diff options
| -rw-r--r-- | src/9p_io.zig | 35 |
1 files changed, 25 insertions, 10 deletions
diff --git a/src/9p_io.zig b/src/9p_io.zig index 393b6210..2ca11a87 100644 --- a/src/9p_io.zig +++ b/src/9p_io.zig @@ -176,6 +176,14 @@ pub const Listener = struct { 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), + /// 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 + /// connection's task as it opens, read with the turn. + accepts: std.atomic.Value(u64) = .init(0), + accepted: [max_conns]std.atomic.Value(u64) = @splat(.init(0)), + /// Connections accepted before this count are being cut by a Restore. + cut_before: u64 = 0, 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, @@ -209,6 +217,16 @@ pub const Listener = struct { /// answers. A request that would change a pane while the editor is out /// in a syscall mid-step is parked in the engine instead, and retried /// when the turn is next given up between steps (`wakeParked`). + fn onOpened(ctx: ?*anyopaque, conn: *Runner.Conn) void { + const l = of(ctx); + l.accepted[conn.index].store(l.accepts.fetchAdd(1, .acq_rel) + 1, .release); + } + + /// A live connection accepted before the Restore under way. + fn cutting(l: *Listener, conn: *Runner.Conn) bool { + return conn.live() and l.accepted[conn.index].load(.acquire) <= l.cut_before and l.cut_before != 0; + } + fn onServe(ctx: ?*anyopaque, conn: *Runner.Conn, req: pardes.ctlfs.Req) void { const l = of(ctx); // Once `stop` has begun the editor is tearing down and may never @@ -221,7 +239,7 @@ pub const Listener = struct { 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) { + if (l.cutting(conn)) { const refused = pardes.ctlfs.Reply.fail(req.tag, pardes.ctlfs.E.IO); return conn.reply(&refused, ""); } @@ -516,18 +534,14 @@ pub const Listener = struct { // 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; - }; + l.cut_before = l.accepts.load(.acquire); 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; + if (!l.cutting(conn)) continue; conn.lock(); const n = conn.engine.output().len; conn.unlock(); @@ -536,16 +550,17 @@ pub const Listener = struct { if (!pending) break; Client.nap(1); } - for (&l.runner.conns) |*conn| if (conn.user != null and conn.live()) conn.close(); + for (&l.runner.conns) |*conn| if (l.cutting(conn)) conn.close(); const deadline = nowMs() +| 2000; while (nowMs() < deadline) { const left = for (&l.runner.conns) |*conn| { - if (conn.user != null and conn.live()) break true; + if (l.cutting(conn)) break true; } else false; if (!left) break; Client.nap(1); } pardes.turn.wake(); + l.cut_before = 0; } l.collectOs(); for (&l.conns) |*conn| conn.accepted_ms = 0; @@ -801,7 +816,7 @@ pub fn listen(io: std.Io, gpa: std.mem.Allocator, core: *pardes.Pardes, named: [ l.runner.init(.{ .io = io, .root = pardes.ctlfs.root, - .handler = .{ .ctx = l, .serve = Listener.onServe }, + .handler = .{ .ctx = l, .serve = Listener.onServe, .opened = Listener.onOpened }, .greet_timeout_ms = Listener.greet_deadline_ms, }); const entry_name = if (named.len != 0) named else fallback; |
