const std = @import("std"); const Io = std.Io; const ArenaAllocator = std.heap.ArenaAllocator; const Writer = Io.Writer; const matcha = @import("root.zig"); const Cmd = matcha.Cmd; const Model = matcha.Model; const Renderer = matcha.Renderer; const Msg = matcha.Msg; const View = matcha.View; const Tty = matcha.Tty; const Terminal = matcha.Terminal; const StreamTerminal = matcha.StreamTerminal; const Queue = matcha.Queue; const Value = std.atomic.Value; const Program = @This(); pub const Options = struct { /// Render frame interval in nanoseconds. Default ~60fps. frame_interval_ns: u64 = 16_666_666, /// Capacity of the message queue. msg_queue_capacity: usize = 256, /// Capacity of the cmd queue. cmd_queue_capacity: usize = 64, // ── Backend injection (mirrors bubbletea program options) ────────────── // Providing `input`+`output` drives the model over an arbitrary stream // (e.g. an SSH channel) instead of the local terminal. Both must be set // together, or neither. /// Injected input byte source. Null → read from the local terminal. /// Mirrors bubbletea's `WithInput`. input: ?*Io.Reader = null, /// Injected sink for rendered frames. Null → write to the local terminal. /// Mirrors bubbletea's `WithOutput`. output: ?*Io.Writer = null, /// Environment for the served session (e.g. `TERM` forwarded over SSH). /// Reserved for color-profile detection. Mirrors `WithEnvironment`. environ: ?[]const []const u8 = null, /// Initial window size, used when the backend can't report one (an /// injected stream never can). Mirrors `WithWindowSize`. window_size: ?Terminal.Size = null, /// Forced color profile, seeded to the model at startup when set. /// Mirrors `WithColorProfile`. color_profile: ?Msg.ColorProfile = null, }; /// Process-wide bundle: gpa, io, permanent arena, environ, args, preopens. /// Forwarded to `Model.process` so callbacks can read env vars / argv / etc. process: std.process.Init, model: *Model, options: Options, /// Per-callback scratch arena. Reset before every model callback; the /// resolved Allocator is published to `model.arena` for callbacks to use. /// Distinct from `process.arena`, which is process-permanent. arena: ArenaAllocator, msgs: Queue(Msg), cmds: Queue(Cmd), // Render-side state: model writes its current frame into `view_buf` while // holding `view_mutex`; render thread snapshots it under the same lock. view_buf: Writer.Allocating, view_meta: View = .{}, view_mutex: Io.Mutex = .init, view_cond: Io.Condition = .init, view_dirty: bool = false, shutdown: Value(bool) = .init(false), stdout_buf: [4096]u8 = undefined, /// Active backend, resolved in `run` to the local TTY or an injected /// stream. Points at storage owned by `run`'s stack frame. terminal: *Terminal = undefined, // Worker tasks, dispatched onto `process.io` (fibers on an evented io, // threads on a Threaded io). Awaited/canceled in `teardown`. input_task: ?Io.Future(void) = null, cmd_task: ?Io.Future(void) = null, render_task: ?Io.Future(void) = null, pub fn init( process: std.process.Init, model: *Model, options: Options, ) !Program { return .{ .process = process, .model = model, .options = options, .arena = .init(process.gpa), .msgs = try .init(process.gpa, process.io, options.msg_queue_capacity), .cmds = try .init(process.gpa, process.io, options.cmd_queue_capacity), .view_buf = .init(process.gpa), }; } pub fn deinit(self: *Program) void { self.msgs.deinit(); self.cmds.deinit(); self.view_buf.deinit(); self.arena.deinit(); self.* = undefined; } /// Inject a message from outside the update loop. Used by a host (e.g. dream) /// to deliver out-of-band events such as SSH window-size changes. Mirrors /// bubbletea's `Program.Send`. Silently drops if the program has stopped. pub fn send(self: *Program, msg: Msg) void { self.msgs.push(msg) catch {}; } /// Ask the program to quit after draining queued messages. Mirrors /// bubbletea's `Program.Quit`. pub fn quit(self: *Program) void { self.send(.quit); } /// Force the program down now, without waiting for a quit msg to drain. /// Mirrors bubbletea's `Program.Kill`. pub fn kill(self: *Program) void { self.shutdown.store(true, .release); self.msgs.close(); } pub fn run(self: *Program) !void { self.model.process = self.process; self.model.gpa = self.process.gpa; self.model.arena = self.arena.allocator(); self.model.io = self.process.io; // Resolve the backend: an injected reader/writer pair (e.g. an SSH // channel) when supplied, otherwise the local terminal. Both `input` and // `output` must be provided together. Storage lives on this stack frame // and outlives the worker threads, which `teardown` joins before return. if ((self.options.input == null) != (self.options.output == null)) return error.StreamRequiresInputAndOutput; var tty: Tty = undefined; var stream: StreamTerminal = undefined; if (self.options.input) |input| { stream = StreamTerminal.init(self.process.gpa, self.options.output.?, input); self.terminal = &stream.terminal; } else { tty = try Tty.init(self.process.gpa, self.process.io, &self.stdout_buf); self.terminal = &tty.terminal; } defer self.terminal.deinit(); var renderer: Renderer = .init(self.process.gpa, self.terminal.writer()); defer renderer.deinit(); defer self.teardown(&renderer); const io = self.process.io; self.input_task = try io.concurrent(inputLoop, .{self}); self.cmd_task = try io.concurrent(cmdLoop, .{self}); self.render_task = try io.concurrent(renderLoop, .{ self, &renderer }); _ = self.arena.reset(.retain_capacity); if (self.model.init()) |cmd| try self.cmds.push(cmd); // Seed the forced color profile (if any) before the first frame. if (self.options.color_profile) |profile| try self.msgs.push(.{ .color_profile = profile }); // Seed the model with the initial terminal dimensions so it can lay out // before the first user resize. A configured `window_size` wins; otherwise // ask the backend (the local TTY reports via ioctl; a stream can't). // Later resizes arrive through the input thread via SIGWINCH (posix) or // WINDOW_BUFFER_SIZE_EVENT (windows), or via injected `window_size` msgs. if (self.options.window_size orelse (self.terminal.windowSize() catch null)) |ws| { try self.msgs.push(.{ .window_size = .{ .width = ws.width, .height = ws.height } }); } try self.renderFrame(); // Main update loop. while (self.msgs.pop()) |msg| { switch (msg) { .quit, .interrupt => break, // Batch is internal: fan out the child cmds and skip the model. .batch => |cmds| { for (cmds) |child| self.cmds.push(child) catch break; self.process.gpa.free(cmds); continue; }, else => {}, } _ = self.arena.reset(.retain_capacity); if (self.model.update(msg)) |cmd| { self.cmds.push(cmd) catch break; } // Free heap-allocated msg payloads owned by the queue (see Msg). switch (msg) { .paste => |t| if (t.len > 0) self.process.gpa.free(t), else => {}, } try self.renderFrame(); } } fn teardown(self: *Program, renderer: *Renderer) void { self.shutdown.store(true, .release); const io = self.process.io; // Unblock each loop's wait so awaiting its task can't deadlock: wake the // input backend (local: breaks the poll; stream: sets the woken flag), // close the queues, and signal the render condition. self.terminal.wake(); self.cmds.close(); self.view_mutex.lockUncancelable(io); self.view_cond.broadcast(io); self.view_mutex.unlock(io); self.msgs.close(); // Cancel the input task rather than plain-await it: a stream read blocked // on an idle client can only be interrupted through io cancellation (the // woken flag can't unblock an in-flight read). `cancel` also awaits. The // cmd/render tasks exit cleanly via the close/broadcast above. if (self.input_task) |*t| t.cancel(io); if (self.cmd_task) |*t| t.await(io); if (self.render_task) |*t| t.await(io); renderer.shutdown() catch {}; self.terminal.writer().writeAll("\r\n") catch {}; self.terminal.writer().flush() catch {}; } fn renderFrame(self: *Program) !void { self.view_mutex.lockUncancelable(self.process.io); defer self.view_mutex.unlock(self.process.io); self.view_buf.writer.end = 0; _ = self.arena.reset(.retain_capacity); self.view_meta = try self.model.view(&self.view_buf.writer); self.view_dirty = true; self.view_cond.signal(self.process.io); } fn inputLoop(self: *Program) void { while (!self.shutdown.load(.acquire)) { const maybe_msg = self.terminal.nextMsg() catch break; const msg = maybe_msg orelse break; self.msgs.push(msg) catch return; } } fn cmdLoop(self: *Program) void { while (self.cmds.pop()) |cmd| if (cmd.run(self.process.gpa, self.process.io)) |msg| self.msgs.push(msg) catch return; } fn renderLoop(self: *Program, renderer: *Renderer) void { while (!self.shutdown.load(.acquire)) { self.view_mutex.lockUncancelable(self.process.io); while (!self.view_dirty and !self.shutdown.load(.acquire)) { self.view_cond.waitUncancelable(self.process.io, &self.view_mutex); } if (self.shutdown.load(.acquire)) { self.view_mutex.unlock(self.process.io); break; } const content = self.process.gpa.dupe(u8, self.view_buf.writer.buffered()) catch { self.view_mutex.unlock(self.process.io); continue; }; var view = self.view_meta; const title_dup: ?[]u8 = if (view.window_title) |t| self.process.gpa.dupe(u8, t) catch null else null; if (title_dup) |t| view.window_title = t; self.view_dirty = false; self.view_mutex.unlock(self.process.io); defer self.process.gpa.free(content); defer if (title_dup) |t| self.process.gpa.free(t); renderer.render(content, view) catch break; Io.sleep(self.process.io, .fromNanoseconds(@intCast(self.options.frame_interval_ns)), .awake) catch break; } } test "run drives a model over an injected stream" { const gpa = std.testing.allocator; var threaded: std.Io.Threaded = .init(gpa, .{}); defer threaded.deinit(); // Minimal model: renders one line, quits on 'q'. const M = struct { model: Model = .{ .vtable = &.{ .init = mInit, .update = mUpdate, .view = mView } }, fn mInit(_: *Model) ?Cmd { return null; } fn mUpdate(_: *Model, msg: Msg) ?Cmd { return switch (msg) { .key_press => |k| if (k.rune == 'q') .quit else null, else => null, }; } fn mView(_: *Model, w: *Writer) Writer.Error!View { try w.writeAll("hello stream"); return .{}; } }; var m: M = .{}; var input: Io.Reader = .fixed("q"); var out: Writer.Allocating = .init(gpa); defer out.deinit(); // Only gpa/io are read off `process`; the rest go unused by `run` and the // test model, so leaving them undefined is safe here. const process: std.process.Init = .{ .minimal = undefined, .arena = undefined, .gpa = gpa, .io = threaded.io(), .environ_map = undefined, .preopens = undefined, }; var program: Program = try .init(process, &m.model, .{ .input = &input, .output = &out.writer, .window_size = .{ .width = 80, .height = 24 }, }); defer program.deinit(); try program.run(); // The model's frame was serialized to the injected output sink. try std.testing.expect( std.mem.indexOf(u8, out.writer.buffered(), "hello stream") != null, ); }