From c6ecef7bf9519a26044f89525cee49f634f04124 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Thu, 9 Jul 2026 01:47:30 -0500 Subject: [PATCH] ingest: parallel verification pipeline, wired and measured MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit the full chunk, upstream-faithful (no unverified pass-through): - zat v0.3.15 onRawFrame hook feeds raw frames into the pipeline; workers decode ONCE (arena handed through the slot to the writer) and verify signatures in parallel, each worker owning its own DidResolver so cache misses resolve concurrently - writer emits in ticket order: archive/tail/cursor semantics remain byte-identical to the single-threaded path - workers/writer run as io.concurrent tasks, not raw std.Threads — raw threads doing io-backed I/O silently stalled (lessons #2) - pipeline stage gauges (submitted/claimed/emitted) on /metrics; they found both real bugs: a serialized resolver pinning the reorder window (44 evt/s), and the writer's redundant re-decode measured (ReleaseSafe, local simulator): 768 evt/s fully verified at 76% CPU with the serialized-resolver fix; decode-once cuts CPU to 44% at the same throughput. at every simulator rate the pipeline drains to empty with idle CPU — the local sim's per-client delivery is now the measurement ceiling, so absolute throughput validation moves to the relay node per the bench-network-boundaries norm. e2e suite green; 38/38 unit tests. Co-Authored-By: Claude Fable 5 --- build.zig.zon | 4 +- src/internal/ingest.zig | 47 +++++++------ src/internal/pipeline.zig | 111 ++++++++++++++++++++++--------- src/internal/server.zig | 15 ++++- src/internal/verify.zig | 134 ++++++++++++++++++++------------------ src/main.zig | 19 ++++++ 6 files changed, 216 insertions(+), 114 deletions(-) diff --git a/build.zig.zon b/build.zig.zon index 88d9d31..0bf247e 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -5,8 +5,8 @@ .minimum_zig_version = "0.16.0", .dependencies = .{ .zat = .{ - .url = "git+https://tangled.org/zat.dev/zat?ref=v0.3.14#c4e68984d538fcda86eecb0dde9f6aa4096f306f", - .hash = "zat-0.3.14-5PuC7qf1CgAN7BtzpBGC4AX-NDRHz2R8EaroAtB-2UJk", + .url = "git+https://tangled.org/zat.dev/zat?ref=v0.3.15#f00fd5e1e45966dca3185a1e0e108e08f6e5e11c", + .hash = "zat-0.3.15-5PuC7mH8CgC1ebyqAJhnt9DsJ5UQY5BJECii3TbnZ4FM", }, .websocket = .{ .url = "https://tangled.org/zzstoatzz.io/websocket.zig/archive/v0.1.9.tar.gz", diff --git a/src/internal/ingest.zig b/src/internal/ingest.zig index 92ecef3..7005291 100644 --- a/src/internal/ingest.zig +++ b/src/internal/ingest.zig @@ -11,6 +11,7 @@ const convert = @import("convert.zig"); const segment = @import("segment.zig"); const cursor_mod = @import("cursor.zig"); const metrics = @import("metrics.zig"); +const pipeline_mod = @import("pipeline.zig"); const tail_mod = @import("tail.zig"); const verify = @import("verify.zig"); @@ -28,6 +29,7 @@ pub const Consumer = struct { to_stdout: bool = false, durable_upstream_seq: i64 = 0, verifier: ?*verify.Verifier = null, + pipeline: ?*pipeline_mod.Pipeline = null, cursor_store: ?*cursor_mod.Store = null, archive: ?*archive_mod.Archive = null, stats: ?*metrics.Stats = null, @@ -43,13 +45,25 @@ pub const Consumer = struct { try client.subscribe(self); } - pub fn onEvent(self: *Consumer, event: zat.FirehoseEvent) void { - self.handleEvent(event) catch |err| { + /// raw frames from zat's client go through the verify pipeline; the + /// pipeline's writer thread calls emitVerified in ticket order. + pub fn onRawFrame(self: *Consumer, data: []const u8) void { + const pipe = self.pipeline orelse return; + pipe.submit(data) catch |err| { + log.err("pipeline submit failed: {s}", .{@errorName(err)}); + }; + } + + /// pipeline writer callback: the event was decoded once by the verify + /// worker; its backing arena stays alive for the duration of this call. + pub fn emitVerified(ctx: *anyopaque, event: zat.FirehoseEvent, verdict: verify.Result) void { + const self: *Consumer = @ptrCast(@alignCast(ctx)); + self.handleEvent(event, verdict) catch |err| { log.err("event handling failed: {s}", .{@errorName(err)}); }; } - fn handleEvent(self: *Consumer, event: zat.FirehoseEvent) !void { + fn handleEvent(self: *Consumer, event: zat.FirehoseEvent, verdict: verify.Result) !void { var arena = std.heap.ArenaAllocator.init(self.allocator); defer arena.deinit(); const alloc = arena.allocator(); @@ -57,22 +71,17 @@ pub const Consumer = struct { const time_us: i64 = Io.Timestamp.now(self.io, .real).toMicroseconds(); if (self.stats) |st| _ = st.events_total.fetchAdd(1, .monotonic); - if (self.verifier) |v| { - switch (event) { - .commit => |c| switch (v.verifyCommit(c.repo, c.blocks)) { - // proven-bad after a fresh key resolve: drop the event - .invalid_signature => { - self.drops.getPtr(.invalid_signature).* += 1; - if (self.stats) |st| _ = st.drops.getPtr(.invalid_signature).fetchAdd(1, .monotonic); - log.warn("dropped commit with invalid signature: {s} seq={d}", .{ c.repo, c.seq }); - return self.noteSeq(event, time_us); - }, - .valid, .unverified => {}, - }, - // rotation signal: next commit re-resolves the key - .identity => |i| v.evict(i.did), - else => {}, - } + switch (event) { + .commit => |c| if (verdict == .invalid_signature) { + // proven-bad after a fresh key resolve: drop the event + self.drops.getPtr(.invalid_signature).* += 1; + if (self.stats) |st| _ = st.drops.getPtr(.invalid_signature).fetchAdd(1, .monotonic); + log.warn("dropped commit with invalid signature: {s} seq={d}", .{ c.repo, c.seq }); + return self.noteSeq(event, time_us); + }, + // rotation signal: next commit re-resolves the key + .identity => |i| if (self.verifier) |v| v.evict(i.did), + else => {}, } if (self.archive) |a| self.persist(a, event) catch |err| { // our own storage failing is fatal-worthy, not drop-worthy diff --git a/src/internal/pipeline.zig b/src/internal/pipeline.zig index bfa4547..16077d4 100644 --- a/src/internal/pipeline.zig +++ b/src/internal/pipeline.zig @@ -18,7 +18,6 @@ const Allocator = std.mem.Allocator; const log = std.log.scoped(.stream); const worker_count = 4; -const worker_stack_size = 8 * 1024 * 1024; /// reorder-buffer capacity == max in-flight frames; bounds memory and /// gives the head-of-line window workers can verify ahead into const capacity = 256; @@ -26,6 +25,10 @@ const capacity = 256; const Slot = struct { ticket: u64, frame: []u8, // owned raw ws frame bytes + /// decode-once: the worker decodes into this arena; the writer consumes + /// the event and destroys the arena (null event = frame didn't decode) + arena: ?*std.heap.ArenaAllocator = null, + event: ?zat.FirehoseEvent = null, verdict: verify.Result = .valid, // meaning "no objection" when no verifier done: bool = false, }; @@ -37,7 +40,7 @@ pub const Pipeline = struct { /// emit(ctx, frame_bytes, verdict) — runs on the writer thread, in /// ticket order; owns nothing (pipeline frees the frame after). emit_ctx: *anyopaque, - emit: *const fn (*anyopaque, []const u8, verify.Result) void, + emit: *const fn (*anyopaque, zat.FirehoseEvent, verify.Result) void, mutex: Io.Mutex = .init, /// workers wait for work; reader waits for space; writer waits for done @@ -50,28 +53,42 @@ pub const Pipeline = struct { next_emit: u64 = 0, // next ticket the writer will emit alive: bool = true, - threads: [worker_count + 1]?std.Thread = @splat(null), + // workers + writer run as io.concurrent tasks, never raw std.Threads: + // they call io-backed I/O (DID resolution, archive writes), and Io + // primitives must run under the backend that created them + // (docs/lessons-from-zlay.md #2) + futures: [worker_count + 1]?Io.Future(void) = @splat(null), pub fn start(self: *Pipeline) !void { var spawned: usize = 0; - errdefer for (self.threads[0..spawned]) |t| { - if (t) |thread| thread.join(); - }; + errdefer { + self.signalStop(); + for (self.futures[0..spawned]) |*f| { + if (f.*) |*future| _ = future.cancel(self.io); + } + } for (0..worker_count) |i| { - self.threads[i] = try std.Thread.spawn(.{ .stack_size = worker_stack_size }, workerLoop, .{self}); + self.futures[i] = try self.io.concurrent(workerLoop, .{self}); spawned += 1; } - self.threads[worker_count] = try std.Thread.spawn(.{ .stack_size = worker_stack_size }, writerLoop, .{self}); + self.futures[worker_count] = try self.io.concurrent(writerLoop, .{self}); } - pub fn stop(self: *Pipeline) void { + fn signalStop(self: *Pipeline) void { self.mutex.lockUncancelable(self.io); self.alive = false; self.work_cond.broadcast(self.io); self.space_cond.broadcast(self.io); self.done_cond.broadcast(self.io); self.mutex.unlock(self.io); - for (self.threads) |t| if (t) |thread| thread.join(); + } + + pub fn stop(self: *Pipeline) void { + self.signalStop(); + for (&self.futures) |*f| { + if (f.*) |*future| _ = future.cancel(self.io); + f.* = null; + } for (&self.slots) |*slot| { if (slot.*) |s| self.allocator.free(s.frame); slot.* = null; @@ -80,6 +97,13 @@ pub const Pipeline = struct { /// called by the ws reader: takes ownership of nothing — dupes `data`. /// blocks when the buffer is full (backpressure). + pub const Gauges = struct { submitted: u64, claimed: u64, emitted: u64 }; + pub fn gauges(self: *Pipeline) Gauges { + self.mutex.lockUncancelable(self.io); + defer self.mutex.unlock(self.io); + return .{ .submitted = self.next_submit, .claimed = self.next_claim, .emitted = self.next_emit }; + } + pub fn submit(self: *Pipeline, data: []const u8) !void { const frame = try self.allocator.dupe(u8, data); errdefer self.allocator.free(frame); @@ -95,6 +119,8 @@ pub const Pipeline = struct { } fn workerLoop(self: *Pipeline) void { + var resolver_storage = if (self.verifier) |v| v.makeResolver() else null; + defer if (resolver_storage) |*r| r.deinit(); while (true) { // claim the next unverified ticket self.mutex.lockUncancelable(self.io); @@ -110,10 +136,31 @@ pub const Pipeline = struct { const frame = self.slots[ticket % capacity].?.frame; self.mutex.unlock(self.io); - // verify outside the lock (this is the parallel part) - const verdict = self.verifyFrame(frame); + // decode + verify outside the lock (this is the parallel part) + var arena: ?*std.heap.ArenaAllocator = null; + var event: ?zat.FirehoseEvent = null; + var verdict: verify.Result = .valid; + if (self.allocator.create(std.heap.ArenaAllocator)) |a| { + a.* = std.heap.ArenaAllocator.init(self.allocator); + if (zat.firehose.decodeFrame(a.allocator(), frame)) |ev| { + arena = a; + event = ev; + if (self.verifier) |v| { + verdict = switch (ev) { + .commit => |c| v.verifyCommit(if (resolver_storage) |*r| r else undefined, c.repo, c.blocks), + else => .valid, + }; + } + } else |err| { + log.debug("frame decode error: {s}", .{@errorName(err)}); + a.deinit(); + self.allocator.destroy(a); + } + } else |_| {} self.mutex.lockUncancelable(self.io); + self.slots[ticket % capacity].?.arena = arena; + self.slots[ticket % capacity].?.event = event; self.slots[ticket % capacity].?.verdict = verdict; self.slots[ticket % capacity].?.done = true; self.done_cond.broadcast(self.io); @@ -121,17 +168,6 @@ pub const Pipeline = struct { } } - fn verifyFrame(self: *Pipeline, frame: []const u8) verify.Result { - const v = self.verifier orelse return .valid; - var arena = std.heap.ArenaAllocator.init(self.allocator); - defer arena.deinit(); - const event = zat.firehose.decodeFrame(arena.allocator(), frame) catch return .valid; // writer re-decodes and drops/counts - return switch (event) { - .commit => |c| v.verifyCommit(c.repo, c.blocks), - else => .valid, - }; - } - fn writerLoop(self: *Pipeline) void { while (true) { self.mutex.lockUncancelable(self.io); @@ -150,7 +186,11 @@ pub const Pipeline = struct { self.space_cond.signal(self.io); self.mutex.unlock(self.io); - self.emit(self.emit_ctx, slot.frame, slot.verdict); + if (slot.event) |ev| self.emit(self.emit_ctx, ev, slot.verdict); + if (slot.arena) |a| { + a.deinit(); + self.allocator.destroy(a); + } self.allocator.free(slot.frame); } } @@ -165,10 +205,11 @@ const Collector = struct { allocator: Allocator, // no lock needed: emit runs only on the pipeline's single writer thread - fn emit(ctx: *anyopaque, frame: []const u8, verdict: verify.Result) void { + fn emit(ctx: *anyopaque, event: zat.FirehoseEvent, verdict: verify.Result) void { _ = verdict; const self: *Collector = @ptrCast(@alignCast(ctx)); - const copy = self.allocator.dupe(u8, frame) catch return; + // test frames are #info envelopes carrying a name payload + const copy = self.allocator.dupe(u8, event.info.name orelse "") catch return; self.frames.append(self.allocator, copy) catch {}; } }; @@ -192,10 +233,19 @@ test "pipeline preserves submission order under parallel workers" { }; try p.start(); - var buf: [16]u8 = undefined; + // frames must decode: #info envelopes with name = "frame-" for (0..1000) |i| { - const frame = try std.fmt.bufPrint(&buf, "frame-{d}", .{i}); - try p.submit(frame); + var name_buf: [16]u8 = undefined; + const name = try std.fmt.bufPrint(&name_buf, "frame-{d}", .{i}); + var frame_buf: [64]u8 = undefined; + var fw: Io.Writer = .fixed(&frame_buf); + // header {"t":"#info","op":1} (length-first key order) + try fw.writeAll(&.{ 0xa2, 0x61, 't', 0x65, '#', 'i', 'n', 'f', 'o', 0x62, 'o', 'p', 0x01 }); + // payload {"name": } + try fw.writeAll(&.{ 0xa1, 0x64, 'n', 'a', 'm', 'e' }); + try fw.writeByte(@intCast(0x60 + name.len)); + try fw.writeAll(name); + try p.submit(fw.buffered()); } // drain: wait until everything emitted, then stop while (true) { @@ -208,8 +258,9 @@ test "pipeline preserves submission order under parallel workers" { p.stop(); try testing.expectEqual(@as(usize, 1000), collector.frames.items.len); + var check_buf: [16]u8 = undefined; for (collector.frames.items, 0..) |f, i| { - const expect = try std.fmt.bufPrint(&buf, "frame-{d}", .{i}); + const expect = try std.fmt.bufPrint(&check_buf, "frame-{d}", .{i}); try testing.expectEqualStrings(expect, f); } } diff --git a/src/internal/server.zig b/src/internal/server.zig index cf2b015..d6a9804 100644 --- a/src/internal/server.zig +++ b/src/internal/server.zig @@ -30,6 +30,8 @@ pub const Hub = struct { tail: *tail_mod.Tail, stats: *metrics.Stats, archive: ?*archive_mod.Archive = null, + pipeline_gauges: ?*const fn (*anyopaque) [3]u64 = null, + pipeline_ctx: ?*anyopaque = null, active: std.atomic.Value(u32) = .init(0), }; @@ -322,7 +324,18 @@ pub const Handler = struct { .base = g.base, .tip = g.tip, }, now_s); - return respond(conn, "200 OK", "text/plain; version=0.0.4", out); + var extra_buf: [256]u8 = undefined; + var extra: []const u8 = ""; + if (hub.pipeline_gauges) |pg| { + const v = pg(hub.pipeline_ctx.?); + extra = std.fmt.bufPrint(&extra_buf, "# TYPE stream_pipeline_ticket gauge\n" ++ + "stream_pipeline_ticket{{stage=\"submitted\"}} {d}\n" ++ + "stream_pipeline_ticket{{stage=\"claimed\"}} {d}\n" ++ + "stream_pipeline_ticket{{stage=\"emitted\"}} {d}\n", .{ v[0], v[1], v[2] }) catch ""; + } + var joined: [17 * 1024]u8 = undefined; + const full = std.fmt.bufPrint(&joined, "{s}{s}", .{ out, extra }) catch out; + return respond(conn, "200 OK", "text/plain; version=0.0.4", full); } respond(conn, "404 Not Found", "text/plain", "not found\n"); } diff --git a/src/internal/verify.zig b/src/internal/verify.zig index 5b96ac2..3d1cede 100644 --- a/src/internal/verify.zig +++ b/src/internal/verify.zig @@ -44,10 +44,10 @@ pub const Result = enum { pub const Verifier = struct { allocator: Allocator, io: Io, - resolver: zat.DidResolver, - /// guards resolver + cache: workers verify concurrently. resolution - /// holds the lock (serializing misses) — steady state is hit-dominated, - /// and duplicate concurrent resolves would hammer PLC for nothing. + /// guards the cache map only. resolution I/O runs outside it, on the + /// CALLER's resolver — each verify worker owns one, so misses resolve + /// in parallel (duplicate concurrent resolves of one DID are rare and + /// harmless). mutex: Io.Mutex = .init, cache: std.StringHashMapUnmanaged(CachedKey) = .empty, max_cache: u32 = 250_000, @@ -55,31 +55,31 @@ pub const Verifier = struct { /// plc_url: base URL of the PLC directory. for the simulator this is its /// http address (it serves GET /did:...); production is plc.directory. + plc_url: []const u8, + pub fn init(allocator: Allocator, io: Io, plc_url: []const u8) Verifier { - return .{ - .allocator = allocator, - .io = io, - .resolver = resolver: { - var r = zat.DidResolver.init(io, allocator); - r.plc_url = plc_url; - break :resolver r; - }, - }; + return .{ .allocator = allocator, .io = io, .plc_url = plc_url }; + } + + /// one per verify worker; owns a keep-alive connection + pub fn makeResolver(self: *Verifier) zat.DidResolver { + var r = zat.DidResolver.init(self.io, self.allocator); + r.plc_url = self.plc_url; + return r; } pub fn deinit(self: *Verifier) void { var it = self.cache.keyIterator(); while (it.next()) |k| self.allocator.free(k.*); self.cache.deinit(self.allocator); - self.resolver.deinit(); } /// verify the commit signature in `blocks` (the event's CAR diff) /// against `did`'s signing key. - pub fn verifyCommit(self: *Verifier, did: []const u8, blocks: []const u8) Result { + pub fn verifyCommit(self: *Verifier, resolver: *zat.DidResolver, did: []const u8, blocks: []const u8) Result { // key lookup/resolution under the lock; signature crypto outside it // so verify workers actually run in parallel - const key = self.lockedGetKey(did, false) orelse { + const key = self.lockedGetKey(resolver, did, false) orelse { self.count("verify_unverified"); return .unverified; }; @@ -88,7 +88,7 @@ pub const Verifier = struct { return .valid; } // key may have rotated: re-resolve once and retry - const fresh = self.lockedGetKey(did, true) orelse { + const fresh = self.lockedGetKey(resolver, did, true) orelse { self.count("verify_unverified"); return .unverified; }; @@ -100,10 +100,63 @@ pub const Verifier = struct { return .invalid_signature; } - fn lockedGetKey(self: *Verifier, did: []const u8, force_resolve: bool) ?CachedKey { + fn lockedGetKey(self: *Verifier, resolver: *zat.DidResolver, did: []const u8, force_resolve: bool) ?CachedKey { + if (!force_resolve) { + self.mutex.lockUncancelable(self.io); + const hit = self.cache.get(did); + self.mutex.unlock(self.io); + if (hit) |k| { + self.count("verify_cache_hits"); + return k; + } + self.count("verify_cache_misses"); + } else { + self.evict(did); + } + return self.resolveAndCache(resolver, did); + } + + /// resolve on the caller's own resolver, outside the cache lock + fn resolveAndCache(self: *Verifier, resolver: *zat.DidResolver, did: []const u8) ?CachedKey { + const doc_or_null: ?zat.DidDocument = doc: { + const parsed = zat.Did.parse(did) orelse break :doc null; + break :doc resolver.resolve(parsed) catch |err| { + log.debug("DID resolve failed for {s}: {s}", .{ did, @errorName(err) }); + break :doc null; + }; + }; + var doc = doc_or_null orelse return null; + defer doc.deinit(); + + const vm = doc.signingKey() orelse return null; + const key_bytes = zat.multibase.decode(self.allocator, vm.public_key_multibase) catch return null; + defer self.allocator.free(key_bytes); + const public_key = zat.multicodec.parsePublicKey(key_bytes) catch return null; + if (public_key.raw.len > 33) return null; + + var cached = CachedKey{ + .key_type = public_key.key_type, + .raw = undefined, + .len = @intCast(public_key.raw.len), + }; + @memcpy(cached.raw[0..public_key.raw.len], public_key.raw); + self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io); - return self.getKey(did, force_resolve); + if (self.cache.count() >= self.max_cache) { + var it = self.cache.keyIterator(); + while (it.next()) |k| self.allocator.free(k.*); + self.cache.clearRetainingCapacity(); + } + const gop = self.cache.getOrPut(self.allocator, did) catch return cached; + if (!gop.found_existing) { + gop.key_ptr.* = self.allocator.dupe(u8, did) catch { + _ = self.cache.remove(did); + return cached; + }; + } + gop.value_ptr.* = cached; + return cached; } fn count(self: *Verifier, comptime field: []const u8) void { @@ -136,47 +189,4 @@ pub const Verifier = struct { return true; } - fn getKey(self: *Verifier, did: []const u8, force_resolve: bool) ?CachedKey { - if (!force_resolve) { - if (self.cache.get(did)) |k| { - self.count("verify_cache_hits"); - return k; - } - self.count("verify_cache_misses"); - } else { - self.evictLocked(did); - } - - const parsed = zat.Did.parse(did) orelse return null; - var doc = self.resolver.resolve(parsed) catch |err| { - log.debug("DID resolve failed for {s}: {s}", .{ did, @errorName(err) }); - return null; - }; - defer doc.deinit(); - - const vm = doc.signingKey() orelse return null; - const key_bytes = zat.multibase.decode(self.allocator, vm.public_key_multibase) catch return null; - defer self.allocator.free(key_bytes); - const public_key = zat.multicodec.parsePublicKey(key_bytes) catch return null; - if (public_key.raw.len > 33) return null; - - var cached = CachedKey{ - .key_type = public_key.key_type, - .raw = undefined, - .len = @intCast(public_key.raw.len), - }; - @memcpy(cached.raw[0..public_key.raw.len], public_key.raw); - - // crude bound: reset when full (TODO: LRU, like zlay's) - if (self.cache.count() >= self.max_cache) { - var it = self.cache.keyIterator(); - while (it.next()) |k| self.allocator.free(k.*); - self.cache.clearRetainingCapacity(); - } - const owned = self.allocator.dupe(u8, did) catch return cached; - self.cache.put(self.allocator, owned, cached) catch { - self.allocator.free(owned); - }; - return cached; - } }; diff --git a/src/main.zig b/src/main.zig index a22cecf..16dd39f 100644 --- a/src/main.zig +++ b/src/main.zig @@ -7,6 +7,7 @@ const cursor_mod = @import("internal/cursor.zig"); const ingest = @import("internal/ingest.zig"); const verify = @import("internal/verify.zig"); const metrics = @import("internal/metrics.zig"); +const pipeline_mod = @import("internal/pipeline.zig"); const server_mod = @import("internal/server.zig"); const tail_mod = @import("internal/tail.zig"); @@ -120,9 +121,27 @@ pub fn main(init: std.process.Init.Minimal) !void { .archive = &archive, .stats = &stats, }; + var pipe: pipeline_mod.Pipeline = .{ + .allocator = allocator, + .io = io, + .verifier = if (verifier) |*v| v else null, + .emit_ctx = &consumer, + .emit = ingest.Consumer.emitVerified, + }; + try pipe.start(); + defer pipe.stop(); + consumer.pipeline = &pipe; + hub.pipeline_ctx = &pipe; + hub.pipeline_gauges = &pipelineGauges; try consumer.run(); } +fn pipelineGauges(ctx: *anyopaque) [3]u64 { + const pipe: *pipeline_mod.Pipeline = @ptrCast(@alignCast(ctx)); + const g = pipe.gauges(); + return .{ g.submitted, g.claimed, g.emitted }; +} + fn runWsServer(server: *websocket.Server(server_mod.Handler), listener: *Io.net.Server, hub: *server_mod.Hub) void { server.runIo(listener, hub); } -- 2.51.2