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;