jetstream client
atproto jetstream client
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392//! replay a Streamplace broadcast's chat, including contiguous rollover records.//!//! run: zig build example-streamplace-chat -- iame.li//! zig build example-streamplace-chat -- at://did:plc:.../place.stream.livestream/3...//!//! the archive endpoints are bearer-gated; set JETSTREAM_ARCHIVE_API_KEY or//! the example skips the archive and replays over the live socket only.
const std = @import("std");const zat = @import("zat");const jetstream_sdk = @import("jetstream");
const rollover_tolerance_us: i64 = 2 * std.time.us_per_s;const ingestion_grace_us: i64 = 5 * std.time.us_per_min;const archive_host = "https://stream.waow.tech";
const Segment = struct { uri: []u8, rkey: []u8, title: []u8, created_at: []u8, ended_at: ?[]u8, url: []u8,
fn deinit(self: *Segment, allocator: std.mem.Allocator) void { allocator.free(self.uri); allocator.free(self.rkey); allocator.free(self.title); allocator.free(self.created_at); if (self.ended_at) |value| allocator.free(value); allocator.free(self.url); }};
const ChatPrinter = struct { allocator: std.mem.Allocator, streamer_did: []const u8, first_stream_rkey: []const u8, ended_at_us: ?i64, count: usize = 0, /// the archive pass and the live replay overlap at the seam by design /// (delivery is at-least-once); "did/rkey" keys dedup the seam seen: std.StringHashMapUnmanaged(void) = .empty,
fn deinit(self: *ChatPrinter) void { var keys = self.seen.keyIterator(); while (keys.next()) |key| self.allocator.free(key.*); self.seen.deinit(self.allocator); }
pub fn onEvent(self: *ChatPrinter, event: zat.JetstreamEvent) void { const commit = switch (event) { .commit => |c| c, else => return, }; if (commit.operation != .create) return; const record = commit.record orelse return;
const streamer = zat.json.getString(record, "streamer") orelse return; if (!std.mem.eql(u8, streamer, self.streamer_did)) return;
// Chat records use TID keys. This excludes messages from broadcasts // before the selected logical session, even if they arrive late. if (std.mem.order(u8, commit.rkey, self.first_stream_rkey) == .lt) return;
const created_at = zat.json.getString(record, "createdAt") orelse return; if (self.ended_at_us) |end| { const created_at_us = parseUtcMicros(created_at) catch return; if (created_at_us > end) return; }
var key_buf: [512]u8 = undefined; const key = std.fmt.bufPrint(&key_buf, "{s}/{s}", .{ commit.did, commit.rkey }) catch return; const entry = self.seen.getOrPut(self.allocator, key) catch return; if (entry.found_existing) return; entry.key_ptr.* = self.allocator.dupe(u8, key) catch { _ = self.seen.remove(key); return; };
const text = zat.json.getString(record, "text") orelse ""; self.count += 1; std.debug.print("[{s}] {s}: {s}\n", .{ created_at, commit.did, text }); }
pub fn onConnect(self: *ChatPrinter, host: []const u8) void { std.debug.print("connected to {s} ({d} observed messages so far)\n", .{ host, self.count }); }
pub fn onError(_: *ChatPrinter, err: anyerror) void { std.debug.print("stream error: {s}\n", .{@errorName(err)}); }};
pub fn main(init: std.process.Init) !void { const allocator = init.gpa; const args = try init.minimal.args.toSlice(init.arena.allocator()); const api_key = init.minimal.environ.getAlloc(init.arena.allocator(), "JETSTREAM_ARCHIVE_API_KEY") catch null; if (args.len != 2) { std.debug.print("usage: zig build example-streamplace-chat -- <handle|livestream-at-uri>\n", .{}); return error.MissingStream; }
var handle_resolver = zat.HandleResolver.init(init.io, allocator); defer handle_resolver.deinit();
var streamer_did: []const u8 = undefined; var livestream_uri: []const u8 = undefined; if (zat.AtUri.parse(args[1])) |uri| { const collection = uri.collection() orelse return error.MissingCollection; if (!std.mem.eql(u8, collection, "place.stream.livestream")) return error.NotLivestreamUri; _ = zat.Tid.parse(uri.rkey() orelse return error.MissingLivestreamRkey) orelse return error.InvalidLivestreamRkey;
const authority = uri.authority(); streamer_did = if (zat.Did.parse(authority) != null) try allocator.dupe(u8, authority) else try handle_resolver.resolve(zat.Handle.parse(authority) orelse return error.InvalidAuthority); livestream_uri = try allocator.dupe(u8, args[1]); } else { const handle = zat.Handle.parse(args[1]) orelse { std.debug.print("expected an ATProto handle or livestream AT-URI: {s}\n", .{args[1]}); return error.InvalidStream; }; streamer_did = try handle_resolver.resolve(handle); livestream_uri = try latestLivestreamUri(init.io, allocator, streamer_did); } defer allocator.free(streamer_did); defer allocator.free(livestream_uri);
const at_uri = zat.AtUri.parse(livestream_uri) orelse return error.InvalidLivestreamUri; const first_rkey = at_uri.rkey() orelse return error.MissingLivestreamRkey; const first_tid = zat.Tid.parse(first_rkey) orelse return error.InvalidLivestreamRkey;
var did_resolver = zat.DidResolver.init(init.io, allocator); defer did_resolver.deinit(); var doc = try did_resolver.resolve(zat.Did.parse(streamer_did).?); defer doc.deinit(); const pds_endpoint = doc.pdsEndpoint() orelse return error.NoPdsEndpoint;
var pds = zat.XrpcClient.initWithUserAgent( init.io, allocator, pds_endpoint, "zat-streamplace-chat-example/1.1", ); defer pds.deinit();
var selected = try getLivestream(&pds, allocator, streamer_did, first_rkey); defer selected.deinit(allocator);
var newer: std.ArrayList(Segment) = .empty; defer { for (newer.items) |*segment| segment.deinit(allocator); newer.deinit(allocator); } if (selected.ended_at != null) { try collectNewerLivestreams(&pds, allocator, streamer_did, first_rkey, &newer); }
var replay_end = selected.ended_at; var previous_end_us = if (selected.ended_at) |end| try parseUtcMicros(end) else null; var previous_url: []const u8 = selected.url; var segment_count: usize = 1;
std.debug.print("logical Streamplace session:\n", .{}); printSegment(&selected);
// listRecords is newest-first. Walking the collected prefix backwards // visits the record immediately after the selected one first. var i = newer.items.len; while (i > 0 and previous_end_us != null) { i -= 1; const candidate = &newer.items[i]; const candidate_start_us = parseUtcMicros(candidate.created_at) catch break; const gap_us = candidate_start_us - previous_end_us.?; const same_target = previous_url.len == 0 or candidate.url.len == 0 or std.mem.eql(u8, previous_url, candidate.url); if (!same_target or gap_us < -rollover_tolerance_us or gap_us > rollover_tolerance_us) break;
segment_count += 1; printSegment(candidate); replay_end = candidate.ended_at; previous_end_us = if (candidate.ended_at) |end| try parseUtcMicros(end) else null; previous_url = candidate.url; }
std.debug.print("\n", .{}); if (segment_count > 1) { std.debug.print("joined {d} contiguous livestream records; Streamplace rolled over the AT-URI mid-session.\n", .{segment_count}); }
var printer = ChatPrinter{ .allocator = allocator, .streamer_did = streamer_did, .first_stream_rkey = first_rkey, .ended_at_us = if (replay_end) |end| try parseUtcMicros(end) else null, }; defer printer.deinit();
// The livestream record's TID is its creation time in microseconds. // Rewind one second so events at the boundary cannot be skipped. const start_us = @as(i64, @intCast(first_tid.timestamp())) - std.time.us_per_s; const end_cursor = if (replay_end) |end| (try parseUtcMicros(end)) + ingestion_grace_us else null;
// Historical chat comes from the archive, not from live-replay retention: // stream.waow.tech archives every event, so any session since its live // capture began replays completely — including messages later deleted. // The witnessed-time window maps to plan seq bounds, which keeps the // fetch to the handful of blocks the session actually touches. const bounds: jetstream_sdk.ArchiveBackfill.SeqBounds = if (api_key != null) try jetstream_sdk.ArchiveBackfill.fetchSeqBounds(init.io, allocator, archive_host, start_us, end_cursor, api_key) else .{ .covered = false, .after_seq = 0, .before_seq = null }; var archive_boundary_us: ?i64 = null; if (api_key == null) { std.debug.print("JETSTREAM_ARCHIVE_API_KEY unset; skipping the archive. Live replay may not include the full session history.\n", .{}); } else if (bounds.covered) { std.debug.print("replaying from the archive ({s})...\n\n", .{archive_host}); const result = try jetstream_sdk.ArchiveBackfill.run(init.io, allocator, .{ .host = archive_host, .collections = &.{"place.stream.chat.message"}, .after_seq = bounds.after_seq, .before_seq = bounds.before_seq, .concurrency = 6, .api_key = api_key, }, &printer); archive_boundary_us = result.last_time_us; std.debug.print("\narchive replay: {d} messages from {d} blocks\n", .{ printer.count, result.blocks_decoded }); } else { std.debug.print("session predates the archive's live capture; only live replay is available.\n", .{}); }
// The newest events live in the archive's unsealed tail, so a complete // replay always finishes over the live socket from the covered boundary // (the seam overlaps; the printer dedups). For a live session this is // also the follow. std.debug.print("{s}\n\n", .{ if (replay_end == null) "catching up over the live socket, then following live..." else "finishing the session tail over the live socket...", }); var client = zat.JetstreamClient.init(init.io, allocator, .{ .wanted_collections = &.{"place.stream.chat.message"}, .cursor = if (archive_boundary_us) |t| t - std.time.us_per_s else start_us, // Stop close to the broadcast end, with room for delayed ingestion. // The idle timeout handles the case where no matching event lands // exactly beyond that boundary. .end_cursor = end_cursor, .replay_idle_timeout_ms = if (end_cursor != null) 10_000 else null, }); defer client.deinit();
try client.subscribe(&printer); std.debug.print("\nreplay complete: {d} observed messages across {d} segment{s}\n", .{ printer.count, segment_count, if (segment_count == 1) "" else "s", });}
fn printSegment(segment: *const Segment) void { std.debug.print("- {s}\n {s} -> {s}\n {s}\n", .{ segment.title, segment.created_at, segment.ended_at orelse "live", segment.uri, });}
fn getLivestream( client: *zat.XrpcClient, allocator: std.mem.Allocator, repo: []const u8, rkey: []const u8,) !Segment { const params = [_]zat.XrpcClient.QueryParam{ .{ .name = "repo", .value = repo }, .{ .name = "collection", .value = "place.stream.livestream" }, .{ .name = "rkey", .value = rkey }, }; const method = zat.Nsid.parse("com.atproto.repo.getRecord").?; var response = try client.queryParams(method, ¶ms); defer response.deinit(); if (!response.ok()) { std.debug.print("getRecord failed ({d}): {s}\n", .{ @backingInt(response.status), response.body }); return error.LivestreamLookupFailed; }
var parsed = try response.json(); defer parsed.deinit(); const uri = zat.json.getString(parsed.value, "uri") orelse return error.MissingLivestreamUri; const value = zat.json.getPath(parsed.value, "value") orelse return error.MissingLivestreamRecord; return copySegment(allocator, uri, value);}
fn collectNewerLivestreams( client: *zat.XrpcClient, allocator: std.mem.Allocator, repo: []const u8, selected_rkey: []const u8, out: *std.ArrayList(Segment),) !void { var cursor: ?[]u8 = null; defer if (cursor) |value| allocator.free(value);
while (true) { var params: std.ArrayList(zat.XrpcClient.QueryParam) = .empty; defer params.deinit(allocator); try params.appendSlice(allocator, &.{ .{ .name = "repo", .value = repo }, .{ .name = "collection", .value = "place.stream.livestream" }, .{ .name = "limit", .value = "100" }, }); if (cursor) |value| try params.append(allocator, .{ .name = "cursor", .value = value });
const method = zat.Nsid.parse("com.atproto.repo.listRecords").?; var response = try client.queryParams(method, params.items); defer response.deinit(); if (!response.ok()) return error.LivestreamListFailed;
var parsed = try response.json(); defer parsed.deinit(); const records = zat.json.getArray(parsed.value, "records") orelse return error.MissingRecords; for (records) |record| { const uri = zat.json.getString(record, "uri") orelse continue; const parsed_uri = zat.AtUri.parse(uri) orelse continue; const rkey = parsed_uri.rkey() orelse continue; if (std.mem.eql(u8, rkey, selected_rkey)) return; const value = zat.json.getPath(record, "value") orelse continue; try out.append(allocator, try copySegment(allocator, uri, value)); }
const next = zat.json.getString(parsed.value, "cursor") orelse return; const next_copy = try allocator.dupe(u8, next); if (cursor) |value| allocator.free(value); cursor = next_copy; }}
fn copySegment(allocator: std.mem.Allocator, uri: []const u8, value: std.json.Value) !Segment { const parsed_uri = zat.AtUri.parse(uri) orelse return error.InvalidLivestreamUri; const rkey = parsed_uri.rkey() orelse return error.MissingLivestreamRkey; const created_at = zat.json.getString(value, "createdAt") orelse return error.MissingCreatedAt; const ended_at = zat.json.getString(value, "endedAt"); return .{ .uri = try allocator.dupe(u8, uri), .rkey = try allocator.dupe(u8, rkey), .title = try allocator.dupe(u8, zat.json.getString(value, "title") orelse "untitled stream"), .created_at = try allocator.dupe(u8, created_at), .ended_at = if (ended_at) |end| try allocator.dupe(u8, end) else null, .url = try allocator.dupe(u8, zat.json.getString(value, "url") orelse ""), };}
fn parseUtcMicros(value: []const u8) !i64 { const dt = zat.Datetime.parse(value) orelse return error.InvalidDatetime; return dt.micros;}
fn latestLivestreamUri( io: std.Io, allocator: std.mem.Allocator, streamer_did: []const u8,) ![]u8 { var url_buffer: [512]u8 = undefined; const url = try std.fmt.bufPrint(&url_buffer, "https://stream.place/api/livestream/{s}", .{streamer_did}); var transport = zat.HttpTransport.initWithUserAgent(io, allocator, "zat-streamplace-chat-example/1.1"); defer transport.deinit(); var response = try transport.fetch(.{ .url = url, .max_response_size = 1024 * 1024 }); defer response.deinit(allocator); if (response.status != .ok) return error.LivestreamLookupFailed;
var parsed = try std.json.parseFromSlice(std.json.Value, allocator, response.body, .{}); defer parsed.deinit(); const uri = zat.json.getString(parsed.value, "uri") orelse return error.MissingLivestreamUri; return try allocator.dupe(u8, uri);}
test "parse Streamplace UTC datetime to microseconds" { try std.testing.expectEqual(@as(i64, 0), try parseUtcMicros("1970-01-01T00:00:00Z")); try std.testing.expectEqual(@as(i64, 1785780238315000), try parseUtcMicros("2026-08-03T18:03:58.315Z")); try std.testing.expectError(error.InvalidDatetime, parseUtcMicros("2026-02-30T00:00:00Z"));}