diff options
Diffstat (limited to 'src/ninep/events.zig')
| -rw-r--r-- | src/ninep/events.zig | 558 |
1 files changed, 558 insertions, 0 deletions
diff --git a/src/ninep/events.zig b/src/ninep/events.zig new file mode 100644 index 00000000..fb49e41b --- /dev/null +++ b/src/ninep/events.zig @@ -0,0 +1,558 @@ +//! Event streams: per-pane `event` records, the editor-wide `log`, and the +//! bounded queues that park a read until something happens. +const std = @import("std"); +const pardes = @import("../pardes.zig"); +const tree = @import("tree.zig"); +const pane_files = @import("pane.zig"); + +const Pardes = pardes.Pardes; +const Pane = pardes.Pane; +const MAX_PANES = pardes.MAX_PANES; +const Req = tree.Req; +const Reply = tree.Reply; +const E = tree.E; + +pub const queue_cap = 64 * 1024; + +/// Wall-clock seconds, or zero where the platform has no clock. +pub fn now() u32 { + if (comptime !pardes.hosted) return 0; + var ts: std.c.timespec = undefined; + if (std.c.clock_gettime(.REALTIME, &ts) != 0) return 0; + return std.math.cast(u32, ts.sec) orelse 0; +} + +pub const Queue = struct { + buf: std.ArrayList(u8) = .empty, + head: usize = 0, + + pub fn deinit(q: *Queue, gpa: std.mem.Allocator) void { + q.buf.deinit(gpa); + q.head = 0; + } + + pub fn push(q: *Queue, gpa: std.mem.Allocator, record: []const u8) void { + if (record.len > std.math.maxInt(u32)) return; + while (q.buf.items.len - q.head + record.len + 4 > queue_cap) { + if (q.peek() == null) return; + q.pop(); + } + q.compact(); + var head: [4]u8 = undefined; + std.mem.writeInt(u32, &head, @intCast(record.len), .little); + q.buf.appendSlice(gpa, &head) catch return; + q.buf.appendSlice(gpa, record) catch { + q.buf.shrinkRetainingCapacity(q.buf.items.len - 4); + return; + }; + } + + pub fn peek(q: *const Queue) ?[]const u8 { + const rest = q.buf.items[@min(q.head, q.buf.items.len)..]; + if (rest.len < 4) return null; + const len = std.mem.readInt(u32, rest[0..4], .little); + if (rest.len < 4 + len) return null; + return rest[4 .. 4 + len]; + } + + pub fn pop(q: *Queue) void { + const record = q.peek() orelse return; + q.head += 4 + record.len; + if (q.head == q.buf.items.len) { + q.buf.clearRetainingCapacity(); + q.head = 0; + } + } + + pub fn popFront(q: *Queue, n: usize) void { + const record = q.peek() orelse return; + if (n >= record.len) return q.pop(); + q.head += n; + std.mem.writeInt(u32, q.buf.items[q.head..][0..4], @intCast(record.len - n), .little); + } + + fn compact(q: *Queue) void { + if (q.head == 0 or q.head * 2 < q.buf.items.len) return; + const rest = q.buf.items.len - q.head; + std.mem.copyForwards(u8, q.buf.items[0..rest], q.buf.items[q.head..]); + q.buf.shrinkRetainingCapacity(rest); + q.head = 0; + } + + pub fn empty(q: *const Queue) bool { + return q.peek() == null; + } + + pub fn clearAndFree(q: *Queue, gpa: std.mem.Allocator) void { + q.buf.clearAndFree(gpa); + q.head = 0; + } +}; + +/// One record per read; `.again` parks the read until a record arrives. +pub fn readQueue(p: *Pardes, req: Req, q: *Queue) Reply { + const record = q.peek() orelse return .{ .tag = req.tag, .status = .again }; + if (req.size < record.len) return Reply.fail(req.tag, E.INVAL); + const out = p.fs.stage(p.gpa); + out.appendSlice(p.gpa, record) catch return Reply.fail(req.tag, E.NOMEM); + q.pop(); + return .{ .tag = req.tag, .payload = .{ .staged = @intCast(out.items.len) } }; +} + +// ---- the editor-wide log ---- + +pub const LogKind = enum { new, del, rename, save }; + +/// A pane installed now is announced by `announce` once the update that made +/// it has finished, when its name is known. +pub fn noteInstall(p: *Pardes, id: usize) void { + if (id < MAX_PANES) p.fs.panes[id].unannounced = true; +} + +/// A pane leaving; one never announced leaves silently. +pub fn noteRetire(p: *Pardes, id: usize, pane: *Pane) void { + if (id >= MAX_PANES) return; + if (p.fs.panes[id].unannounced) { + p.fs.panes[id].unannounced = false; + return; + } + noteLog(p, .del, pane); +} + +/// Called at the end of every update: reports the panes installed by it. +pub fn announce(p: *Pardes) void { + for (p.panes, 0..) |slot, id| { + const pf = &p.fs.panes[id]; + if (!pf.unannounced) continue; + pf.unannounced = false; + if (slot) |pane| noteLog(p, .new, pane); + } +} + +/// Records `<kind> <serial> <name>` while a reader holds /log open. +pub fn noteLog(p: *Pardes, kind: LogKind, pane: *Pane) void { + if (p.fs.log_readers == 0) return; + var buf: [4096 + 64]u8 = undefined; + const name = pane_files.nameOf(pane); + const record = std.fmt.bufPrint(&buf, "{s} {d} {s}\n", .{ @tagName(kind), pane.serial, name[0..@min(name.len, 4096)] }) catch return; + p.fs.log.push(p.gpa, record); +} + +// ---- per-pane event records ---- + +pub const max_record_text = 256; + +pub const Action = enum(u8) { + body_delete = 'D', + tag_delete = 'd', + body_insert = 'I', + tag_insert = 'i', + body_look = 'L', + tag_look = 'l', + body_exec = 'X', + tag_exec = 'x', + + pub fn char(a: Action) u8 { + return @intFromEnum(a); + } + + pub fn fromChar(c: u8) ?Action { + return std.enums.fromInt(Action, c); + } + + pub fn onTag(a: Action) bool { + return @intFromEnum(a) >= 'a'; + } +}; + +pub const flag_builtin: u32 = 1; +pub const flag_expansion: u32 = 2; +pub const flag_filename: u32 = 4; +pub const flag_chorded: u32 = 8; + +pub fn formatRecord( + buf: []u8, + origin: u8, + action: Action, + q0: u32, + q1: u32, + flag: u32, + text: []const u8, +) []const u8 { + const sent = if (text.len >= max_record_text) text[0..0] else text; + return std.fmt.bufPrint(buf, "{c}{c}{d} {d} {d} {d} {s}\n", .{ + origin, + action.char(), + q0, + q1, + flag, + sent.len, + sent, + }) catch buf[0..0]; +} + +pub const Span = struct { at: u32, removed: u32, inserted: u32 }; + +pub fn diffSpan(old: []const u8, new: []const u8) Span { + const both = @min(old.len, new.len); + const stride = 64; + var head: usize = 0; + while (head + stride <= both and + std.mem.eql(u8, old[head..][0..stride], new[head..][0..stride])) head += stride; + while (head < both and old[head] == new[head]) head += 1; + var tail: usize = 0; + const rest = both - head; + while (tail + stride <= rest and std.mem.eql( + u8, + old[old.len - tail - stride ..][0..stride], + new[new.len - tail - stride ..][0..stride], + )) tail += stride; + while (tail < rest and old[old.len - 1 - tail] == new[new.len - 1 - tail]) tail += 1; + return .{ + .at = @intCast(head), + .removed = @intCast(old.len - tail - head), + .inserted = @intCast(new.len - tail - head), + }; +} + +pub fn noteReplace(p: *Pardes, id: usize, on_tag: bool, old: []const u8, new: []const u8) void { + if (!p.fs.scripted(id)) return; + const span = diffSpan(old, new); + if (span.removed == 0 and span.inserted == 0) return; + if (span.removed > 0) _ = noteAction( + p, + id, + if (on_tag) .tag_delete else .body_delete, + span.at, + span.at + span.removed, + 0, + "", + ); + if (span.inserted > 0) _ = noteAction( + p, + id, + if (on_tag) .tag_insert else .body_insert, + span.at, + span.at + span.inserted, + 0, + new[span.at..][0..span.inserted], + ); +} + +pub fn noteAction( + p: *Pardes, + id: usize, + action: Action, + q0: u32, + q1: u32, + flag: u32, + text: []const u8, +) bool { + if (!p.fs.scripted(id)) return false; + var buf: [max_record_text + 64]u8 = undefined; + const record = formatRecord(&buf, p.fs.origin, action, q0, q1, flag, text); + p.fs.panes[id].events.push(p.gpa, record); + return true; +} + +pub fn notePtyOutput(p: *Pardes, id: usize, bytes: []const u8) void { + if (id >= MAX_PANES or bytes.len == 0) return; + const pf = &p.fs.panes[id]; + if (pf.pty_readers == 0) return; + var off: usize = 0; + while (off < bytes.len) { + const n = @min(bytes.len - off, queue_cap / 2); + pf.pty_out.push(p.gpa, bytes[off..][0..n]); + off += n; + } +} + +const EventRecord = struct { action: Action, q0: u32, q1: u32 }; + +const EventReader = struct { + data: []const u8, + i: usize = 0, + + fn next(er: *EventReader) ?EventRecord { + if (er.i >= er.data.len) return null; + var i = er.i; + if (i + 2 > er.data.len) return null; + i += 1; + const action = Action.fromChar(er.data[i]) orelse return null; + i += 1; + const q0 = scanNumber(er.data, &i) orelse return null; + const q1 = scanNumber(er.data, &i) orelse return null; + while (i < er.data.len and er.data[i] == ' ') i += 1; + if (i >= er.data.len or er.data[i] != '\n') return null; + er.i = i + 1; + return .{ .action = action, .q0 = q0, .q1 = q1 }; + } +}; + +fn scanNumber(data: []const u8, i: *usize) ?u32 { + while (i.* < data.len and data[i.*] == ' ') i.* += 1; + const s = i.*; + var n: u64 = 0; + while (i.* < data.len and data[i.*] >= '0' and data[i.*] <= '9') : (i.* += 1) + n = @min(n * 10 + (data[i.*] - '0'), std.math.maxInt(u32)); + if (i.* == s) return null; + return @intCast(n); +} + +/// Writing a Look or Exec record back performs the action it names. +pub fn writeEvent(p: *Pardes, req: Req, id: usize) Reply { + const pane0 = p.panes[id] orelse return Reply.fail(req.tag, E.NOENT); + const serial = pane0.serial; + { + const body = pane_files.bodyOf(pane0); + const tag = pane_files.tagOf(p, pane0); + var check: EventReader = .{ .data = req.data }; + while (check.next()) |r| { + switch (r.action) { + .body_look, .tag_look, .body_exec, .tag_exec => {}, + else => return Reply.fail(req.tag, E.INVAL), + } + const n = if (r.action.onTag()) tag.len else body.len; + if (r.q0 > r.q1 or r.q1 > n) return Reply.fail(req.tag, E.INVAL); + } + if (check.i != req.data.len) return Reply.fail(req.tag, E.INVAL); + } + var run: EventReader = .{ .data = req.data }; + while (run.next()) |r| { + const live = p.paneBySerial(serial) orelse break; + const pane = p.panes[live].?; + const whole = if (r.action.onTag()) pane_files.tagOf(p, pane) else pane_files.bodyOf(pane); + const lo = @min(@as(usize, r.q0), whole.len); + const hi = @max(lo, @min(@as(usize, r.q1), whole.len)); + const text = p.scratch.allocator().dupe(u8, whole[lo..hi]) catch continue; + switch (r.action) { + .body_exec, .tag_exec => _ = p.execute(live, text), + .body_look, .tag_look => p.lookAt(live, text), + else => unreachable, + } + } + return .{ .tag = req.tag, .written = @intCast(req.data.len) }; +} + +// ---- tests ---- + +const testing = std.testing; +const th = @import("testing.zig"); +const call = th.call; +const rd = th.rd; +const wr = th.wr; +const withFile = th.withFile; +const serialOf = th.serialOf; +const Node = tree.Node; +const Status = tree.Status; + +test "event records are acme's bytes, one per read, and .again when empty" { + const gpa = testing.allocator; + const p = try withFile(gpa, "Msg fs-ran\n"); + defer p.deinit(); + const serial = serialOf(p); + const event = Node.of(serial, .event); + + _ = noteAction(p, 0, .body_exec, 1, 4, flag_builtin, "sg "); + try testing.expect(p.fs.panes[0].events.empty()); + + const h = call(p, .{ .tag = 10, .op = .open, .node = event }); + try testing.expect(h.reply.handle != 0); + try testing.expectEqual(@as(u16, 1), p.fs.listeners); + + try testing.expectEqual(Status.again, rd(p, event, 0, 4096).reply.status); + + p.fs.origin = 'M'; + _ = noteAction(p, 0, .body_exec, 1, 4, flag_builtin, "ell"); + _ = noteAction(p, 0, .body_delete, 0, 3, 0, ""); + try testing.expectEqualStrings("MX1 4 1 3 ell\n", rd(p, event, 0, 4096).bytes); + try testing.expectEqualStrings("MD0 3 0 0 \n", rd(p, event, 0, 4096).bytes); + try testing.expectEqual(Status.again, rd(p, event, 0, 4096).reply.status); + + _ = noteAction(p, 0, .body_look, 0, 3, flag_filename, "one"); + try testing.expectEqual(E.INVAL, rd(p, event, 0, 4).errno()); + try testing.expectEqualStrings("ML0 3 4 3 one\n", rd(p, event, 0, 4096).bytes); + + const big = "z" ** max_record_text; + _ = noteAction(p, 0, .body_exec, 0, max_record_text, 0, big); + try testing.expectEqualStrings("MX0 256 0 0 \n", rd(p, event, 0, 4096).bytes); + + _ = call(p, .{ .tag = 11, .op = .release, .node = event, .handle = h.reply.handle }); + try testing.expectEqual(@as(u16, 0), p.fs.listeners); +} + +fn drainEvents(p: *Pardes, node: u64, store: []u8, out: [][]const u8) [][]const u8 { + var used: usize = 0; + var n: usize = 0; + while (n < out.len) { + const a = rd(p, node, 0, 4096); + if (a.reply.status != .ok) break; + @memcpy(store[used..][0..a.bytes.len], a.bytes); + out[n] = store[used..][0..a.bytes.len]; + used += a.bytes.len; + n += 1; + } + return out[0..n]; +} + +test "a write through the filesystem is reported once, attributed to the file it came through" { + const gpa = testing.allocator; + const p = try withFile(gpa, "one\ntwo\n"); + defer p.deinit(); + const serial = serialOf(p); + const event = Node.of(serial, .event); + _ = call(p, .{ .tag = 40, .op = .open, .node = event }); + var store: [4096]u8 = undefined; + var slots: [16][]const u8 = undefined; + _ = drainEvents(p, event, &store, &slots); + + _ = wr(p, Node.of(serial, .body), "three\n"); + const body_recs = drainEvents(p, event, &store, &slots); + try testing.expect(body_recs.len >= 1); + try testing.expectEqualStrings("EI8 14 0 6 three\n\n", body_recs[0]); + for (body_recs[1..]) |r| { + try testing.expectEqual(@as(u8, 'E'), r[0]); + try testing.expect(Action.fromChar(r[1]).?.onTag()); + } + + _ = wr(p, Node.of(serial, .addr), "1"); + _ = wr(p, Node.of(serial, .data), "ONE\n"); + const data_recs = drainEvents(p, event, &store, &slots); + try testing.expectEqual(@as(usize, 2), data_recs.len); + try testing.expectEqualStrings("FD0 3 0 0 \n", data_recs[0]); + try testing.expectEqualStrings("FI0 3 0 3 ONE\n", data_recs[1]); + try testing.expectEqual(Status.again, rd(p, event, 0, 4096).reply.status); +} + +test "two event readers each count once, and the second closing leaves the first" { + const gpa = testing.allocator; + const p = try withFile(gpa, "x\n"); + defer p.deinit(); + const event = Node.of(serialOf(p), .event); + + _ = call(p, .{ .tag = 12, .op = .open, .node = event }); + _ = call(p, .{ .tag = 13, .op = .open, .node = event }); + try testing.expectEqual(@as(u16, 2), p.fs.panes[0].readers); + try testing.expectEqual(@as(u16, 2), p.fs.listeners); + + _ = call(p, .{ .tag = 14, .op = .release, .node = event }); + try testing.expectEqual(@as(u16, 1), p.fs.panes[0].readers); + try testing.expect(p.fs.scripted(0)); + + _ = call(p, .{ .tag = 15, .op = .release, .node = Node.of(serialOf(p), .body) }); + try testing.expectEqual(@as(u16, 1), p.fs.listeners); + + _ = call(p, .{ .tag = 16, .op = .release, .node = event }); + try testing.expectEqual(@as(u16, 0), p.fs.listeners); + try testing.expect(!p.fs.scripted(0)); + _ = call(p, .{ .tag = 17, .op = .release, .node = event }); + try testing.expectEqual(@as(u16, 0), p.fs.listeners); +} + +test "a pane deleted while its event file is open leaves no suppression behind" { + const gpa = testing.allocator; + const p = try withFile(gpa, "x\n"); + defer p.deinit(); + const serial = serialOf(p); + const event = Node.of(serial, .event); + _ = try th.newPane(p); + + const a = call(p, .{ .tag = 18, .op = .open, .node = event }); + const b = call(p, .{ .tag = 19, .op = .open, .node = event }); + try testing.expectEqual(@as(u16, 2), p.fs.listeners); + + _ = wr(p, Node.of(serial, .ctl), "exec Del\n"); + try testing.expect(p.paneBySerial(serial) == null); + try testing.expectEqual(@as(u16, 0), p.fs.listeners); + + _ = call(p, .{ .tag = 20, .op = .release, .node = event, .handle = a.reply.handle }); + _ = call(p, .{ .tag = 21, .op = .release, .node = event, .handle = b.reply.handle }); + try testing.expectEqual(@as(u16, 0), p.fs.listeners); + + try testing.expectEqual(E.NOENT, rd(p, event, 0, 64).errno()); + try testing.expectEqual(E.NOENT, rd(p, Node.of(serial, .body), 0, 64).errno()); + try testing.expectEqual(E.NOENT, wr(p, Node.of(serial, .ctl), "clean\n").errno()); + try testing.expectEqual(E.NOENT, call(p, .{ .tag = 22, .op = .open, .node = event }).errno()); +} + +test "writing an event record back performs the action it names" { + const gpa = testing.allocator; + const p = try withFile(gpa, "Msg fs-ran\n"); + defer p.deinit(); + const serial = serialOf(p); + const event = Node.of(serial, .event); + const pane = p.panes[0].?; + + const w = wr(p, event, "FX0 10\n"); + try testing.expectEqual(Status.ok, w.reply.status); + try testing.expectEqualStrings("fs-ran", pane.msg[0..pane.msg_len]); + + pane.msg_len = 0; + try testing.expectEqual(Status.ok, wr(p, event, "FX0 10\nFX0 10\n").reply.status); + try testing.expectEqualStrings("fs-ran", pane.msg[0..pane.msg_len]); + + pane.msg_len = 0; + for ([_][]const u8{ + "FX0 10\nFQ0 1\n", // unknown type character + "FX0 999\n", // out of range + "FX0 10", // no newline + "FX5 1\n", // q0 > q1 + "FD0 3\n", // a report, not a request + "F\n", + }) |bad| { + try testing.expectEqual(E.INVAL, wr(p, event, bad).errno()); + try testing.expectEqual(@as(usize, 0), pane.msg_len); + } + + _ = call(p, .{ .tag = 23, .op = .open, .node = event }); + p.fs.origin = 'K'; + _ = wr(p, event, "KX0 10\n"); + try testing.expectEqual(@as(u8, 'F'), p.fs.origin); +} + +test "the log parks until a pane is created, renamed, saved or deleted" { + const gpa = testing.allocator; + const p = try withFile(gpa, "logged\n"); + defer p.deinit(); + const log = @intFromEnum(tree.TopFile.log); + _ = try th.newPane(p); + try testing.expect(p.fs.log.empty()); + + const opened = call(p, .{ .tag = 1, .op = .open, .node = log }); + try testing.expectEqual(Status.ok, opened.reply.status); + try testing.expectEqual(@as(u16, 1), p.fs.log_readers); + try testing.expectEqual(Status.again, rd(p, log, 0, 4096).reply.status); + + const serial = try th.newPane(p); + const id = p.paneBySerial(serial).?; + var expected: [4200]u8 = undefined; + try testing.expectEqualStrings( + try std.fmt.bufPrint(&expected, "new {d} {s}\n", .{ serial, pane_files.nameOf(p.panes[id].?) }), + rd(p, log, 0, 4096).bytes, + ); + try testing.expectEqual(Status.ok, wr(p, Node.of(serial, .name), "/tmp/logged.txt\n").reply.status); + try testing.expectEqualStrings( + try std.fmt.bufPrint(&expected, "rename {d} /tmp/logged.txt\n", .{serial}), + rd(p, log, 0, 4096).bytes, + ); + const saving = wr(p, Node.of(serial, .ctl), "exec Save\n"); + try testing.expectEqual(Status.ok, saving.reply.status); + try testing.expect(saving.saved); + p.perform(.{ .save_file = .{ .pane = @intCast(id) } }); + try testing.expectEqualStrings( + try std.fmt.bufPrint(&expected, "save {d} /tmp/logged.txt\n", .{serial}), + rd(p, log, 0, 4096).bytes, + ); + try testing.expectEqual(Status.ok, wr(p, Node.of(serial, .ctl), "exec Del\n").reply.status); + try testing.expectEqualStrings( + try std.fmt.bufPrint(&expected, "del {d} /tmp/logged.txt\n", .{serial}), + rd(p, log, 0, 4096).bytes, + ); + try testing.expectEqual(Status.again, rd(p, log, 0, 4096).reply.status); + + _ = call(p, .{ .tag = 2, .op = .release, .node = log, .handle = opened.reply.handle }); + try testing.expectEqual(@as(u16, 0), p.fs.log_readers); + _ = try th.newPane(p); + try testing.expect(p.fs.log.empty()); + try testing.expectEqual(@as(usize, 0), p.fs.log.buf.capacity); +} |
