From 173a890d1394946b5d7623c66cd34bcd36d8eeb8 Mon Sep 17 00:00:00 2001 From: Amp Date: Wed, 23 Sep 2026 23:02:10 +0000 Subject: [PATCH] loop: recover read failures and interrupt blocked input on stop An input error silently ends the reader, and stopping it can hang while it waits for queue capacity or a terminal response. Widget dispatch also holds the queue mutex while invoking application handlers. Retry recoverable reads with bounded backoff and expose permanent errors through queue closure after buffered events are consumed. Keep key text and partial Windows input alive across reader failures. Close the queue and cancel the task on stop, using a dedicated event to wake native Windows console waits without terminal cooperation. Dispatch a bounded batch of events outside the queue lock so continued input cannot indefinitely postpone frames and timers. Preserve lossless blocking posts by default, and release paste data if enqueueing fails. Add regression coverage for retries, closure, shutdown, and reentrant event dispatch. Linux tests and Windows/macOS cross-builds pass; native Windows runtime testing remains necessary. Amp-Thread-ID: T-01a0cd7f-3c80-709d-9be2-84643b76cfc3 Co-authored-by: Tim Culverhouse Fixes: #367 --- src/Loop.zig | 272 +++++++++++++++++++++++++++++++++++++++-------- src/queue.zig | 92 +++++++++++++++- src/tty.zig | 103 +++++++++++++++++- src/vxfw/App.zig | 145 +++++++++++++++++++------ 4 files changed, 529 insertions(+), 83 deletions(-) diff --git a/src/Loop.zig b/src/Loop.zig index c3c0176..15cefb8 100644 --- a/src/Loop.zig +++ b/src/Loop.zig @@ -22,7 +22,8 @@ pub fn Loop(comptime T: type) type { queue: Queue(T, 512), thread: ?std.Io.Future(void) = null, - should_quit: bool = false, + // Queued key text must outlive the input task, including on failure. + cache: GraphemeCache = .{}, resize_handler_installed: bool = false, /// Initialize the event loop. This is an intrusive init so that we have @@ -70,6 +71,9 @@ pub fn Loop(comptime T: type) type { /// spawns the input thread to read input from the tty pub fn start(self: *Self) !void { if (self.thread) |_| return; + if (builtin.os.tag == .windows and !builtin.is_test) try self.tty.resetInput(); + self.queue.reopen(); + errdefer self.queue.close(error.Closed); self.thread = try self.io.concurrent(Self.ttyRun, .{ self, self.vaxis.opts.system_clipboard_allocator, @@ -80,18 +84,19 @@ pub fn Loop(comptime T: type) type { pub fn stop(self: *Self) void { // If we don't have a thread, we have nothing to stop if (self.thread == null) return; - self.should_quit = true; - // trigger a read - self.vaxis.deviceStatusReport(self.tty.writer()) catch {}; + self.queue.close(error.Closed); + if (builtin.os.tag == .windows and !builtin.is_test) self.tty.interruptInput(); if (self.thread) |*thread| { - thread.await(self.io); + // Interrupt POSIX reads, retry sleeps, and queue waits as well as + // joining. Native Windows console waits use input_stop above. + thread.cancel(self.io); self.thread = null; - self.should_quit = false; } } - /// returns the next available event, blocking until one is available + /// Returns the next event, blocking until available. After buffered events + /// are drained, returns Closed on stop or the input task's failure. pub fn nextEvent(self: *Self) !T { return try self.queue.pop(); } @@ -129,54 +134,43 @@ pub fn Loop(comptime T: type) type { } } - const TtyRunError = error{ - AccessDenied, - Canceled, - CodepointTooLarge, - ConnectionResetByPeer, - EndOfStream, - InputOutput, - InvalidCharacter, - InvalidColorSpec, - InvalidPadding, - InvalidUTF8, - IoctlError, - IsDir, - LockViolation, - NoSpaceLeft, - NotOpenForReading, - OutOfMemory, - Overflow, - SocketUnconnected, - SystemResources, - Unexpected, - Utf8CannotEncodeSurrogateHalf, - WouldBlock, - }; - /// read input from the tty. This is run in a separate thread fn ttyRun(self: *Self, paste_allocator: ?std.mem.Allocator) void { - self._ttyRun(paste_allocator) catch {}; + self._ttyRun(paste_allocator) catch |err| self.inputFailed(err); + } + + fn inputFailed(self: *Self, err: anyerror) void { + if (err != error.Canceled and err != error.Closed) + log.warn("input stopped: {s}", .{@errorName(err)}); + self.queue.close(if (err == error.Canceled) error.Closed else err); + } + + fn runWindows(self: *Self, reader: anytype, paste_allocator: ?std.mem.Allocator) !void { + var parser: Parser = .{ .cursor_position_requests = &self.vaxis.cursor_position_requests }; + var retry: ReadRetry = .{}; + while (true) { + try self.io.checkCancel(); + const event = reader.nextEvent(&parser, paste_allocator) catch |err| { + if (malformedInput(err)) continue; + try retry.wait(self.io, err); + continue; + }; + retry = .{}; + try handleEventGeneric(self, self.vaxis, &self.cache, Event, event, paste_allocator); + } } fn _ttyRun( self: *Self, paste_allocator: ?std.mem.Allocator, - ) TtyRunError!void { + ) !void { // Return early if we're in test mode to avoid infinite loops if (builtin.is_test) return; - // initialize a grapheme cache - var cache: GraphemeCache = .{}; var parser: Parser = .{ .cursor_position_requests = &self.vaxis.cursor_position_requests }; switch (builtin.os.tag) { - .windows => { - while (!self.should_quit) { - const event = try self.tty.nextEvent(&parser, paste_allocator); - try handleEventGeneric(self, self.vaxis, &cache, Event, event, paste_allocator); - } - }, + .windows => try self.runWindows(self.tty, paste_allocator), else => { // get our initial winsize const winsize = try self.tty.getWinsize(); @@ -187,9 +181,16 @@ pub fn Loop(comptime T: type) type { // initialize the read buffer var buf: [1024]u8 = undefined; var read_start: usize = 0; + var retry: ReadRetry = .{}; // read loop - read_loop: while (!self.should_quit) { - const bytes_read = try self.tty.read(buf[read_start..]); + read_loop: while (true) { + try self.io.checkCancel(); + const bytes_read = self.tty.read(buf[read_start..]) catch |err| { + try retry.wait(self.io, err); + continue; + }; + if (bytes_read == 0) return error.EndOfStream; + retry = .{}; const n = read_start + bytes_read; var seq_start: usize = 0; while (seq_start < n) { @@ -208,7 +209,13 @@ pub fn Loop(comptime T: type) type { } } - const result = try parser.parse(buf[seq_start..n], paste_allocator); + const result = parser.parse(buf[seq_start..n], paste_allocator) catch |err| { + if (!malformedInput(err)) return err; + // There is no consumed length on parse errors. Discard + // this batch rather than retrying the same bad bytes. + read_start = 0; + continue :read_loop; + }; if (result.n == 0) { // copy the read to the beginning. We don't use memcpy because // this could be overlapping, and it's also rare @@ -223,7 +230,7 @@ pub fn Loop(comptime T: type) type { seq_start += result.n; const event = result.event orelse continue; - try handleEventGeneric(self, self.vaxis, &cache, Event, event, paste_allocator); + try handleEventGeneric(self, self.vaxis, &self.cache, Event, event, paste_allocator); } } }, @@ -232,6 +239,41 @@ pub fn Loop(comptime T: type) type { }; } +const ReadRetry = struct { + attempts: u8 = 0, + + fn delay(self: *ReadRetry, err: anyerror) !std.Io.Duration { + switch (err) { + error.InputInterrupted, error.WouldBlock, error.InputOutput, error.SystemResources => {}, + else => return err, + } + if (self.attempts == 8) return err; + const ms = @min(@as(u32, 10) << @intCast(self.attempts), 250); + self.attempts += 1; + return .fromMilliseconds(ms); + } + + fn wait(self: *ReadRetry, io: std.Io, err: anyerror) !void { + const duration = try self.delay(err); + if (self.attempts == 1) log.warn("input read failed: {s}; retrying", .{@errorName(err)}); + try io.sleep(duration, .awake); + } +}; + +fn malformedInput(err: anyerror) bool { + return switch (err) { + error.InvalidCharacter, + error.InvalidColorSpec, + error.InvalidPadding, + error.InvalidUTF8, + error.Utf8CannotEncodeSurrogateHalf, + error.CodepointTooLarge, + error.Overflow, + => true, + else => false, + }; +} + // Use return on the self.postEvent's so it can either return error union or void pub fn handleEventGeneric(self: anytype, vx: *Vaxis, cache: *GraphemeCache, Event: type, event: anytype, paste_allocator: ?std.mem.Allocator) !void { switch (event) { @@ -251,6 +293,7 @@ pub fn handleEventGeneric(self: anytype, vx: *Vaxis, cache: *GraphemeCache, Even }, .paste => |text| { if (@hasField(Event, "paste")) { + errdefer if (paste_allocator) |allocator| allocator.free(text); return self.postEvent(.{ .paste = text }); } if (paste_allocator) |allocator| allocator.free(text); @@ -580,6 +623,143 @@ test "paste dispatch preserves boundaries and transfers or frees clipboard text" try testing.expectEqual(@as(?KeyEvent, null), try keys.tryEvent()); } +test "read retries back off, cap attempts, and reject permanent errors" { + const testing = std.testing; + var retry: ReadRetry = .{}; + for ([_]i64{ 10, 20, 40, 80, 160, 250, 250, 250 }) |ms| { + const duration = try retry.delay(error.InputInterrupted); + try testing.expectEqual(ms, duration.toMilliseconds()); + } + try testing.expectError(error.InputInterrupted, retry.delay(error.InputInterrupted)); + retry = .{}; + for ([_]anyerror{ error.Canceled, error.AccessDenied, error.InvalidHandle, error.OutOfMemory, error.EndOfStream }) |err| { + try testing.expectError(err, retry.delay(err)); + } + try testing.expectEqual(0, retry.attempts); +} + +test "Windows reader recovers, resets retries, and reports failure after queued keys" { + const testing = std.testing; + const Event = union(enum) { key_press: vaxis.Key }; + const Reader = struct { + calls: usize = 0, + + fn nextEvent(self: *@This(), _: *Parser, _: ?std.mem.Allocator) !vaxis.Event { + defer self.calls += 1; + return switch (self.calls) { + 0...7, 9 => error.InputInterrupted, + 8 => .{ .key_press = .{ .codepoint = 'a', .text = "abc" } }, + 10 => error.InvalidUTF8, + 11 => .{ .key_press = .{ .codepoint = 'z', .text = "zy" } }, + else => error.AccessDenied, + }; + } + + fn run(self: *@This(), loop: *Loop(Event)) void { + loop.runWindows(self, null) catch |err| loop.inputFailed(err); + } + }; + var vx: Vaxis = undefined; + var loop: Loop(Event) = .init(testing.io, undefined, &vx); + var reader: Reader = .{}; + var task = try testing.io.concurrent(Reader.run, .{ &reader, &loop }); + task.await(testing.io); + try testing.expectEqual(13, reader.calls); + const first = (try loop.nextEvent()).key_press; + const second = (try loop.nextEvent()).key_press; + try testing.expectEqual('a', first.codepoint); + try testing.expectEqual('z', second.codepoint); + try testing.expectEqualStrings("abc", first.text.?); + try testing.expectEqualStrings("zy", second.text.?); + try testing.expect(first.text.?.ptr == loop.cache.buf[0..].ptr); + try testing.expectError(error.AccessDenied, loop.nextEvent()); + try testing.expectError(error.AccessDenied, loop.tryEvent()); + try testing.expectError(error.AccessDenied, loop.pollEvent()); +} + +test "stop interrupts a reader posting to a full queue" { + const testing = std.testing; + const Event = union(enum) { focus_in }; + const Reader = struct { + ready: std.Io.Event = .unset, + calls: usize = 0, + + fn nextEvent(self: *@This(), _: *Parser, _: ?std.mem.Allocator) !vaxis.Event { + self.calls += 1; + self.ready.set(testing.io); + return .focus_in; + } + + fn run(self: *@This(), loop: *Loop(Event)) void { + loop.runWindows(self, null) catch |err| loop.inputFailed(err); + } + }; + var vx: Vaxis = undefined; + var loop: Loop(Event) = .init(testing.io, undefined, &vx); + for (0..512) |_| try loop.postEvent(.focus_in); + var reader: Reader = .{}; + loop.thread = try testing.io.concurrent(Reader.run, .{ &reader, &loop }); + defer loop.stop(); + try reader.ready.wait(testing.io); + try testing.io.sleep(.fromMilliseconds(10), .awake); + loop.stop(); + try testing.expectEqual(1, reader.calls); + try testing.expect(loop.thread == null); + for (0..512) |_| try testing.expectEqual(Event.focus_in, try loop.nextEvent()); + try testing.expectError(error.Closed, loop.nextEvent()); +} + +test "stop cancels an idle POSIX read without a terminal response" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + const testing = std.testing; + const Event = union(enum) { focus_in }; + var tty = try Tty.init(testing.io, &.{}); + defer tty.deinit(); + const Reader = struct { + tty: @import("tty.zig").PosixTty, + ready: std.Io.Event = .unset, + result: anyerror!usize = undefined, + + fn run(self: *@This(), loop: *Loop(Event)) void { + var buf: [16]u8 = undefined; + self.ready.set(testing.io); + self.result = self.tty.read(&buf); + _ = self.result catch |err| { + loop.inputFailed(err); + return; + }; + } + }; + var reader: Reader = .{ .tty = undefined }; + reader.tty.io = testing.io; + reader.tty.fd = .{ .handle = tty.pipe_read, .flags = .{ .nonblocking = false } }; + var vx: Vaxis = undefined; + var loop: Loop(Event) = .init(testing.io, &tty, &vx); + loop.thread = try testing.io.concurrent(Reader.run, .{ &reader, &loop }); + defer loop.stop(); + try reader.ready.wait(testing.io); + try testing.io.sleep(.fromMilliseconds(10), .awake); + loop.stop(); + try testing.expectError(error.Canceled, reader.result); + try testing.expectError(error.Closed, loop.tryEvent()); +} + +test "failed paste enqueue frees its allocation" { + const Event = union(enum) { paste: []const u8 }; + var vx: Vaxis = undefined; + var loop: Loop(Event) = .init(std.testing.io, undefined, &vx); + loop.queue.close(error.Closed); + const text = try std.testing.allocator.dupe(u8, "owned paste"); + try std.testing.expectError(error.Closed, handleEventGeneric( + &loop, + &vx, + &loop.cache, + Event, + @as(vaxis.Event, .{ .paste = text }), + std.testing.allocator, + )); +} + test { std.testing.refAllDecls(@This()); } diff --git a/src/queue.zig b/src/queue.zig index 35cf854..3f7cf4b 100644 --- a/src/queue.zig +++ b/src/queue.zig @@ -12,6 +12,7 @@ pub fn Queue( read_index: usize = 0, write_index: usize = 0, + closed: ?anyerror = null, io: std.Io, mutex: std.Io.Mutex = .init, @@ -26,11 +27,30 @@ pub fn Queue( return .{ .io = io }; } + /// Reject further writes and wake all waiters. Buffered items remain readable; + /// once drained, readers receive the first close reason. + pub fn close(self: *Self, reason: anyerror) void { + self.mutex.lockUncancelable(self.io); + defer self.mutex.unlock(self.io); + if (self.closed != null) return; + self.closed = reason; + self.not_full.broadcast(self.io); + self.not_empty.broadcast(self.io); + } + + /// Reopen after the previous producer has stopped. Retains buffered items. + pub fn reopen(self: *Self) void { + self.mutex.lockUncancelable(self.io); + defer self.mutex.unlock(self.io); + self.closed = null; + } + /// Pop an item from the queue. Blocks until an item is available. pub fn pop(self: *Self) !T { try self.mutex.lock(self.io); defer self.mutex.unlock(self.io); while (self.isEmptyLH()) { + if (self.closed) |err| return err; try self.not_empty.wait(self.io, &self.mutex); } std.debug.assert(!self.isEmptyLH()); @@ -43,8 +63,10 @@ pub fn Queue( try self.mutex.lock(self.io); defer self.mutex.unlock(self.io); while (self.isFullLH()) { + if (self.closed) |err| return err; try self.not_full.wait(self.io, &self.mutex); } + if (self.closed) |err| return err; std.debug.assert(!self.isFullLH()); self.pushAndSignalLH(item); } @@ -55,6 +77,7 @@ pub fn Queue( pub fn tryPush(self: *Self, item: T) !bool { try self.mutex.lock(self.io); defer self.mutex.unlock(self.io); + if (self.closed) |err| return err; if (self.isFullLH()) return false; self.pushAndSignalLH(item); return true; @@ -65,7 +88,7 @@ pub fn Queue( pub fn tryPop(self: *Self) !?T { try self.mutex.lock(self.io); defer self.mutex.unlock(self.io); - if (self.isEmptyLH()) return null; + if (self.isEmptyLH()) return if (self.closed) |err| err else null; return self.popAndSignalLH(); } @@ -74,6 +97,7 @@ pub fn Queue( try self.mutex.lock(self.io); defer self.mutex.unlock(self.io); while (self.isEmptyLH()) { + if (self.closed) |err| return err; try self.not_empty.wait(self.io, &self.mutex); } std.debug.assert(!self.isEmptyLH()); @@ -371,6 +395,72 @@ test "2 writers" { try t2.await(io); } +test "close preserves buffered items and the first failure" { + var q: Queue(u8, 2) = .init(testing.io); + try q.push(17); + try q.push(29); + q.close(error.InputOutput); + q.close(error.Closed); + try testing.expectError(error.InputOutput, q.push(31)); + try testing.expectError(error.InputOutput, q.tryPush(31)); + try q.poll(); + try testing.expectEqual(17, try q.pop()); + try testing.expectEqual(29, (try q.tryPop()).?); + try testing.expectError(error.InputOutput, q.pop()); + try testing.expectError(error.InputOutput, q.tryPop()); + try testing.expectError(error.InputOutput, q.poll()); + // A closed queue rejects writes even when it has free capacity. + try testing.expectError(error.InputOutput, q.push(31)); + try testing.expectError(error.InputOutput, q.tryPush(31)); + q.reopen(); + try testing.expectEqual(null, try q.tryPop()); + try q.push(31); + try testing.expectEqual(31, try q.pop()); +} + +test "close wakes all blocked writers and readers" { + const Waiter = struct { + ready: std.Io.Event = .unset, + + fn write(self: *@This(), q: *Queue(u8, 1)) !void { + self.ready.set(q.io); + try testing.expectError(error.Closed, q.push(99)); + } + + fn read(self: *@This(), q: *Queue(u8, 1)) !void { + self.ready.set(q.io); + try testing.expectError(error.Closed, q.pop()); + } + + fn poll(self: *@This(), q: *Queue(u8, 1)) !void { + self.ready.set(q.io); + try testing.expectError(error.Closed, q.poll()); + } + }; + const io = testing.io; + var full: Queue(u8, 1) = .init(io); + var empty: Queue(u8, 1) = .init(io); + try full.push(7); + var waiters: [4]Waiter = @splat(.{}); + var a = try io.concurrent(Waiter.write, .{ &waiters[0], &full }); + defer a.cancel(io) catch {}; + var b = try io.concurrent(Waiter.write, .{ &waiters[1], &full }); + defer b.cancel(io) catch {}; + var c = try io.concurrent(Waiter.read, .{ &waiters[2], &empty }); + defer c.cancel(io) catch {}; + var d = try io.concurrent(Waiter.poll, .{ &waiters[3], &empty }); + defer d.cancel(io) catch {}; + for (&waiters) |*waiter| try waiter.ready.wait(io); + try io.sleep(.fromMilliseconds(10), .awake); + full.close(error.Closed); + empty.close(error.Closed); + try a.await(io); + try b.await(io); + try c.await(io); + try d.await(io); + try testing.expectEqual(7, try full.pop()); +} + test { std.testing.refAllDecls(@This()); } diff --git a/src/tty.zig b/src/tty.zig index 8e98754..620a873 100644 --- a/src/tty.zig +++ b/src/tty.zig @@ -307,6 +307,7 @@ pub const PosixTty = struct { pub const WindowsTty = struct { stdin: windows.HANDLE, stdout: windows.HANDLE, + input_stop: windows.HANDLE, initial_codepage: c_uint, initial_input_mode: CONSOLE_MODE_INPUT, @@ -350,6 +351,9 @@ pub const WindowsTty = struct { pub fn init(io: std.Io, buffer: []u8) !WindowsTty { const stdin: std.Io.File = .stdin(); const stdout: std.Io.File = .stdout(); + const input_stop = CreateEventW(null, .TRUE, .FALSE, null) orelse + return windows.unexpectedError(windows.GetLastError()); + errdefer windows.CloseHandle(input_stop); // get initial modes const initial_output_codepage = GetConsoleOutputCP(); @@ -368,6 +372,7 @@ pub const WindowsTty = struct { var self: WindowsTty = .{ .stdin = stdin.handle, .stdout = stdout.handle, + .input_stop = input_stop, .initial_codepage = initial_output_codepage, .initial_input_mode = initial_input_mode, .initial_output_mode = initial_output_mode, @@ -391,6 +396,7 @@ pub const WindowsTty = struct { _ = SetConsoleOutputCP(self.initial_codepage); setConsoleMode(self.stdin, self.initial_input_mode) catch {}; setConsoleMode(self.stdout, self.initial_output_mode) catch {}; + windows.CloseHandle(self.input_stop); windows.CloseHandle(self.stdin); windows.CloseHandle(self.stdout); } @@ -442,15 +448,53 @@ pub const WindowsTty = struct { return posix.read(self.fd, buf); } + pub fn resetInput(self: *WindowsTty) !void { + if (ResetEvent(self.input_stop) == .FALSE) + return windows.unexpectedError(windows.GetLastError()); + self.event_state = .{}; + } + + pub fn interruptInput(self: *WindowsTty) void { + // The handle is owned by this Tty and remains open until deinit. + if (SetEvent(self.input_stop) == .FALSE) @panic("invalid console stop event"); + } + + fn inputError(code: windows.Win32Error) anyerror { + return switch (code) { + .INVALID_HANDLE => error.InvalidHandle, + .ACCESS_DENIED => error.AccessDenied, + .OPERATION_ABORTED => error.InputInterrupted, + .NOT_READY, .BUSY, .RETRY => error.WouldBlock, + .NOT_ENOUGH_MEMORY, .NO_SYSTEM_RESOURCES => error.SystemResources, + else => { + std.log.scoped(.vaxis).warn("console input failed: Win32 error {d}", .{@intFromEnum(code)}); + return error.Unexpected; + }, + }; + } + pub fn nextEvent(self: *WindowsTty, parser: *Parser, paste_allocator: ?std.mem.Allocator) !Event { - // We use a loop so we can ignore certain events + // Keep partial ANSI and UTF-16 input across transient console read errors. while (true) { + // A console handle is signaled while input is available. Put the stop + // event first so shutdown wins even during a continuous input stream. + // This Tty must be the sole reader of the console input buffer. + const handles = [_]windows.HANDLE{ self.input_stop, self.stdin }; + switch (WaitForMultipleObjects(handles.len, &handles, .FALSE, 0xffffffff)) { + 0 => return error.Canceled, + 1 => {}, + else => return inputError(windows.GetLastError()), + } var event_count: u32 = 0; var input_record: INPUT_RECORD = undefined; if (ReadConsoleInputW(self.stdin, &input_record, 1, &event_count) == .FALSE) - return windows.unexpectedError(windows.GetLastError()); + return inputError(windows.GetLastError()); - if (try self.eventFromRecord(&input_record, &self.event_state, parser, paste_allocator)) |ev| { + const event = self.eventFromRecord(&input_record, &self.event_state, parser, paste_allocator) catch |err| { + self.event_state = .{}; + return err; + }; + if (event) |ev| { return ev; } } @@ -811,7 +855,7 @@ pub const WindowsTty = struct { // the size directly when we get this event var console_info: CONSOLE_SCREEN_BUFFER_INFO = undefined; if (GetConsoleScreenBufferInfo(self.stdout, &console_info) == .FALSE) { - return windows.unexpectedError(windows.GetLastError()); + return inputError(windows.GetLastError()); } const window_rect = console_info.srWindow; const width = window_rect.Right - window_rect.Left + 1; @@ -936,6 +980,10 @@ pub const WindowsTty = struct { pub const PINPUT_RECORD = *INPUT_RECORD; + extern "kernel32" fn CreateEventW(?*windows.SECURITY_ATTRIBUTES, windows.BOOL, windows.BOOL, ?windows.LPCWSTR) callconv(.winapi) ?windows.HANDLE; + extern "kernel32" fn SetEvent(windows.HANDLE) callconv(.winapi) windows.BOOL; + extern "kernel32" fn ResetEvent(windows.HANDLE) callconv(.winapi) windows.BOOL; + extern "kernel32" fn WaitForMultipleObjects(windows.DWORD, [*]const windows.HANDLE, windows.BOOL, windows.DWORD) callconv(.winapi) windows.DWORD; pub extern "kernel32" fn ReadConsoleInputW(hConsoleInput: windows.HANDLE, lpBuffer: PINPUT_RECORD, nLength: windows.DWORD, lpNumberOfEventsRead: *windows.DWORD) callconv(.winapi) windows.BOOL; pub extern "kernel32" fn GetConsoleOutputCP() callconv(.winapi) windows.UINT; pub extern "kernel32" fn GetConsoleMode(kConsoleHandle: windows.HANDLE, lpMode: *windows.DWORD) callconv(.winapi) windows.BOOL; @@ -1078,6 +1126,7 @@ const WindowsInputTest = struct { tty: WindowsTty = .{ .stdin = undefined, .stdout = undefined, + .input_stop = undefined, .initial_codepage = 0, .initial_input_mode = .{}, .initial_output_mode = .{}, @@ -1152,6 +1201,52 @@ test "Windows paste delimiters wrapped in win32-input-mode" { } } +test "console errors distinguish interruption from permanent failures" { + try std.testing.expectEqual(error.InputInterrupted, WindowsTty.inputError(.OPERATION_ABORTED)); + try std.testing.expectEqual(error.InvalidHandle, WindowsTty.inputError(.INVALID_HANDLE)); + try std.testing.expectEqual(error.AccessDenied, WindowsTty.inputError(.ACCESS_DENIED)); + try std.testing.expectEqual(error.SystemResources, WindowsTty.inputError(.NO_SYSTEM_RESOURCES)); +} + +test "Windows input wait is interruptible without a console response" { + if (builtin.os.tag != .windows) return error.SkipZigTest; + const io = std.testing.io; + // An unsignaled event stands in for an idle console handle. Cancellation + // must return without ever reaching ReadConsoleInputW. + const input = WindowsTty.CreateEventW(null, .TRUE, .FALSE, null) orelse return error.Unexpected; + defer windows.CloseHandle(input); + const stop = WindowsTty.CreateEventW(null, .TRUE, .FALSE, null) orelse return error.Unexpected; + defer windows.CloseHandle(stop); + var tty: WindowsTty = undefined; + tty.stdin = input; + tty.input_stop = stop; + const Reader = struct { + fn run(t: *WindowsTty, ready: *std.Io.Event) !void { + var parser: Parser = .{}; + ready.set(std.testing.io); + try std.testing.expectError(error.Canceled, t.nextEvent(&parser, null)); + } + }; + var ready: std.Io.Event = .unset; + var task = try io.concurrent(Reader.run, .{ &tty, &ready }); + defer { + tty.interruptInput(); + task.cancel(io) catch {}; + } + try ready.wait(io); + try io.sleep(.fromMilliseconds(10), .awake); + tty.interruptInput(); + try task.await(io); + + try tty.resetInput(); + try std.testing.expectEqual(@as(u32, 258), WindowsTty.WaitForMultipleObjects(1, &.{stop}, .FALSE, 0)); + // Shutdown also wins if input and stop are both already signaled. + try std.testing.expect(WindowsTty.SetEvent(input) != .FALSE); + tty.interruptInput(); + var parser: Parser = .{}; + try std.testing.expectError(error.Canceled, tty.nextEvent(&parser, null)); +} + test { std.testing.refAllDecls(@This()); } diff --git a/src/vxfw/App.zig b/src/vxfw/App.zig index 57a2d7c..1cf122d 100644 --- a/src/vxfw/App.zig +++ b/src/vxfw/App.zig @@ -123,38 +123,7 @@ pub fn run(self: *App, widget: vxfw.Widget, opts: Options) anyerror!void { } try self.checkTimers(&ctx); - - { - try loop.queue.lock(); - defer loop.queue.unlock(); - while (loop.queue.drain()) |event| { - defer resetEventState(&ctx); - switch (event) { - .key_press => { - try focus_handler.handleEvent(&ctx, event); - try self.handleCommand(&ctx.cmds); - }, - .focus_out => { - try mouse_handler.mouseExit(self, &ctx); - try focus_handler.handleEvent(&ctx, .focus_out); - try self.handleCommand(&ctx.cmds); - }, - .focus_in => { - try focus_handler.handleEvent(&ctx, .focus_in); - try self.handleCommand(&ctx.cmds); - }, - .mouse => |mouse| try mouse_handler.handleMouse(self, &ctx, mouse), - .winsize => |ws| { - try vx.resize(self.allocator, tty.writer(), ws); - ctx.redraw = true; - }, - else => { - try focus_handler.handleEvent(&ctx, event); - try self.handleCommand(&ctx.cmds); - }, - } - } - } + try self.dispatchEvents(&loop, &ctx, &mouse_handler, &focus_handler); // If we have a focus change, handle that event before we layout if (self.wants_focus) |wants_focus| { @@ -201,6 +170,46 @@ pub fn run(self: *App, widget: vxfw.Widget, opts: Options) anyerror!void { } } +fn dispatchEvents( + self: *App, + loop: *EventLoop, + ctx: *vxfw.EventContext, + mouse_handler: *MouseHandler, + focus_handler: *FocusHandler, +) !void { + // Bound the batch so continuously arriving input cannot starve frames + // or timers. tryEvent releases the queue mutex before handlers run. + for (0..loop.queue.buf.len) |_| { + const event = try loop.tryEvent() orelse break; + defer resetEventState(ctx); + switch (event) { + .key_press => { + try focus_handler.handleEvent(ctx, event); + try self.handleCommand(&ctx.cmds); + }, + .focus_out => { + try mouse_handler.mouseExit(self, ctx); + try focus_handler.handleEvent(ctx, .focus_out); + try self.handleCommand(&ctx.cmds); + }, + .focus_in => { + try focus_handler.handleEvent(ctx, .focus_in); + try self.handleCommand(&ctx.cmds); + }, + .mouse => |mouse| try mouse_handler.handleMouse(self, ctx, mouse), + .winsize => |ws| { + try self.vx.resize(self.allocator, self.tty.writer(), ws); + ctx.redraw = true; + }, + else => { + try focus_handler.handleEvent(ctx, event); + try self.handleCommand(&ctx.cmds); + }, + } + if (ctx.quit) return; + } +} + fn doLayout( self: *App, widget: vxfw.Widget, @@ -728,6 +737,78 @@ test "FocusHandler: removed focus falls back to root without calling the removed } } +test "event dispatch unlocks before handlers and yields a continuously replenished queue" { + const testing = std.testing; + const TestWidget = struct { + loop: *EventLoop, + keys: usize = 0, + ticks: usize = 0, + + fn handle(userdata: *anyopaque, ctx: *vxfw.EventContext, event: vxfw.Event) anyerror!void { + const self: *@This() = @ptrCast(@alignCast(userdata)); + switch (event) { + .key_press => { + // Fail rather than deadlock if dispatch still holds the lock. + try testing.expect(self.loop.queue.mutex.tryLock()); + self.loop.queue.mutex.unlock(testing.io); + self.keys += 1; + if (self.keys > 512) return error.UnboundedBatch; + try self.loop.postEvent(event); + ctx.redraw = true; + }, + .tick => self.ticks += 1, + else => {}, + } + } + + fn draw(_: *anyopaque, _: vxfw.DrawContext) Allocator.Error!vxfw.Surface { + unreachable; + } + }; + var app: App = .{ + .io = testing.io, + .allocator = testing.allocator, + .tty = undefined, + .vx = undefined, + .timers = .empty, + .wants_focus = null, + }; + defer app.timers.deinit(testing.allocator); + var loop: EventLoop = .init(testing.io, &app.tty, &app.vx); + var test_widget: TestWidget = .{ .loop = &loop }; + const widget: Widget = .{ + .userdata = &test_widget, + .eventHandler = TestWidget.handle, + .drawFn = TestWidget.draw, + }; + var mouse = MouseHandler.init(widget); + defer mouse.deinit(testing.allocator); + var focus = FocusHandler.init(testing.allocator, widget); + defer focus.deinit(testing.allocator); + try focus.path_to_focused.append(testing.allocator, widget); + var ctx: vxfw.EventContext = .{ + .io = testing.io, + .alloc = testing.allocator, + .phase = .capturing, + .cmds = .empty, + .consume_event = false, + .redraw = false, + .quit = false, + }; + defer ctx.cmds.deinit(testing.allocator); + try loop.postEvent(.{ .key_press = .{ .codepoint = 'x' } }); + try app.dispatchEvents(&loop, &ctx, &mouse, &focus); + try testing.expectEqual(512, test_widget.keys); + try testing.expect(ctx.redraw); + try testing.expectEqual('x', (try loop.tryEvent()).?.key_press.codepoint); + try app.timers.append(testing.allocator, .{ + .deadline = std.Io.Timestamp.now(testing.io, .awake).addDuration(.fromMilliseconds(-1)), + .widget = widget, + }); + try app.checkTimers(&ctx); + try testing.expectEqual(1, test_widget.ticks); +} + test "timer consume does not leak to the next event" { const testing = std.testing; -- 2.51.2