jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294//! dev tool: write a sample active segment to the path in argv[1]//! (used to differentially verify the writer against upstream's//! inspect-segment; see docs/jss-format-v1.md)
const std = @import("std");const archive_mod = @import("internal/storage/archive.zig");const segment = @import("internal/storage/segment.zig");const writer_mod = @import("internal/storage/segment_writer.zig");
pub fn main(init: std.process.Init.Minimal) !void { const allocator = std.heap.c_allocator; var threaded: std.Io.Threaded = .init(allocator, .{}); defer threaded.deinit(); const io = threaded.io();
var args = init.args.iterate(); _ = args.next(); const path = args.next() orelse return error.MissingPath;
// --parse <car-file>: run the backfill CAR-to-rows conversion on a file if (std.mem.eql(u8, path, "--parse")) { const car_path = args.next() orelse return error.MissingPath; const repos = @import("internal/ingest/backfill/repos.zig"); const seg = @import("internal/storage/segment.zig"); var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const bytes = try std.Io.Dir.cwd().readFileAlloc(io, car_path, arena.allocator(), .limited(1 << 30)); var drops: repos.DropCounts = .{}; const Counter = struct { rows: u64 = 0, fn emit(self: *@This(), _: seg.Event) anyerror!void { self.rows += 1; } }; var counter: Counter = .{}; const rev = repos.carToRows(arena.allocator(), bytes, "", 0, .create, &drops, &counter, Counter.emit) catch |err| { std.debug.print("carToRows failed: {t}\n", .{err}); return err; }; std.debug.print("ok rev={s} rows={d} drops={d}/{d}/{d}\n", .{ rev, counter.rows, drops.invalid_collection, drops.invalid_rkey, drops.field_too_long, }); return; }
// --fetch <did> <out>: fetch a repo CAR via zat XrpcClient (backfill diagnostics) if (std.mem.eql(u8, path, "--fetch")) { const did = args.next() orelse return error.MissingPath; const out_path = args.next() orelse return error.MissingPath; const zat = @import("zat"); var client = zat.XrpcClient.init(io, allocator, "http://localhost:7777"); defer client.deinit(); var p2 = std.StringHashMap([]const u8).init(allocator); 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 }); var f = try std.Io.Dir.cwd().createFile(io, out_path, .{}); defer f.close(io); try f.writeStreamingAll(io, resp.body); return; }
// --status <dir>: seed real repo/host metadata for the HTTP contract. if (std.mem.eql(u8, path, "--status")) { const dir = args.next() orelse return error.MissingPath; const meta_store = @import("internal/storage/meta_store.zig"); const repo_store = @import("internal/ingest/backfill/repo_store.zig"); var meta = try meta_store.Store.open(allocator, dir); defer meta.deinit(); var store = try repo_store.Store.init(allocator, io, &meta); defer store.deinit(); const now_us = std.Io.Timestamp.now(io, .real).toMicroseconds(); var rows: [12]repo_store.Transition = undefined; for (&rows, 0..) |*row, i| row.* = .{ .did = try std.fmt.allocPrint(allocator, "did:plc:large{d}", .{i}), .state = .{ .status = .complete, .host = "pds-large.example", .last_attempt_us = now_us - 2 * std.time.us_per_s }, }; defer for (rows) |row| allocator.free(row.did); try store.putMany(&rows); try store.put("did:plc:failing", .{ .status = .failed, .host = "pds-failing.example", .last_error = "<script>alert(1)</script>", .error_class = .http_5xx, .attempts = 1, .last_attempt_us = now_us - std.time.us_per_s, }); try store.put("did:plc:failing", .{ .status = .failed, .host = "pds-failing.example", .last_error = "HTTP 503 again", .error_class = .http_5xx, .attempts = 2, .last_attempt_us = now_us, }); return; }
// --retry-status <dir> <same|unique> <count>: seed durable failed rows // for the production retry-control receipt. if (std.mem.eql(u8, path, "--retry-status")) { const dir = args.next() orelse return error.MissingPath; const mode = args.next() orelse return error.MissingMode; const count = try std.fmt.parseInt(usize, args.next() orelse return error.MissingCount, 10); if (!std.mem.eql(u8, mode, "same") and !std.mem.eql(u8, mode, "unique")) return error.InvalidMode; const meta_store = @import("internal/storage/meta_store.zig"); const repo_store = @import("internal/ingest/backfill/repo_store.zig"); var meta = try meta_store.Store.open(allocator, dir); defer meta.deinit(); var store = try repo_store.Store.init(allocator, io, &meta); defer store.deinit(); for (0..count) |i| { var did_buf: [96]u8 = undefined; const did = try std.fmt.bufPrint(&did_buf, "did:plc:retryconfig{d}", .{i}); var host_buf: [96]u8 = undefined; const host = if (std.mem.eql(u8, mode, "same")) "same.retry.invalid" else try std.fmt.bufPrint(&host_buf, "host{d}.retry.invalid", .{i}); try store.put(did, .{ .status = .failed, .active = true, .host = host, .last_error = "seeded retry failure", .error_class = .http_5xx, .attempts = 1, }); } return; }
// --retry-assert <dir> <count> <min-delay-us> <max-delay-us>: reopen // RocksDB after the process exits and prove one persisted retry failure // per seed plus the configured jitter/cap window. if (std.mem.eql(u8, path, "--retry-assert")) { const dir = args.next() orelse return error.MissingPath; const count = try std.fmt.parseInt(usize, args.next() orelse return error.MissingCount, 10); const min_delay_us = try std.fmt.parseInt(i64, args.next() orelse return error.MissingDelay, 10); const max_delay_us = try std.fmt.parseInt(i64, args.next() orelse return error.MissingDelay, 10); const meta_store = @import("internal/storage/meta_store.zig"); const repo_store = @import("internal/ingest/backfill/repo_store.zig"); var meta = try meta_store.Store.open(allocator, dir); defer meta.deinit(); var store = try repo_store.Store.init(allocator, io, &meta); defer store.deinit(); var saw_jitter = false; for (0..count) |i| { var did_buf: [96]u8 = undefined; const did = try std.fmt.bufPrint(&did_buf, "did:plc:retryconfig{d}", .{i}); const state = (try store.get(allocator, did)) orelse return error.MissingRetryState; defer store.freeState(allocator, state); if (state.status != .failed or state.retry_count != 1 or state.attempts != 2) return error.InvalidRetryState; const delay_us = state.next_attempt_us - state.last_attempt_us; if (delay_us < min_delay_us or delay_us > max_delay_us) return error.InvalidRetryDelay; saw_jitter = saw_jitter or delay_us > min_delay_us; } if (count > 1 and !saw_jitter) return error.MissingRetryJitter; std.debug.print("retry state: PASS ({d} rows, delay {d}..{d}us)\n", .{ count, min_delay_us, max_delay_us }); return; }
// --archive <dir>: drive the Archive instead (rotation + seal + active) if (std.mem.eql(u8, path, "--archive")) { const dir = args.next() orelse return error.MissingPath; var a = try archive_mod.Archive.init(allocator, io, dir); defer a.deinit(); a.max_events_per_block = 500; a.writer.max_events_per_block = 500; a.max_segment_bytes = 8 * 1024; var n: u64 = 1; while (n <= 5000) : (n += 1) { _ = try a.append(.{ .seq = 0, .witnessed_at = @intCast(1_700_000_000_000_000 + n * 1000), .indexed_at = 0, .kind = if (n % 100 == 0) .identity else .create, .did = if (n % 2 == 0) "did:plc:streamsampleaccount" else "did:plc:otherstreamaccount", .collection = if (n % 100 == 0) "" else "app.bsky.feed.post", .rkey = if (n % 100 == 0) "" else "3l3qo2vuowo2b", .rev = if (n % 100 == 0) "" else "3l3qo2vutsw2b", .payload = "\xa1\x64\x74\x65\x78\x74\x62\x68\x69", }, @intCast(n)); } a.close(); return; }
// --plan-config-archive <dir>: three real sealed segments with a 3/4 // collection density and two disjoint matching block ranges per segment. // This lets the production HTTP contract distinguish planner threshold, // pagination, and unlimited-page configuration without fixture verdicts. if (std.mem.eql(u8, path, "--plan-config-archive")) { const dir = args.next() orelse return error.MissingPath; var a = try archive_mod.Archive.init(allocator, io, dir); defer a.deinit(); a.max_events_per_block = 1; a.writer.max_events_per_block = 1; a.max_segment_bytes = 1024 * 1024; var seq: u64 = 1; for (0..3) |_| { for (0..4) |at| { _ = try a.append(.{ .seq = 0, .witnessed_at = @intCast(1_700_000_000_000_000 + seq * 1000), .indexed_at = 0, .kind = .create, .did = "did:plc:planconfigaccount", .collection = if (at == 1) "app.bsky.feed.like" else "app.bsky.feed.post", .rkey = "3l3qo2vuowo2b", .rev = "3l3qo2vutsw2b", .payload = "\xa1\x64\x74\x65\x78\x74\x62\x68\x69", }, @intCast(seq)); seq += 1; } try a.rotate(); } a.close(); return; }
// --compaction-config-archive <dir>: two sealed create/delete pairs. // cap=1 must compact them as two chunks; cap=0 must use one unlimited // chunk. Two physical segments also expose the rewrite-worker bound. if (std.mem.eql(u8, path, "--compaction-config-archive")) { const dir = args.next() orelse return error.MissingPath; var a = try archive_mod.Archive.init(allocator, io, dir); defer a.deinit(); a.max_events_per_block = 1; a.writer.max_events_per_block = 1; a.max_segment_bytes = 1024 * 1024; const post = "app.bsky.feed.post"; for (0..2) |i| { const seq: u64 = @intCast(i * 2 + 1); const rkey = if (i == 0) "3l3qo2vuowo2a" else "3l3qo2vuowo2b"; _ = try a.append(.{ .seq = 0, .witnessed_at = @intCast(1_700_000_000_000_000 + seq * 1000), .indexed_at = 0, .kind = .create, .did = "did:plc:compactionconfigaccount", .collection = post, .rkey = rkey, .rev = "3l3qo2vutsw2b", .payload = "\xa1\x64\x74\x65\x78\x74\x62\x68\x69", }, @intCast(seq)); _ = try a.append(.{ .seq = 0, .witnessed_at = @intCast(1_700_000_000_000_000 + (seq + 1) * 1000), .indexed_at = 0, .kind = .delete, .did = "did:plc:compactionconfigaccount", .collection = post, .rkey = rkey, .rev = "3l3qo2vutsw2b", .payload = "", }, @intCast(seq + 1)); try a.rotate(); } a.close(); return; }
var w = try writer_mod.ActiveWriter.init(allocator); defer w.deinit(); w.max_events_per_block = 1000;
var seq: u64 = 1; while (seq <= 2500) : (seq += 1) { try w.append(.{ .seq = seq, .witnessed_at = @intCast(1_700_000_000_000_000 + seq * 1000), .indexed_at = 0, .kind = if (seq % 100 == 0) .identity else if (seq % 10 == 0) .delete else .create, .did = if (seq % 2 == 0) "did:plc:streamsampleaccount" else "did:plc:otherstreamaccount", .collection = if (seq % 100 == 0) "" else if (seq % 3 == 0) "app.bsky.feed.like" else "app.bsky.feed.post", .rkey = if (seq % 100 == 0) "" else "3l3qo2vuowo2b", .rev = if (seq % 100 == 0) "" else "3l3qo2vutsw2b", .payload = if (seq % 10 == 0 and seq % 100 != 0) "" else "\xa1\x64\x74\x65\x78\x74\x62\x68\x69", }); } try w.seal();
var f = try std.Io.Dir.cwd().createFile(io, path, .{}); defer f.close(io); try f.writeStreamingAll(io, w.bytes());}