//! 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;