diff options
Diffstat (limited to 'src/9p_io.zig')
| -rw-r--r-- | src/9p_io.zig | 63 |
1 files changed, 53 insertions, 10 deletions
diff --git a/src/9p_io.zig b/src/9p_io.zig index d6eef9d4..393b6210 100644 --- a/src/9p_io.zig +++ b/src/9p_io.zig @@ -174,6 +174,8 @@ pub const Listener = struct { core: *pardes.Pardes, runner: Runner = undefined, stopping: std.atomic.Value(bool) = .init(false), + /// A connection wrote the Restore and has yet to answer it. + restore_writer: std.atomic.Value(bool) = .init(false), tcp_address: ?std.Io.net.IpAddress = null, quic: if (quic_enabled) ?quic.Listener else void = if (quic_enabled) null else {}, quic_address: ?std.Io.net.IpAddress = null, @@ -217,6 +219,12 @@ pub const Listener = struct { } const quiet = pardes.turn.take(); defer pardes.turn.give(); + // A connection a Restore is cutting (`reset`) is the old editor's + // client: its fids name the old editor's panes and opens. + if (conn.user != null) { + const refused = pardes.ctlfs.Reply.fail(req.tag, pardes.ctlfs.E.IO); + return conn.reply(&refused, ""); + } if (!quiet and pardes.ctlfs.needsQuiet(req)) { pardes.turn.parked = true; const later: pardes.ctlfs.Reply = .{ .tag = req.tag, .status = .again }; @@ -267,13 +275,20 @@ pub const Listener = struct { // and `deinit` lets the answer out before it cuts the connections. if (core.quit) return conn.reply(&reply, ""); 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); + defer if (restoring) l.restore_writer.store(false, .release); if (core.effects_len != 0) pardes.turn.awaitSettled(epoch); + // The Restore's own write is answered once it is done, and before + // `reset` hangs every connection up, this one too, so the writer + // hears that it happened rather than a cut; a failed one leaves the + // connection and says why on the message row and in the log. + if (restoring) { + if (l.core == core) pardes.turn.awaitRestored(restores); + return conn.reply(&reply, ""); + } // Once a wait returns `core` may be gone: a Restore meanwhile put a - // replacement in (`reset`) and is hanging this connection up, which - // is a Restore's answer. Only one that failed leaves this connection - // here to answer. - if (l.core != core) return; - if (restoring) pardes.turn.awaitRestored(restores); + // replacement in and is hanging this connection up. if (l.core != core) return; conn.reply(&reply, ""); } @@ -484,9 +499,9 @@ pub const Listener = struct { /// Hangs every connection up, so that `replacement` starts with no /// client holding anything. From here requests are answered from the /// replacement -- a client that connects meanwhile is a client of the - /// replacement -- and every task waiting on the old core is let go: it - /// finds the core changed and answers nothing, its connection being - /// hung up. The connections pay their releases on their own tasks, so + /// replacement -- and every task waiting on the old core is let go: the + /// Restore's own writer answers, the rest find the core changed and + /// answer nothing, their connections being hung up. The connections pay their releases on their own tasks, so /// the editor rests while they do. pub fn reset(l: *Listener, replacement: *pardes.Pardes) void { for (0..quic_slots) |i| l.drop(@intCast(i)); @@ -496,12 +511,40 @@ pub const Listener = struct { replacement.fs.tcp_address = l.tcp_address; replacement.fs.quic_address = l.quic_address; if (comptime supported) { - l.runner.closeAll(); + // The Restore's writer answers as it wakes (`onServe`), and its + // answer gets 200 ms to leave before the cut; every other request + // from the old editor's clients meanwhile is refused, since the + // fids it names are the old editor's. A client that connects now + // is the replacement's, and stays. + // `user` marks the connections to cut: cloud9 leaves it to us and + // clears it for a new connection in the slot. + for (&l.runner.conns) |*conn| if (conn.live()) { + conn.user = l; + }; pardes.turn.settle(); pardes.turn.restoreSettled(); pardes.turn.rest(); + const flush_by = nowMs() +| 200; + while (nowMs() < flush_by) { + const pending = l.restore_writer.load(.acquire) or for (&l.runner.conns) |*conn| { + if (conn.user == null or !conn.live()) continue; + conn.lock(); + const n = conn.engine.output().len; + conn.unlock(); + if (n != 0) break true; + } else false; + if (!pending) break; + Client.nap(1); + } + for (&l.runner.conns) |*conn| if (conn.user != null and conn.live()) conn.close(); const deadline = nowMs() +| 2000; - while (l.runner.count() != 0 and nowMs() < deadline) Client.nap(1); + while (nowMs() < deadline) { + const left = for (&l.runner.conns) |*conn| { + if (conn.user != null and conn.live()) break true; + } else false; + if (!left) break; + Client.nap(1); + } pardes.turn.wake(); } l.collectOs(); |
