diff --git a/build.zig.zon b/build.zig.zon index e05a795..19466a4 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -5,12 +5,12 @@ .minimum_zig_version = "0.16.0", .dependencies = .{ .zat = .{ - .url = "git+https://tangled.org/zat.dev/zat#03383118cd30e6e40fbb115eeb69c7a0c724a7ad", - .hash = "zat-0.3.16-5PuC7o8pCwAgpjdP8uWtJJVX5YU7P3NnoVIIRiqgKSeD", + .url = "git+https://tangled.org/zat.dev/zat#4803be926a03db7b0f099c54da9aeb95ce323f99", + .hash = "zat-0.3.16-5PuC7o8pCwCi5sqcgqR47OZlACzgHbmPYxFjYMW1scqm", }, .websocket = .{ - .url = "https://tangled.org/zzstoatzz.io/websocket.zig/archive/ac5d16e.tar.gz", - .hash = "websocket-0.1.9-ZPISdWJuBAD1NNhJiyV9_BwI9ksRWFOnzpzB4WqFM93N", + .url = "https://tangled.org/zzstoatzz.io/websocket.zig/archive/bfd23c6.tar.gz", + .hash = "websocket-0.1.9-ZPISdT59BADf5MaPXEGp96Ypq45Icgi_cZFNY2ppaUpJ", }, .zstd = .{ .url = "https://github.com/facebook/zstd/releases/download/v1.5.7/zstd-1.5.7.tar.gz", diff --git a/docs/semantic-parity.md b/docs/semantic-parity.md index ad85799..2e7515b 100644 --- a/docs/semantic-parity.md +++ b/docs/semantic-parity.md @@ -15,7 +15,7 @@ this document or `bootstrap-semantic-parity.md` is open. | Bootstrap, merge, retry | Detailed row-by-row contract and adversity receipt | See `bootstrap-semantic-parity.md`; adversity remains open | | JSS v1 storage | Upstream-produced fixture, reciprocal segment parsing, checksum/index/gloom vectors | closed for sealed format; rerun reciprocal corpus before experiment | | Archive XRPC | Official Go client plus Stream conformance replay over listSegments/getSegment/getBlock/planBackfill | implementation present; fresh full suite required | -| Subscribe v1/v2 | Wire/filter/cursor/zstd unit and local e2e coverage | Cursor parsing/resolution now precedes upgrade; both endpoints use the seq/time-us magnitude split, v1 clamps below-floor seqs, and v2 rejects them. Hot/cold scan, skip, encode, options-update, oversize-parser, clean-peer-close, cursor-resolution, and sustained adversarial-rate boundaries feed the canonical metrics, with deterministic offline tests derived from upstream. **Still open:** official-client replay, 30-second ping/write deadlines, graceful-shutdown close accounting, v2 dictionary negotiation/download, and pre-upgrade timestamp translation fault classification. | +| Subscribe v1/v2 | Wire/filter/cursor/zstd unit and local e2e coverage | Cursor parsing/resolution now precedes upgrade; both endpoints use the seq/time-us magnitude split, v1 clamps below-floor seqs, and v2 rejects them. Hot/cold scan, skip, encode, options-update, oversize-parser, clean peer/server close, cursor-resolution, and sustained adversarial-rate boundaries feed the canonical metrics, with deterministic offline tests derived from upstream. Subscriber transport sends a ping every 30 seconds, applies a kernel-enforced five-second deadline to every frame write, and sends close code 1001 before intentional server shutdown while keeping transport failure out of clean-disconnect accounting. A failed delivery, cold read, ping, or adversarial-rate check interrupts the server reader and removes the connection instead of leaving a ping-only zombie. **Still open:** official-client replay, v2 dictionary negotiation/download, and pre-upgrade timestamp translation fault classification. | | Sync 1.1 live verification | Upstream requires durable per-DID chain/hosting state, MST inversion, op-CID consistency, rev replay/future guards, default acceptance of legacy-shaped commits, transparent whole-repo repair, and Atmos delivery scheduling | **closed at the implementation boundary.** Diff-CAR/MST inversion (including upstream's narrow default lenient carve-out), post-state op-CID proof, signed inner/outer consistency, exact decimal size/future-rev/replay gates, exact legacy-shape detection with default `LegacyAccept`, two-phase durable chain/hosting state, and account/identity replay ratchets are implemented. Recoverable decode/inversion/duplicate-path/op-CID/chain failures route to repair; signature failures and outer/inner producer-integrity mismatches bypass it, and no failed event is archived. `prepareRepo` retains and authenticates the canonical signed complete head without a second multi-gigabyte parse/walk. The live path now uses Atmos's 32 worker slots, one worker-held FIFO chain per DID, a 64-pending drop-oldest boundary, completion-order batches of 50/500 ms across DIDs, and `min(inflight)-1` cursor watermarks. Dropped and silent events leave the inflight set without inventing a delivery batch; cursor advancement remains coupled to durable archive boundaries. Offline receipts use real encoded frames, real RocksDB verifier state, and actual worker blocking rather than fixture verdicts or scheduler mocks. See `live-scheduler.md`. | | Sync 1.1 resync ordering | Upstream serializes per DID, drops stale async repairs, and emits a sync tombstone plus authoritative replacements | **closed at the implementation boundary.** The live coordinator has 32 real fetch workers, a 64-job queue, a five-minute fetch budget plus the 1,000 B/s-for-30-seconds slow-transfer guard, a bounded 16,384-entry 5/minute-per-DID limiter, 2,048 pending commits with drop-oldest overflow, fetch outside the DID lane, authenticated apply with older/equal-contradictory head rejection, ordered pending replay, and trigger-ticket-gated outbox delivery. Commit divergence schedules async repair; `#sync` divergence attempts inline repair and queues one retry only for transient failure, preserving the original relay cursor only on inline success. The writer emits `sync` then bounded 1,024-row `create_resync` batches and stages the fetched/pending chain checkpoints under the archive lock before any matching fsync. Loopback-only tests exercise a real HTTP getRepo mmap, real DID resolution/signature, buffering while the response is held open, slow-transfer cancellation, durable RocksDB promotion/cursor coupling, v2 tail publication, and sealed JSS row order. | | Timestamp import | Canonical dashboard and API include timestamp import lifecycle/rewrite metrics and behavior | **Closed at the implementation boundary.** The format preserves distinct `witnessed_at` and sentinel-zero `indexed_at` columns. Subscriber v1/v2 encoding applies upstream's `indexed_at != 0 ? indexed_at : witnessed_at` display rule while ranges and timestamp cursors remain anchored to immutable witness envelopes. The strict seekable RFC 4180 parser matches upstream's header, row rejection, byte-offset, CID, RFC3339, 64 KiB, and bounded-sampling contracts in a forward-only allocation-free pass. A separate Bloom-filtered RocksDB rule store bulk-loads sorted external SST chunks with CSV-order last-write-wins, reconstructs its resident collection gate on open, and implements specific-CID precedence with all-version fallback. Its central archive-lock hook stamps live, repair, bootstrap, and merge materializations before sequence assignment, and feeds the mutated rows to immediate delivery so hot and cold encodings agree. The sealed-JSS patch primitive preserves block topology/envelopes, clean frames, and the opaque bloom/collection footer while atomically changing only `indexed_at`; it is exercised against an upstream-produced segment. A concrete sealed-segment bloom catalog, bounded DID/FD bucketer, fsynced packed offsets, positioned revalidation, and collision-safe per-segment patch plan implement Phase B/C with specific-CID precedence and idempotent disk-to-disk coverage. A dedicated archive rewrite mutex serializes the real delete compactor with timestamp patching while leaving live append available. The unified metadata DB synchronously persists the current-job pointer, complete job record, and per-segment done set. The real runner activates rules before force-seal, durably crosses the bucketed handoff, checkpoints only after patch fsync/rename/dir-sync, retains partial cancellation progress, and resumes after full archive/rule/metadata reopen. The background manager canonicalizes symlinks, confines regular files, enforces one durable nonterminal job, and adopts it on startup. Bearer-gated import/status XRPC matches the pinned lexicons, fixed-width digest comparison, error names/statuses, steady-state admission, and current/by-id status behavior; a real loopback HTTP run exercised it against the local simulator. Every canonical import metric is produced at its real parse, route, patch, or durable-terminal boundary; cancellation remains a resumable pause and is not counted as terminal failure. | diff --git a/src/internal/homepage.zig b/src/internal/homepage.zig index 3855723..11b9a49 100644 --- a/src/internal/homepage.zig +++ b/src/internal/homepage.zig @@ -227,7 +227,6 @@ pub fn renderStatus(buf: []u8, args: StatusArgs) []const u8 { " 8/8 lifecycle crashpoints recovered\n" ++ " 16 throttled archive downloads stayed below 40 MiB RSS\n\n" ++ "still open\n" ++ - " stalled WebSocket writers have no Go-equivalent deadline\n" ++ " permessage-deflate is not negotiated (custom zstd works)\n" ++ " {s}\n\n" ++ "build {s} · https://tangled.org/zat.dev/stream\n", @@ -339,7 +338,6 @@ pub fn renderStatusHtml(buf: []u8, args: StatusArgs) []const u8 { \\
  • Failed-repo healing{s}
  • \\
  • Commit verification{d} valid · {d} invalid
  • \\
    not pretending

    Still open

    \\
    touch the system

    Public surfaces

    diff --git a/src/internal/server.zig b/src/internal/server.zig index f10a09d..b2b3084 100644 --- a/src/internal/server.zig +++ b/src/internal/server.zig @@ -26,9 +26,12 @@ const log = std.log.scoped(.stream); /// 8 MB, matching the Io.Threaded stacks (docs/lessons-from-zlay.md #1) const subscriber_stack_size: usize = 8 * 1024 * 1024; +const ping_stack_size: usize = 256 * 1024; const slow_window_us: i64 = 60 * std.time.us_per_s; const slow_lag_threshold: u64 = 100_000; const slow_min_rate: u64 = 5; +const ping_interval: Io.Clock.Duration = .{ .raw = Io.Duration.fromSeconds(30), .clock = .awake }; +const frame_write_timeout_ms: u32 = 5_000; const SlowDetector = struct { streak_start_us: i64 = 0, @@ -113,14 +116,17 @@ const Subscriber = struct { /// /subscribe-v2: superset wire (seq, record_cbor, #sync), strict cursors v2: bool = false, thread: ?std.Thread = null, + ping_thread: ?std.Thread = null, + ping_stop: Io.Event = .unset, slow: SlowDetector = .{}, /// the pull loop. cleanup happens in Handler.close, which joins this /// thread — the subscriber never frees itself, so the handler can always /// safely signal it. fn run(self: *Subscriber) void { - if (self.cold_from_seq) |from_seq| self.coldPhaseSeq(from_seq); - if (self.cold_from_us) |from_us| self.coldPhase(from_us); + defer self.workerDone(); + if (self.cold_from_seq) |from_seq| if (!self.coldPhaseSeq(from_seq)) return; + if (self.cold_from_us) |from_us| if (!self.coldPhase(from_us)) return; while (self.alive.load(.acquire)) { const json = self.hub.tail.next(self.hub.allocator, &self.idx, &self.stopped, self.v2, self, hotWantsPred) catch return orelse return; defer self.hub.allocator.free(json); @@ -143,6 +149,45 @@ const Subscriber = struct { } } + fn pingLoop(self: *Subscriber) void { + while (self.alive.load(.acquire)) { + self.ping_stop.waitTimeout(self.hub.io, .{ .duration = ping_interval }) catch |err| switch (err) { + error.Timeout => {}, + error.Canceled => return, + }; + if (!self.alive.load(.acquire)) return; + self.write_mutex.lockUncancelable(self.hub.io); + if (!self.alive.load(.acquire)) { + self.write_mutex.unlock(self.hub.io); + return; + } + var payload: [0]u8 = .{}; + self.conn.writePing(&payload) catch { + const interrupt = self.alive.swap(false, .acq_rel); + self.stopped.store(true, .release); + self.write_mutex.unlock(self.hub.io); + self.ping_stop.set(self.hub.io); + self.hub.tail.wake(); + if (interrupt) self.conn.interruptRead(); + return; + }; + self.write_mutex.unlock(self.hub.io); + } + } + + /// A subscriber-worker failure must end the connection, not leave a + /// ping-only zombie. Half-shutdown wakes the server reader; that reader + /// remains the sole owner of final descriptor cleanup. + fn workerDone(self: *Subscriber) void { + self.write_mutex.lockUncancelable(self.hub.io); + const interrupt = self.alive.swap(false, .acq_rel); + self.stopped.store(true, .release); + self.write_mutex.unlock(self.hub.io); + self.ping_stop.set(self.hub.io); + self.hub.tail.wake(); + if (interrupt) self.conn.interruptRead(); + } + fn write(self: *Subscriber, json: []const u8) !cold.WriteResult { self.write_mutex.lockUncancelable(self.hub.io); defer self.write_mutex.unlock(self.hub.io); @@ -170,9 +215,9 @@ const Subscriber = struct { return if (self.conn.compression != null) .deflate else .none; } - fn coldPhase(self: *Subscriber, from_us: i64) void { + fn coldPhase(self: *Subscriber, from_us: i64) bool { const hub = self.hub; - const archive = hub.archive orelse return; + const archive = hub.archive orelse return true; const until_us = hub.tail.oldestTime() orelse std.math.maxInt(i64); _ = hub.stats.subscribe_cold_reads_total.fetchAdd(1, .monotonic); const n = cold.replay( @@ -188,14 +233,15 @@ const Subscriber = struct { until_us, ) catch |err| { log.warn("cold replay failed: {s}", .{@errorName(err)}); - return; + return false; }; log.debug("cold replay served {d} frames", .{n}); + return true; } - fn coldPhaseSeq(self: *Subscriber, from_seq: u64) void { + fn coldPhaseSeq(self: *Subscriber, from_seq: u64) bool { const hub = self.hub; - const archive = hub.archive orelse return; + const archive = hub.archive orelse return true; _ = hub.stats.subscribe_cold_reads_total.fetchAdd(1, .monotonic); const n = cold.replaySeq( hub.allocator, @@ -210,9 +256,10 @@ const Subscriber = struct { from_seq, ) catch |err| { log.warn("seq cold replay failed: {s}", .{@errorName(err)}); - return; + return false; }; log.debug("seq cold replay served {d} frames", .{n}); + return true; } fn coldWrite(self: *Subscriber, json: []const u8) anyerror!cold.WriteResult { @@ -287,6 +334,7 @@ pub const Handler = struct { awaiting_hello: bool, compress: bool, v2: bool, + clean_disconnect_counted: bool = false, pub fn init(handshake: *const websocket.Handshake, conn: *websocket.Conn, hub: *Hub) !Handler { const url = handshake.url; @@ -340,6 +388,7 @@ pub const Handler = struct { if (compress) if (handshake.headers.get("sec-websocket-extensions")) |extensions| { if (offersPerMessageDeflate(extensions)) return error.InvalidRequest; }; + try conn.writeTimeout(frame_write_timeout_ms); return .{ .hub = hub, @@ -399,6 +448,16 @@ pub const Handler = struct { self.subscriber = null; return err; }; + sub.ping_thread = std.Thread.spawn(.{ .stack_size = ping_stack_size }, Subscriber.pingLoop, .{sub}) catch |err| { + self.signalStop(sub); + sub.thread.?.join(); + sub.filter.deinit(); + _ = hub.active.fetchSub(1, .monotonic); + _ = hub.stats.subscribe_active.getPtr(scheme).fetchSub(1, .monotonic); + hub.allocator.destroy(sub); + self.subscriber = null; + return err; + }; } /// client → server messages: SubscriberSourcedMessage envelope. @@ -496,16 +555,24 @@ pub const Handler = struct { pub fn clientClose(self: *Handler, _: []const u8) !void { if (self.subscriber) |sub| self.signalStop(sub); - _ = self.hub.stats.subscribe_clean_disconnects_total.fetchAdd(1, .monotonic); + self.noteCleanDisconnect(); self.conn.close(.{}) catch {}; } + /// Called only by the websocket server's intentional shutdown path. + pub fn serverClose(self: *Handler) void { + if (self.subscriber) |sub| self.signalStop(sub); + self.conn.writeClose(.{ .code = 1001, .reason = "server shutting down" }) catch {}; + self.noteCleanDisconnect(); + } + /// always called by the server when the connection ends — the single /// place the subscriber thread is joined and freed. pub fn close(self: *Handler) void { if (self.subscriber) |sub| { self.signalStop(sub); if (sub.thread) |t| t.join(); + if (sub.ping_thread) |t| t.join(); if (sub.compressor) |*comp| comp.deinit(); sub.filter.deinit(); _ = self.hub.active.fetchSub(1, .monotonic); @@ -521,9 +588,16 @@ pub const Handler = struct { sub.alive.store(false, .release); sub.stopped.store(true, .release); sub.write_mutex.unlock(self.hub.io); + sub.ping_stop.set(self.hub.io); self.hub.tail.wake(); } + fn noteCleanDisconnect(self: *Handler) void { + if (self.clean_disconnect_counted) return; + self.clean_disconnect_counted = true; + _ = self.hub.stats.subscribe_clean_disconnects_total.fetchAdd(1, .monotonic); + } + /// plain-HTTP requests on the websocket port: /metrics and /healthz pub fn httpFallback( conn: *websocket.Conn,