jetstream client
atproto jetstream client
Something went wrong. Try again.
Zig
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061//! segment export: one sealed jss segment on local disk → one TSV row per//! event on stdout, for loading into a columnar store (duckdb, parquet).//!//! columns: seq, time_us, did, kind, op, collection, rkey. time_us is the//! archive's witnessed time (jetstream v2 semantics), not record creation//! time. record bodies are neither decoded nor emitted (decode_records = false).//!//! run: zig build segment-export -- <path/to/seg_xxxxxxxxxx.jss>
const std = @import("std");const jetstream_sdk = @import("jetstream");
const Exporter = struct { out: *std.Io.Writer, rows: u64 = 0, failed: ?anyerror = null,
pub fn onRow(self: *Exporter, event: jetstream_sdk.Event) bool { self.write(event) catch |err| { self.failed = err; return false; }; return true; }
fn write(self: *Exporter, event: jetstream_sdk.Event) !void { switch (event.payload) { .commit => |c| try self.out.print("{d}\t{d}\t{s}\tcommit\t{s}\t{s}\t{s}\n", .{ event.seq, event.time_us, event.did, @tagName(c.operation), c.collection, c.rkey, }), inline else => |_, tag| try self.out.print("{d}\t{d}\t{s}\t{s}\t\t\t\n", .{ event.seq, event.time_us, event.did, @tagName(tag), }), } self.rows += 1; }};
pub fn main(init: std.process.Init) !void { const allocator = init.gpa; const args = try init.minimal.args.toSlice(init.arena.allocator()); if (args.len != 2) { std.debug.print("usage: zig build segment-export -- <segment.jss>\n", .{}); return error.MissingSegment; } const body = try std.Io.Dir.cwd().readFileAlloc(init.io, args[1], allocator, .limited(1 << 31)); defer allocator.free(body);
var buffer: [1 << 16]u8 = undefined; var stdout = std.Io.File.stdout().writer(init.io, &buffer); var exporter: Exporter = .{ .out = &stdout.interface }; const started = std.Io.Timestamp.now(init.io, .awake).nanoseconds; _ = jetstream_sdk.ArchiveBackfill.replaySegment(allocator, body, .{ .host = "", .decode_records = false }, &exporter) catch |err| switch (err) { error.Stopped => return exporter.failed orelse err, else => return err, }; try stdout.interface.flush(); const elapsed_ms = @divTrunc(std.Io.Timestamp.now(init.io, .awake).nanoseconds - started, std.time.ns_per_ms); std.debug.print("{d} rows from {d} bytes in {d} ms\n", .{ exporter.rows, body.len, elapsed_ms });}