//! A `Terminal` backed by an injected reader/writer pair rather than the //! local TTY. This is the seam that lets a `Program` render a model into //! something other than the controlling terminal — e.g. an SSH channel //! (dream). Mirrors bubbletea's `WithInput`/`WithOutput`. //! //! Input bytes are pulled from `reader` and decoded through the same //! `InputBuffer` the TTY backends use, so key/mouse/paste parsing is //! identical. Window-size events are *not* produced here — over a remote //! transport they arrive out-of-band, so the host injects them with //! `Program.send(.{ .window_size = ... })`. //! //! Cancellation contract: `wake` makes the *next* `nextMsg` return null, but //! it cannot interrupt a `readSliceShort` already blocked on `reader`. The //! owner of the stream is responsible for making a blocked read return (by //! closing its end) when it wants the input loop to stop; a closed source //! surfaces as end-of-stream, which `nextMsg` reports as null. const std = @import("std"); const Io = std.Io; const Allocator = std.mem.Allocator; const Terminal = @import("Terminal.zig"); const InputBuffer = @import("InputBuffer.zig"); const Msg = @import("msg.zig").Msg; const StreamTerminal = @This(); terminal: Terminal = .{ .vtable = &vtable }, reader: *Io.Reader, out: *Io.Writer, input_buf: InputBuffer, woken: std.atomic.Value(bool) = .init(false), read_buf: [256]u8 = undefined, pub fn init(gpa: Allocator, out: *Io.Writer, reader: *Io.Reader) StreamTerminal { return .{ .reader = reader, .out = out, .input_buf = .init(gpa), }; } pub fn deinit(self: *StreamTerminal) void { self.input_buf.deinit(); } pub fn nextMsg(self: *StreamTerminal) !?Msg { while (true) { if (self.woken.load(.acquire)) return null; // Drain any complete msg already buffered before pulling more bytes. if (try self.input_buf.next()) |msg| return msg; // Single short read: return whatever bytes are available now rather // than blocking for a full buffer (`readSliceShort` would). EOF and // read failures (incl. a canceled read on teardown) end the loop. var data: [1][]u8 = .{&self.read_buf}; const n = self.reader.readVec(&data) catch return null; if (n > 0) self.input_buf.append(self.read_buf[0..n]); } } pub fn writer(self: *StreamTerminal) *Io.Writer { return self.out; } pub fn windowSize(_: *StreamTerminal) !Terminal.Size { // Unknown for a raw stream — the Program falls back to its configured // initial size, and later sizes arrive via injected window_size msgs. return error.NoWindowSize; } pub fn wake(self: *StreamTerminal) void { self.woken.store(true, .release); } // ─── Terminal vtable ───────────────────────────────────────────────────── const vtable: Terminal.VTable = .{ .nextMsg = vNextMsg, .writer = vWriter, .windowSize = vWindowSize, .wake = vWake, .deinit = vDeinit, }; fn vNextMsg(t: *Terminal) anyerror!?Msg { const self: *StreamTerminal = @fieldParentPtr("terminal", t); return self.nextMsg(); } fn vWriter(t: *Terminal) *Io.Writer { const self: *StreamTerminal = @fieldParentPtr("terminal", t); return self.writer(); } fn vWindowSize(t: *Terminal) anyerror!Terminal.Size { const self: *StreamTerminal = @fieldParentPtr("terminal", t); return self.windowSize(); } fn vWake(t: *Terminal) void { const self: *StreamTerminal = @fieldParentPtr("terminal", t); self.wake(); } fn vDeinit(t: *Terminal) void { const self: *StreamTerminal = @fieldParentPtr("terminal", t); self.deinit(); } test "decodes keys pulled from an injected reader" { var input: Io.Reader = .fixed("ab\x1b[A"); var sink: Io.Writer = .fixed(&.{}); var term: StreamTerminal = .init(std.testing.allocator, &sink, &input); defer term.deinit(); try std.testing.expectEqual( Msg.Key.Kind.rune, (try term.nextMsg()).?.key_press.kind, ); try std.testing.expectEqual( @as(u21, 'b'), (try term.nextMsg()).?.key_press.rune, ); try std.testing.expectEqual( Msg.Key.Kind.up, (try term.nextMsg()).?.key_press.kind, ); // Source drained → EOF → null. try std.testing.expectEqual(@as(?Msg, null), try term.nextMsg()); } test "wake makes nextMsg return null" { var input: Io.Reader = .fixed("aaaa"); var sink: Io.Writer = .fixed(&.{}); var term: StreamTerminal = .init(std.testing.allocator, &sink, &input); defer term.deinit(); term.wake(); try std.testing.expectEqual(@as(?Msg, null), try term.nextMsg()); }