diff --git a/docs/compression-review.md b/docs/compression-review.md index 59d98ad..699ca36 100644 --- a/docs/compression-review.md +++ b/docs/compression-review.md @@ -62,3 +62,35 @@ disconnects/errors, and increasing message counts in every interval (the reporter prints 0/s for its initial baseline). Upstream sequence advanced from 33,455,790,952 to 33,455,832,751. Observed RSS was 4.19 GB; drop and compaction-error counters were zero. + +## Final live-decoder audit, 2026-09-08 + +References: official Jetstream `58c4d7f7a9130e53b40348ad3d1f7aafed0e4843`, +its Atmos v0.4.0 dependency, and klauspost/compress v1.19.2. Atmos supplies +HTTP retry policy; Jetstream live.go and klauspost's DecodeAll/framedec.go +supply the live compression contract. + +Two additional regressions failed first: concatenated frames were rejected, +and a 2 KiB window with four output bytes was accepted under a 1 KiB limit. +The decoder now walks all zstd/skippable frames, bounds aggregate output, +and checks every window against min(readLimit, 512 MiB), with upstream's +1 KiB minimum effective window. This supersedes the earlier 448-byte-limit +acceptance experiment: upstream also rejects a large window even when actual +output is small. Omitted-size frames remain supported within both limits. + +A loopback test verifies malformed compressed messages invoke onError: +accepting delivers the next valid event on the same connection; declining +stops without delivering it. That callback behavior needed no code change. +All 63 tests pass. + +Raw-record scope is not a gap: upstream options.go explicitly restricts +WithRawRecords to archive CBOR decoding. Live events remain JSON. The SDK's +archive-only decode_records=false option follows that scope; callback-borrowed +bytes must be copied when retained. Existing malformed-record/delete coverage +passes. No additional public raw-record API is required for this release. + +The separate HTTP-library follow-up must compare Atmos v0.4.0's complete +header contract: RateLimit-Reset takes precedence over Retry-After and is a +signed Unix timestamp. The current HTTP helper prioritizes Retry-After and +also interprets small reset values as relative delays. This is broader than +the previously recorded negative-value edge and is not fixed in this release. diff --git a/src/loopback_test.zig b/src/loopback_test.zig index 372c5c4..8d6a09a 100644 --- a/src/loopback_test.zig +++ b/src/loopback_test.zig @@ -55,6 +55,7 @@ const Instance = struct { stalled_plan: bool = false, archive_fault: enum { none, fetch, decode, segment } = .none, compression: enum { none, valid, rotate, same, invalid } = .none, + bad_compressed_first: bool = false, dict_calls: usize = 0, compressed_sessions: usize = 0, reject_plan: bool = false, @@ -317,6 +318,7 @@ const Instance = struct { try w.print("HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: {s}\r\nSec-WebSocket-Protocol: xrpc.v1.json\r\n\r\n", .{accept}); try w.flush(); + if (compressed and self.bad_compressed_first) try writeFrame(w, "invalid zstd", 0x82); if (!self.silent_ws) { const cursor = queryCursor(target); for (self.live) |row| { @@ -842,3 +844,37 @@ test "explicitly disabled compression makes no dictionary request" { try testing.expectEqual(@as(usize, 0), host.dict_calls); try testing.expectEqual(@as(usize, 0), host.compressed_sessions); } + +test "malformed compressed live frame follows caller continue or stop policy" { + inline for (.{ true, false }) |accept| { + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var host = try Instance.init(io, testing.allocator, &.{}, &.{.{ .seq = 1, .t = 1 }}); + host.compression = .valid; + host.bad_compressed_first = true; + var future = try io.concurrent(Instance.serve, .{&host}); + defer host.deinit(); + defer _ = future.cancel(io) catch {}; + defer host.quit(); + const Handler = struct { + errors: usize = 0, + events: usize = 0, + pub fn onError(self: *@This(), err: anyerror) bool { + testing.expectEqual(error.DecompressFailed, err) catch @panic("unexpected live error"); + self.errors += 1; + return accept; + } + pub fn onEvent(self: *@This(), _: livedecode.Event) bool { + self.events += 1; + return false; + } + }; + var handler = Handler{}; + var url: [64]u8 = undefined; + try client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{host.hostUrl(&url)} }, &handler); + try testing.expectEqual(@as(usize, 1), handler.errors); + try testing.expectEqual(@as(usize, if (accept) 1 else 0), handler.events); + try testing.expectEqual(@as(u32, 1), host.ws_sessions); + } +} diff --git a/src/zstd.zig b/src/zstd.zig index 17dd4fd..e18b5eb 100644 --- a/src/zstd.zig +++ b/src/zstd.zig @@ -14,6 +14,19 @@ 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 @@ -92,8 +105,29 @@ pub const DictDecoder = struct { } pub fn decompressAlloc(self: *DictDecoder, allocator: Allocator, frame: []const u8) DecompressError![]u8 { - if (frame.len == 0) return error.DecompressFailed; - const size = c.ZSTD_getFrameContentSize(frame.ptr, frame.len); + // 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; @@ -167,8 +201,8 @@ fn unknownSizeFrame(buf: *[21]u8) []u8 { return buf; } -test "live dictionary accepts omitted size with output exactly at limit" { - var decoder = try DictDecoder.init(@embedFile("testdata/live.dict"), 4); +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)); @@ -182,7 +216,7 @@ test "live dictionary limits unknown output and preserves checksum validation" { var buf: [21]u8 = undefined; const frame = unknownSizeFrame(&buf); try testing.expectError(error.BlockOversize, decoder.decompressAlloc(testing.allocator, frame)); - decoder.max_out = 4; + decoder.max_out = 1024; frame[20] ^= 1; try testing.expectError(error.DecompressFailed, decoder.decompressAlloc(testing.allocator, frame)); frame[20] ^= 1; @@ -198,7 +232,7 @@ test "live dictionary decodes upstream Go streaming frame without declared size" // 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"), 448); + 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); @@ -206,3 +240,34 @@ test "live dictionary decodes upstream Go streaming frame without declared size" 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)); +}