summaryrefslogtreecommitdiff
path: root/src/9p_io.zig
diff options
context:
space:
mode:
Diffstat (limited to 'src/9p_io.zig')
-rw-r--r--src/9p_io.zig104
1 files changed, 98 insertions, 6 deletions
diff --git a/src/9p_io.zig b/src/9p_io.zig
index 0b3db274..18c51ece 100644
--- a/src/9p_io.zig
+++ b/src/9p_io.zig
@@ -180,6 +180,10 @@ pub const Listener = struct {
restore_writer: std.atomic.Value(bool) = .init(false),
/// Clients turned away with every slot taken, since the log last said so.
refused: std.atomic.Value(u32) = .init(0),
+ /// Writes being served right now, a retry of a parked one included: a
+ /// held Edit write a retry has taken out of its park is not flushed
+ /// (watchEdit).
+ writes_out: std.atomic.Value(u32) = .init(0),
/// Connections accepted so far, and the count each slot's connection
/// was accepted at: a Restore cuts those accepted before it (`reset`),
/// and a new client in a freed slot has a later stamp. Written on the
@@ -241,6 +245,10 @@ pub const Listener = struct {
fn onServe(ctx: ?*anyopaque, conn: *Runner.Conn, req: pardes.ctlfs.Req) void {
const l = of(ctx);
+ if (req.op == .write) _ = l.writes_out.fetchAdd(1, .acq_rel);
+ defer if (req.op == .write) {
+ _ = l.writes_out.fetchSub(1, .acq_rel);
+ };
// Once `stop` has begun the editor is tearing down and may never
// rest again: answer without waiting on it.
if (l.stopping.load(.acquire)) {
@@ -268,7 +276,17 @@ pub const Listener = struct {
const changes = pardes.ctlfs.changesPane(core, req);
const fills = pardes.ctlfs.handsOff(core, req);
const reply = core.serveFs(req);
+ const edit_started = core.fs.edit_started;
+ core.fs.edit_started = false;
+ // The held Edit write asked again by a retry, its Edit still
+ // running: parked again, as it was.
+ if (reply.status == .again and req.op == .write) return conn.reply(&reply, "");
if (req.op == .release and quiet) l.collectOs();
+ // A close's held last line runs now when no step is out, so its err
+ // is logged before the close is answered, as before; with a step
+ // out (perhaps through a mount this close came by) it waits for
+ // that step, and the close is answered at once all the same.
+ if (req.op == .release and quiet) pardes.ctlfs.runClosedLines(core);
// A read with nothing yet stays parked in the engine, and the core
// keeps its ticket to answer it when what it waits on has something.
if (reply.status == .again and req.op == .read) {
@@ -308,14 +326,38 @@ pub const Listener = struct {
// mount: 9ns answers nothing else on it while a clunk is out, and
// the editor's step may be out opening a file through that mount.
if (fills) return conn.reply(&reply, core.fsPayload(reply));
- // A refusal answers at once: its text may be in the core's one
- // buffer for it, which a request run while this one waited would
- // write over.
+ // A refusal that left saves going (Putall's one refused among the
+ // rest) is answered once they have landed, its words copied out of
+ // the core's one buffer for them, which a request run meanwhile
+ // could write over.
+ if (reply.status == .err and core.effects_len != 0 and !core.quit) {
+ var kept: [320]u8 = undefined;
+ const words = kept[0..@min(reply.ename.len, kept.len)];
+ @memcpy(words, reply.ename[0..words.len]);
+ var waited = reply;
+ waited.ename = words;
+ pardes.turn.awaitSettled(epoch);
+ return conn.reply(&waited, "");
+ }
+ // Any other refusal answers at once, for the same buffer's sake.
if (reply.status == .err) return conn.reply(&reply, "");
// A write that quits the editor (Kill) is answered now: the editor
// is on its way out and will settle nothing this could wait for,
// and `deinit` lets the answer out before it cuts the connections.
if (core.quit) return conn.reply(&reply, "");
+ // An Edit whose `<`, `|` or `>` commands run off the loop: the write
+ // is held, as a read with nothing yet is, and answered when the Edit
+ // is done (answerHeld), so the connection serves on meanwhile -- a
+ // filter that reads this session through a mount among what it
+ // serves. A flush of it, or a hang-up, stops the Edit (watchEdit).
+ if (edit_started and core.pipe.edit_run != null) {
+ const later: pardes.ctlfs.Reply = .{ .tag = req.tag, .status = .again };
+ const ticket = conn.hold(&later) orelse return;
+ core.fs.edit_hold = .{ .asker = conn, .slot = ticket.slot, .seq = ticket.seq, .tag = req.tag, .node = req.node, .handle = req.handle, .written = reply.written };
+ const watcher = std.Thread.spawn(.{}, watchEdit, .{ l, conn, ticket }) catch return;
+ watcher.detach();
+ return;
+ }
const restoring = core.restore_req != null;
// Set with the turn still held, so `reset` sees it and waits.
if (restoring) l.restore_writer.store(true, .release);
@@ -395,6 +437,46 @@ pub const Listener = struct {
const conn: *Runner.Conn = @ptrCast(@alignCast(held.asker));
if (!conn.answerWith(held.ticket, Held{ .core = core, .req = held.req }, Held.make)) o.held = null;
}
+ // An Edit's held write, its Edit done. Not waiting: a retry has it
+ // out, and its asking answers it (Pardes.serveFs); or it was
+ // flushed, and watchEdit lets the record go.
+ if (core.fs.edit_hold) |h| if (h.done) {
+ const conn: *Runner.Conn = @ptrCast(@alignCast(h.asker));
+ _ = conn.answerWith(.{ .slot = h.slot, .seq = h.seq }, core, struct {
+ fn make(c: *pardes.Pardes) ?Runner.Conn.Answer {
+ return .{ .reply = pardes.filesystem.EditHold.answer(&c.fs) };
+ }
+ }.make);
+ };
+ }
+
+ /// While an Edit's write is held: when nobody waits for it any more --
+ /// flushed (an interrupted writer), its fid clunked, its connection
+ /// hung up -- the Edit is stopped, its commands killed, nothing
+ /// changed (edit_cmd.cancel). The engine tells the backend nothing of a
+ /// flush, so it is looked for, every 100 ms; a park out for a retry is
+ /// not waiting either, so only three looks in a row with no write being
+ /// served count.
+ fn watchEdit(l: *Listener, conn: *Runner.Conn, ticket: cloud9.fs.Ticket) void {
+ const io = pardes.turn.io orelse return;
+ var misses: u8 = 0;
+ while (!l.stopping.load(.acquire)) {
+ io.sleep(.fromMilliseconds(100), .awake) catch return;
+ _ = pardes.turn.take();
+ defer pardes.turn.give();
+ const core = l.core;
+ const h = core.fs.edit_hold orelse return;
+ if (h.slot != ticket.slot or h.seq != ticket.seq or h.asker != @as(*anyopaque, conn)) return;
+ if (conn.waiting(ticket) or l.writes_out.load(.acquire) != 0) {
+ misses = 0;
+ continue;
+ }
+ misses += 1;
+ if (misses < 3) continue;
+ if (h.done) core.fs.edit_hold = null else @import("edit_cmd.zig").cancel(core);
+ l.kick();
+ return;
+ }
}
/// A held read asked again, to make its answer while its park waits.
@@ -1910,6 +1992,16 @@ pub fn start(io: std.Io, gpa: std.mem.Allocator, core: *pardes.Pardes) ?*Listene
/// open to scripts and to mounted shells — while withholding `PARDES_PID`,
/// which is the whole of what `--nested` means.
pub fn exportPaneEnv(listener: ?*const Listener, serial: u32, adopts: bool) void {
+ exportEnv(if (listener) |l| l.path() else null, serial, adopts);
+}
+
+/// exportPaneEnv for a command run off the loop (an Edit's filter), with
+/// this process's own socket: what a pane's command would be given.
+pub fn exportSessionEnv(serial: u32, adopts: bool) void {
+ exportEnv(if (own_socket_len > 0) own_socket[0..own_socket_len] else null, serial, adopts);
+}
+
+fn exportEnv(socket: ?[]const u8, serial: u32, adopts: bool) void {
var announced = false;
if (adopts) announcing: {
var buf: [16]u8 = undefined;
@@ -1919,13 +2011,13 @@ pub fn exportPaneEnv(listener: ?*const Listener, serial: u32, adopts: bool) void
}
if (!announced) _ = unsetenv("PARDES_PID");
- if (listener) |l| exporting: {
+ if (socket) |own| exporting: {
var sock: [sun_path_len]u8 = undefined;
- const path = std.fmt.bufPrintSentinel(&sock, "{s}", .{l.path()}, 0) catch break :exporting;
+ const path = std.fmt.bufPrintSentinel(&sock, "{s}", .{own}, 0) catch break :exporting;
var buf: [16]u8 = undefined;
const id = std.fmt.bufPrintSentinel(&buf, "{d}", .{serial}, 0) catch break :exporting;
if (setenv("PARDES_9P", path, 1) != 0) break :exporting;
- exportMount(l.path());
+ exportMount(own);
if (setenv("PARDES_PANE", id, 1) == 0) return;
}
_ = unsetenv("PARDES_9P");