diff --git a/docs/grafana-dashboard.md b/docs/grafana-dashboard.md index 4a547ec..55e8a2a 100644 --- a/docs/grafana-dashboard.md +++ b/docs/grafana-dashboard.md @@ -41,7 +41,11 @@ counters. Values come from `getrusage` plus Linux procfs or Darwin `proc_pidinfo`. Unavailable measurements are omitted, never reported as zero. `just process-metrics-contract` launches the production binary with only loopback dependencies and validates types, values, invariants, monotonicity, -and the absence of `go_*` series across independent HTTP scrapes. +and the absence of `go_*` series across independent HTTP scrapes. The same +receipt proves real startup RocksDB reads, point writes, and batch commits +populate the canonical store histogram. A real RocksDB unit receipt covers +point-read `ok`/`notfound`, durable `set`/`delete`, and atomic `batch_commit` +boundaries; registered error series remain zero until RocksDB reports one. `just http-metrics-contract` seeds a real sealed archive and exercises the ReleaseSafe binary through ordinary HTTP, Range/conditional getBlock requests, and a real WebSocket upgrade/close. It proves the same public-mux middleware @@ -102,6 +106,7 @@ The following dashboard families already have non-placeholder producers: - `jetstream_manifest_block_index_cache_hits_total` - `jetstream_manifest_block_index_cache_misses_total` - `jetstream_manifest_block_index_load_seconds` +- `jetstream_store_op_duration_seconds` - `jetstream_verifier_failures_total` - `jetstream_backfill_completed_total` - `jetstream_backfill_failed_total` diff --git a/docs/semantic-parity.md b/docs/semantic-parity.md index 956a158..e1508ef 100644 --- a/docs/semantic-parity.md +++ b/docs/semantic-parity.md @@ -20,7 +20,7 @@ this document or `bootstrap-semantic-parity.md` is open. | 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. | | Status/diagnostics | Repo attempts/error class/final PDS and per-host aggregates survive restart and are operator-queryable | **closed at the implementation boundary.** Durable v5 repo rows distinguish initial Backfill.Rev from latest Rev/UpdatedAt, retain resolved Handle/PDS and reserved record/byte fields, and atomically maintain `handle/`. Merge source cursor commits refresh latest revisions without mutating backfill watermarks. Same-batch `host/` totals, active/status counts, cumulative error classes, and five bounded recent samples survive restart; host moves and active flips are atomic with the source repo row. `/status?tab=hosts` reads those rows without a whole-network scan and matches upstream ordering and error presentation. Account queries preserve resolver-first handle semantics, the legacy `did`/`handle` aliases, and missing-identity hydration, then reconstruct records last-writer-wins from checksum-verified, DID-bloom-pruned sealed JSS plus rotation-safe active/pending and bootstrap-live snapshots. The resulting canonical MST is compared with the PDS commit root through real Sync 1.1 `getLatestCommit`/`getBlocks`; the page renders the same match fields and error state as upstream. `HEAD` performs lookup without verification; successful status responses are `no-store` and carry a generation timestamp. Expensive verification uses upstream's exact 4/source-IP/minute fixed-window limiter, 4,096-entry bound, stale pruning, oldest eviction, and explicit disable option. The offline receipt uses a real HTTP Stream route, real RocksDB row, real JSS archive, signed commit, and loopback PLC/PDS; four separate connections verify successfully, a fifth changing only its ephemeral port is blocked without another upstream request, HEAD remains side-effect-free, and disabling the limiter performs a new authoritative verification. | -| Prometheus/Grafana | Every canonical dashboard query maps to a real, same-semantics producer; runtime-specific panels use honest Zig/process metrics | exact checksum-pinned upstream dashboard is vendored with an offline fail-closed scrape contract; canonical build/live-firehose/archive-append/drop/verifier/backfill/timestamp-import producers are present. The complete canonical subscriber family set is exported from real scan, delivery, rejection, cursor, and disconnect boundaries; `none`, `deflate`, and `zstd` labels now arise only from the actual negotiated delivery path. Live frame inspection distinguishes gaps, bounded-label server errors, forward-compatible unknowns, and malformed CBOR at the raw transport boundary; a Zat reconnect callback measures attempts rather than connection failures. The only dashboard divergence is now a deterministic, checksum-pinned process-runtime row: same-semantics start/CPU/RSS/FD measurements, real VM/thread/page-fault/context-switch signals, no `go_*` aliases, and no host-network counters mislabeled as process I/O. A production-binary loopback receipt proves values and cumulative monotonicity on independent scrapes. The hot readable log now matches upstream's durability pin: rows above the archive watermark may exceed the byte budget but cannot be evicted, publication releases their pinned-byte accounting, and only then may retention advance the floor; all three canonical readable-log gauges come from that exact state. Data-directory free bytes are collected fresh from the archive directory's filesystem at scrape time, and the canonical seal histogram measures only successful flush-through-footer/header-fsync work with upstream's exact 15 slow-latency buckets. HTTP duration uses upstream's public-mux boundary and stable labels, including inline WebSocket lifetime and the single `xrpc/` subtree; debug and unmatched routes remain unobserved. getBlock outcomes, duration, and served bytes come from actual response completion, with bytes counted only for complete 200 responses. A ReleaseSafe offline receipt proves 200/206/304/400/404/416 behavior and WebSocket close accounting. Sealed headers and block indexes are now genuinely resident: startup and verified refresh latency, loaded count, and lookup hits/misses use upstream's exact boundaries. Refcounted snapshots remain memory-safe across refresh, while fresh-file offsets prevent cross-generation splicing; the pinned official client proves cold replay produces hits. **Remaining resident bloom/collection planning, store, orchestrator, backfill, and compaction families open.** | +| Prometheus/Grafana | Every canonical dashboard query maps to a real, same-semantics producer; runtime-specific panels use honest Zig/process metrics | exact checksum-pinned upstream dashboard is vendored with an offline fail-closed scrape contract; canonical build/live-firehose/archive-append/drop/verifier/backfill/timestamp-import producers are present. The complete canonical subscriber family set is exported from real scan, delivery, rejection, cursor, and disconnect boundaries; `none`, `deflate`, and `zstd` labels now arise only from the actual negotiated delivery path. Live frame inspection distinguishes gaps, bounded-label server errors, forward-compatible unknowns, and malformed CBOR at the raw transport boundary; a Zat reconnect callback measures attempts rather than connection failures. The only dashboard divergence is now a deterministic, checksum-pinned process-runtime row: same-semantics start/CPU/RSS/FD measurements, real VM/thread/page-fault/context-switch signals, no `go_*` aliases, and no host-network counters mislabeled as process I/O. A production-binary loopback receipt proves values and cumulative monotonicity on independent scrapes. The hot readable log now matches upstream's durability pin: rows above the archive watermark may exceed the byte budget but cannot be evicted, publication releases their pinned-byte accounting, and only then may retention advance the floor; all three canonical readable-log gauges come from that exact state. Data-directory free bytes are collected fresh from the archive directory's filesystem at scrape time, and the canonical seal histogram measures only successful flush-through-footer/header-fsync work with upstream's exact 15 slow-latency buckets. HTTP duration uses upstream's public-mux boundary and stable labels, including inline WebSocket lifetime and the single `xrpc/` subtree; debug and unmatched routes remain unobserved. getBlock outcomes, duration, and served bytes come from actual response completion, with bytes counted only for complete 200 responses. A ReleaseSafe offline receipt proves 200/206/304/400/404/416 behavior and WebSocket close accounting. Sealed headers and block indexes are now genuinely resident: startup and verified refresh latency, loaded count, and lookup hits/misses use upstream's exact boundaries. Refcounted snapshots remain memory-safe across refresh, while fresh-file offsets prevent cross-generation splicing; the pinned official client proves cold replay produces hits. RocksDB point reads, durable point writes/deletes, and atomic batch commits populate upstream's exact `{op,status}` store histogram and fast-latency buckets; real-database and production-startup receipts prove the non-error paths. **Remaining resident bloom/collection planning, orchestrator, backfill, and compaction families open.** | | Oracle | Event-log equivalence, final-state convergence, crash/power-loss, verifier repair, hostile input, and anti-vacuity receipts | **open — current crash matrix is necessary but substantially smaller than upstream's oracle** | ## Experiment admission diff --git a/src/internal/meta_store.zig b/src/internal/meta_store.zig index ea56486..8039396 100644 --- a/src/internal/meta_store.zig +++ b/src/internal/meta_store.zig @@ -7,8 +7,10 @@ const std = @import("std"); const rocksdb = @import("rocksdb"); +const metrics = @import("metrics.zig"); const Allocator = std.mem.Allocator; +const Io = std.Io; const log = std.log.scoped(.stream_meta); pub const subdir = "meta.rocksdb"; @@ -17,8 +19,18 @@ pub const Store = struct { allocator: Allocator, db: rocksdb.DB, default_cf: rocksdb.ColumnFamilyHandle, + io: ?Io = null, + stats: ?*metrics.Stats = null, pub fn open(allocator: Allocator, data_dir: []const u8) !Store { + return openInner(allocator, data_dir, null, null); + } + + pub fn openWithStats(allocator: Allocator, io: Io, data_dir: []const u8, stats: *metrics.Stats) !Store { + return openInner(allocator, data_dir, io, stats); + } + + fn openInner(allocator: Allocator, data_dir: []const u8, io: ?Io, stats: ?*metrics.Stats) !Store { const path = try std.fs.path.join(allocator, &.{ data_dir, subdir }); defer allocator.free(path); @@ -42,7 +54,13 @@ pub const Store = struct { if (families.len != 1) return error.InvalidMetadataStore; log.info("metadata store opened at {s}", .{path}); - return .{ .allocator = allocator, .db = db, .default_cf = families[0].handle }; + return .{ + .allocator = allocator, + .db = db, + .default_cf = families[0].handle, + .io = io, + .stats = stats, + }; } pub fn deinit(self: *Store) void { @@ -94,28 +112,44 @@ pub const Store = struct { } pub fn getAlloc(self: *Store, allocator: Allocator, key: []const u8) !?[]u8 { + const started_at = self.started(); var err_data: ?rocksdb.Data = null; defer if (err_data) |e| e.deinit(); const value = self.db.get(self.default_cf, key, &err_data) catch |err| { + self.observe(.get, .error_result, started_at); if (err_data) |e| log.err("get {s}: {s}", .{ key, e.data }); return err; - } orelse return null; + } orelse { + self.observe(.get, .notfound, started_at); + return null; + }; + self.observe(.get, .ok, started_at); defer value.deinit(); return try allocator.dupe(u8, value.data); } pub fn putDurable(self: *Store, key: []const u8, value: []const u8) !void { - var writes = self.batch(); - defer writes.deinit(); - writes.put(key, value); - try writes.commitDurable(); + const started_at = self.started(); + var err_data: ?rocksdb.Data = null; + defer if (err_data) |e| e.deinit(); + self.db.putWithOptions(self.default_cf, key, value, .{ .sync = true }, &err_data) catch |err| { + self.observe(.set, .error_result, started_at); + if (err_data) |e| log.err("durable put {s}: {s}", .{ key, e.data }); + return err; + }; + self.observe(.set, .ok, started_at); } pub fn deleteDurable(self: *Store, key: []const u8) !void { - var writes = self.batch(); - defer writes.deinit(); - writes.delete(key); - try writes.commitDurable(); + const started_at = self.started(); + var err_data: ?rocksdb.Data = null; + defer if (err_data) |e| e.deinit(); + self.db.deleteWithOptions(self.default_cf, key, .{ .sync = true }, &err_data) catch |err| { + self.observe(.delete, .error_result, started_at); + if (err_data) |e| log.err("durable delete {s}: {s}", .{ key, e.data }); + return err; + }; + self.observe(.delete, .ok, started_at); } pub fn iterator(self: *Store, start: ?[]const u8) rocksdb.Iterator { @@ -125,6 +159,19 @@ pub const Store = struct { pub fn batch(self: *Store) Batch { return .{ .store = self, .inner = rocksdb.WriteBatch.init() }; } + + fn started(self: *const Store) ?Io.Timestamp { + const io = self.io orelse return null; + return Io.Timestamp.now(io, .awake); + } + + fn observe(self: *Store, op: metrics.StoreOp, status: metrics.StoreStatus, started_at: ?Io.Timestamp) void { + const stats = self.stats orelse return; + const io = self.io orelse return; + const started_value = started_at orelse return; + const elapsed = started_value.durationTo(Io.Timestamp.now(io, .awake)).toMicroseconds(); + stats.observeStore(op, status, @intCast(@max(elapsed, 0))); + } }; pub const Batch = struct { @@ -144,12 +191,15 @@ pub const Batch = struct { } pub fn commitDurable(self: *Batch) !void { + const started_at = self.store.started(); var err_data: ?rocksdb.Data = null; defer if (err_data) |e| e.deinit(); self.store.db.writeWithOptions(self.inner, .{ .sync = true }, &err_data) catch |err| { + self.store.observe(.batch_commit, .error_result, started_at); if (err_data) |e| log.err("durable batch commit: {s}", .{e.data}); return err; }; + self.store.observe(.batch_commit, .ok, started_at); } }; @@ -181,3 +231,35 @@ test "durable point reads and atomic batches survive reopen" { try testing.expectEqualStrings("complete\x003lrev", repo); } } + +test "real RocksDB operations populate canonical store metric outcomes" { + const testing = std.testing; + var tmp = testing.tmpDir(.{}); + defer tmp.cleanup(); + var path_buf: [std.fs.max_path_bytes]u8 = undefined; + const path_len = try tmp.dir.realPath(testing.io, &path_buf); + var stats: metrics.Stats = .{}; + var store = try Store.openWithStats(testing.allocator, testing.io, path_buf[0..path_len], &stats); + defer store.deinit(); + + try store.putDurable("point", "value"); + const found = (try store.getAlloc(testing.allocator, "point")).?; + testing.allocator.free(found); + try testing.expectEqual(null, try store.getAlloc(testing.allocator, "missing")); + try store.deleteDurable("point"); + var writes = store.batch(); + defer writes.deinit(); + writes.put("batch/a", "1"); + writes.put("batch/b", "2"); + try writes.commitDurable(); + + try testing.expectEqual(1, stats.store_op_duration_count[storeMetricIndex(.get, .ok)].load(.monotonic)); + try testing.expectEqual(1, stats.store_op_duration_count[storeMetricIndex(.get, .notfound)].load(.monotonic)); + try testing.expectEqual(1, stats.store_op_duration_count[storeMetricIndex(.set, .ok)].load(.monotonic)); + try testing.expectEqual(1, stats.store_op_duration_count[storeMetricIndex(.delete, .ok)].load(.monotonic)); + try testing.expectEqual(1, stats.store_op_duration_count[storeMetricIndex(.batch_commit, .ok)].load(.monotonic)); +} + +fn storeMetricIndex(op: metrics.StoreOp, status: metrics.StoreStatus) usize { + return @intFromEnum(op) * @typeInfo(metrics.StoreStatus).@"enum".fields.len + @intFromEnum(status); +} diff --git a/src/internal/metrics.zig b/src/internal/metrics.zig index 9e6d535..c5cbf98 100644 --- a/src/internal/metrics.zig +++ b/src/internal/metrics.zig @@ -30,9 +30,13 @@ pub const HttpHandler = enum { pub const HttpMethod = enum { get, head, post, put, patch, delete, options, connect, trace, other }; pub const HttpCode = enum { c200, c201, c206, c304, c400, c401, c404, c405, c409, c416, c500, c503, other }; pub const GetBlockResult = enum { ok, not_found, bad_request, error_result }; +pub const StoreOp = enum { get, set, delete, batch_commit }; +pub const StoreStatus = enum { ok, notfound, error_result }; const http_series_count = @typeInfo(HttpHandler).@"enum".fields.len * @typeInfo(HttpMethod).@"enum".fields.len * @typeInfo(HttpCode).@"enum".fields.len; +const store_series_count = @typeInfo(StoreOp).@"enum".fields.len * + @typeInfo(StoreStatus).@"enum".fields.len; pub const Stats = struct { start_time_s: i64 = 0, @@ -60,6 +64,9 @@ pub const Stats = struct { manifest_block_index_load_count: Counter = .init(0), manifest_block_index_load_sum_us: Counter = .init(0), manifest_block_index_load_buckets: [8]Counter = @splat(.init(0)), + store_op_duration_count: [store_series_count]Counter = @splat(.init(0)), + store_op_duration_sum_us: [store_series_count]Counter = @splat(.init(0)), + store_op_duration_buckets: [store_series_count][15]Counter = @splat(@splat(.init(0))), upstream_seq: std.atomic.Value(i64) = .init(0), drops: std.enums.EnumArray(convert.DropReason, Counter) = .initFill(.init(0)), verify_valid: Counter = .init(0), @@ -200,8 +207,23 @@ pub const Stats = struct { _ = self.manifest_block_index_load_buckets[index].fetchAdd(1, .monotonic); } } + + pub fn observeStore(self: *Stats, op: StoreOp, status: StoreStatus, duration_us: u64) void { + const index = storeIndex(op, status); + _ = self.store_op_duration_count[index].fetchAdd(1, .monotonic); + _ = self.store_op_duration_sum_us[index].fetchAdd(duration_us, .monotonic); + inline for (0..15) |bucket_index| { + const upper_us = @as(u64, 100) << @intCast(bucket_index); + if (duration_us <= upper_us) + _ = self.store_op_duration_buckets[index][bucket_index].fetchAdd(1, .monotonic); + } + } }; +fn storeIndex(op: StoreOp, status: StoreStatus) usize { + return @intFromEnum(op) * @typeInfo(StoreStatus).@"enum".fields.len + @intFromEnum(status); +} + fn httpIndex(handler: HttpHandler, method: HttpMethod, code: HttpCode) usize { const method_count = @typeInfo(HttpMethod).@"enum".fields.len; const code_count = @typeInfo(HttpCode).@"enum".fields.len; @@ -517,6 +539,37 @@ pub fn format( ) catch {}; w.print("jetstream_manifest_block_index_load_seconds_bucket{{le=\"+Inf\"}} {d}\n", .{manifest_load_count}) catch {}; + w.print("# TYPE jetstream_store_op_duration_seconds histogram\n", .{}) catch {}; + const store_bounds = [_][]const u8{ + "0.0001", "0.0002", "0.0004", "0.0008", "0.0016", + "0.0032", "0.0064", "0.0128", "0.0256", "0.0512", + "0.1024", "0.2048", "0.4096", "0.8192", "1.6384", + }; + inline for (comptime std.enums.values(StoreOp)) |op| { + inline for (comptime std.enums.values(StoreStatus)) |status| { + if (op != .get and status == .notfound) continue; + const index = storeIndex(op, status); + const count = stats.store_op_duration_count[index].load(.monotonic); + const sum_us = stats.store_op_duration_sum_us[index].load(.monotonic); + inline for (store_bounds, 0..) |upper, bucket_index| w.print( + "jetstream_store_op_duration_seconds_bucket{{op=\"{s}\",status=\"{s}\",le=\"{s}\"}} {d}\n", + .{ storeOpLabel(op), storeStatusLabel(status), upper, stats.store_op_duration_buckets[index][bucket_index].load(.monotonic) }, + ) catch {}; + w.print( + "jetstream_store_op_duration_seconds_bucket{{op=\"{s}\",status=\"{s}\",le=\"+Inf\"}} {d}\n", + .{ storeOpLabel(op), storeStatusLabel(status), count }, + ) catch {}; + w.print( + "jetstream_store_op_duration_seconds_sum{{op=\"{s}\",status=\"{s}\"}} {d}.{d:0>6}\n", + .{ storeOpLabel(op), storeStatusLabel(status), sum_us / std.time.us_per_s, sum_us % std.time.us_per_s }, + ) catch {}; + w.print( + "jetstream_store_op_duration_seconds_count{{op=\"{s}\",status=\"{s}\"}} {d}\n", + .{ storeOpLabel(op), storeStatusLabel(status), count }, + ) catch {}; + } + } + // Canonical upstream names. These are aliases only where Stream owns the // same lifecycle boundary; absent families stay absent so the dashboard // contract cannot be satisfied by placeholder zeroes. @@ -871,6 +924,23 @@ fn getBlockResultLabel(value: GetBlockResult) []const u8 { }; } +fn storeOpLabel(value: StoreOp) []const u8 { + return switch (value) { + .get => "get", + .set => "set", + .delete => "delete", + .batch_commit => "batch_commit", + }; +} + +fn storeStatusLabel(value: StoreStatus) []const u8 { + return switch (value) { + .ok => "ok", + .notfound => "notfound", + .error_result => "error", + }; +} + // === tests === test "http handler labels match upstream public middleware boundaries" { @@ -956,6 +1026,9 @@ test "format renders counters" { _ = stats.manifest_block_index_cache_hits_total.fetchAdd(4, .monotonic); _ = stats.manifest_block_index_cache_misses_total.fetchAdd(5, .monotonic); stats.observeManifestLoad(400); + stats.observeStore(.get, .ok, 200); + stats.observeStore(.get, .notfound, 400); + stats.observeStore(.set, .error_result, 800); var buf: [64 * 1024]u8 = undefined; const out = format(&buf, &stats, 2, .{ @@ -1011,6 +1084,10 @@ test "format renders counters" { try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_manifest_block_index_cache_misses_total 5") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_manifest_block_index_load_seconds_bucket{le=\"0.0004\"} 1") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_manifest_block_index_load_seconds_count 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_store_op_duration_seconds_bucket{op=\"get\",status=\"ok\",le=\"0.0002\"} 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_store_op_duration_seconds_count{op=\"get\",status=\"notfound\"} 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_store_op_duration_seconds_count{op=\"set\",status=\"error\"} 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_store_op_duration_seconds_count{op=\"delete\",status=\"error\"} 0") != null); try std.testing.expect(std.mem.indexOf(u8, out, "stream_backfill_repos_total{result=\"complete\"} 11") != null); try std.testing.expect(std.mem.indexOf(u8, out, "stream_backfill_workers_active 4") != null); try std.testing.expect(std.mem.indexOf(u8, out, "stream_backfill_queue_depth 23") != null); diff --git a/src/main.zig b/src/main.zig index dcfaec2..e446fe9 100644 --- a/src/main.zig +++ b/src/main.zig @@ -185,7 +185,7 @@ pub fn main(init: std.process.Init.Minimal) !void { var xrpc_api: xrpcapi.Api = .{ .allocator = allocator, .io = io, .archive = &archive }; hub.xrpc = &xrpc_api; - var meta = try meta_store.Store.open(allocator, data_dir); + var meta = try meta_store.Store.openWithStats(allocator, io, data_dir, &stats); defer meta.deinit(); try meta.migrateLegacy(io, data_dir); var repo_diagnostics = try repo_store.Store.init(allocator, io, &meta); diff --git a/tests/process_metrics_contract.py b/tests/process_metrics_contract.py index 0ecb5a6..3f47faa 100644 --- a/tests/process_metrics_contract.py +++ b/tests/process_metrics_contract.py @@ -13,6 +13,14 @@ import urllib.request BASE = os.environ.get("STREAM_BASE_URL", "http://127.0.0.1:6008") TYPE_RE = re.compile(r"^# TYPE (\S+) (\S+)$") SAMPLE_RE = re.compile(r"^(\S+?)(?:\{[^}]*\})?\s+([^\s]+)$") +STORE_COUNT_RE = re.compile( + r'^jetstream_store_op_duration_seconds_count\{op="([^"]+)",status="([^"]+)"\}\s+([^\s]+)$', + re.MULTILINE, +) +STORE_BUCKET_RE = re.compile( + r'^jetstream_store_op_duration_seconds_bucket\{op="([^"]+)",status="([^"]+)",le="([^"]+)"\}\s+([^\s]+)$', + re.MULTILINE, +) EXPECTED_TYPES = { "process_start_time_seconds": "gauge", @@ -51,6 +59,36 @@ def main() -> None: assert first_types.get(name) == kind, (name, first_types.get(name), kind) assert name in first and math.isfinite(first[name]) and first[name] >= 0, (name, first.get(name)) + assert first_types.get("jetstream_store_op_duration_seconds") == "histogram" + store_counts = { + (op, status): float(value) + for op, status, value in STORE_COUNT_RE.findall(first_raw) + } + expected_store_series = { + ("get", "ok"), + ("get", "notfound"), + ("get", "error"), + ("set", "ok"), + ("set", "error"), + ("delete", "ok"), + ("delete", "error"), + ("batch_commit", "ok"), + ("batch_commit", "error"), + } + assert set(store_counts) == expected_store_series, store_counts + assert store_counts[("get", "notfound")] > 0 + assert store_counts[("set", "ok")] > 0 + assert store_counts[("batch_commit", "ok")] > 0 + buckets: dict[tuple[str, str], list[tuple[str, float]]] = {} + for op, status, upper, value in STORE_BUCKET_RE.findall(first_raw): + buckets.setdefault((op, status), []).append((upper, float(value))) + assert set(buckets) == expected_store_series + for labels, observations in buckets.items(): + assert len(observations) == 16, (labels, observations) + assert observations[-1] == ("+Inf", store_counts[labels]), (labels, observations[-1]) + values = [value for _, value in observations] + assert values == sorted(values), (labels, values) + assert 0 < first["process_start_time_seconds"] <= time.time() assert 0 < first["process_resident_memory_bytes"] <= first["process_virtual_memory_bytes"] assert first["stream_process_threads"] >= 1 @@ -71,7 +109,10 @@ def main() -> None: f"rss={int(second['process_resident_memory_bytes'])} " f"vm={int(second['process_virtual_memory_bytes'])} " f"threads={int(second['stream_process_threads'])} " - f"fds={int(second['process_open_fds'])}/{int(second['process_max_fds'])}" + f"fds={int(second['process_open_fds'])}/{int(second['process_max_fds'])} " + f"store=get-notfound:{int(store_counts[('get', 'notfound')])}," + f"set:{int(store_counts[('set', 'ok')])}," + f"batch:{int(store_counts[('batch_commit', 'ok')])}" )