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.zig63
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();