summaryrefslogtreecommitdiff
path: root/src/ninep/events.zig
diff options
context:
space:
mode:
Diffstat (limited to 'src/ninep/events.zig')
-rw-r--r--src/ninep/events.zig558
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);
+}