jetstream client
atproto jetstream client
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307//! minimal libzstd decode bindings//!//! backed by the vendored facebook/zstd v1.5.7 source (pinned by url+hash in//! build.zig.zon, decode-only subset compiled in build.zig — no system//! library). two reasons this exists instead of std.compress.zstd://! measured ~20x faster on jss archive blocks, and libzstd verifies the//! frame content checksum (present on every jss block), which std does not//! implement.
const std = @import("std");const Allocator = std.mem.Allocator;
const c = struct { extern fn ZSTD_getFrameContentSize(src: [*]const u8, src_size: usize) c_ulonglong; extern fn ZSTD_decompress(dst: [*]u8, dst_cap: usize, src: [*]const u8, src_len: usize) usize; extern fn ZSTD_isError(code: usize) c_uint; const FrameHeader = extern struct { frame_content_size: c_ulonglong, window_size: c_ulonglong, block_size_max: c_uint, frame_type: c_uint, header_size: c_uint, dict_id: c_uint, checksum_flag: c_uint, reserved1: c_uint, reserved2: c_uint, }; extern fn ZSTD_getFrameHeader(header: *FrameHeader, src: [*]const u8, src_size: usize) usize; extern fn ZSTD_findFrameCompressedSize(src: [*]const u8, src_size: usize) usize; extern fn ZSTD_decompressBound(src: [*]const u8, src_size: usize) c_ulonglong; extern fn ZSTD_getErrorCode(code: usize) c_uint; const error_dst_size_too_small: c_uint = 70; // ZSTD_error_dstSize_tooSmall, zstd_errors.h extern fn ZSTD_createDCtx() ?*anyopaque; extern fn ZSTD_freeDCtx(dctx: ?*anyopaque) usize; extern fn ZSTD_createDDict(dict: [*]const u8, dict_size: usize) ?*anyopaque; extern fn ZSTD_freeDDict(ddict: ?*anyopaque) usize; extern fn ZSTD_decompress_usingDDict(dctx: ?*anyopaque, dst: [*]u8, dst_cap: usize, src: [*]const u8, src_len: usize, ddict: ?*anyopaque) usize; extern fn ZSTD_getDictID_fromDict(dict: [*]const u8, dict_size: usize) c_uint;
const contentsize_unknown: c_ulonglong = std.math.maxInt(c_ulonglong); // (ull)-1 const contentsize_error: c_ulonglong = std.math.maxInt(c_ulonglong) - 1; // (ull)-2};
/// jss spec cap on a decompressed block; also the zstd-bomb guardpub const max_decompressed: usize = 1 << 30;
pub const DecompressError = error{ /// malformed frame, size lie, or content checksum mismatch DecompressFailed, BlockOversize, OutOfMemory,};
/// Decompress an archive frame, including small upstream blocks that omit/// content size. Bound output allocation and preserve checksum verification.pub fn decompressAlloc(allocator: Allocator, frame: []const u8) DecompressError![]u8 { if (frame.len == 0) return error.DecompressFailed; const size = c.ZSTD_getFrameContentSize(frame.ptr, frame.len); if (size == c.contentsize_unknown) { var header: c.FrameHeader = undefined; if (c.ZSTD_getFrameHeader(&header, frame.ptr, frame.len) != 0) return error.DecompressFailed; // Match upstream's default 512 MiB window ceiling as well as the // archive's 1 GiB output ceiling. A bound can overestimate output. if (header.window_size > 512 << 20) return error.BlockOversize; const bound = c.ZSTD_decompressBound(frame.ptr, frame.len); if (bound == c.contentsize_error) return error.DecompressFailed; const dst = try allocator.alloc(u8, @intCast(@min(bound, max_decompressed))); errdefer allocator.free(dst); const written = c.ZSTD_decompress(dst.ptr, dst.len, frame.ptr, frame.len); if (c.ZSTD_isError(written) != 0) { if (c.ZSTD_getErrorCode(written) == c.error_dst_size_too_small) return error.BlockOversize; return error.DecompressFailed; } return try allocator.realloc(dst, written); } if (size == c.contentsize_error) return error.DecompressFailed; if (size > max_decompressed) return error.BlockOversize;
const dst = try allocator.alloc(u8, @intCast(size)); errdefer allocator.free(dst); const written = c.ZSTD_decompress(dst.ptr, dst.len, frame.ptr, frame.len); if (c.ZSTD_isError(written) != 0 or written != dst.len) return error.DecompressFailed; return dst;}
/// dictionary ID declared in a dictionary blob's header (zero = not a/// structured dictionary). the live tail pins this ID in ?zstdDictionary=/// and the server refuses unknown IDs pre-upgrade.pub fn dictId(dict: []const u8) u32 { if (dict.len == 0) return 0; return @intCast(c.ZSTD_getDictID_fromDict(dict.ptr, dict.len));}
/// dictionary-seeded decoder for live subscribeEvents binary frames./// the decompressed-size cap mirrors the connection's read limit: an/// uncompressed frame must fit the limit on the wire, so its compressed/// twin may not expand past it either (upstream newZstdDecoder).pub const DictDecoder = struct { dctx: *anyopaque, ddict: *anyopaque, id: u32, max_out: usize,
pub fn init(dict: []const u8, max_out: usize) error{InvalidDictionary}!DictDecoder { const id = dictId(dict); if (id == 0) return error.InvalidDictionary; const dctx = c.ZSTD_createDCtx() orelse return error.InvalidDictionary; const ddict = c.ZSTD_createDDict(dict.ptr, dict.len) orelse { _ = c.ZSTD_freeDCtx(dctx); return error.InvalidDictionary; }; return .{ .dctx = dctx, .ddict = ddict, .id = id, .max_out = max_out }; }
pub fn deinit(self: *DictDecoder) void { _ = c.ZSTD_freeDDict(self.ddict); _ = c.ZSTD_freeDCtx(self.dctx); self.* = undefined; }
pub fn decompressAlloc(self: *DictDecoder, allocator: Allocator, frame: []const u8) DecompressError![]u8 { // A WebSocket message may contain multiple zstd/skippable frames. // Validate every window and sum declared output before allocating. var remaining = frame; var size: c_ulonglong = 0; while (remaining.len > 0) { var header: c.FrameHeader = undefined; if (c.ZSTD_getFrameHeader(&header, remaining.ptr, remaining.len) != 0) return error.DecompressFailed; const consumed = c.ZSTD_findFrameCompressedSize(remaining.ptr, remaining.len); if (c.ZSTD_isError(consumed) != 0 or consumed == 0 or consumed > remaining.len) return error.DecompressFailed; if (header.frame_type == 0) { // Match klauspost's default 512 MiB window cap and the // upstream client's WithDecoderMaxMemory(readLimit). const window_limit = @min(self.max_out, 512 << 20); if (@max(header.window_size, 1024) > window_limit) return error.BlockOversize; if (header.frame_content_size == c.contentsize_unknown) { size = c.contentsize_unknown; } else if (size != c.contentsize_unknown) { size = std.math.add(c_ulonglong, size, header.frame_content_size) catch return error.BlockOversize; if (size > self.max_out) return error.BlockOversize; } } remaining = remaining[consumed..]; } if (size == c.contentsize_unknown) { const bound = c.ZSTD_decompressBound(frame.ptr, frame.len); if (bound == c.contentsize_error) return error.DecompressFailed; // The bound may exceed actual output; cap the allocation, not the // frame's window. libzstd checks whether decoded bytes fit. const capacity: usize = @intCast(@min(bound, self.max_out)); const dst = try allocator.alloc(u8, capacity); errdefer allocator.free(dst); const written = c.ZSTD_decompress_usingDDict(self.dctx, dst.ptr, dst.len, frame.ptr, frame.len, self.ddict); if (c.ZSTD_isError(written) != 0) { if (c.ZSTD_getErrorCode(written) == c.error_dst_size_too_small) return error.BlockOversize; return error.DecompressFailed; } // Return an allocation with the actual length so callers can free // it normally, including allocators that cannot shrink in place. return try allocator.realloc(dst, written); } if (size == c.contentsize_error) return error.DecompressFailed; if (size > self.max_out) return error.BlockOversize; const dst = try allocator.alloc(u8, @intCast(size)); errdefer allocator.free(dst); const written = c.ZSTD_decompress_usingDDict(self.dctx, dst.ptr, dst.len, frame.ptr, frame.len, self.ddict); if (c.ZSTD_isError(written) != 0 or written != dst.len) return error.DecompressFailed; return dst; }};
// === tests ===
const testing = std.testing;
// hand-built frame: magic, FHD (single-segment | checksum), content size 4,// one raw last-block of "zat!", XXH64(seed 0) low 32 bits as content checksumfn testFrame(buf: *[17]u8, corrupt_checksum: bool) []const u8 { const content = "zat!"; buf[0..4].* = .{ 0x28, 0xb5, 0x2f, 0xfd }; buf[4] = 0x20 | 0x04; // single_segment + content_checksum buf[5] = content.len; buf[6..9].* = .{ 0x21, 0x00, 0x00 }; // last block, raw, size 4 buf[9..13].* = content.*; const check: u32 = @truncate(std.hash.XxHash64.hash(0, content)); std.mem.writeInt(u32, buf[13..17], if (corrupt_checksum) check ^ 1 else check, .little); return buf;}
test "decompressAlloc round-trips a checksummed frame" { var buf: [17]u8 = undefined; const out = try decompressAlloc(testing.allocator, testFrame(&buf, false)); defer testing.allocator.free(out); try testing.expectEqualStrings("zat!", out);}
test "decompressAlloc rejects a corrupted content checksum" { var buf: [17]u8 = undefined; try testing.expectError(error.DecompressFailed, decompressAlloc(testing.allocator, testFrame(&buf, true)));}
test "decompressAlloc rejects garbage" { try testing.expectError(error.DecompressFailed, decompressAlloc(testing.allocator, "not a zstd frame")); try testing.expectError(error.DecompressFailed, decompressAlloc(testing.allocator, ""));}
// A complete dictionary frame with a window descriptor but no content size.// Its 1 KiB window is deliberately larger than the four-byte output limit.fn unknownSizeFrame(buf: *[21]u8) []u8 { buf[0..6].* = .{ 0x28, 0xb5, 0x2f, 0xfd, 0x07, 0x00 }; std.mem.writeInt(u32, buf[6..10], dictId(@embedFile("testdata/live.dict")), .little); buf[10..13].* = .{ 0x21, 0, 0 }; buf[13..17].* = "zat!".*; std.mem.writeInt(u32, buf[17..21], @truncate(std.hash.XxHash64.hash(0, "zat!")), .little); return buf;}
test "live dictionary accepts omitted size within window and output limits" { var decoder = try DictDecoder.init(@embedFile("testdata/live.dict"), 1024); defer decoder.deinit(); var buf: [21]u8 = undefined; const out = try decoder.decompressAlloc(testing.allocator, unknownSizeFrame(&buf)); defer testing.allocator.free(out); try testing.expectEqualStrings("zat!", out);}
test "live dictionary limits unknown output and preserves checksum validation" { var decoder = try DictDecoder.init(@embedFile("testdata/live.dict"), 3); defer decoder.deinit(); var buf: [21]u8 = undefined; const frame = unknownSizeFrame(&buf); try testing.expectError(error.BlockOversize, decoder.decompressAlloc(testing.allocator, frame)); decoder.max_out = 1024; frame[20] ^= 1; try testing.expectError(error.DecompressFailed, decoder.decompressAlloc(testing.allocator, frame)); frame[20] ^= 1; try testing.expectError(error.DecompressFailed, decoder.decompressAlloc(testing.allocator, frame[0..20])); const out = try decoder.decompressAlloc(testing.allocator, frame); defer testing.allocator.free(out); try testing.expectEqualStrings("zat!", out);}
test "live dictionary decodes upstream Go streaming frame without declared size" { // klauspost/compress v1.19.2: encoder dictionary + Write, Flush, Close. // Verified with upstream's DecodeAll options; the window is 8 MiB, // but actual output is only 448 bytes. const frame = [_]u8{ 0x28, 0xb5, 0x2f, 0xfd, 0x07, 0x68, 0xcb, 0x27, 0x35, 0x01, 0x0c, 0x01, 0x00, 0x34, 0x01, 0x7b, 0x70, 0x6f, 0x73, 0x74, 0x22, 0x68, 0x65, 0x6c, 0x6c, 0x6f, 0x20, 0x77, 0x6f, 0x72, 0x6c, 0x64, 0x22, 0x7d, 0x03, 0x00, 0x85, 0x5b, 0xa9, 0x3f, 0x45, 0x34, 0xdd, 0x60, 0x16, 0x1b, 0x01, 0x00, 0x00, 0x01, 0xa7, 0xb1, 0x7d }; try testing.expectEqual(c.contentsize_unknown, c.ZSTD_getFrameContentSize(&frame, frame.len)); var decoder = try DictDecoder.init(@embedFile("testdata/live.dict"), 8 << 20); defer decoder.deinit(); const out = try decoder.decompressAlloc(testing.allocator, &frame); defer testing.allocator.free(out); const record = "{\"collection\":\"app.bsky.feed.post\",\"text\":\"hello world\"}"; try testing.expectEqual(record.len * 8, out.len); var records = std.mem.window(u8, out, record.len, record.len); while (records.next()) |r| try testing.expectEqualStrings(record, r); decoder.max_out = 447; try testing.expectError(error.BlockOversize, decoder.decompressAlloc(testing.allocator, &frame));}
test "live decoder accepts concatenated and skippable frames with aggregate limits" { var decoder = try DictDecoder.init(@embedFile("testdata/live.dict"), 1024); defer decoder.deinit(); var buf: [17]u8 = undefined; const frame = testFrame(&buf, false); const skip = [_]u8{ 0x50, 0x2a, 0x4d, 0x18, 1, 0, 0, 0, 42 }; var joined: [43]u8 = undefined; @memcpy(joined[0..17], frame); @memcpy(joined[17..26], &skip); @memcpy(joined[26..], frame); const out = try decoder.decompressAlloc(testing.allocator, &joined); defer testing.allocator.free(out); try testing.expectEqualStrings("zat!zat!", out); const empty = try decoder.decompressAlloc(testing.allocator, &skip); defer testing.allocator.free(empty); try testing.expectEqual(@as(usize, 0), empty.len); const too_large = try testing.allocator.alloc(u8, frame.len * 257); defer testing.allocator.free(too_large); for (0..257) |i| @memcpy(too_large[i * frame.len ..][0..frame.len], frame); try testing.expectError(error.BlockOversize, decoder.decompressAlloc(testing.allocator, too_large));}
test "live decoder rejects oversized window even for small output" { var decoder = try DictDecoder.init(@embedFile("testdata/live.dict"), 1024); defer decoder.deinit(); var buf: [21]u8 = undefined; const frame = unknownSizeFrame(&buf); frame[5] = 8; // 2 KiB window, four bytes output. try testing.expectError(error.BlockOversize, decoder.decompressAlloc(testing.allocator, frame));}
test "archive frames accept omitted content size and still verify checksums" { var buf: [17]u8 = undefined; _ = testFrame(&buf, false); buf[4] = 0x04; // checksum, no content size, separate window descriptor buf[5] = 0x00; // 1 KiB window const out = try decompressAlloc(testing.allocator, &buf); defer testing.allocator.free(out); try testing.expectEqualStrings("zat!", out); buf[16] ^= 1; try testing.expectError(error.DecompressFailed, decompressAlloc(testing.allocator, &buf)); buf[16] ^= 1; try testing.expectError(error.DecompressFailed, decompressAlloc(testing.allocator, buf[0..16])); buf[5] = 0xa0; // oversized 1 GiB window, rejected before allocation try testing.expectError(error.BlockOversize, decompressAlloc(testing.allocator, &buf));}