const std = @import("std"); const sdk = @import("jetstream"); const LoadHandler = struct { io: std.Io, samples: []i64, count: usize = 0, failed: bool = false, delay_ns: i64, rate: i64, epoch: i64 = 0, pub fn onEvent(self: *@This(), event: sdk.Event) bool { if (event.seq != self.count + 1) { self.failed = true; return false; } var scheduled = std.fmt.parseInt(i64, event.payload.commit.rkey, 10) catch { self.failed = true; return false; }; if (self.count == 0) self.epoch = scheduled; scheduled = self.epoch + @divTrunc(@as(i64, @intCast(self.count)) * std.time.ns_per_s, self.rate); self.samples[self.count] = @intCast(std.Io.Clock.real.now(self.io).toNanoseconds() - scheduled); self.count += 1; if (self.delay_ns > 0) self.io.sleep(std.Io.Duration.fromNanoseconds(self.delay_ns), .awake) catch { self.failed = true; return false; }; return self.count < self.samples.len; } pub fn onError(self: *@This(), _: anyerror) bool { self.failed = true; return false; } }; const CatchupHandler = struct { target: usize, count: usize = 0, digest: u64 = 14695981039346656037, failed: bool = false, pub fn onEvent(self: *@This(), event: sdk.Event) bool { self.count += 1; if (event.seq != self.count) { self.failed = true; return false; } for (event.payload.commit.record_cbor) |byte| self.digest = (self.digest ^ byte) *% 1099511628211; return self.count < self.target; } pub fn onError(self: *@This(), _: anyerror) bool { self.failed = true; return false; } }; const Handler = struct { target: usize, count: usize = 0, stop_on_error: bool, error_count: usize = 0, max_errors: usize = 0, pub fn onEvent(self: *@This(), event: sdk.Event) bool { const c = event.payload.commit; std.debug.print("EVENT\t{d}\t{s}\t{s}\t{s}\t{s}\n", .{ event.seq, event.did, c.collection, c.rkey, @tagName(c.operation) }); self.count += 1; return self.count < self.target; } pub fn onError(self: *@This(), err: anyerror) bool { std.debug.print("ERROR\nDETAIL {s}\n", .{@errorName(err)}); self.error_count += 1; return !self.stop_on_error and (self.max_errors == 0 or self.error_count < self.max_errors); } }; pub fn main(init: std.process.Init) !void { const args = try init.minimal.args.toSlice(init.arena.allocator()); if (args.len == 4 and std.mem.eql(u8, args[2], "catchup")) { const count = try std.fmt.parseInt(usize, args[3], 10); if (count == 0 or count > 1_000_000) return error.InvalidCount; var handler = CatchupHandler{ .target = count + 1 }; std.debug.print("READY\n", .{}); const start = std.Io.Clock.awake.now(init.io); try sdk.subscribe(init.io, init.gpa, .{ .hosts = &.{args[1]}, .after_seq = 0, .zstd_compression = false, .concurrency = 4 }, &handler); if (handler.failed or handler.count != count + 1) return error.InvalidDelivery; const elapsed = start.durationTo(std.Io.Clock.awake.now(init.io)).nanoseconds; std.debug.print("CATCHUP {d} {d} {d}\n", .{ handler.count, handler.digest, elapsed }); return; } if (args.len == 7 and std.mem.eql(u8, args[2], "load")) { const count = try std.fmt.parseInt(usize, args[3], 10); if (count == 0 or count > 1_000_000) return error.InvalidCount; const samples = try init.gpa.alloc(i64, count); defer init.gpa.free(samples); std.debug.print("READY\n", .{}); var handler = LoadHandler{ .io = init.io, .samples = samples, .delay_ns = try std.fmt.parseInt(i64, args[5], 10), .rate = try std.fmt.parseInt(i64, args[6], 10) }; try sdk.subscribe(init.io, init.gpa, .{ .hosts = &.{args[1]}, .zstd_compression = std.mem.eql(u8, args[4], "compressed") }, &handler); if (handler.failed or handler.count != count) return error.InvalidDelivery; std.mem.sort(i64, samples, {}, std.sort.asc(i64)); std.debug.print("LOAD {d} {d} {d} {d}\n", .{ count, samples[count / 2], samples[@min(count - 1, count * 99 / 100)], samples[count - 1] }); return; } if (args.len != 3) return error.ExpectedHostAndScenario; const mode = args[2]; const exact = std.mem.eql(u8, mode, "exact"); const wildcard = std.mem.eql(u8, mode, "wildcard"); const resuming = std.mem.eql(u8, mode, "resume"); const kinds = std.mem.startsWith(u8, mode, "archive-kinds"); const many_dids = std.mem.eql(u8, mode, "many-dids"); var did_storage: [120][32]u8 = undefined; var dids: [120][]const u8 = undefined; const alphabet = "abcdefghijklmnopqrstuvwxyz234567"; for (&did_storage, 0..) |*did, i| { @memcpy(did[0..8], "did:plc:"); @memset(did[8..30], 'a'); did[30] = alphabet[i / 32]; did[31] = alphabet[i % 32]; dids[i] = did; } var handler = Handler{ .target = if (exact or resuming or kinds) 2 else 3, .stop_on_error = std.mem.endsWith(u8, mode, "-stop"), .max_errors = if (many_dids) 3 else 0 }; try sdk.subscribe(init.io, init.gpa, .{ .hosts = &.{args[1]}, .collections = if (exact) &.{"app.bsky.feed.post"} else if (wildcard) &.{"app.bsky.feed.*"} else &.{}, .dids = if (many_dids) &dids else &.{}, .kinds = if (kinds) &.{.commit} else &.{}, .backoff_min_ns = std.time.ns_per_ms, .backoff_max_ns = std.time.ns_per_ms, .live_cursor = if (resuming) 1 else null, .after_seq = if (std.mem.startsWith(u8, mode, "archive-")) 0 else null, .snapshot_only = std.mem.eql(u8, mode, "archive-pages") or kinds, .zstd_compression = !std.mem.eql(u8, mode, "plain") and !std.mem.startsWith(u8, mode, "archive-") and !std.mem.startsWith(u8, mode, "terminal-") and !many_dids, }, &handler); }