diff options
| author | Gabriel Schneider <[email protected]> | 2026-09-29 19:49:19 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-10-01 00:12:17 -0300 |
| commit | d1e93d4cdc2e2625d5a1724a7167a5c3e64fe6a2 (patch) | |
| tree | 32ec6983d7ef0c8b828a041447dce53928961d11 | |
| parent | 156891c9346ee15aea5e10edbc80850f0cf36797 (diff) | |
| download | pardes-d1e93d4cdc2e2625d5a1724a7167a5c3e64fe6a2.tar.gz pardes-d1e93d4cdc2e2625d5a1724a7167a5c3e64fe6a2.zip | |
Over a slow link the tty keeps one frame in flight: it asks the terminal to answer after each, and draws the latest state when it does
| -rw-r--r-- | src/tty/tty.zig | 172 | ||||
| -rw-r--r-- | test/e2e_harness.zig | 28 |
2 files changed, 194 insertions, 6 deletions
diff --git a/src/tty/tty.zig b/src/tty/tty.zig index df1e5d02..f5ae078a 100644 --- a/src/tty/tty.zig +++ b/src/tty/tty.zig @@ -132,7 +132,14 @@ fn readInput(loop: *Loop, tty: anytype, cache: *vaxis.GraphemeCache) !void { if (carried == buf.len) return error.InputSequenceTooLong; const received = try tty.read(buf[carried..]); if (received == 0) return; - const end = carried + received; + // The terminal's answers to frames (FrameAck) come out first, + // wherever they sit: never events, and never glued to a key (an + // Escape pressed just before one would read as Alt on it). + const end = FrameAck.take(buf[0 .. carried + received], loop); + if (end == 0) { + carried = 0; + continue; + } if (kitty_shm_probing.load(.acquire)) if (std.mem.indexOf(u8, buf[0..end], std.fmt.comptimePrint("\x1b_Gi={d};", .{kitty_shm_probe_id}))) |at| { kitty_shm_probing.store(false, .release); kitty_shm.store(std.mem.startsWith(u8, buf[at..end], std.fmt.comptimePrint("\x1b_Gi={d};OK", .{kitty_shm_probe_id})), .release); @@ -153,7 +160,8 @@ fn readInput(loop: *Loop, tty: anytype, cache: *vaxis.GraphemeCache) !void { // its writes whole, so no sequence arrives cut right after its ESC. // Waiting on it cost every Escape press 50 ms on a legacy terminal // (tmux, most ssh sessions), and 1 s before kitty's flags were on. - var waited = carried == 0 and received == 1; + // (A terminal's answer to a frame is one write too: FrameAck.) + var waited = carried == 0 and end == 1; while (consumed < parse_end) { // A lone ESC that ends a longer read is the Escape key, or the // first byte of a sequence whose rest the read boundary held @@ -162,6 +170,9 @@ fn readInput(loop: *Loop, tty: anytype, cache: *vaxis.GraphemeCache) !void { // kitty keyboard protocol Escape is `CSI 27 u`, so there the ESC // can only start a sequence and the wait can be long. const rest = buf[consumed..parse_end]; + // A frame's answer cut short by the read waits for its rest. + if (parse_end == end and rest.len > 1 and rest.len < FrameAck.answer.len and + std.mem.startsWith(u8, FrameAck.answer, rest)) break; if (!waited and parse_end == end and rest[0] == 0x1b and (rest.len == 1 or (rest.len == 2 and rest[1] == 0x1b))) { waited = true; if (try inputFollows(tty, loop.io, if (loop.vaxis.caps.kitty_keyboard) 1000 else 50)) break; @@ -191,6 +202,63 @@ fn readInput(loop: *Loop, tty: anytype, cache: *vaxis.GraphemeCache) !void { } } +/// One frame in flight at most. Each frame the tty draws ends with a device +/// status request (`CSI 5 n`), and no other frame is drawn until the +/// terminal's `CSI 0 n` comes back: over ssh the frames a burst of scrolling +/// or typing would queue behind each other (sshd's window holds megabytes) +/// are never drawn, and the burst ends one frame and one round trip after +/// its last input. On a local terminal the answer is back in well under a +/// millisecond. A terminal that never answers turns this off for good after +/// `timeout_ns`. +const FrameAck = struct { + const request = "\x1b[5n"; + const answer = "\x1b[0n"; + const timeout_ns: u64 = std.time.ns_per_s; + var waiting = std.atomic.Value(bool).init(false); + var off = std.atomic.Value(bool).init(false); + /// When the request went out (the loop's thread only). + var sent_ns: u64 = 0; + + /// Take every answer out of `bytes`, closing up the rest; the new + /// length. The loop is woken for the frame that waited on one. + fn take(bytes: []u8, loop: *Loop) usize { + var len = bytes.len; + var found = false; + while (std.mem.indexOf(u8, bytes[0..len], answer)) |at| { + std.mem.copyForwards(u8, bytes[at..], bytes[at + answer.len .. len]); + len -= answer.len; + found = true; + } + if (found and waiting.swap(false, .acq_rel)) loop.postEvent(.nop) catch {}; + return len; + } + + /// Whether a new frame must wait (and the loop has a wake at the + /// timeout); past the timeout the terminal is taken to never answer. + fn holds(now_ns: u64, io: std.Io) bool { + if (off.load(.acquire) or !waiting.load(.acquire)) return false; + const waited_ns = now_ns -| sent_ns; + if (waited_ns >= timeout_ns) { + off.store(true, .release); + waiting.store(false, .release); + return false; + } + requestWake(io, @intCast((timeout_ns - waited_ns + std.time.ns_per_ms - 1) / std.time.ns_per_ms)); + return true; + } + + /// After a frame that wrote something: ask for the answer. + fn ask(w: *std.Io.Writer, now_ns: u64) void { + if (off.load(.acquire)) return; + // Waiting before the request goes: a local terminal can answer + // before the write returns, and an answer nobody waited for is lost. + sent_ns = now_ns; + waiting.store(true, .release); + w.writeAll(request) catch return waiting.store(false, .release); + w.flush() catch return waiting.store(false, .release); + } +}; + /// Whether more input is readable within `ms`, polled in short slices so /// that a cancel of the input thread is still seen. A test reader answers /// for itself. @@ -521,6 +589,98 @@ test "an ESC read by itself is the Escape key at once" { vx.caps.kitty_keyboard = false; } +test "a frame's answer is taken wherever the reads cut it, and never typed" { + if (comptime builtin.os.tag == .windows) return error.SkipZigTest; + const gpa = std.testing.allocator; + const io = std.testing.io; + var env = try std.testing.environ.createMap(gpa); + defer env.deinit(); + var vx = try vaxis.init(io, gpa, &env, .{}); + var output: std.Io.Writer.Allocating = .init(gpa); + defer output.deinit(); + defer vx.deinit(gpa, &output.writer); + var tty: vaxis.Tty = undefined; + var loop: Loop = .init(io, &tty, &vx); + var cache: vaxis.GraphemeCache = .{}; + const Reader = struct { + parts: [2][]const u8, + next: usize = 0, + fn getWinsize(_: *@This()) !vaxis.Winsize { + return .{ .rows = 24, .cols = 80, .x_pixel = 0, .y_pixel = 0 }; + } + fn inputFollows(self: *@This(), _: u32) bool { + return self.next < self.parts.len; + } + fn read(self: *@This(), buf: []u8) !usize { + while (self.next < self.parts.len) { + const part = self.parts[self.next]; + self.next += 1; + if (part.len == 0) continue; + @memcpy(buf[0..part.len], part); + return part.len; + } + return 0; + } + }; + const all = "x" ++ FrameAck.answer ++ "y"; + defer FrameAck.waiting.store(false, .release); + for (0..all.len) |cut| { + FrameAck.waiting.store(true, .release); + var reader: Reader = .{ .parts = .{ all[0..cut], all[cut..] } }; + try readInput(&loop, &reader, &cache); + try std.testing.expect(!FrameAck.waiting.load(.acquire)); + var typed: std.ArrayList(u8) = .empty; + defer typed.deinit(gpa); + var woke = false; + while (try loop.tryEvent()) |event| switch (event) { + .winsize => {}, + .nop => woke = true, + .key_press => |key| try typed.appendSlice(gpa, key.text orelse return error.UnexpectedInputKey), + else => return error.UnexpectedInputEvent, + }; + std.testing.expectEqualStrings("xy", typed.items) catch |err| { + std.debug.print("cut at {d}\n", .{cut}); + return err; + }; + try std.testing.expect(woke); + } + // An Escape pressed just before an answer, in the same read, is the + // Escape key, not Alt on the answer. + FrameAck.waiting.store(true, .release); + var glued: Reader = .{ .parts = .{ "\x1b" ++ FrameAck.answer, "j" } }; + try readInput(&loop, &glued, &cache); + try std.testing.expectEqual(.winsize, std.meta.activeTag((try loop.tryEvent()).?)); + try std.testing.expectEqual(.nop, std.meta.activeTag((try loop.tryEvent()).?)); + const escape = (try loop.tryEvent()).?.key_press; + try std.testing.expectEqual(vaxis.Key.escape, escape.codepoint); + try std.testing.expect(!escape.mods.alt); + try std.testing.expectEqual(@as(u21, 'j'), (try loop.tryEvent()).?.key_press.codepoint); + try std.testing.expect(try loop.tryEvent() == null); +} + +test "a terminal that never answers a frame stops being asked after one timeout" { + const io = std.testing.io; + defer { + FrameAck.off.store(false, .release); + FrameAck.waiting.store(false, .release); + FrameAck.sent_ns = 0; + } + var out: std.Io.Writer.Allocating = .init(std.testing.allocator); + defer out.deinit(); + const t0: u64 = 5 * std.time.ns_per_s; + FrameAck.ask(&out.writer, t0); + try std.testing.expectEqualStrings(FrameAck.request, out.written()); + // Unanswered: the next frame waits, until the timeout. + try std.testing.expect(FrameAck.holds(t0 + 10 * std.time.ns_per_ms, io)); + try std.testing.expect(FrameAck.holds(t0 + 999 * std.time.ns_per_ms, io)); + try std.testing.expect(!FrameAck.holds(t0 + FrameAck.timeout_ns, io)); + // Then never again: no request, no wait. + out.clearRetainingCapacity(); + FrameAck.ask(&out.writer, t0 + 2 * std.time.ns_per_s); + try std.testing.expectEqual(@as(usize, 0), out.written().len); + try std.testing.expect(!FrameAck.holds(t0 + 2 * std.time.ns_per_s, io)); +} + test "terminal input cancellation joins blocked reads and queued EOF" { if (comptime builtin.os.tag == .windows) return error.SkipZigTest; const gpa = std.testing.allocator; @@ -1438,6 +1598,12 @@ const Shell = struct { const s = of(ctx); const vx = s.vx; s.tracks = canonical.panelTracks(); + // The last frame is not answered yet: this one stays owed, and the + // answer's wake draws the state as it is then. + if (FrameAck.holds(host_io.monotonicNs(), s.io)) { + s.core.present_skipped = true; + return; + } _ = s.frame.reset(.retain_capacity); const surface = panel_compositor.compose( s.frame.allocator(), @@ -1523,7 +1689,9 @@ const Shell = struct { } if (surface.cursor) |cur| paintCursor(win, cur.x, cur.y, cur.bar); const tz_render = tracy.zone(@src(), "vx.render"); + const before = s.tty.tty_writer.pos + s.tty.writer().end; vx.render(s.tty.writer()) catch {}; + if (s.tty.tty_writer.pos + s.tty.writer().end != before) FrameAck.ask(s.tty.writer(), host_io.monotonicNs()); tz_render.end(); } diff --git a/test/e2e_harness.zig b/test/e2e_harness.zig index eaa19b43..9e67c49f 100644 --- a/test/e2e_harness.zig +++ b/test/e2e_harness.zig @@ -64,6 +64,8 @@ pub const Harness = struct { /// when true, print the captured screen state after each pump/waitFor and on /// every assertion, so live test runs can be inspected (zig build test -Dtrace). trace: bool = false, + /// How much of a `CSI 5 n` the output has ended on (`feed`). + status_request: u8 = 0, /// every raw byte the app has emitted (accumulated in pump). Lets tests assert /// on control sequences the emulator consumes and never renders (e.g. OSC 52 /// clipboard writes). Capture is bounded explicitly so a runaway child cannot @@ -147,6 +149,24 @@ pub const Harness = struct { /// Read pty output and feed it to our ghostty terminal for `ms` ms. After /// this, the grid reflects everything the app rendered so far. + /// The app's output into the emulator. The one query answered is the + /// device status request (`CSI 5 n` → `CSI 0 n`), as every terminal + /// does: the tty asks it after each frame and draws no other until the + /// answer (FrameAck in src/tty/tty.zig). + fn feed(self: *Harness, bytes: []const u8) void { + self.stream.nextSlice(bytes); + const request = "\x1b[5n"; + for (bytes) |b| { + if (b == request[self.status_request]) { + self.status_request += 1; + } else self.status_request = if (b == 0x1b) 1 else 0; + if (self.status_request == request.len) { + self.status_request = 0; + _ = libc.write(self.master, "\x1b[0n", 4); + } + } + } + pub fn pump(self: *Harness, ms: i64) !void { const deadline = nowMs() + ms; var buf: [4096]u8 = undefined; @@ -157,7 +177,7 @@ pub const Harness = struct { const n = posix.read(self.master, &buf) catch break; if (n == 0) break; self.recordRaw(buf[0..n]); - self.stream.nextSlice(buf[0..n]); + self.feed(buf[0..n]); } } if (self.trace) self.traceScreen("pump"); @@ -176,7 +196,7 @@ pub const Harness = struct { const n = posix.read(self.master, &buf) catch return false; if (n == 0) return false; self.recordRaw(buf[0..n]); - self.stream.nextSlice(buf[0..n]); + self.feed(buf[0..n]); return true; } @@ -243,7 +263,7 @@ pub const Harness = struct { const n = posix.read(self.master, &buf) catch break; if (n == 0) break; self.recordRaw(buf[0..n]); - self.stream.nextSlice(buf[0..n]); + self.feed(buf[0..n]); } const text = try self.screenText(); defer self.gpa.free(text); @@ -270,7 +290,7 @@ pub const Harness = struct { const n = posix.read(self.master, &buf) catch break; if (n == 0) break; self.recordRaw(buf[0..n]); - self.stream.nextSlice(buf[0..n]); + self.feed(buf[0..n]); if (std.mem.indexOf(u8, self.raw.items, needle) != null) return true; } return false; |
