Something went wrong. Try again.
A charm-like tui library
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302//! An accepted SSH session. Mirrors charmbracelet/ssh.Session.//!//! Phase 1 architecture: the session-fiber owns the underlying libssh session//! (calls `ssh_event_dopoll`, drains output, etc). The handler runs on a//! separate fiber and never touches libssh directly — it reads from this//! Session's input buffer and writes into its output buffer. Coordination is//! via std.Io.Mutex + Condition; window-change events arrive through a//! bounded queue with drop-oldest semantics.const Session = @This();
const log = std.log.scoped(.dream);
pub const window_queue_capacity: usize = 8;
allocator: Allocator,io: Io,
ssh_session: *c.ssh_session_struct,channel: *c.ssh_channel_struct,
user_name: []const u8 = "",command_argv: []const []const u8 = &.{},pty_value: ?Pty = null,last_window: Window = .{ .width = 0, .height = 0 },
window_changes: Queue(Window),
exit_status: u32 = 0,
in_mutex: Io.Mutex = .init,in_cond: Io.Condition = .init,in_buf: std.ArrayList(u8) = .empty,in_eof: bool = false,
out_mutex: Io.Mutex = .init,out_buf: std.ArrayList(u8) = .empty,err_buf: std.ArrayList(u8) = .empty,closed: bool = false,
// std.Io.Reader/Writer views over the input/output buffers, so a matcha// Program (or any std.Io consumer) can drive the session as a plain stream.// The reader has an empty buffer (reads land straight in the caller's dest);// the writer buffers in `out_writer_buf` and drains into `out_buf`.in_reader: Io.Reader = undefined,out_writer: Io.Writer = undefined,out_writer_buf: []u8 = &.{},/// Underlying error stashed when a vtable call must report ReadFailed//// WriteFailed (which carry no payload).in_err: ?anyerror = null,out_err: ?anyerror = null,
pub fn init( io: Io, allocator: Allocator, ssh_session: *c.ssh_session_struct, channel: *c.ssh_channel_struct,) !Session { const out_buf = try allocator.alloc(u8, 4096); errdefer allocator.free(out_buf); return .{ .allocator = allocator, .io = io, .ssh_session = ssh_session, .channel = channel, .window_changes = try .init(allocator, io, window_queue_capacity), .in_reader = .{ .vtable = &.{ .stream = inStream, .readVec = inReadVec }, .buffer = &.{}, .seek = 0, .end = 0, }, .out_writer = .{ .vtable = &.{ .drain = outDrain }, .buffer = out_buf, .end = 0, }, .out_writer_buf = out_buf, };}
pub fn deinit(self: *Session) void { self.in_buf.deinit(self.allocator); self.out_buf.deinit(self.allocator); self.err_buf.deinit(self.allocator); self.allocator.free(self.out_writer_buf); self.window_changes.deinit(); self.* = undefined;}
/// `std.Io.Reader` view over the client→server byte stream. Valid once the/// Session is at its final address (the vtable recovers `*Session` from this/// field via `@fieldParentPtr`).pub fn reader(self: *Session) *Io.Reader { return &self.in_reader;}
/// `std.Io.Writer` view over the server→client byte stream. Bytes land in/// `out_buf`; the session-fiber flushes them to libssh via `drainOutput`.pub fn writer(self: *Session) *Io.Writer { return &self.out_writer;}
// ─── accessors (charmbracelet/ssh.Session parity) ───────────────────────────
pub fn user(self: *const Session) []const u8 { return self.user_name;}
pub fn command(self: *const Session) []const []const u8 { return self.command_argv;}
/// Returns the PTY metadata if the client requested one, else null.pub fn pty(self: *const Session) ?Pty { return self.pty_value;}
/// Latest window size known to the server. Updated as window-change events/// arrive; `window_changes` queue also yields each event.pub fn window(self: *const Session) Window { return self.last_window;}
/// Block until the next window-change event. Returns null once window events/// have been stopped (handler finishing) or the session has closed.pub fn nextWindow(self: *Session) ?Window { return self.window_changes.pop();}
/// Stop delivering window-change events; unblocks any fiber parked in/// `nextWindow`. Idempotent. A handler that spawned a window watcher should/// call this before awaiting that watcher.pub fn stopWindowEvents(self: *Session) void { self.window_changes.close();}
// ─── handler-side I/O ───────────────────────────────────────────────────────
/// Read up to `buf.len` bytes from the channel. Blocks until at least one/// byte is available or EOF is received. Returns 0 on EOF.pub fn read(self: *Session, buf: []u8) !usize { self.in_mutex.lockUncancelable(self.io); defer self.in_mutex.unlock(self.io); // Cancelable wait: a Program tearing down (model quit) cancels its input // task, which surfaces here as error.Canceled even with an idle client. // `Condition.wait` re-locks the mutex before returning, so the `defer // unlock` above stays correct on the cancel path. while (self.in_buf.items.len == 0 and !self.in_eof) { try self.in_cond.wait(self.io, &self.in_mutex); } if (self.in_buf.items.len == 0) return 0; const n = @min(buf.len, self.in_buf.items.len); @memcpy(buf[0..n], self.in_buf.items[0..n]); // Shift the remainder forward. Cheap enough; ring buffer if profiling // later says it matters. const rest = self.in_buf.items.len - n; if (rest > 0) std.mem.copyForwards(u8, self.in_buf.items[0..rest], self.in_buf.items[n..]); self.in_buf.shrinkRetainingCapacity(rest); return n;}
pub fn write(self: *Session, bytes: []const u8) !usize { self.out_mutex.lockUncancelable(self.io); defer self.out_mutex.unlock(self.io); if (self.closed) return error.SessionClosed; try self.out_buf.appendSlice(self.allocator, bytes); return bytes.len;}
pub fn writeStderr(self: *Session, bytes: []const u8) !usize { self.out_mutex.lockUncancelable(self.io); defer self.out_mutex.unlock(self.io); if (self.closed) return error.SessionClosed; try self.err_buf.appendSlice(self.allocator, bytes); return bytes.len;}
pub fn exit(self: *Session, status: u32) void { self.exit_status = status;}
// ─── server-side helpers (called only from session-fiber) ───────────────────
/// Append bytes received from the client to the input buffer. Returns the/// number of bytes consumed (always all of them at the moment).pub fn pushInput(self: *Session, bytes: []const u8) !usize { self.in_mutex.lockUncancelable(self.io); defer self.in_mutex.unlock(self.io); try self.in_buf.appendSlice(self.allocator, bytes); self.in_cond.broadcast(self.io); return bytes.len;}
pub fn markEof(self: *Session) void { self.in_mutex.lockUncancelable(self.io); defer self.in_mutex.unlock(self.io); self.in_eof = true; self.in_cond.broadcast(self.io);}
/// Drain queued output to libssh. Called by the session-fiber between event/// polls. Caller owns the libssh session (no internal lock here).pub fn drainOutput(self: *Session) void { var stdout_owned: []u8 = &.{}; var stderr_owned: []u8 = &.{}; { self.out_mutex.lockUncancelable(self.io); defer self.out_mutex.unlock(self.io); if (self.out_buf.items.len > 0) { stdout_owned = self.out_buf.toOwnedSlice(self.allocator) catch &.{}; } if (self.err_buf.items.len > 0) { stderr_owned = self.err_buf.toOwnedSlice(self.allocator) catch &.{}; } } defer if (stdout_owned.len > 0) self.allocator.free(stdout_owned); defer if (stderr_owned.len > 0) self.allocator.free(stderr_owned);
if (stdout_owned.len > 0) { const w = c.ssh_channel_write(self.channel, stdout_owned.ptr, @intCast(stdout_owned.len)); if (w < 0) log.warn("ssh_channel_write failed", .{}); } if (stderr_owned.len > 0) { const w = c.ssh_channel_write_stderr(self.channel, stderr_owned.ptr, @intCast(stderr_owned.len)); if (w < 0) log.warn("ssh_channel_write_stderr failed", .{}); }}
pub fn setWindow(self: *Session, win: Window) void { self.last_window = win; self.window_changes.pushLatest(win) catch {};}
pub fn close(self: *Session) void { self.out_mutex.lockUncancelable(self.io); self.closed = true; self.out_mutex.unlock(self.io); self.markEof(); self.window_changes.close();}
// ─── std.Io.Reader / Writer vtables ─────────────────────────────────────────
fn inStream(io_r: *Io.Reader, io_w: *Io.Writer, limit: Io.Limit) Io.Reader.StreamError!usize { const dest = limit.slice(try io_w.writableSliceGreedy(1)); var data: [1][]u8 = .{dest}; const n = try inReadVec(io_r, &data); io_w.advance(n); return n;}
fn inReadVec(io_r: *Io.Reader, data: [][]u8) Io.Reader.Error!usize { const self: *Session = @fieldParentPtr("in_reader", io_r); // `read` fills one contiguous slice, so serve the first non-empty caller // slice directly. (The reader carries no buffer of its own, so we don't go // through `writableVector`, which could otherwise hand back an empty dest.) for (data) |slice| { if (slice.len == 0) continue; const n = self.read(slice) catch |e| { self.in_err = e; return error.ReadFailed; }; if (n == 0) return error.EndOfStream; return n; } return 0;}
fn outDrain(io_w: *Io.Writer, data: []const []const u8, splat: usize) Io.Writer.Error!usize { const self: *Session = @fieldParentPtr("out_writer", io_w); self.out_mutex.lockUncancelable(self.io); defer self.out_mutex.unlock(self.io); if (self.closed) { self.out_err = error.SessionClosed; return error.WriteFailed; } // Flush the writer's own buffered bytes, then the vectored `data`; the // last slice repeats `splat` times (0 is legal). Return counts bytes // consumed from `data` only — buffered bytes are already accounted for. self.out_buf.appendSlice(self.allocator, io_w.buffered()) catch return error.WriteFailed; io_w.end = 0; var consumed: usize = 0; for (data[0 .. data.len - 1]) |slice| { self.out_buf.appendSlice(self.allocator, slice) catch return error.WriteFailed; consumed += slice.len; } const pattern = data[data.len - 1]; var i: usize = 0; while (i < splat) : (i += 1) { self.out_buf.appendSlice(self.allocator, pattern) catch return error.WriteFailed; } return consumed + pattern.len * splat;}
const Allocator = std.mem.Allocator;const Io = std.Io;const std = @import("std");const c = @import("libssh");const Pty = @import("Pty.zig").Pty;const Window = @import("Pty.zig").Window;const Queue = @import("queue.zig").Queue;