//! 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; } }; /// What the next read would answer, which is the size a stat reports: a /// client can see there is something waiting without parking on it. pub fn pending(q: *const Queue) u64 { const record = q.peek() orelse return 0; return record.len; } /// 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 ` ` 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 tree.failText(req.tag, E.INVAL, tree.e_bad_event), } const n = if (r.action.onTag()) tag.len else body.len; if (r.q0 > r.q1 or r.q1 > n) return tree.failText(req.tag, E.INVAL, tree.e_bad_event); } if (check.i != req.data.len) return tree.failText(req.tag, E.INVAL, tree.e_bad_event); } 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, .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, .dirty), "0\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, .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, .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); }