From de45ca73ea3bb6df33ab8937f04f3b530f4bb45e Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Fri, 2 Oct 2026 18:29:17 -0700 Subject: [PATCH] zig 0.17: port to the 0.17.0 release Builds, tests, and formats on zig 0.17.0. zig build test: 365 passed, 1 skipped (the linux-only process metrics test), same as 0.16 on main. Dependencies: zat v0.6.0-alpha.2, websocket v0.2.0-alpha.2, and the 0.17 ports of rocksdb-zig and otel-zig. CI and the Dockerfile install zig 0.17.0. errdefer no longer captures the error, so the 19 traced scopes that recorded failures with `errdefer |err| trace.fail(err)` now run their covered code in a function and pass its result through Span.recordError, which fails the span on error and returns the result unchanged. Each wrapper covers exactly the scope the errdefer did, so the same errors are recorded and the surrounding defers run in the same order. Other 0.17 changes: - array repetition (`**`) becomes @splat or literals; zstd tests build their repeated strings with a comptime helper - @typeInfo enum `.fields` becomes `.field_names`; std.meta.fields is gone (enum counts, cli exposures, crashpoint parsing) - @cImport is gone: disk_space declares fstatvfs and struct statvfs for darwin, glibc, and musl - Uri.getHost is gone; host extraction reads the raw URI host, as getHost did, because HostName.fromUri rejects IPv6 literals - EnumSet initFull/initEmpty, Crc32Iscsi, and dupeZ renames - zig fmt rewrites @intFromEnum/@enumFromInt to @backingInt and @fromBackingInt zat 0.4 -> 0.6: - FetchOptions.redirect_behavior is on_redirect, response headers are captured by name - the repo export transport identifies as stream instead of the transport library's default user agent - cbor.Value gained .float (zat 0.4.2 accepts finite 64-bit floats like atmos), so the wire encoder writes it as a JSON number, and the undecodable-record hot tail test uses a 32-bit float, which is still rejected build.zig poisons the configure cache where it shells out to git: 0.17 caches the configure result, so the version command and stream_build_info kept reporting the previous commit after HEAD moved. Also fixes the bench step, which passed a removed argument to wire.encodeCommitV2 and did not compile on 0.16 either. Co-Authored-By: Claude Opus 5.5 (1M context) --- .tangled/workflows/ci.yml | 6 +- Dockerfile | 4 +- build.zig | 9 +- build.zig.zon | 18 +- src/bench.zig | 2 +- src/internal/compact/steady.zig | 7 +- src/internal/ingest/backfill/engine.zig | 193 +++++++++++------- src/internal/ingest/backfill/fetched_car.zig | 10 +- src/internal/ingest/backfill/host_store.zig | 30 +-- src/internal/ingest/backfill/lifecycle.zig | 109 +++++++--- src/internal/ingest/backfill/merge.zig | 74 +++++-- .../ingest/backfill/prepared_fetch.zig | 2 +- src/internal/ingest/backfill/repo_store.zig | 57 +++--- src/internal/ingest/backfill/repos.zig | 4 +- src/internal/ingest/backfill/resync.zig | 40 +++- src/internal/ingest/convert.zig | 2 +- src/internal/ingest/ingest.zig | 28 ++- src/internal/ingest/pipeline.zig | 2 +- src/internal/ingest/repair.zig | 2 +- .../ingest/repair_integration_test.zig | 6 +- src/internal/ingest/verify.zig | 10 +- src/internal/runtime/cli.zig | 12 +- src/internal/runtime/crashpoint.zig | 5 +- src/internal/runtime/logging.zig | 12 +- src/internal/runtime/metrics.zig | 32 +-- src/internal/runtime/observability.zig | 7 + src/internal/serve/filter.zig | 4 +- src/internal/serve/repo_export.zig | 9 +- src/internal/serve/server.zig | 6 +- src/internal/serve/wire.zig | 1 + src/internal/storage/archive.zig | 39 +++- src/internal/storage/cold.zig | 2 +- src/internal/storage/cold_cache.zig | 2 +- src/internal/storage/disk_space.zig | 33 ++- src/internal/storage/gloom.zig | 6 +- src/internal/storage/meta_store.zig | 2 +- src/internal/storage/segment.zig | 4 +- src/internal/storage/segment_fixture_test.zig | 2 +- src/internal/storage/segment_writer.zig | 4 +- src/internal/storage/zstd.zig | 13 +- src/internal/timestamp/jobs.zig | 2 +- src/internal/timestamp/parse.zig | 6 +- src/internal/timestamp/rules.zig | 8 +- src/internal/timestamp/runner.zig | 7 +- src/main.zig | 108 +++++++--- src/write_sample.zig | 2 +- 46 files changed, 616 insertions(+), 327 deletions(-) diff --git a/.tangled/workflows/ci.yml b/.tangled/workflows/ci.yml index 89156e7..c1fdc90 100644 --- a/.tangled/workflows/ci.yml +++ b/.tangled/workflows/ci.yml @@ -12,10 +12,10 @@ dependencies: - xz steps: - - name: install zig 0.16 + - name: install zig 0.17 command: | - curl -sSL https://ziglang.org/download/0.16.0/zig-x86_64-linux-0.16.0.tar.xz | tar -xJ -C /tangled/workspace - mv /tangled/workspace/zig-x86_64-linux-0.16.0 /tangled/workspace/.zig + curl -sSL https://ziglang.org/download/0.17.0/zig-x86_64-linux-0.17.0.tar.xz | tar -xJ -C /tangled/workspace + mv /tangled/workspace/zig-x86_64-linux-0.17.0 /tangled/workspace/.zig - name: check formatting command: | diff --git a/Dockerfile b/Dockerfile index 8c7c61c..336b049 100644 --- a/Dockerfile +++ b/Dockerfile @@ -2,8 +2,8 @@ # target; zstd/xxhash are vendored so no system dev packages are needed) FROM --platform=linux/amd64 debian:bookworm-slim AS builder RUN apt-get update && apt-get install -y --no-install-recommends curl xz-utils ca-certificates git && rm -rf /var/lib/apt/lists/* -RUN curl -fSL https://ziglang.org/download/0.16.0/zig-x86_64-linux-0.16.0.tar.xz | tar xJ -C /opt -ENV PATH=/opt/zig-x86_64-linux-0.16.0:$PATH +RUN curl -fSL https://ziglang.org/download/0.17.0/zig-x86_64-linux-0.17.0.tar.xz | tar xJ -C /opt +ENV PATH=/opt/zig-x86_64-linux-0.17.0:$PATH WORKDIR /build # fetch dependencies first (cacheable — only changes with build.zig.zon) diff --git a/build.zig b/build.zig index 24696b1..7ad7af5 100644 --- a/build.zig +++ b/build.zig @@ -100,6 +100,9 @@ pub fn build(b: *std.Build) void { build_options.addOption([]const u8, "version", b.option([]const u8, "version", "Release version reported by the version command") orelse "dev"); build_options.addOption([]const u8, "build_date", b.option([]const u8, "build-date", "RFC3339 build date reported by the version command") orelse "unknown"); build_options.addOption([]const u8, "git_sha", git_sha: { + // The configure result is cached and nothing it can track changes + // when HEAD moves; without this the binary reports a stale commit. + b.graph.poisonCache(); var code: u8 = 0; const result = b.runAllowFail(&.{ "git", "rev-parse", "--short", "HEAD" }, &code, .ignore); if (result) |output| { @@ -127,7 +130,7 @@ pub fn build(b: *std.Build) void { b.installArtifact(exe); const run_exe = b.addRunArtifact(exe); - if (b.args) |args| run_exe.addArgs(args); + run_exe.addPassthruArgs(); const run_step = b.step("run", "run stream"); run_step.dependOn(&run_exe.step); @@ -142,7 +145,7 @@ pub fn build(b: *std.Build) void { linkVendoredC(sample_mod, b, target, optimize); const write_sample = b.addExecutable(.{ .name = "write-sample-segment", .root_module = sample_mod }); const run_sample = b.addRunArtifact(write_sample); - if (b.args) |args| run_sample.addArgs(args); + run_sample.addPassthruArgs(); b.step("write-sample", "write a sample active segment (writer dev tool)").dependOn(&run_sample.step); const bench_mod = b.createModule(.{ @@ -156,7 +159,7 @@ pub fn build(b: *std.Build) void { linkVendoredC(bench_mod, b, target, optimize); const bench = b.addExecutable(.{ .name = "bench", .root_module = bench_mod }); const run_bench = b.addRunArtifact(bench); - if (b.args) |args| run_bench.addArgs(args); + run_bench.addPassthruArgs(); b.step("bench", "hot-path timings per subsystem (comparative, not a gate)").dependOn(&run_bench.step); const test_step = b.step("test", "run unit tests"); diff --git a/build.zig.zon b/build.zig.zon index 505ac35..c0047fc 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -2,15 +2,15 @@ .name = .stream, .version = "0.0.1", .fingerprint = 0xf0e9be1c73517a70, // Changing this has security and trust implications. - .minimum_zig_version = "0.16.0", + .minimum_zig_version = "0.17.0", .dependencies = .{ .zat = .{ - .url = "https://tangled.org/zat.dev/zat/archive/v0.4.0.tar.gz", - .hash = "zat-0.4.0-5PuC7iFkDAD8OoBB-bDlexUoSI24rveRWfINE90Oxebs", + .url = "https://tangled.org/zat.dev/zat/archive/v0.6.0-alpha.2.tar.gz", + .hash = "zat-0.6.0-alpha.2-5PuC7knODABZ9UsoxRvI94zsmsE3XQPWdBIBn6zDduW2", }, .websocket = .{ - .url = "https://github.com/zzstoatzz/websocket.zig/archive/73429df.tar.gz", - .hash = "websocket-0.1.10-ZPISdUowBQAXIl_X9VpZozIxBLVf_o5wYjOMMdobrXF5", + .url = "https://tangled.org/zzstoatzz.io/websocket.zig/archive/v0.2.0-alpha.2.tar.gz", + .hash = "websocket-0.2.0-alpha.2-ZPISdd9ZBQAlweSmUfRIOl4m_BrcbeAmDT_5H6kA0aXC", }, .zstd = .{ .url = "https://github.com/facebook/zstd/releases/download/v1.5.7/zstd-1.5.7.tar.gz", @@ -21,12 +21,12 @@ .hash = "N-V-__8AABgCUwB4fhP90zPC1u1U84-BEirnJ0dxOPi5ZM7Z", }, .rocksdb = .{ - .url = "https://github.com/zzstoatzz/rocksdb-zig/archive/9d2ebd8.tar.gz", - .hash = "rocksdb-9.7.4-z_CUTu7bAADWOmsU-JQTz1V91P7c6I8BzMS2_fkGtyLu", + .url = "https://github.com/zzstoatzz/rocksdb-zig/archive/30cf93da3ea6c0814ac6c490c8e3bed9214c5377.tar.gz", + .hash = "rocksdb-9.7.4-z_CUTvbbAAD2p1ScTnhTREoFY8lzJh9-EMLnYjZA1MZd", }, .otel = .{ - .url = "git+https://github.com/zzstoatzz/otel-zig.git#158b32d87762be4f6b65a7d723d08b20d74d1294", - .hash = "otel-0.0.1-9q9ZWueJFADcwU-zVrOU0CdC9j7Zukfhse1Q5ALaSi3t", + .url = "git+https://github.com/zzstoatzz/otel-zig.git#0339b99d358a230e2468b4a28e5e9d82f6f30f6c", + .hash = "otel-0.0.1-9q9ZWu27FAAQFjhWND9nXirEt5_mdQiVDEB2-e7oQClw", }, }, .paths = .{ diff --git a/src/bench.zig b/src/bench.zig index 648bdd6..d476788 100644 --- a/src/bench.zig +++ b/src/bench.zig @@ -295,7 +295,7 @@ fn benchWire(alloc: Allocator, io: Io, comptime v2: bool) !Result { var buf: std.Io.Writer.Allocating = .init(alloc); defer buf.deinit(); if (v2) { - try wire.encodeCommitV2(alloc, &buf.writer, op, n, payload); + try wire.encodeCommitV2(alloc, &buf.writer, op, n); } else { try wire.encodeCommit(alloc, &buf.writer, op); } diff --git a/src/internal/compact/steady.zig b/src/internal/compact/steady.zig index 97cc1f6..b670a4d 100644 --- a/src/internal/compact/steady.zig +++ b/src/internal/compact/steady.zig @@ -207,9 +207,12 @@ pub const Compactor = struct { ); defer trace.deinit(); defer trace.succeed(); - errdefer |trace_err| trace.fail(trace_err); + return trace.recordError(self.runObservedPass(trace.context())); + } + + fn runObservedPass(self: *Compactor, trace_context: observability.Context) !pass.Stats { const started = Io.Timestamp.now(self.io, .awake); - const result = self.runOnceInner(trace.context()) catch |err| { + const result = self.runOnceInner(trace_context) catch |err| { if (self.stats) |s| s.observeCompactionPass(.error_result, elapsedUs(self.io, started)); self.publishLastPassNow(); return @errorCast(err); diff --git a/src/internal/ingest/backfill/engine.zig b/src/internal/ingest/backfill/engine.zig index d625e51..5db3667 100644 --- a/src/internal/ingest/backfill/engine.zig +++ b/src/internal/ingest/backfill/engine.zig @@ -697,7 +697,17 @@ pub fn runTraced( var trace = observability.Span.start(allocator, "ingest/backfill", "Run", trace_context, &.{}); defer trace.deinit(); defer trace.succeed(); - errdefer |trace_err| trace.fail(trace_err); + return trace.recordError(runBackfill(allocator, io, config, archive, store, trace.context())); +} + +fn runBackfill( + allocator: Allocator, + io: Io, + config: Config, + archive: *archive_mod.Archive, + store: *repo_store.Store, + trace_context: observability.Context, +) !Stats { if (config.selected_repos.len > 0 and config.max_repos > 0) return error.ConflictingBackfillSelection; if (config.selected_repos.len > 0 and config.plc_url == null) return error.SelectedBackfillRequiresIdentityResolver; var scratch_buf: [Io.Dir.max_path_bytes]u8 = undefined; @@ -712,7 +722,7 @@ pub fn runTraced( // parallel whole-network engine. var work_config = config; if (config.selected_repos.len > 0) work_config.workers = 1; - var work = try Work.init(allocator, io, work_config, archive, trace.context()); + var work = try Work.init(allocator, io, work_config, archive, trace_context); defer work.deinit(); try work.start(); var work_completed = false; @@ -1253,84 +1263,115 @@ fn downloadOne( ); defer trace.deinit(); defer trace.succeed(); - errdefer |err| trace.fail(err); - const handle_started = Io.Timestamp.now(io, .awake); - var local: Stats = .{}; - const now_us: i64 = Io.Timestamp.now(io, .real).toMicroseconds(); - const Sink = struct { - archive: *archive_mod.Archive, - now_us: i64, - batch: [1024]segment.Event = undefined, - batch_len: usize = 0, - rows: u64 = 0, - last_seq: u64 = 0, - trace_context: observability.Context, - - fn flush(self: *@This()) anyerror!void { - if (self.batch_len == 0) return; - const first_seq = (try self.archive.appendBatchTraced( - self.batch[0..self.batch_len], - self.now_us, - self.trace_context, - )).?; - self.last_seq = first_seq + self.batch_len - 1; - self.batch_len = 0; - } - - fn emit(self: *@This(), event: segment.Event) anyerror!void { - self.batch[self.batch_len] = event; - self.batch_len += 1; - self.rows += 1; - if (self.batch_len == self.batch.len) try self.flush(); - } - }; - var sink: Sink = .{ - .archive = archive, - .now_us = now_us, - .trace_context = trace.context(), - }; - const head_rev = repos.emitPrepared(alloc, &prepared.repo, did, now_us, .create, &local.drops, &sink, Sink.emit) catch |err| switch (err) { - error.OutOfMemory, error.ArchiveAppendFailed => return @errorCast(err), - error.MissingBlock, error.IncompleteCar => { - if (transient_attempt >= transient_retries) return failureResult(result_allocator, did, &final_host, attempts, started_us, .car, @errorName(err)); - try sleepBackoff(io, transient_attempt, 30); - transient_attempt += 1; - continue; - }, - else => return failureResult(result_allocator, did, &final_host, attempts, started_us, .car, @errorName(err)), - }; - sink.flush() catch |err| switch (err) { - error.OutOfMemory, error.ArchiveAppendFailed => return @errorCast(err), - else => return error.ArchiveAppendFailed, - }; - if (instrumentation) |m| m.observeBackfillHandleRepo(@intCast(@max( - handle_started.durationTo(Io.Timestamp.now(io, .awake)).toMicroseconds(), - 0, - ))); - const owned_rev = try result_allocator.dupe(u8, head_rev); - errdefer result_allocator.free(owned_rev); - const owned_latest_rev = try result_allocator.dupe(u8, head_rev); - errdefer result_allocator.free(owned_latest_rev); - const owned_host = try result_allocator.dupe(u8, final_host.slice()); - errdefer if (owned_host.len > 0) result_allocator.free(owned_host); - return .{ - .did = did, - .status = .complete, - .rev = owned_rev, - .latest_rev = owned_latest_rev, - .host = owned_host, - .attempts = attempts, - .started_us = started_us, - .rows = sink.rows, - .last_seq = sink.last_seq, - .completion_requires_durability = true, - .drops = local.drops, - }; + if (try trace.recordError(handlePrepared( + result_allocator, + alloc, + io, + did, + prepared, + archive, + instrumentation, + &final_host, + attempts, + started_us, + &transient_attempt, + trace.context(), + ))) |result| return result; }, } } } +/// One downloaded repository's emit and flush. Null asks the caller to retry +/// the download. +fn handlePrepared( + result_allocator: Allocator, + alloc: Allocator, + io: Io, + did: []const u8, + prepared: *prepared_fetch.Prepared, + archive: *archive_mod.Archive, + instrumentation: ?*metrics_mod.Stats, + final_host: *const fetched_car.FinalHost, + attempts: u32, + started_us: i64, + transient_attempt: *u8, + trace_context: observability.Context, +) !?Result { + const handle_started = Io.Timestamp.now(io, .awake); + var local: Stats = .{}; + const now_us: i64 = Io.Timestamp.now(io, .real).toMicroseconds(); + const Sink = struct { + archive: *archive_mod.Archive, + now_us: i64, + batch: [1024]segment.Event = undefined, + batch_len: usize = 0, + rows: u64 = 0, + last_seq: u64 = 0, + trace_context: observability.Context, + + fn flush(self: *@This()) anyerror!void { + if (self.batch_len == 0) return; + const first_seq = (try self.archive.appendBatchTraced( + self.batch[0..self.batch_len], + self.now_us, + self.trace_context, + )).?; + self.last_seq = first_seq + self.batch_len - 1; + self.batch_len = 0; + } + + fn emit(self: *@This(), event: segment.Event) anyerror!void { + self.batch[self.batch_len] = event; + self.batch_len += 1; + self.rows += 1; + if (self.batch_len == self.batch.len) try self.flush(); + } + }; + var sink: Sink = .{ + .archive = archive, + .now_us = now_us, + .trace_context = trace_context, + }; + const head_rev = repos.emitPrepared(alloc, &prepared.repo, did, now_us, .create, &local.drops, &sink, Sink.emit) catch |err| switch (err) { + error.OutOfMemory, error.ArchiveAppendFailed => return @errorCast(err), + error.MissingBlock, error.IncompleteCar => { + if (transient_attempt.* >= transient_retries) return try failureResult(result_allocator, did, final_host, attempts, started_us, .car, @errorName(err)); + try sleepBackoff(io, transient_attempt.*, 30); + transient_attempt.* += 1; + return null; + }, + else => return try failureResult(result_allocator, did, final_host, attempts, started_us, .car, @errorName(err)), + }; + sink.flush() catch |err| switch (err) { + error.OutOfMemory, error.ArchiveAppendFailed => return @errorCast(err), + else => return error.ArchiveAppendFailed, + }; + if (instrumentation) |m| m.observeBackfillHandleRepo(@intCast(@max( + handle_started.durationTo(Io.Timestamp.now(io, .awake)).toMicroseconds(), + 0, + ))); + const owned_rev = try result_allocator.dupe(u8, head_rev); + errdefer result_allocator.free(owned_rev); + const owned_latest_rev = try result_allocator.dupe(u8, head_rev); + errdefer result_allocator.free(owned_latest_rev); + const owned_host = try result_allocator.dupe(u8, final_host.slice()); + errdefer if (owned_host.len > 0) result_allocator.free(owned_host); + return .{ + .did = did, + .status = .complete, + .rev = owned_rev, + .latest_rev = owned_latest_rev, + .host = owned_host, + .attempts = attempts, + .started_us = started_us, + .rows = sink.rows, + .last_seq = sink.last_seq, + .completion_requires_durability = true, + .drops = local.drops, + }; +} + fn successResult(allocator: Allocator, did: []const u8, status: repo_store.Status, host: *const fetched_car.FinalHost, attempts: u32, started_us: i64) !Result { return .{ .did = did, @@ -1799,7 +1840,7 @@ test "successful real CAR handling records only the handler boundary" { var tree = zat.mst.Mst.init(fixture_alloc); try tree.put("app.bsky.feed.post/3k2abcdefghij", record_cid); const data_cid = try tree.rootCid(); - const keypair = try zat.Keypair.fromSecretKey(.p256, .{17} ** 32); + const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(17)); const embedded_did = try keypair.did(fixture_alloc); const signed = try zat.signCommit(fixture_alloc, .{ .did = embedded_did, diff --git a/src/internal/ingest/backfill/fetched_car.zig b/src/internal/ingest/backfill/fetched_car.zig index a6914af..44ebd51 100644 --- a/src/internal/ingest/backfill/fetched_car.zig +++ b/src/internal/ingest/backfill/fetched_car.zig @@ -333,15 +333,15 @@ fn fetchWithTimeout( pub fn captureFinalHost(uri: std.Uri, out: *FinalHost) !void { var host_buf: [std.Io.net.HostName.max_len]u8 = undefined; - const host = try uri.getHost(&host_buf); + const host = try (uri.host orelse return error.UriMissingHost).toRaw(&host_buf); const default_port: ?u16 = if (std.ascii.eqlIgnoreCase(uri.scheme, "https")) 443 else if (std.ascii.eqlIgnoreCase(uri.scheme, "http")) 80 else null; const include_port = uri.port != null and uri.port != default_port; const normalized = if (include_port) - try std.fmt.bufPrint(&out.buf, "{s}:{d}", .{ host.bytes, uri.port.? }) + try std.fmt.bufPrint(&out.buf, "{s}:{d}", .{ host, uri.port.? }) else blk: { - if (host.bytes.len > out.buf.len) return error.NoSpaceLeft; - for (host.bytes, out.buf[0..host.bytes.len]) |c, *dst| dst.* = std.ascii.toLower(c); - break :blk out.buf[0..host.bytes.len]; + if (host.len > out.buf.len) return error.NoSpaceLeft; + for (host, out.buf[0..host.len]) |c, *dst| dst.* = std.ascii.toLower(c); + break :blk out.buf[0..host.len]; }; if (include_port) { for (normalized) |*c| c.* = std.ascii.toLower(c.*); diff --git a/src/internal/ingest/backfill/host_store.zig b/src/internal/ingest/backfill/host_store.zig index 0ac28a1..6a05119 100644 --- a/src/internal/ingest/backfill/host_store.zig +++ b/src/internal/ingest/backfill/host_store.zig @@ -48,15 +48,15 @@ pub const Status = struct { host: []const u8 = "", total: u64 = 0, active: u64 = 0, - by_status: [std.meta.fields(RepoStatus).len]u64 = @splat(0), + by_status: [@typeInfo(RepoStatus).@"enum".field_names.len]u64 = @splat(0), last_attempt_us: i64 = 0, latest_error: []const u8 = "", latest_error_class: ErrorClass = .unknown, - error_counts: [std.meta.fields(ErrorClass).len]u64 = @splat(0), + error_counts: [@typeInfo(ErrorClass).@"enum".field_names.len]u64 = @splat(0), recent_errors: []ErrorSample = &.{}, pub fn count(self: Status, status: RepoStatus) u64 { - return self.by_status[@intFromEnum(status)]; + return self.by_status[@backingInt(status)]; } pub fn deinit(self: *Status, allocator: Allocator) void { @@ -92,21 +92,21 @@ pub fn applyTransition(status: *Status, first_in_bucket: bool, old_active: bool, if (first_in_bucket) { status.total += 1; if (new_active) status.active += 1; - status.by_status[@intFromEnum(next)] += 1; + status.by_status[@backingInt(next)] += 1; return; } if (old_active != new_active) { if (new_active) status.active += 1 else status.active -|= 1; } if (old == next) return; - status.by_status[@intFromEnum(old)] -|= 1; - status.by_status[@intFromEnum(next)] += 1; + status.by_status[@backingInt(old)] -|= 1; + status.by_status[@backingInt(next)] += 1; } pub fn remove(status: *Status, active: bool, lifecycle: RepoStatus) void { status.total -|= 1; if (active) status.active -|= 1; - status.by_status[@intFromEnum(lifecycle)] -|= 1; + status.by_status[@backingInt(lifecycle)] -|= 1; } pub fn recordError(status: *Status, allocator: Allocator, did: []const u8, attempted_us: i64, class: ErrorClass, raw_message: []const u8) !void { @@ -124,7 +124,7 @@ pub fn recordError(status: *Status, allocator: Allocator, did: []const u8, attem if (status.latest_error.len > 0) allocator.free(@constCast(status.latest_error)); status.latest_error = next_latest; status.latest_error_class = class; - status.error_counts[@intFromEnum(class)] += 1; + status.error_counts[@backingInt(class)] += 1; next[0] = .{ .did = next_did, .attempted_us = attempted_us, @@ -157,7 +157,7 @@ pub fn encode(allocator: Allocator, status: Status) ![]u8 { writeInt(u64, out, &at, status.active); for (status.by_status) |value| writeInt(u64, out, &at, value); writeInt(i64, out, &at, status.last_attempt_us); - out[at] = @intFromEnum(status.latest_error_class); + out[at] = @backingInt(status.latest_error_class); at += 1; for (status.error_counts) |value| writeInt(u64, out, &at, value); writeInt(u16, out, &at, @intCast(status.latest_error.len)); @@ -167,7 +167,7 @@ pub fn encode(allocator: Allocator, status: Status) ![]u8 { at += status.latest_error.len; for (status.recent_errors) |sample| { writeInt(i64, out, &at, sample.attempted_us); - out[at] = @intFromEnum(sample.class); + out[at] = @backingInt(sample.class); at += 1; writeInt(u16, out, &at, @intCast(sample.did.len)); writeInt(u16, out, &at, @intCast(sample.message.len)); @@ -190,8 +190,8 @@ pub fn decode(allocator: Allocator, host: []const u8, bytes: []const u8) !Status out.active = try readInt(u64, bytes, &at); for (&out.by_status) |*value| value.* = try readInt(u64, bytes, &at); out.last_attempt_us = try readInt(i64, bytes, &at); - if (at >= bytes.len or bytes[at] >= std.meta.fields(ErrorClass).len) return error.InvalidHostStatus; - out.latest_error_class = @enumFromInt(bytes[at]); + if (at >= bytes.len or bytes[at] >= @typeInfo(ErrorClass).@"enum".field_names.len) return error.InvalidHostStatus; + out.latest_error_class = @fromBackingInt(@intCast(bytes[at])); at += 1; for (&out.error_counts) |*value| value.* = try readInt(u64, bytes, &at); const latest_len = try readInt(u16, bytes, &at); @@ -215,8 +215,8 @@ pub fn decode(allocator: Allocator, host: []const u8, bytes: []const u8) !Status } while (initialized < sample_count) : (initialized += 1) { const attempted_us = try readInt(i64, bytes, &at); - if (at >= bytes.len or bytes[at] >= std.meta.fields(ErrorClass).len) return error.InvalidHostStatus; - const class: ErrorClass = @enumFromInt(bytes[at]); + if (at >= bytes.len or bytes[at] >= @typeInfo(ErrorClass).@"enum".field_names.len) return error.InvalidHostStatus; + const class: ErrorClass = @fromBackingInt(@intCast(bytes[at])); at += 1; const did_len = try readInt(u16, bytes, &at); const message_len = try readInt(u16, bytes, &at); @@ -261,7 +261,7 @@ test "host status round trips bounded cumulative diagnostics" { try testing.expectEqual(@as(u64, 1), decoded.total); try testing.expectEqual(@as(u64, 1), decoded.active); try testing.expectEqual(@as(u64, 1), decoded.count(.failed)); - try testing.expectEqual(@as(u64, 7), decoded.error_counts[@intFromEnum(ErrorClass.http_5xx)]); + try testing.expectEqual(@as(u64, 7), decoded.error_counts[@backingInt(ErrorClass.http_5xx)]); try testing.expectEqual(@as(usize, 5), decoded.recent_errors.len); try testing.expectEqual(@as(i64, 6), decoded.recent_errors[0].attempted_us); try testing.expectEqual(@as(i64, 2), decoded.recent_errors[4].attempted_us); diff --git a/src/internal/ingest/backfill/lifecycle.zig b/src/internal/ingest/backfill/lifecycle.zig index 2102f76..9b2efc6 100644 --- a/src/internal/ingest/backfill/lifecycle.zig +++ b/src/internal/ingest/backfill/lifecycle.zig @@ -136,8 +136,7 @@ pub fn runTraced( ); defer bootstrap_trace.deinit(); defer bootstrap_trace.succeed(); - errdefer |err| bootstrap_trace.fail(err); - try runBootstrap( + try bootstrap_trace.recordError(runBootstrap( allocator, io, opts, @@ -149,7 +148,7 @@ pub fn runTraced( stats, &repo_st, bootstrap_trace.context(), - ); + )); phase = .merging; bootstrap_trace.succeed(); } @@ -163,7 +162,36 @@ pub fn runTraced( ); defer merge_trace.deinit(); defer merge_trace.succeed(); - errdefer |err| merge_trace.fail(err); + try merge_trace.recordError(runMerge( + allocator, + io, + opts, + live_data_dir, + dst, + meta, + &phase_store, + &wm, + meta_dir, + &repo_st, + stats, + &merge_trace, + )); +} + +fn runMerge( + allocator: Allocator, + io: Io, + opts: Options, + live_data_dir: []const u8, + dst: *archive_mod.Archive, + meta: *meta_store.Store, + phase_store: *phase_mod.Store, + wm: *watermark_mod.Store, + meta_dir: Io.Dir, + repo_st: *repo_store_mod.Store, + stats: *metrics.Stats, + merge_trace: *observability.Span, +) !void { const merge_started = Io.Timestamp.now(io, .awake); var merge_observed = false; defer if (!merge_observed) stats.observeOrchestratorState(.merge, elapsedUs(io, merge_started)); @@ -179,7 +207,7 @@ pub fn runTraced( var bf_dir = try Io.Dir.cwd().openDir(io, bf_path, .{}); defer bf_dir.close(io); - const m = try merge.runTraced(allocator, io, live_data_dir, dst, &repo_st, stats, merge_trace.context()); + const m = try merge.runTraced(allocator, io, live_data_dir, dst, repo_st, stats, merge_trace.context()); log.info("merge: {d} kept, {d} dropped across {d} sources", .{ m.kept, m.dropped, m.sources }); try dst.flushBlock(); crashpoint.hit(.after_merge_drain_before_pending); @@ -207,25 +235,8 @@ pub fn runTraced( ); defer compaction_trace.deinit(); defer compaction_trace.succeed(); - errdefer |err| compaction_trace.fail(err); - if (opts.compaction_enabled) { - const compaction_started = Io.Timestamp.now(io, .awake); - var seg_path_buf: [256]u8 = undefined; - const seg_path = try std.fmt.bufPrint(&seg_path_buf, "{s}/segments", .{opts.data_dir}); - var seg_dir = try Io.Dir.cwd().openDir(io, seg_path, .{ .iterate = true }); - defer seg_dir.close(io); - const cs = compact_pass.run(allocator, io, seg_dir, &wm, null, .{ - .instrumentation = stats, - .measure_pass = false, - }) catch |err| { - stats.observeCompactionPass(.error_result, elapsedUs(io, compaction_started)); - return @errorCast(err); - }; - stats.observeCompactionPass(.ok, elapsedUs(io, compaction_started)); - log.info("merge-tail compaction: watermark {d}, {d} chunks, {d}/{d} segments rewritten, {d} rows dropped", .{ - cs.target_watermark, cs.chunks, cs.segments_rewritten, cs.segments_examined, cs.rows_dropped, - }); - } + if (opts.compaction_enabled) + try compaction_trace.recordError(compactMergeTail(allocator, io, opts.data_dir, wm, stats)); compaction_trace.succeed(); // Merge-tail rewriting is deliberately manifest-oblivious. Before // serving can ungate, every resident header/index must match disk. @@ -244,7 +255,7 @@ pub fn runTraced( // outage here fails the merge (upstream couples them the same // way) — re-entry is idempotent. if (!opts.skip_merge_discovery) { - try runDiscovery(allocator, io, opts.relay_http, &repo_st, stats, merge_trace.context()); + try runDiscovery(allocator, io, opts.relay_http, repo_st, stats, merge_trace.context()); } else { log.info("merge discovery: skipped by configuration", .{}); } @@ -310,6 +321,31 @@ pub fn runTraced( crashpoint.hit(.after_steady_phase_before_steady_run); } +fn compactMergeTail( + allocator: Allocator, + io: Io, + data_dir: []const u8, + wm: *watermark_mod.Store, + stats: *metrics.Stats, +) !void { + const compaction_started = Io.Timestamp.now(io, .awake); + var seg_path_buf: [256]u8 = undefined; + const seg_path = try std.fmt.bufPrint(&seg_path_buf, "{s}/segments", .{data_dir}); + var seg_dir = try Io.Dir.cwd().openDir(io, seg_path, .{ .iterate = true }); + defer seg_dir.close(io); + const cs = compact_pass.run(allocator, io, seg_dir, wm, null, .{ + .instrumentation = stats, + .measure_pass = false, + }) catch |err| { + stats.observeCompactionPass(.error_result, elapsedUs(io, compaction_started)); + return @errorCast(err); + }; + stats.observeCompactionPass(.ok, elapsedUs(io, compaction_started)); + log.info("merge-tail compaction: watermark {d}, {d} chunks, {d}/{d} segments rewritten, {d} rows dropped", .{ + cs.target_watermark, cs.chunks, cs.segments_rewritten, cs.segments_examined, cs.rows_dropped, + }); +} + /// Upstream's RunPendingRepoRetryPass: the same bounded retry runner with /// eligibleStatus flipped to pending. Reusing it rather than hand-rolling a /// second pass is what makes the durable bookkeeping match — worker pool, @@ -357,7 +393,16 @@ fn runDiscovery( ); defer trace.deinit(); defer trace.succeed(); - errdefer |trace_err| trace.fail(trace_err); + try trace.recordError(discover(allocator, io, relay_http, repo_st, stats)); +} + +fn discover( + allocator: Allocator, + io: Io, + relay_http: []const u8, + repo_st: *repo_store_mod.Store, + stats: *metrics.Stats, +) !void { const start = (try engine.loadBootstrapLastListReposCursor(repo_st.meta, allocator)) orelse { log.info("merge discovery: no bootstrap cursor, skipping", .{}); return; @@ -578,7 +623,16 @@ fn runBootstrap( ); defer finish_trace.deinit(); defer finish_trace.succeed(); - errdefer |err| finish_trace.fail(err); + try finish_trace.recordError(finishBootstrap(io, stats, &pipe, &live_archive)); + live_sealed = true; +} + +fn finishBootstrap( + io: Io, + stats: *metrics.Stats, + pipe: *pipeline_mod.Pipeline, + live_archive: *archive_mod.Archive, +) !void { const drain_started = Io.Timestamp.now(io, .awake); stats.observeOrchestratorState(.drain_bootstrap, elapsedUs(io, drain_started)); // cancel joins the websocket reader, so no new submissions can race this @@ -597,7 +651,6 @@ fn runBootstrap( return @errorCast(err); }; stats.observeOrchestratorState(.seal_bootstrap, elapsedUs(io, seal_started)); - live_sealed = true; } fn orchestratorPhase(phase: phase_mod.Phase) metrics.OrchestratorPhase { diff --git a/src/internal/ingest/backfill/merge.zig b/src/internal/ingest/backfill/merge.zig index cdf78f4..b7c35bd 100644 --- a/src/internal/ingest/backfill/merge.zig +++ b/src/internal/ingest/backfill/merge.zig @@ -98,7 +98,18 @@ pub fn runTraced( var trace = observability.Span.start(allocator, "ingest/orchestrator", "run", trace_context, &.{}); defer trace.deinit(); defer trace.succeed(); - errdefer |trace_err| trace.fail(trace_err); + return trace.recordError(mergeSources(allocator, io, live_data_dir, dst, store, metric_stats, trace.context())); +} + +fn mergeSources( + allocator: Allocator, + io: Io, + live_data_dir: []const u8, + dst: *archive_mod.Archive, + store: *repo_store.Store, + metric_stats: *metrics.Stats, + trace_context: observability.Context, +) !Stats { // seal guard: a crash between live-consumer teardown and seal leaves // an active trailing source. Archive recovery seals torn actives; // sealAndClose then deletes the fresh empty one it opened. @@ -137,40 +148,65 @@ pub fn runTraced( allocator, "ingest/orchestrator", "processSourceSegment", - trace.context(), + trace_context, &.{}, ); defer source_trace.deinit(); defer source_trace.succeed(); - errdefer |err| source_trace.fail(err); - const revs = try drainOne( + try source_trace.recordError(mergeSource( source_arena.allocator(), io, src_dir, idx, dst, + store, &cache, &stats, metric_stats, - source_trace.context(), - ); - // durability ordering: dst rows fsynced before the cursor advances - try dst.flushBlockTraced(source_trace.context()); - source_trace.succeed(); - crashpoint.hit(.after_merge_dst_flush_before_source_commit); - const updates = try source_arena.allocator().alloc(repo_store.RevUpdate, revs.count()); - var update_it = revs.iterator(); - var update_index: usize = 0; - while (update_it.next()) |entry| : (update_index += 1) - updates[update_index] = .{ .did = entry.key_ptr.*, .rev = entry.value_ptr.* }; - try store.commitMergeSource(cursor_key, idx + 1, updates, Io.Timestamp.now(io, .real).toMicroseconds()); - stats.sources += 1; - _ = metric_stats.orchestrator_merge_segments_consumed_total.fetchAdd(1, .monotonic); - _ = metric_stats.orchestrator_merge_repo_revs_updated_total.fetchAdd(updates.len, .monotonic); + &source_trace, + )); } return stats; } +fn mergeSource( + arena: Allocator, + io: Io, + src_dir: Io.Dir, + idx: u64, + dst: *archive_mod.Archive, + store: *repo_store.Store, + cache: *RepoStatusCache, + stats: *Stats, + metric_stats: *metrics.Stats, + source_trace: *observability.Span, +) !void { + const revs = try drainOne( + arena, + io, + src_dir, + idx, + dst, + cache, + stats, + metric_stats, + source_trace.context(), + ); + // durability ordering: dst rows fsynced before the cursor advances + try dst.flushBlockTraced(source_trace.context()); + source_trace.succeed(); + crashpoint.hit(.after_merge_dst_flush_before_source_commit); + const updates = try arena.alloc(repo_store.RevUpdate, revs.count()); + var update_it = revs.iterator(); + var update_index: usize = 0; + while (update_it.next()) |entry| : (update_index += 1) + updates[update_index] = .{ .did = entry.key_ptr.*, .rev = entry.value_ptr.* }; + try store.commitMergeSource(cursor_key, idx + 1, updates, Io.Timestamp.now(io, .real).toMicroseconds()); + stats.sources += 1; + _ = metric_stats.orchestrator_merge_segments_consumed_total.fetchAdd(1, .monotonic); + _ = metric_stats.orchestrator_merge_repo_revs_updated_total.fetchAdd(updates.len, .monotonic); +} + fn drainOne( allocator: Allocator, io: Io, diff --git a/src/internal/ingest/backfill/prepared_fetch.zig b/src/internal/ingest/backfill/prepared_fetch.zig index 3c2b519..684b504 100644 --- a/src/internal/ingest/backfill/prepared_fetch.zig +++ b/src/internal/ingest/backfill/prepared_fetch.zig @@ -153,7 +153,7 @@ test "prepared repo emission owns the CAR after its source is destroyed" { var tree = zat.mst.Mst.init(source); try tree.put("app.bsky.feed.post/3k2abcdefghij", record_cid); const data_cid = try tree.rootCid(); - const keypair = try zat.Keypair.fromSecretKey(.p256, .{29} ** 32); + const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(29)); const did = try keypair.did(source); const signed = try zat.signCommit(source, .{ .did = did, diff --git a/src/internal/ingest/backfill/repo_store.zig b/src/internal/ingest/backfill/repo_store.zig index f0936c6..905b6bf 100644 --- a/src/internal/ingest/backfill/repo_store.zig +++ b/src/internal/ingest/backfill/repo_store.zig @@ -49,12 +49,13 @@ pub fn hostFromPDS(allocator: Allocator, endpoint: []const u8) ![]u8 { if (!std.ascii.eqlIgnoreCase(uri.scheme, "http") and !std.ascii.eqlIgnoreCase(uri.scheme, "https")) return allocator.dupe(u8, "invalid-pds"); var host_buf: [std.Io.net.HostName.max_len]u8 = undefined; - const host = uri.getHost(&host_buf) catch return allocator.dupe(u8, "invalid-pds"); - if (host.bytes.len == 0) return allocator.dupe(u8, "invalid-pds"); - const bare_host = if (host.bytes.len >= 2 and host.bytes[0] == '[' and host.bytes[host.bytes.len - 1] == ']') - host.bytes[1 .. host.bytes.len - 1] + const host_component = uri.host orelse return allocator.dupe(u8, "invalid-pds"); + const host = host_component.toRaw(&host_buf) catch return allocator.dupe(u8, "invalid-pds"); + if (host.len == 0) return allocator.dupe(u8, "invalid-pds"); + const bare_host = if (host.len >= 2 and host[0] == '[' and host[host.len - 1] == ']') + host[1 .. host.len - 1] else - host.bytes; + host; const default_port: ?u16 = if (std.ascii.eqlIgnoreCase(uri.scheme, "https")) 443 else 80; const raw = if (uri.port != null and uri.port != default_port) if (std.mem.indexOfScalar(u8, bare_host, ':') != null) @@ -235,9 +236,9 @@ pub const Store = struct { if (old_bytes) |bytes| { old = try decode(bytes); had_old = true; - if (counts[@intFromEnum(old.status)] > 0) counts[@intFromEnum(old.status)] -= 1; + if (counts[@backingInt(old.status)] > 0) counts[@backingInt(old.status)] -= 1; } - counts[@intFromEnum(transition.state.status)] += 1; + counts[@backingInt(transition.state.status)] += 1; const old_host = if (had_old) try host_store.normalizeAuthority(self.allocator, old.host) else ""; defer if (old_host.len > 0) self.allocator.free(old_host); @@ -414,7 +415,7 @@ pub const Store = struct { pub fn countByStatus(self: *Store, status: Status) u32 { const counts = self.loadCounts() catch return 0; - return @intCast(@min(counts[@intFromEnum(status)], std.math.maxInt(u32))); + return @intCast(@min(counts[@backingInt(status)], std.math.maxInt(u32))); } fn rebuildCounts(self: *Store) !void { @@ -426,7 +427,7 @@ pub const Store = struct { while (try it.next(&err_data)) |entry| { if (!std.mem.startsWith(u8, entry[0].data, repo_prefix)) break; const state = try decode(entry[1].data); - counts[@intFromEnum(state.status)] += 1; + counts[@backingInt(state.status)] += 1; } var bytes: [@sizeOf(Counts)]u8 = undefined; encodeCounts(&bytes, counts); @@ -534,7 +535,7 @@ pub const Store = struct { } }; -const Counts = [std.meta.fields(Status).len]u64; +const Counts = [@typeInfo(Status).@"enum".field_names.len]u64; fn encodeCounts(out: *[@sizeOf(Counts)]u8, counts: Counts) void { for (counts, 0..) |n, i| std.mem.writeInt(u64, out[i * 8 ..][0..8], n, .little); @@ -574,7 +575,7 @@ fn encode(allocator: Allocator, state: RepoState) ![]u8 { return error.RepoStateFieldTooLong; const out = try allocator.alloc(u8, header_len + state.rev.len + state.latest_rev.len + state.host.len + state.handle.len + state.pds.len + state.last_error.len); out[0] = encoding_version; - out[1] = @intFromEnum(state.status); + out[1] = @backingInt(state.status); out[2] = @intFromBool(state.active); std.mem.writeInt(u32, out[3..7], state.retry_count, .little); std.mem.writeInt(i64, out[7..15], state.next_attempt_us, .little); @@ -591,7 +592,7 @@ fn encode(allocator: Allocator, state: RepoState) ![]u8 { std.mem.writeInt(u16, out[77..79], @intCast(state.handle.len), .little); std.mem.writeInt(u16, out[79..81], @intCast(state.pds.len), .little); std.mem.writeInt(u16, out[81..83], @intCast(state.last_error.len), .little); - out[83] = @intFromEnum(state.error_class); + out[83] = @backingInt(state.error_class); var offset: usize = header_len; @memcpy(out[offset..][0..state.rev.len], state.rev); offset += state.rev.len; @@ -609,9 +610,9 @@ fn encode(allocator: Allocator, state: RepoState) ![]u8 { fn decode(bytes: []const u8) !RepoState { if (bytes.len >= legacy_header_len and bytes[0] == 2) { - if (bytes[1] >= std.meta.fields(Status).len) return error.InvalidRepoState; + if (bytes[1] >= @typeInfo(Status).@"enum".field_names.len) return error.InvalidRepoState; return .{ - .status = @enumFromInt(bytes[1]), + .status = @fromBackingInt(@intCast(bytes[1])), .active = true, .retry_count = std.mem.readInt(u32, bytes[2..6], .little), .next_attempt_us = std.mem.readInt(i64, bytes[6..14], .little), @@ -620,9 +621,9 @@ fn decode(bytes: []const u8) !RepoState { }; } if (bytes.len >= 15 and bytes[0] == 3) { - if (bytes[1] >= std.meta.fields(Status).len) return error.InvalidRepoState; + if (bytes[1] >= @typeInfo(Status).@"enum".field_names.len) return error.InvalidRepoState; return .{ - .status = @enumFromInt(bytes[1]), + .status = @fromBackingInt(@intCast(bytes[1])), .active = bytes[2] != 0, .retry_count = std.mem.readInt(u32, bytes[3..7], .little), .next_attempt_us = std.mem.readInt(i64, bytes[7..15], .little), @@ -631,8 +632,8 @@ fn decode(bytes: []const u8) !RepoState { }; } if (bytes.len >= v4_header_len and bytes[0] == 4) { - if (bytes[1] >= std.meta.fields(Status).len) return error.InvalidRepoState; - if (bytes[51] >= std.meta.fields(ErrorClass).len) return error.InvalidRepoState; + if (bytes[1] >= @typeInfo(Status).@"enum".field_names.len) return error.InvalidRepoState; + if (bytes[51] >= @typeInfo(ErrorClass).@"enum".field_names.len) return error.InvalidRepoState; const rev_len = std.mem.readInt(u32, bytes[43..47], .little); const host_len = std.mem.readInt(u16, bytes[47..49], .little); const error_len = std.mem.readInt(u16, bytes[49..51], .little); @@ -641,13 +642,13 @@ fn decode(bytes: []const u8) !RepoState { const rev_end = v4_header_len + rev_len; const host_end = rev_end + host_len; return .{ - .status = @enumFromInt(bytes[1]), + .status = @fromBackingInt(@intCast(bytes[1])), .active = bytes[2] != 0, .rev = bytes[v4_header_len..rev_end], .latest_rev = bytes[v4_header_len..rev_end], .host = bytes[rev_end..host_end], .last_error = bytes[host_end..total], - .error_class = @enumFromInt(bytes[51]), + .error_class = @fromBackingInt(@intCast(bytes[51])), .retry_count = std.mem.readInt(u32, bytes[3..7], .little), .next_attempt_us = std.mem.readInt(i64, bytes[7..15], .little), .attempts = std.mem.readInt(u32, bytes[15..19], .little), @@ -657,8 +658,8 @@ fn decode(bytes: []const u8) !RepoState { }; } if (bytes.len < header_len or bytes[0] != encoding_version) return error.InvalidRepoState; - if (bytes[1] >= std.meta.fields(Status).len) return error.InvalidRepoState; - if (bytes[83] >= std.meta.fields(ErrorClass).len) return error.InvalidRepoState; + if (bytes[1] >= @typeInfo(Status).@"enum".field_names.len) return error.InvalidRepoState; + if (bytes[83] >= @typeInfo(ErrorClass).@"enum".field_names.len) return error.InvalidRepoState; const rev_len = std.mem.readInt(u32, bytes[67..71], .little); const latest_rev_len = std.mem.readInt(u32, bytes[71..75], .little); const host_len = std.mem.readInt(u16, bytes[75..77], .little); @@ -673,7 +674,7 @@ fn decode(bytes: []const u8) !RepoState { const handle_end = host_end + handle_len; const pds_end = handle_end + pds_len; return .{ - .status = @enumFromInt(bytes[1]), + .status = @fromBackingInt(@intCast(bytes[1])), .active = bytes[2] != 0, .rev = bytes[header_len..rev_end], .latest_rev = bytes[rev_end..latest_rev_end], @@ -681,7 +682,7 @@ fn decode(bytes: []const u8) !RepoState { .handle = bytes[host_end..handle_end], .pds = bytes[handle_end..pds_end], .last_error = bytes[pds_end..total], - .error_class = @enumFromInt(bytes[83]), + .error_class = @fromBackingInt(@intCast(bytes[83])), .retry_count = std.mem.readInt(u32, bytes[3..7], .little), .next_attempt_us = std.mem.readInt(i64, bytes[7..15], .little), .attempts = std.mem.readInt(u32, bytes[15..19], .little), @@ -812,7 +813,7 @@ test "repo transitions atomically maintain durable upstream-shaped host diagnost try testing.expectEqual(@as(u64, 0), hosts[0].active); try testing.expectEqual(@as(u64, 1), hosts[0].count(.complete)); try testing.expectEqual(@as(u64, 0), hosts[0].count(.failed)); - try testing.expectEqual(@as(u64, 2), hosts[0].error_counts[@intFromEnum(ErrorClass.http_5xx)]); + try testing.expectEqual(@as(u64, 2), hosts[0].error_counts[@backingInt(ErrorClass.http_5xx)]); try testing.expectEqual(@as(usize, 2), hosts[0].recent_errors.len); try testing.expectEqualStrings("HTTP 503 second", hosts[0].recent_errors[0].message); @@ -832,7 +833,7 @@ test "repo transitions atomically maintain durable upstream-shaped host diagnost } else return error.MissingNewHost; try testing.expectEqual(@as(u64, 0), a.total); try testing.expectEqual(@as(u64, 0), a.count(.complete)); - try testing.expectEqual(@as(u64, 2), a.error_counts[@intFromEnum(ErrorClass.http_5xx)]); + try testing.expectEqual(@as(u64, 2), a.error_counts[@backingInt(ErrorClass.http_5xx)]); try testing.expectEqual(@as(u64, 1), b.total); try testing.expectEqual(@as(u64, 1), b.count(.complete)); @@ -938,12 +939,12 @@ test "v4 rows migrate latest revision without losing lifecycle state" { const testing = std.testing; var bytes: [v4_header_len + 3 + 4 + 5]u8 = @splat(0); bytes[0] = 4; - bytes[1] = @intFromEnum(Status.complete); + bytes[1] = @backingInt(Status.complete); bytes[2] = 1; std.mem.writeInt(u32, bytes[43..47], 3, .little); std.mem.writeInt(u16, bytes[47..49], 4, .little); std.mem.writeInt(u16, bytes[49..51], 5, .little); - bytes[51] = @intFromEnum(ErrorClass.timeout); + bytes[51] = @backingInt(ErrorClass.timeout); @memcpy(bytes[52..55], "rev"); @memcpy(bytes[55..59], "host"); @memcpy(bytes[59..64], "error"); diff --git a/src/internal/ingest/backfill/repos.zig b/src/internal/ingest/backfill/repos.zig index 218db56..3e90533 100644 --- a/src/internal/ingest/backfill/repos.zig +++ b/src/internal/ingest/backfill/repos.zig @@ -424,7 +424,7 @@ test "allocation failure never masquerades as a malformed repository" { var tree = zat.mst.Mst.init(fixture); try tree.put("app.bsky.feed.post/3k2abcdefghij", record_cid); const data_cid = try tree.rootCid(); - const keypair = try zat.Keypair.fromSecretKey(.p256, .{17} ** 32); + const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(17)); const did = try keypair.did(fixture); const signed = try zat.signCommit(fixture, .{ .did = did, @@ -488,7 +488,7 @@ test "listRepos DID remains authoritative over a valid CAR's embedded DID" { var tree = zat.mst.Mst.init(allocator); try tree.put("app.bsky.feed.post/3k2abcdefghij", record_cid); const data_cid = try tree.rootCid(); - const keypair = try zat.Keypair.fromSecretKey(.p256, .{13} ** 32); + const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(13)); const embedded_did = try keypair.did(allocator); const signed = try zat.signCommit(allocator, .{ .did = embedded_did, diff --git a/src/internal/ingest/backfill/resync.zig b/src/internal/ingest/backfill/resync.zig index aa8f1a6..3bc3553 100644 --- a/src/internal/ingest/backfill/resync.zig +++ b/src/internal/ingest/backfill/resync.zig @@ -141,11 +141,39 @@ pub fn resyncRepo( var trace = observability.Span.start(allocator, "ingest/backfill", "handleRepo", trace_context, &attributes); defer trace.deinit(); defer trace.succeed(); - errdefer |trace_err| trace.fail(trace_err); + return trace.recordError(appendReplacement( + alloc, + io, + pds, + did, + &prepared_fetch_result.repo, + archive, + store, + drops, + details, + previous, + active, + trace.context(), + )); +} +fn appendReplacement( + alloc: Allocator, + io: Io, + pds: []const u8, + did: []const u8, + prepared_repo: *repos.PreparedRepo, + archive: *archive_mod.Archive, + store: *repo_store.Store, + drops: *repos.DropCounts, + details: *AttemptDetails, + previous: ?repo_store.RepoState, + active: bool, + trace_context: observability.Context, +) !Stats { const now_us: i64 = Io.Timestamp.now(io, .real).toMicroseconds(); - const head_rev = prepared_fetch_result.repo.rev; + const head_rev = prepared_repo.rev; const sync_ev = try syncTombstoneEvent(alloc, did, head_rev, now_us); @@ -174,13 +202,13 @@ pub fn resyncRepo( .archive = archive, .allocator = alloc, .now_us = now_us, - .trace_context = trace.context(), + .trace_context = trace_context, }; defer sink.batch.deinit(alloc); try sink.emit(sync_ev); - _ = try repos.emitPrepared(alloc, &prepared_fetch_result.repo, did, now_us, .create_resync, drops, &sink, Sink.emit); + _ = try repos.emitPrepared(alloc, prepared_repo, did, now_us, .create_resync, drops, &sink, Sink.emit); try sink.flush(); - try archive.flushBlockTraced(trace.context()); + try archive.flushBlockTraced(trace_context); // durability ordering: every replacement batch is durable before the // completion row that describes it. @@ -409,7 +437,7 @@ fn buildRepoCar(alloc: Allocator, secret: u8) !struct { car: []const u8, did: [] var tree = zat.mst.Mst.init(alloc); try tree.put("app.bsky.feed.post/3k2abcdefghij", record_cid); const data_cid = try tree.rootCid(); - const keypair = try zat.Keypair.fromSecretKey(.p256, .{secret} ** 32); + const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(secret)); const did = try keypair.did(alloc); const signed = try zat.signCommit(alloc, .{ .did = did, diff --git a/src/internal/ingest/convert.zig b/src/internal/ingest/convert.zig index 0da5906..a19c5af 100644 --- a/src/internal/ingest/convert.zig +++ b/src/internal/ingest/convert.zig @@ -283,7 +283,7 @@ const testing = std.testing; test "unencodable record is isolated instead of failing the live batch" { // A CID whose base32 form overflows cidString's fixed buffer makes // writeCborValue fail with WriteFailed -- reachable from remote input. - const long_cid: zat.cbor.Cid = .{ .raw = &[_]u8{0x55} ** 100 }; + const long_cid: zat.cbor.Cid = .{ .raw = &@as([100]u8, @splat(0x55)) }; const bad_record: zat.cbor.Value = .{ .map = &.{ .{ .key = "subject", .value = .{ .cid = long_cid } }, } }; diff --git a/src/internal/ingest/ingest.zig b/src/internal/ingest/ingest.zig index 20e7e86..e9ee978 100644 --- a/src/internal/ingest/ingest.zig +++ b/src/internal/ingest/ingest.zig @@ -268,7 +268,15 @@ pub const Consumer = struct { var trace = observability.Span.start(self.allocator, "ingest/backfill", "handleRepo", trace_context, &attributes); defer trace.deinit(); defer trace.succeed(); - errdefer |trace_err| trace.fail(trace_err); + try trace.recordError(self.applyRepair(archive, completion, trace.context())); + } + + fn applyRepair( + self: *Consumer, + archive: *archive_mod.Archive, + completion: *repair.Completion, + trace_context: observability.Context, + ) !void { var arena = std.heap.ArenaAllocator.init(self.allocator); defer arena.deinit(); const alloc = arena.allocator(); @@ -339,7 +347,7 @@ pub const Consumer = struct { .fetched_chain = completion.fetched_chain, .did = completion.did, .upstream_seq = completion.upstream_seq, - .trace_context = trace.context(), + .trace_context = trace_context, }; defer sink.batch.deinit(alloc); const sync_event = if (completion.sync_payload) |payload| @@ -359,8 +367,8 @@ pub const Consumer = struct { try writes.commitDurable(); self.drops.getPtr(.invalid_rev).* += 1; if (self.stats) |stats| _ = stats.drops.getPtr(.invalid_rev).fetchAdd(1, .monotonic); - for (completion.pending.items) |item| try self.handleEvent(item.event, item.verdict, true, trace.context()); - try archive.flushBlockTraced(trace.context()); + for (completion.pending.items) |item| try self.handleEvent(item.event, item.verdict, true, trace_context); + try archive.flushBlockTraced(trace_context); return; } @@ -386,8 +394,8 @@ pub const Consumer = struct { } sink.stage_chain_on_flush = true; try sink.flush(); - for (completion.pending.items) |item| try self.handleEvent(item.event, item.verdict, true, trace.context()); - try archive.flushBlockTraced(trace.context()); + for (completion.pending.items) |item| try self.handleEvent(item.event, item.verdict, true, trace_context); + try archive.flushBlockTraced(trace_context); } /// Atmos advances its cursor only after a non-empty delivery batch @@ -951,9 +959,11 @@ test "repair hot tail skips an undecodable archived record without failing inges defer consumer.deinit(); // A complete repository can contain opaque record bytes outside - // DAG-CBOR. This exact float decoder error killed the experiment while a - // verifier repair was filling the hot tail from its fetched repository. - const float_record = "\xa1\x61x\xfb\x3f\xf0\x00\x00\x00\x00\x00\x00"; + // DAG-CBOR. A float decoder error killed the experiment while a verifier + // repair was filling the hot tail from its fetched repository. zat 0.4.2 + // accepts finite 64-bit floats, so this uses a 32-bit float, which + // DAG-CBOR still forbids. + const float_record = "\xa1\x61x\xfa\x3f\x80\x00\x00"; try appendArchivedTailRow(&consumer, .{ .seq = 42, .witnessed_at = 100, diff --git a/src/internal/ingest/pipeline.zig b/src/internal/ingest/pipeline.zig index c5437a9..878f0e4 100644 --- a/src/internal/ingest/pipeline.zig +++ b/src/internal/ingest/pipeline.zig @@ -857,7 +857,7 @@ test "pipeline preserves same-key submission order under parallel workers" { try testing.expectEqual(@as(u64, 0), stats.verify_queue_dropped_total.load(.monotonic)); try testing.expectEqual(@as(usize, 1000), collector.frames.items.len); // per-key submission order must survive parallel workers - var next_per_key = [_]usize{0} ** 100; + var next_per_key: [100]usize = @splat(0); for (collector.frames.items) |f| { var it = std.mem.splitScalar(u8, f, '-'); _ = it.next(); // "frame" diff --git a/src/internal/ingest/repair.zig b/src/internal/ingest/repair.zig index afc16b9..cc11f36 100644 --- a/src/internal/ingest/repair.zig +++ b/src/internal/ingest/repair.zig @@ -747,7 +747,7 @@ test "offline coordinator fetches and authenticates a real complete repo" { defer arena.deinit(); const alloc = arena.allocator(); - const keypair = try zat.Keypair.fromSecretKey(.p256, .{21} ** 32); + const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(21)); const did = "did:plc:repairfixture"; const did_key = try keypair.did(alloc); const multibase = did_key["did:key:".len..]; diff --git a/src/internal/ingest/repair_integration_test.zig b/src/internal/ingest/repair_integration_test.zig index 0cd5d73..a8ee1b2 100644 --- a/src/internal/ingest/repair_integration_test.zig +++ b/src/internal/ingest/repair_integration_test.zig @@ -59,7 +59,7 @@ test "repair completion durably emits tombstone then replacement offline" { var tree = zat.mst.Mst.init(alloc); try tree.put("app.bsky.feed.post/3k2abcdefghij", record_cid); const data_cid = try tree.rootCid(); - const keypair = try zat.Keypair.fromSecretKey(.p256, .{31} ** 32); + const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(31)); const signed = try zat.signCommit(alloc, .{ .did = did, .rev = "3k2abcdefghij", @@ -165,7 +165,7 @@ test "large repair keeps wire encoding scratch bounded" { } const alloc = prepared_arena.allocator(); const did = "did:plc:large-durable-repair"; - const body = [_]u8{'x'} ** 512; + const body: [512]u8 = @splat('x'); const record = try zat.cbor.encodeAlloc(alloc, .{ .map = &.{ .{ .key = "$type", .value = .{ .text = "app.bsky.feed.post" } }, .{ .key = "text", .value = .{ .text = &body } }, @@ -177,7 +177,7 @@ test "large repair keeps wire encoding scratch bounded" { try tree.put(key, record_cid); } const data_cid = try tree.rootCid(); - const keypair = try zat.Keypair.fromSecretKey(.p256, .{29} ** 32); + const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(29)); const signed = try zat.signCommit(alloc, .{ .did = did, .rev = "3k2abcdefghij", diff --git a/src/internal/ingest/verify.zig b/src/internal/ingest/verify.zig index 7a15f57..0eb82c3 100644 --- a/src/internal/ingest/verify.zig +++ b/src/internal/ingest/verify.zig @@ -674,7 +674,7 @@ test "op CID proof rejects an envelope that disagrees with post-state MST" { var tree = zat.mst.Mst.init(allocator); try tree.put("app.bsky.feed.post/3k2abcdefghij", record_cid); const data_cid = try tree.rootCid(); - const keypair = try zat.Keypair.fromSecretKey(.p256, .{7} ** 32); + const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(7)); const did = try keypair.did(allocator); const signed = try zat.signCommit(allocator, .{ .did = did, @@ -728,7 +728,7 @@ test "commit block limit is upstream decimal two million bytes" { defer meta.deinit(); var verifier = Verifier.init(testing.allocator, threaded.io(), "http://offline.invalid", &meta); defer verifier.deinit(); - const keypair = try zat.Keypair.fromSecretKey(.p256, .{29} ** 32); + const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(29)); const public_key = try keypair.publicKey(); var cached: CachedKey = .{ .key_type = .p256, .raw = undefined, .len = @intCast(public_key.len) }; @memcpy(cached.raw[0..public_key.len], &public_key); @@ -760,7 +760,7 @@ test "signed chain failures follow upstream repair policy and legacy default" { var tree = zat.mst.Mst.init(allocator); const data_cid = try tree.rootCid(); - const keypair = try zat.Keypair.fromSecretKey(.p256, .{11} ** 32); + const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(11)); const did = try keypair.did(allocator); const signed = try zat.signCommit(allocator, .{ .did = did, @@ -847,7 +847,7 @@ test "full-repo repair authenticates the signed complete repository head" { var tree = zat.mst.Mst.init(allocator); try tree.put("app.bsky.feed.post/3k2abcdefghij", record_cid); const data_cid = try tree.rootCid(); - const keypair = try zat.Keypair.fromSecretKey(.p256, .{9} ** 32); + const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(9)); const did = try keypair.did(allocator); const signed = try zat.signCommit(allocator, .{ .did = did, @@ -884,7 +884,7 @@ test "full-repo repair authenticates the signed complete repository head" { repos.prepareRepo(allocator, incomplete_car), ); - const wrong_keypair = try zat.Keypair.fromSecretKey(.p256, .{7} ** 32); + const wrong_keypair = try zat.Keypair.fromSecretKey(.p256, @splat(7)); const wrong_public_key = try wrong_keypair.publicKey(); var wrong: CachedKey = .{ .key_type = .p256, .raw = undefined, .len = @intCast(wrong_public_key.len) }; @memcpy(wrong.raw[0..wrong_public_key.len], &wrong_public_key); diff --git a/src/internal/runtime/cli.zig b/src/internal/runtime/cli.zig index 9d16333..2728e12 100644 --- a/src/internal/runtime/cli.zig +++ b/src/internal/runtime/cli.zig @@ -68,12 +68,12 @@ fn overrideFor(comptime T: type, comptime field: []const u8) ?Exposure { pub fn exposures(comptime T: type) []const Exposure { return comptime blk: { @setEvalBranchQuota(20_000); - const fields = std.meta.fields(T); - var out: [fields.len]Exposure = undefined; - for (fields, 0..) |f, i| { - out[i] = overrideFor(T, f.name) orelse .{ - .flag = flagName(f.name), - .env = envName(f.name), + const names = @typeInfo(T).@"struct".field_names; + var out: [names.len]Exposure = undefined; + for (names, 0..) |name, i| { + out[i] = overrideFor(T, name) orelse .{ + .flag = flagName(name), + .env = envName(name), }; } const frozen = out; diff --git a/src/internal/runtime/crashpoint.zig b/src/internal/runtime/crashpoint.zig index 7448c9a..56d23fc 100644 --- a/src/internal/runtime/crashpoint.zig +++ b/src/internal/runtime/crashpoint.zig @@ -87,9 +87,8 @@ pub fn pointName(point: Point) []const u8 { } fn parse(name: []const u8) ?Point { - inline for (@typeInfo(Point).@"enum".fields) |field| { - const point: Point = @enumFromInt(field.value); - if (std.mem.eql(u8, name, pointName(point)) or std.mem.eql(u8, name, field.name)) return point; + inline for (comptime std.enums.values(Point)) |point| { + if (std.mem.eql(u8, name, pointName(point)) or std.mem.eql(u8, name, @tagName(point))) return point; } return null; } diff --git a/src/internal/runtime/logging.zig b/src/internal/runtime/logging.zig index 203f0fb..afd840a 100644 --- a/src/internal/runtime/logging.zig +++ b/src/internal/runtime/logging.zig @@ -7,14 +7,14 @@ pub const Format = enum(u8) { json, }; -var configured_level: std.atomic.Value(u8) = .init(@intFromEnum(std.log.Level.info)); -var configured_format: std.atomic.Value(u8) = .init(@intFromEnum(Format.json)); +var configured_level: std.atomic.Value(u8) = .init(@backingInt(std.log.Level.info)); +var configured_format: std.atomic.Value(u8) = .init(@backingInt(Format.json)); pub fn configure(level_raw: []const u8, format_raw: []const u8) !void { const level = try parseLevel(level_raw); const format = try parseFormat(format_raw); - configured_level.store(@intFromEnum(level), .release); - configured_format.store(@intFromEnum(format), .release); + configured_level.store(@backingInt(level), .release); + configured_format.store(@backingInt(format), .release); } pub fn parseLevel(raw: []const u8) !std.log.Level { @@ -39,7 +39,7 @@ pub fn logFn( comptime format: []const u8, args: anytype, ) void { - if (@intFromEnum(level) > configured_level.load(.acquire)) return; + if (@backingInt(level) > configured_level.load(.acquire)) return; var message: Io.Writer.Allocating = .init(std.heap.page_allocator); defer message.deinit(); @@ -55,7 +55,7 @@ pub fn logFn( const terminal = std.debug.lockStderr(&buffer).terminal(); defer std.debug.unlockStderr(); - const selected: Format = @enumFromInt(configured_format.load(.acquire)); + const selected: Format = @fromBackingInt(@intCast(configured_format.load(.acquire))); writeRecord( terminal.writer, selected, diff --git a/src/internal/runtime/metrics.zig b/src/internal/runtime/metrics.zig index 5f5f8ab..d7cbaf5 100644 --- a/src/internal/runtime/metrics.zig +++ b/src/internal/runtime/metrics.zig @@ -71,11 +71,11 @@ pub const CompactionPassResult = enum { ok, error_result }; pub const CompactionReason = enum { record, sync, account }; pub const PipelineFailureStage = enum { scheduler, emit, repair, cursor }; -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; -const orchestrator_phase_count = @typeInfo(OrchestratorPhase).@"enum".fields.len; +const http_series_count = @typeInfo(HttpHandler).@"enum".field_names.len * + @typeInfo(HttpMethod).@"enum".field_names.len * @typeInfo(HttpCode).@"enum".field_names.len; +const store_series_count = @typeInfo(StoreOp).@"enum".field_names.len * + @typeInfo(StoreStatus).@"enum".field_names.len; +const orchestrator_phase_count = @typeInfo(OrchestratorPhase).@"enum".field_names.len; pub const Stats = struct { start_time_s: i64 = 0, @@ -110,9 +110,9 @@ pub const Stats = struct { store_op_duration_buckets: [store_series_count][15]Counter = @splat(@splat(.init(0))), orchestrator_phase: std.atomic.Value(usize) = .init(0), orchestrator_phase_transitions: [orchestrator_phase_count * orchestrator_phase_count]Counter = @splat(.init(0)), - orchestrator_state_duration_count: [@typeInfo(OrchestratorState).@"enum".fields.len]Counter = @splat(.init(0)), - orchestrator_state_duration_sum_us: [@typeInfo(OrchestratorState).@"enum".fields.len]Counter = @splat(.init(0)), - orchestrator_state_duration_buckets: [@typeInfo(OrchestratorState).@"enum".fields.len][14]Counter = @splat(@splat(.init(0))), + orchestrator_state_duration_count: [@typeInfo(OrchestratorState).@"enum".field_names.len]Counter = @splat(.init(0)), + orchestrator_state_duration_sum_us: [@typeInfo(OrchestratorState).@"enum".field_names.len]Counter = @splat(.init(0)), + orchestrator_state_duration_buckets: [@typeInfo(OrchestratorState).@"enum".field_names.len][14]Counter = @splat(@splat(.init(0))), orchestrator_merge_events_kept_total: Counter = .init(0), orchestrator_merge_events_dropped_total: Counter = .init(0), orchestrator_merge_segments_consumed_total: Counter = .init(0), @@ -347,7 +347,7 @@ pub const Stats = struct { } pub fn setOrchestratorPhase(self: *Stats, phase: OrchestratorPhase) void { - self.orchestrator_phase.store(@intFromEnum(phase) + 1, .monotonic); + self.orchestrator_phase.store(@backingInt(phase) + 1, .monotonic); } pub fn observeOrchestratorTransition(self: *Stats, from: OrchestratorPhase, to: OrchestratorPhase) void { @@ -355,7 +355,7 @@ pub const Stats = struct { } pub fn observeOrchestratorState(self: *Stats, state: OrchestratorState, duration_us: u64) void { - const index = @intFromEnum(state); + const index = @backingInt(state); _ = self.orchestrator_state_duration_count[index].fetchAdd(1, .monotonic); _ = self.orchestrator_state_duration_sum_us[index].fetchAdd(duration_us, .monotonic); inline for (0..14) |bucket_index| { @@ -404,17 +404,17 @@ fn observeSlowHistogram(count: *Counter, sum_us: *Counter, buckets: *[15]Counter } fn storeIndex(op: StoreOp, status: StoreStatus) usize { - return @intFromEnum(op) * @typeInfo(StoreStatus).@"enum".fields.len + @intFromEnum(status); + return @backingInt(op) * @typeInfo(StoreStatus).@"enum".field_names.len + @backingInt(status); } fn orchestratorTransitionIndex(from: OrchestratorPhase, to: OrchestratorPhase) usize { - return @intFromEnum(from) * orchestrator_phase_count + @intFromEnum(to); + return @backingInt(from) * orchestrator_phase_count + @backingInt(to); } 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; - return (@intFromEnum(handler) * method_count + @intFromEnum(method)) * code_count + @intFromEnum(code); + const method_count = @typeInfo(HttpMethod).@"enum".field_names.len; + const code_count = @typeInfo(HttpCode).@"enum".field_names.len; + return (@backingInt(handler) * method_count + @backingInt(method)) * code_count + @backingInt(code); } fn httpCode(code: u16) HttpCode { @@ -866,7 +866,7 @@ pub fn format( "1.28", "2.56", "5.12", "10.24", "20.48", "40.96", "81.92", }; inline for (comptime std.enums.values(OrchestratorState)) |state| { - const index = @intFromEnum(state); + const index = @backingInt(state); const count = stats.orchestrator_state_duration_count[index].load(.monotonic); if (count != 0) { const sum_us = stats.orchestrator_state_duration_sum_us[index].load(.monotonic); diff --git a/src/internal/runtime/observability.zig b/src/internal/runtime/observability.zig index 3c7f0ba..e2dd771 100644 --- a/src/internal/runtime/observability.zig +++ b/src/internal/runtime/observability.zig @@ -331,6 +331,13 @@ pub const Span = struct { self.ended = true; } + /// Fails the span with `result`'s error, if it carries one, and returns + /// `result` unchanged for the caller to `try` or `return`. + pub fn recordError(self: *Span, result: anytype) @TypeOf(result) { + if (result) |_| {} else |err| self.fail(err); + return result; + } + pub fn deinit(self: *Span) void { if (!self.ended) self.succeed(); if (self.owned_context) |owned| diff --git a/src/internal/serve/filter.zig b/src/internal/serve/filter.zig index 7059a08..cabc448 100644 --- a/src/internal/serve/filter.zig +++ b/src/internal/serve/filter.zig @@ -40,7 +40,7 @@ pub const Filter = struct { /// wanted DIDs; null means match-all dids: ?std.StringHashMapUnmanaged(void) = null, /// v2 kinds predicate; full set means match-all (v1 always full) - kinds: KindSet = KindSet.initFull(), + kinds: KindSet = .full, arena: std.heap.ArenaAllocator, pub fn init(allocator: Allocator) Filter { @@ -79,7 +79,7 @@ pub const Filter = struct { errdefer f.deinit(); var saw_kinds = false; var saw_collections = false; - var kinds = KindSet.initEmpty(); + var kinds = KindSet.empty; var it = std.mem.splitScalar(u8, query, '&'); while (it.next()) |pair| { if (pair.len == 0) continue; diff --git a/src/internal/serve/repo_export.zig b/src/internal/serve/repo_export.zig index 0891fd8..a707eb7 100644 --- a/src/internal/serve/repo_export.zig +++ b/src/internal/serve/repo_export.zig @@ -12,6 +12,7 @@ const archive_mod = @import("../storage/archive.zig"); const gloom = @import("../storage/gloom.zig"); const segment = @import("../storage/segment.zig"); const writer_mod = @import("../storage/segment_writer.zig"); +const cli = @import("../runtime/cli.zig"); const Io = std.Io; const Allocator = std.mem.Allocator; @@ -61,9 +62,9 @@ pub fn verify( const latest_url = try xrpcUrlAlloc(allocator, endpoint, "com.atproto.sync.getLatestCommit", &.{.{ .name = "did", .value = did }}); defer allocator.free(latest_url); - var transport = zat.HttpTransport.init(io, allocator); + var transport = zat.HttpTransport.initWithUserAgent(io, allocator, cli.userAgent()); defer transport.deinit(); - const latest_response = try transport.fetch(.{ .url = latest_url, .max_response_size = 64 * 1024, .redirect_behavior = .not_allowed }); + const latest_response = try transport.fetch(.{ .url = latest_url, .max_response_size = 64 * 1024, .on_redirect = .refuse }); defer allocator.free(latest_response.body); if (latest_response.status.class() != .success) return error.GetLatestCommitFailed; const latest_json = try std.json.parseFromSlice(std.json.Value, allocator, latest_response.body, .{}); @@ -82,7 +83,7 @@ pub fn verify( const blocks_response = try transport.fetch(.{ .url = blocks_url, .max_response_size = 8 * (1024 * 1024 + 128) + 64 * 1024, - .redirect_behavior = .not_allowed, + .on_redirect = .refuse, }); defer allocator.free(blocks_response.body); if (blocks_response.status.class() != .success) return error.GetCommitBlockFailed; @@ -553,7 +554,7 @@ test "verify compares a real JSS reconstruction with real PLC and PDS Sync respo var authoritative_tree = zat.mst.Mst.init(alloc); try authoritative_tree.put("app.bsky.feed.post/fixture", record_cid); const data_root = try authoritative_tree.rootCid(); - const keypair = try zat.Keypair.fromSecretKey(.p256, .{41} ** 32); + const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(41)); const signed = try zat.signCommit(alloc, .{ .did = did, .rev = rev, .data = data_root }, &keypair); const commit_cid_string = try zat.multibase.base32lower.encode(alloc, signed.cid.raw); const decoy_data = "decoy"; diff --git a/src/internal/serve/server.zig b/src/internal/serve/server.zig index 7175692..e22b5b3 100644 --- a/src/internal/serve/server.zig +++ b/src/internal/serve/server.zig @@ -2196,10 +2196,10 @@ test "XRPC readiness gate returns the upstream JSON service error" { var transport = zat.HttpTransport.init(io, testing.allocator); transport.keep_alive = false; defer transport.deinit(); - var response = try transport.fetch(.{ .url = url, .max_response_size = 4096, .capture_response_headers = true }); + var response = try transport.fetch(.{ .url = url, .max_response_size = 4096, .capture_headers = &.{"content-type"} }); defer response.deinit(testing.allocator); try testing.expectEqual(std.http.Status.service_unavailable, response.status); - try testing.expectEqualStrings("application/json", response.oauth.content_type.?); + try testing.expectEqualStrings("application/json", response.headers.get("content-type").?); const parsed = try std.json.parseFromSlice(std.json.Value, testing.allocator, response.body, .{}); defer parsed.deinit(); try testing.expectEqualStrings("ServiceUnavailable", parsed.value.object.get("error").?.string); @@ -2370,7 +2370,7 @@ test "accounts HTTP route verifies real archive state and rate limits by port-fr var authoritative_tree = zat.mst.Mst.init(alloc); try authoritative_tree.put("app.bsky.feed.post/fixture", record_cid); const data_root = try authoritative_tree.rootCid(); - const keypair = try zat.Keypair.fromSecretKey(.p256, .{43} ** 32); + const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(43)); const signed = try zat.signCommit(alloc, .{ .did = did, .rev = rev, .data = data_root }, &keypair); const commit_cid_string = try zat.multibase.base32lower.encode(alloc, signed.cid.raw); const get_blocks_car = try zat.car.writeAlloc(alloc, .{ diff --git a/src/internal/serve/wire.zig b/src/internal/serve/wire.zig index 170058d..6b94341 100644 --- a/src/internal/serve/wire.zig +++ b/src/internal/serve/wire.zig @@ -194,6 +194,7 @@ fn writeCborValue(s: *std.json.Stringify, value: cbor.Value) !void { switch (value) { .unsigned => |v| try s.write(v), .negative => |v| try s.write(v), + .float => |v| try s.write(v), // record text can be any bytes a PDS chose to write; the same // array-instead-of-string hazard applies to it as to the envelope .text => |v| try writeString(s, v), diff --git a/src/internal/storage/archive.zig b/src/internal/storage/archive.zig index 4856ea7..e8f2762 100644 --- a/src/internal/storage/archive.zig +++ b/src/internal/storage/archive.zig @@ -528,7 +528,20 @@ pub const Archive = struct { ) !u64 { var flush_trace: ?observability.Span = null; defer if (flush_trace) |*trace| trace.deinit(); - errdefer |err| if (flush_trace) |*trace| trace.fail(err); + const result = self.appendBatchSyncLocked(events, now_us, hook_ctx, before_durable, trace_context, &flush_trace); + if (flush_trace) |*trace| return trace.recordError(result); + return result; + } + + fn appendBatchSyncLocked( + self: *Archive, + events: []segment.Event, + now_us: i64, + hook_ctx: ?*anyopaque, + before_durable: ?*const fn (*anyopaque, u64, u64) anyerror!void, + trace_context: observability.Context, + flush_trace: *?observability.Span, + ) !u64 { self.mu.lockUncancelable(self.io); defer self.mu.unlock(self.io); errdefer self.poisoned = true; @@ -539,7 +552,7 @@ pub const Archive = struct { } if (before_durable) |hook| try hook(hook_ctx.?, first_seq, self.next_seq - 1); if (self.writer.pending.items.len > 0 and now_us - self.last_sync_us >= self.flush_interval_us) { - flush_trace = observability.Span.start( + flush_trace.* = observability.Span.start( self.allocator, "ingest", "flushAndRotateLocked", @@ -549,8 +562,8 @@ pub const Archive = struct { try self.writer.flush(); // partial block: durability over block size } if (self.writer.bytes().len > self.file_written) { - if (flush_trace == null) - flush_trace = observability.Span.start( + if (flush_trace.* == null) + flush_trace.* = observability.Span.start( self.allocator, "ingest", "flushAndRotateLocked", @@ -560,9 +573,9 @@ pub const Archive = struct { try self.syncToDisk(null); self.last_sync_us = now_us; if (self.writer.bytes().len >= self.max_segment_bytes) - try self.rotateLockedTraced(flush_trace.?.context()); + try self.rotateLockedTraced(flush_trace.*.?.context()); } - if (flush_trace) |*trace| trace.succeed(); + if (flush_trace.*) |*trace| trace.succeed(); return first_seq; } @@ -675,8 +688,10 @@ pub const Archive = struct { var trace = observability.Span.start(self.allocator, "ingest", "rotateIfFull", trace_context, &.{}); defer trace.deinit(); defer trace.succeed(); - errdefer |trace_err| trace.fail(trace_err); + try trace.recordError(self.rotateDrained(pipeline, trace.context())); + } + fn rotateDrained(self: *Archive, pipeline: *AsyncFlush, trace_context: observability.Context) !void { pipeline.admission_mu.lockUncancelable(self.io); defer pipeline.admission_mu.unlock(self.io); pipeline.drain(); @@ -684,7 +699,7 @@ pub const Archive = struct { defer self.mu.unlock(self.io); errdefer self.poisoned = true; if (self.poisoned) return error.ArchivePoisoned; - if (self.writer.bytes().len >= self.max_segment_bytes) try self.rotateLockedTraced(trace.context()); + if (self.writer.bytes().len >= self.max_segment_bytes) try self.rotateLockedTraced(trace_context); } fn appendStampedLocked(self: *Archive, event: *segment.Event) !u64 { @@ -767,7 +782,10 @@ pub const Archive = struct { var trace = observability.Span.start(self.allocator, "ingest", "flushAndRotateLocked", trace_context, &.{}); defer trace.deinit(); defer trace.succeed(); - errdefer |err| trace.fail(err); + try trace.recordError(self.flushBlockSync()); + } + + fn flushBlockSync(self: *Archive) !void { self.mu.lockUncancelable(self.io); defer self.mu.unlock(self.io); errdefer self.poisoned = true; @@ -874,8 +892,7 @@ pub const Archive = struct { var trace = observability.Span.start(self.allocator, "ingest", "rotateLocked", trace_context, &.{}); defer trace.deinit(); defer trace.succeed(); - errdefer |err| trace.fail(err); - try self.rotateLocked(); + try trace.recordError(self.rotateLocked()); trace.setAttribute(.{ .key = "active_idx", .value = .{ .int = @intCast(self.seg_index) } }); } diff --git a/src/internal/storage/cold.zig b/src/internal/storage/cold.zig index 63ba949..4e59c0a 100644 --- a/src/internal/storage/cold.zig +++ b/src/internal/storage/cold.zig @@ -1590,7 +1590,7 @@ test "timestamp cursor surfaces corrupt selected block" { var entry_buf: [segment.BlockIndexEntry.encoded_size]u8 = undefined; _ = try file.readPositionalAll(io, &entry_buf, header.block_index_offset + segment.BlockIndexEntry.encoded_size); const entry = segment.BlockIndexEntry.decodePublic(&entry_buf); - try file.writePositionalAll(io, &([_]u8{0} ** 8), entry.offset); + try file.writePositionalAll(io, &@as([8]u8, @splat(0)), entry.offset); try file.sync(io); file.close(io); diff --git a/src/internal/storage/cold_cache.zig b/src/internal/storage/cold_cache.zig index 9822b79..d5dcbae 100644 --- a/src/internal/storage/cold_cache.zig +++ b/src/internal/storage/cold_cache.zig @@ -75,7 +75,7 @@ pub const Handle = struct { comptime compute: fn (@TypeOf(ctx), Allocator, segment.Event) anyerror!?[]u8, ) !?[]const u8 { if (event_index >= self.item.block.events.len) return error.BadEventIndex; - const at = event_index * 4 + @intFromEnum(lane); + const at = event_index * 4 + @backingInt(lane); self.item.memo_mu.lockUncancelable(self.cache.io); defer self.item.memo_mu.unlock(self.cache.io); const memo = &self.item.memos[at]; diff --git a/src/internal/storage/disk_space.zig b/src/internal/storage/disk_space.zig index ce06b07..5d7aef8 100644 --- a/src/internal/storage/disk_space.zig +++ b/src/internal/storage/disk_space.zig @@ -1,17 +1,40 @@ //! Scrape-time filesystem capacity for the archive data directory. const std = @import("std"); +const builtin = @import("builtin"); -const c = @cImport({ - @cInclude("sys/statvfs.h"); -}); +const FsCount = if (builtin.os.tag.isDarwin()) + c_uint +else if (builtin.abi.isMusl()) + u64 +else + c_ulong; + +/// `struct statvfs` through `f_namemax`; `reserved` covers the trailing +/// fields Linux libcs append, which are never read. +const Statvfs = extern struct { + f_bsize: c_ulong, + f_frsize: c_ulong, + f_blocks: FsCount, + f_bfree: FsCount, + f_bavail: FsCount, + f_files: FsCount, + f_ffree: FsCount, + f_favail: FsCount, + f_fsid: c_ulong, + f_flag: c_ulong, + f_namemax: c_ulong, + reserved: [8]c_int, +}; + +extern "c" fn fstatvfs(fd: std.c.fd_t, buf: *Statvfs) c_int; /// Bytes available to this unprivileged process on the filesystem containing /// dir. Failure is represented as unavailable so /metrics never lies with a /// zero that looks like a full disk. pub fn freeBytes(dir: std.Io.Dir) ?u64 { - var stat: c.struct_statvfs = undefined; - if (c.fstatvfs(dir.handle, &stat) != 0) return null; + var stat: Statvfs = undefined; + if (fstatvfs(dir.handle, &stat) != 0) return null; return std.math.mul(u64, @intCast(stat.f_bavail), @intCast(stat.f_frsize)) catch null; } diff --git a/src/internal/storage/gloom.zig b/src/internal/storage/gloom.zig index 83e0754..9778a5f 100644 --- a/src/internal/storage/gloom.zig +++ b/src/internal/storage/gloom.zig @@ -162,7 +162,7 @@ pub const Filter = struct { for (self.words, 0..) |w, i| { std.mem.writeInt(u64, out[header_size + i * 8 ..][0..8], w, .little); } - const crc = std.hash.crc.Crc32Iscsi.hash(out[0 .. total - 4]); + const crc = std.hash.crc.@"CRC-32/ISCSI".hash(out[0 .. total - 4]); std.mem.writeInt(u32, out[total - 4 ..][0..4], crc, .little); return out; } @@ -240,7 +240,7 @@ test "membership: added keys always found" { pub fn unmarshal(allocator: Allocator, blob: []const u8) !Filter { if (blob.len < header_size + 4) return error.Truncated; const crc_stored = std.mem.readInt(u32, blob[blob.len - 4 ..][0..4], .little); - if (std.hash.crc.Crc32Iscsi.hash(blob[0 .. blob.len - 4]) != crc_stored) return error.BadChecksum; + if (std.hash.crc.@"CRC-32/ISCSI".hash(blob[0 .. blob.len - 4]) != crc_stored) return error.BadChecksum; if (blob[0] != 1) return error.UnsupportedVersion; const k = std.mem.readInt(u32, blob[1..5], .little); if (k < min_k or k > max_k) return error.UnsupportedK; @@ -262,7 +262,7 @@ pub fn unmarshal(allocator: Allocator, blob: []const u8) !Filter { pub fn validateMarshaled(blob: []const u8) !void { if (blob.len < header_size + 4) return error.Truncated; const crc_stored = std.mem.readInt(u32, blob[blob.len - 4 ..][0..4], .little); - if (std.hash.crc.Crc32Iscsi.hash(blob[0 .. blob.len - 4]) != crc_stored) return error.BadChecksum; + if (std.hash.crc.@"CRC-32/ISCSI".hash(blob[0 .. blob.len - 4]) != crc_stored) return error.BadChecksum; _ = try dimensions(blob); } diff --git a/src/internal/storage/meta_store.zig b/src/internal/storage/meta_store.zig index 1b634f0..318269e 100644 --- a/src/internal/storage/meta_store.zig +++ b/src/internal/storage/meta_store.zig @@ -396,5 +396,5 @@ test "armed durable store fault fires once at matching batch ordinal" { } fn storeMetricIndex(op: metrics.StoreOp, status: metrics.StoreStatus) usize { - return @intFromEnum(op) * @typeInfo(metrics.StoreStatus).@"enum".fields.len + @intFromEnum(status); + return @backingInt(op) * @typeInfo(metrics.StoreStatus).@"enum".field_names.len + @backingInt(status); } diff --git a/src/internal/storage/segment.zig b/src/internal/storage/segment.zig index c4759b8..bcb3fc6 100644 --- a/src/internal/storage/segment.zig +++ b/src/internal/storage/segment.zig @@ -457,7 +457,7 @@ pub fn decodeBlock(allocator: Allocator, buffer: []u8) Error!Block { .seq = readInt(u64, buffer, seqs + i * 8), .witnessed_at = @bitCast(readInt(u64, buffer, witnessed + i * 8)), .indexed_at = @bitCast(readInt(u64, buffer, indexed + i * 8)), - .kind = @enumFromInt(kind_raw), + .kind = @fromBackingInt(@intCast(kind_raw)), .collection = buffer[c..][0..clen], .did = buffer[d..][0..dlen], .rkey = buffer[rk..][0..rklen], @@ -489,7 +489,7 @@ inline fn readInt(comptime T: type, bytes: []const u8, offset: usize) T { const testing = std.testing; test "header rejects bad magic, active files, wrong version" { - var buf = [_]u8{0} ** header_size; + var buf: [header_size]u8 = @splat(0); try testing.expectError(error.BadMagic, Header.decode(&buf)); @memcpy(buf[0..4], magic); try testing.expectError(error.ActiveSegment, Header.decode(&buf)); diff --git a/src/internal/storage/segment_fixture_test.zig b/src/internal/storage/segment_fixture_test.zig index ebbca64..891526b 100644 --- a/src/internal/storage/segment_fixture_test.zig +++ b/src/internal/storage/segment_fixture_test.zig @@ -140,7 +140,7 @@ test "fixture: collection index matches upstream inspect output" { } // bitmask cross-check: recount per-collection events by walking blocks - var counts = [_]u32{0} ** 7; + var counts: [7]u32 = @splat(0); for (0..sealed.block_index.len) |i| { var block = try sealed.readBlock(testing.allocator, i); defer block.deinit(testing.allocator); diff --git a/src/internal/storage/segment_writer.zig b/src/internal/storage/segment_writer.zig index df98af0..6bc6ac2 100644 --- a/src/internal/storage/segment_writer.zig +++ b/src/internal/storage/segment_writer.zig @@ -86,7 +86,7 @@ pub fn encodeBlock(allocator: Allocator, events: []const segment.Event) Error![] off += n * 8; for (events, 0..) |ev, i| std.mem.writeInt(i64, out[off + i * 8 ..][0..8], ev.indexed_at, .little); off += n * 8; - for (events, 0..) |ev, i| out[off + i] = @intFromEnum(ev.kind); + for (events, 0..) |ev, i| out[off + i] = @backingInt(ev.kind); off += n; for (events, 0..) |ev, i| out[off + i] = @intCast(ev.collection.len); off += n; @@ -433,7 +433,7 @@ test "block encode/decode round-trip" { test "validateEvent enforces column widths" { var ev = sampleEvent(1); - ev.rkey = "x" ** 256; + ev.rkey = &@as([256]u8, @splat('x')); try testing.expectError(error.FieldTooLong, validateEvent(ev)); ev = sampleEvent(0); try testing.expectError(error.InvalidEvent, validateEvent(ev)); diff --git a/src/internal/storage/zstd.zig b/src/internal/storage/zstd.zig index 4263fd7..b141162 100644 --- a/src/internal/storage/zstd.zig +++ b/src/internal/storage/zstd.zig @@ -54,8 +54,17 @@ pub fn decompressExact(allocator: std.mem.Allocator, frame: []const u8, expected // === tests === +fn repeated(comptime unit: []const u8, comptime count: usize) *const [unit.len * count]u8 { + return comptime blk: { + var out: [unit.len * count]u8 = undefined; + for (0..count) |i| @memcpy(out[i * unit.len ..][0..unit.len], unit); + const final = out; + break :blk &final; + }; +} + test "round-trip through std decompress" { - const src = "the quick brown fox jumps over the lazy dog" ** 100; + const src = repeated("the quick brown fox jumps over the lazy dog", 100); const frame = try compress(std.testing.allocator, src); defer std.testing.allocator.free(frame); try std.testing.expect(frame.len < src.len); @@ -69,7 +78,7 @@ test "round-trip through std decompress" { } test "libzstd decompress verifies the frame content checksum" { - const src = "checksum-protected block" ** 100; + const src = repeated("checksum-protected block", 100); const frame = try compress(std.testing.allocator, src); defer std.testing.allocator.free(frame); diff --git a/src/internal/timestamp/jobs.zig b/src/internal/timestamp/jobs.zig index 02edfcf..86166f2 100644 --- a/src/internal/timestamp/jobs.zig +++ b/src/internal/timestamp/jobs.zig @@ -47,7 +47,7 @@ pub const Record = struct { rows_total: u64 = 0, rows_valid: u64 = 0, rows_rejected: u64 = 0, - rejects_by_reason: [@typeInfo(parse.RejectReason).@"enum".fields.len]u64 = @splat(0), + rejects_by_reason: [@typeInfo(parse.RejectReason).@"enum".field_names.len]u64 = @splat(0), segments_examined: u64 = 0, segments_patched: u64 = 0, rows_mutated: u64 = 0, diff --git a/src/internal/timestamp/parse.zig b/src/internal/timestamp/parse.zig index 6ed53e7..4d37b2e 100644 --- a/src/internal/timestamp/parse.zig +++ b/src/internal/timestamp/parse.zig @@ -73,11 +73,11 @@ pub const Stats = struct { rows_total: u64 = 0, rows_valid: u64 = 0, rows_rejected: u64 = 0, - rejects_by_reason: [@typeInfo(RejectReason).@"enum".fields.len]u64 = @splat(0), + rejects_by_reason: [@typeInfo(RejectReason).@"enum".field_names.len]u64 = @splat(0), reject_sample: std.ArrayList(OwnedReject) = .empty, pub fn rejected(self: *const Stats, reason: RejectReason) u64 { - return self.rejects_by_reason[@intFromEnum(reason)]; + return self.rejects_by_reason[@backingInt(reason)]; } pub fn deinit(self: *Stats, allocator: std.mem.Allocator) void { @@ -360,7 +360,7 @@ fn noteReject( }; stats.rows_total += 1; stats.rows_rejected += 1; - stats.rejects_by_reason[@intFromEnum(reject.reason)] += 1; + stats.rejects_by_reason[@backingInt(reject.reason)] += 1; if (opts.on_reject) |callback| callback(opts.context, bounded); const limit = if (opts.reject_sample_limit == 0) diff --git a/src/internal/timestamp/rules.zig b/src/internal/timestamp/rules.zig index 264b80a..ffbe398 100644 --- a/src/internal/timestamp/rules.zig +++ b/src/internal/timestamp/rules.zig @@ -21,7 +21,7 @@ const key_specific: u8 = 's'; const key_has_specific: u8 = 'p'; const key_collection: u8 = 'c'; const key_sep: u8 = 0; -const marker_value: [8]u8 = .{1} ++ .{0} ** 7; +const marker_value: [8]u8 = .{ 1, 0, 0, 0, 0, 0, 0, 0 }; pub const ImportResult = struct { parse: parse_mod.Stats = .{}, @@ -56,7 +56,7 @@ pub const Store = struct { if (data_dir.len == 0) return error.MissingDataDir; const path = try std.fs.path.join(allocator, &.{ data_dir, subdir }); defer allocator.free(path); - const path_z = try allocator.dupeZ(u8, path); + const path_z = try allocator.dupeSentinel(u8, path, 0); defer allocator.free(path_z); const options = rdb.rocksdb_options_create() orelse return error.OutOfMemory; @@ -336,7 +336,7 @@ const Builder = struct { ); errdefer self.allocator.free(path); errdefer Io.Dir.cwd().deleteFile(self.store.io, path) catch {}; - const path_z = try self.allocator.dupeZ(u8, path); + const path_z = try self.allocator.dupeSentinel(u8, path, 0); defer self.allocator.free(path_z); const env_options = rdb.rocksdb_envoptions_create() orelse return error.OutOfMemory; @@ -385,7 +385,7 @@ const Builder = struct { rdb.rocksdb_ingestexternalfileoptions_set_allow_global_seqno(options, 1); rdb.rocksdb_ingestexternalfileoptions_set_allow_blocking_flush(options, 1); for (self.paths.items) |path| { - const path_z = try self.allocator.dupeZ(u8, path); + const path_z = try self.allocator.dupeSentinel(u8, path, 0); defer self.allocator.free(path_z); const paths = [_][*c]const u8{path_z.ptr}; var err_ptr: ?[*:0]u8 = null; diff --git a/src/internal/timestamp/runner.zig b/src/internal/timestamp/runner.zig index 0b0ad40..51a0604 100644 --- a/src/internal/timestamp/runner.zig +++ b/src/internal/timestamp/runner.zig @@ -49,7 +49,10 @@ pub fn run(config: Config, id: []const u8) !void { ); defer trace.deinit(); defer trace.succeed(); - errdefer |trace_err| trace.fail(trace_err); + try trace.recordError(runJob(config, id, trace.context())); +} + +fn runJob(config: Config, id: []const u8, trace_context: observability.Context) !void { const started_us = Io.Timestamp.now(config.io, .real).toMicroseconds(); defer if (config.stats) |stats| stats.import_phase.store(0, .release); var parsed = (try config.jobs.get(id)) orelse return error.ImportJobNotFound; @@ -57,7 +60,7 @@ pub fn run(config: Config, id: []const u8) !void { var record = parsed.value; if (record.state.terminal()) return error.ImportJobTerminal; - runInner(config, &record, trace.context()) catch |err| { + runInner(config, &record, trace_context) catch |err| { if (err == error.ImportCancelled) { try config.jobs.put(record); return @errorCast(err); diff --git a/src/main.zig b/src/main.zig index b1f84ee..06d5d60 100644 --- a/src/main.zig +++ b/src/main.zig @@ -496,8 +496,17 @@ pub fn main(init: std.process.Init.Minimal) !void { ); defer orchestrator_trace.deinit(); defer orchestrator_trace.succeed(); - errdefer |err| orchestrator_trace.fail(err); + try orchestrator_trace.recordError(serve(allocator, io, cfg, backfill_repos, &tracing, &orchestrator_trace)); +} +fn serve( + allocator: std.mem.Allocator, + io: Io, + cfg: *cli.ServeConfig, + backfill_repos: []const []const u8, + tracing: *observability.Runtime, + orchestrator_trace: *observability.Span, +) !void { var hot_tail = tail_mod.Tail.init(allocator, io, cfg.subscribe_read_log_retention_bytes); defer hot_tail.deinit(); var stats: metrics.Stats = .{ @@ -827,7 +836,7 @@ pub fn main(init: std.process.Init.Minimal) !void { archive_closed = true; try listener_shutdowns.await(io); orchestrator_trace.succeed(); - try shutdownTracing(&tracing, cfg.shutdown_timeout); + try shutdownTracing(tracing, cfg.shutdown_timeout); return; } lifecycle_future.await(io); @@ -851,18 +860,63 @@ pub fn main(init: std.process.Init.Minimal) !void { ); defer steady_trace.deinit(); defer steady_trace.succeed(); - errdefer |err| steady_trace.fail(err); + try steady_trace.recordError(runSteadyState( + allocator, + io, + cfg, + tracing, + orchestrator_trace, + &steady_trace, + &hub, + &hot_tail, + &stats, + &archive, + &archive_closed, + &meta, + &cursor_store, + &verifier, + &import_manager, + &import_manager_err_owned, + &import_ready, + live_repair_scratch, + &server_future, + &debug_future, + )); +} + +fn runSteadyState( + allocator: std.mem.Allocator, + io: Io, + cfg: *cli.ServeConfig, + tracing: *observability.Runtime, + orchestrator_trace: *observability.Span, + steady_trace: *observability.Span, + hub: *server_mod.Hub, + hot_tail: *tail_mod.Tail, + stats: *metrics.Stats, + archive: *archive_mod.Archive, + archive_closed: *bool, + meta: *meta_store.Store, + cursor_store: *cursor_mod.Store, + verifier: *?verify.Verifier, + import_manager: *timestamp_manager.Manager, + import_manager_err_owned: *bool, + import_ready: *std.atomic.Value(bool), + live_repair_scratch: []const u8, + server_future: *Io.Future(void), + debug_future: *?Io.Future(void), +) !void { import_manager.setTraceContext(steady_trace.context()); // steady compactor: live tombstone set fed from the archive's append // hook, periodic + cap-triggered passes. interval 0 disables all of it // (run() returns immediately, no hook, no rebuild). - var compactor = steady_mod.Compactor.init(allocator, io, &archive, &meta, cfg.data_dir); + var compactor = steady_mod.Compactor.init(allocator, io, archive, meta, cfg.data_dir); defer compactor.deinit(); compactor.interval_ns = cfg.compaction_interval_ns; compactor.tombstone_cap = cfg.compaction_tombstone_cap; compactor.rewrite_workers = cfg.compaction_rewrite_workers; - compactor.stats = &stats; + compactor.stats = stats; compactor.trace_context = steady_trace.context(); if (cfg.compaction_interval_ns != 0) { try compactor.rebuild(); @@ -888,9 +942,9 @@ pub fn main(init: std.process.Init.Minimal) !void { var rebloom_sweep: rebloom_mod.Sweep = .{ .allocator = allocator, .io = io, - .archive = &archive, + .archive = archive, .data_dir = cfg.data_dir, - .stats = &stats, + .stats = stats, .enabled = cfg.rebloom_sweep and cfg.compaction_interval_ns == 0, }; var rebloom_future = try io.concurrent(rebloom_mod.Sweep.run, .{&rebloom_sweep}); @@ -902,13 +956,13 @@ pub fn main(init: std.process.Init.Minimal) !void { // failed-repo retry loop: periodic resync of repos whose backfill // download failed (or that discovery queued). appends ride the // archive writer mutex beside the live consumer. - var retry = retry_mod.Retry.init(allocator, io, &archive, &meta, cfg.data_dir, cfg.relay_http); + var retry = retry_mod.Retry.init(allocator, io, archive, meta, cfg.data_dir, cfg.relay_http); defer retry.deinit(); retry.interval_ns = cfg.retry_interval_ns; retry.workers = cfg.retry_workers; retry.host_workers = cfg.retry_host_workers; retry.max_delay_ns = cfg.retry_max_delay_ns; - retry.stats = &stats; + retry.stats = stats; retry.trace_context = steady_trace.context(); retry.shutdown = &shutdown_requested; if (cfg.retry_interval_ns != 0) log.info("failed-repo retry loop on (interval {d}ns, workers {d}, host workers {d}, max delay {d}ns)", .{ @@ -926,14 +980,14 @@ pub fn main(init: std.process.Init.Minimal) !void { var consumer: ingest.Consumer = .{ .allocator = allocator, .io = io, - .tail = &hot_tail, + .tail = hot_tail, .upstream = cfg.upstream, .cursor = cfg.cursor, .to_stdout = cfg.to_stdout, - .verifier = if (verifier) |*v| v else null, - .cursor_store = &cursor_store, - .archive = &archive, - .stats = &stats, + .verifier = if (verifier.*) |*v| v else null, + .cursor_store = cursor_store, + .archive = archive, + .stats = stats, .shutdown = &shutdown_requested, .slow_upstream = .{ .min_rate = cfg.upstream_slow_min_rate }, }; @@ -941,17 +995,17 @@ pub fn main(init: std.process.Init.Minimal) !void { // The archive's durability hook points into consumer. Keep the fallback // close above consumer's deinit in the unwind stack so every error path // flushes while the callback context is still alive. - defer if (!archive_closed) { + defer if (!archive_closed.*) { archive.close(); - archive_closed = true; + archive_closed.* = true; }; archive.on_durable_ctx = &consumer; archive.on_durable = ingest.Consumer.onArchiveDurable; var pipe: pipeline_mod.Pipeline = .{ .allocator = allocator, .io = io, - .verifier = if (verifier) |*v| v else null, - .stats = &stats, + .verifier = if (verifier.*) |*v| v else null, + .stats = stats, .emit_ctx = &consumer, .emit = ingest.Consumer.emitVerified, .emit_repair = ingest.Consumer.emitRepair, @@ -959,13 +1013,13 @@ pub fn main(init: std.process.Init.Minimal) !void { .cursor_watermark = cfg.cursor orelse 0, .trace_context = steady_trace.context(), }; - var repair_coordinator: ?repair_mod.Coordinator = if (verifier) |*v| .{ + var repair_coordinator: ?repair_mod.Coordinator = if (verifier.*) |*v| .{ .allocator = allocator, .io = io, .verifier = v, .relay_http = cfg.relay_http, .scratch_dir = live_repair_scratch, - .stats = &stats, + .stats = stats, .publish_ctx = &pipe, .publish = pipeline_mod.Pipeline.publishRepair, } else null; @@ -978,7 +1032,7 @@ pub fn main(init: std.process.Init.Minimal) !void { defer pipe.stop(); consumer.pipeline = &pipe; defer import_manager.deinit(); - import_manager_err_owned = false; + import_manager_err_owned.* = false; import_ready.store(true, .release); if (try import_manager.resumeIncomplete()) log.info("resumed incomplete timestamp import", .{}); @@ -1004,8 +1058,8 @@ pub fn main(init: std.process.Init.Minimal) !void { // preserved while storage/pipeline shutdown proceeds concurrently. var listener_shutdowns: Io.Group = .init; errdefer listener_shutdowns.cancel(io); - try listener_shutdowns.concurrent(io, cancelWsServer, .{ &server_future, io }); - if (debug_future) |*future| + try listener_shutdowns.concurrent(io, cancelWsServer, .{ server_future, io }); + if (debug_future.*) |*future| try listener_shutdowns.concurrent(io, cancelWsServer, .{ future, io }); var shutdown_err: ?anyerror = null; @@ -1020,21 +1074,21 @@ pub fn main(init: std.process.Init.Minimal) !void { // happen here while consumer is still alive—not in archive's older // outer defer after consumer teardown. archive.close(); - archive_closed = true; + archive_closed.* = true; try listener_shutdowns.await(io); if (shutdown_err) |err| return @errorCast(err); steady_trace.succeed(); orchestrator_trace.succeed(); - try shutdownTracing(&tracing, cfg.shutdown_timeout); + try shutdownTracing(tracing, cfg.shutdown_timeout); return; } try consumer_future.await(io); try pipe.drainAndStop(); archive.close(); - archive_closed = true; + archive_closed.* = true; steady_trace.succeed(); orchestrator_trace.succeed(); - try shutdownTracing(&tracing, cfg.shutdown_timeout); + try shutdownTracing(tracing, cfg.shutdown_timeout); } fn shutdownTracing(runtime: *observability.Runtime, configured: Io.Duration) !void { diff --git a/src/write_sample.zig b/src/write_sample.zig index 2affbfa..01f0fbb 100644 --- a/src/write_sample.zig +++ b/src/write_sample.zig @@ -58,7 +58,7 @@ pub fn main(init: std.process.Init.Minimal) !void { try p2.put("did", did); var resp = try client.query(zat.Nsid.parse("com.atproto.sync.getRepo").?, p2); defer resp.deinit(); - std.debug.print("status={d} body={d} bytes\n", .{ @intFromEnum(resp.status), resp.body.len }); + std.debug.print("status={d} body={d} bytes\n", .{ @backingInt(resp.status), resp.body.len }); var f = try std.Io.Dir.cwd().createFile(io, out_path, .{}); defer f.close(io); try f.writeStreamingAll(io, resp.body); -- 2.51.2