diff --git a/src/terminal/snapshot/checkpoint.zig b/src/terminal/snapshot/checkpoint.zig index 7d8f0f097..dabf37773 100644 --- a/src/terminal/snapshot/checkpoint.zig +++ b/src/terminal/snapshot/checkpoint.zig @@ -30,7 +30,7 @@ const record = @import("record.zig"); const Blake3 = std.crypto.hash.Blake3; /// The prefix digest exchanged by checkpoint codecs and the snapshot driver. -pub const Digest = [Blake3.digest_length]u8; +pub const Digest = record.PrefixDigest; comptime { std.debug.assert(@sizeOf(Digest) == 32); @@ -52,21 +52,20 @@ pub const Kind = enum { pub const EncodeError = std.Io.Writer.Error || record.Writer.FinishError; -/// Append one checkpoint covering every byte already in `destination`. +/// Append one checkpoint covering every byte already emitted by `stream`. /// -/// If writing fails, the incomplete checkpoint is removed while all covered -/// prefix bytes remain unchanged. +/// Finalizing does not consume the running hasher. READY is therefore included +/// as the stream continues toward FINISH. pub fn encode( kind: Kind, - destination: *std.Io.Writer.Allocating, + stream: *record.Writer, ) EncodeError!void { - var digest: Digest = undefined; - Blake3.hash(destination.written(), &digest, .{}); + const digest = stream.prefixDigest(); - var record_writer = try record.Writer.init(destination, kind.tag()); - errdefer record_writer.cancel(); - try record_writer.payloadWriter().writeAll(&digest); - try record_writer.finish(); + const payload = stream.begin(kind.tag()); + errdefer stream.cancel(); + try payload.writeAll(&digest); + try stream.finish(); } pub const DecodeError = record.Reader.InitError || @@ -82,15 +81,20 @@ pub const DecodeError = record.Reader.InitError || TrailingData, }; -/// Decode and validate one checkpoint against its preceding-byte digest. +/// Decode and validate one checkpoint against the running prefix digest. /// -/// `expected` must be finalized before any byte of this record is included in -/// the caller's running digest. FINISH additionally requires end-of-file. +/// READY is consumed through the hashing reader so FINISH covers it. FINISH is +/// consumed from the underlying source so neither checkpoint includes itself. pub fn decode( kind: Kind, - expected: Digest, - source: *std.Io.Reader, + stream: *record.StreamReader, ) DecodeError!void { + const expected = stream.prefixDigest(); + const source = if (kind == .finish) + stream.source() + else + stream.reader(); + var record_reader: record.Reader = undefined; try record_reader.init(source); if (record_reader.header.tag != kind.tag()) { @@ -120,8 +124,13 @@ test "READY golden encoding and BLAKE3-256 registry" { var snapshot: std.Io.Writer.Allocating = .init(std.testing.allocator); defer snapshot.deinit(); - try snapshot.writer.writeAll(prefix); - try encode(.ready, &snapshot); + var stream: record.Writer = .init( + std.testing.allocator, + &snapshot.writer, + ); + defer stream.deinit(); + try stream.writer().writeAll(prefix); + try encode(.ready, &stream); try test_fixture.expectEqual( .bytes, @@ -158,8 +167,10 @@ test "READY golden encoding and BLAKE3-256 registry" { ); // Decode the checked-in record rather than the generated candidate. - var source: std.Io.Reader = .fixed(test_ready_fixture[prefix.len..]); - try decode(.ready, actual, &source); + var source: std.Io.Reader = .fixed(&test_ready_fixture); + var source_stream: record.StreamReader = .init(&source); + try source_stream.reader().discardAll(prefix.len); + try decode(.ready, &source_stream); } test "READY and FINISH checkpoint coverage" { @@ -167,30 +178,38 @@ test "READY and FINISH checkpoint coverage" { var snapshot: std.Io.Writer.Allocating = .init(testing.allocator); defer snapshot.deinit(); - try snapshot.writer.writeAll("prefix"); + var stream: record.Writer = .init( + testing.allocator, + &snapshot.writer, + ); + defer stream.deinit(); + try stream.writer().writeAll("prefix"); // READY covers only the bytes that precede its own record. const ready_offset = snapshot.written().len; - var ready_digest: Digest = undefined; - Blake3.hash(snapshot.written(), &ready_digest, .{}); - try encode(.ready, &snapshot); + const ready_digest = stream.prefixDigest(); + try encode(.ready, &stream); // History follows READY and is included by FINISH. - try snapshot.writer.writeAll("history"); + try stream.writer().writeAll("history"); const finish_offset = snapshot.written().len; - var finish_digest: Digest = undefined; - Blake3.hash(snapshot.written(), &finish_digest, .{}); - try encode(.finish, &snapshot); + const finish_digest = stream.prefixDigest(); + try encode(.finish, &stream); - var ready_source: std.Io.Reader = .fixed( - snapshot.written()[ready_offset..], - ); - try decode(.ready, ready_digest, &ready_source); + // The running stream reaches both checkpoints without rehashing its prefix. + var source: std.Io.Reader = .fixed(snapshot.written()); + var source_stream: record.StreamReader = .init(&source); + try source_stream.reader().discardAll(ready_offset); + try testing.expectEqual(ready_digest, source_stream.prefixDigest()); + try decode(.ready, &source_stream); + try source_stream.reader().discardAll("history".len); + try testing.expectEqual(finish_digest, source_stream.prefixDigest()); + try decode(.finish, &source_stream); - var finish_source: std.Io.Reader = .fixed( - snapshot.written()[finish_offset..], + try testing.expectEqual( + snapshot.written().len - finish_offset, + record.Header.len + @sizeOf(Digest), ); - try decode(.finish, finish_digest, &finish_source); } test "checkpoint rejects wrong tags, digests, and FINISH trailing data" { @@ -198,42 +217,68 @@ test "checkpoint rejects wrong tags, digests, and FINISH trailing data" { var snapshot: std.Io.Writer.Allocating = .init(testing.allocator); defer snapshot.deinit(); - try snapshot.writer.writeAll("prefix"); + var stream: record.Writer = .init( + testing.allocator, + &snapshot.writer, + ); + defer stream.deinit(); + try stream.writer().writeAll("prefix"); const checkpoint_offset = snapshot.written().len; - var digest: Digest = undefined; - Blake3.hash(snapshot.written(), &digest, .{}); - try encode(.ready, &snapshot); + try encode(.ready, &stream); - var wrong_tag: std.Io.Reader = .fixed( + var wrong_tag_source: std.Io.Reader = .fixed( snapshot.written()[checkpoint_offset..], ); + var wrong_tag: record.StreamReader = .init(&wrong_tag_source); try testing.expectError( error.UnexpectedRecordTag, - decode(.finish, digest, &wrong_tag), + decode(.finish, &wrong_tag), ); - var invalid_digest = digest; - invalid_digest[0] ^= 1; - var invalid_digest_source: std.Io.Reader = .fixed( - snapshot.written()[checkpoint_offset..], + // Build a correctly framed checkpoint containing an unrelated digest. + var invalid: std.Io.Writer.Allocating = .init(testing.allocator); + defer invalid.deinit(); + var invalid_stream: record.Writer = .init( + testing.allocator, + &invalid.writer, ); + defer invalid_stream.deinit(); + try invalid_stream.writer().writeAll("prefix"); + var invalid_digest = invalid_stream.prefixDigest(); + invalid_digest[0] ^= 1; + const invalid_payload = invalid_stream.begin(.ready); + errdefer invalid_stream.cancel(); + try invalid_payload.writeAll(&invalid_digest); + try invalid_stream.finish(); + + var invalid_digest_source: std.Io.Reader = .fixed(invalid.written()); + var invalid_digest_stream: record.StreamReader = .init( + &invalid_digest_source, + ); + try invalid_digest_stream.reader().discardAll("prefix".len); try testing.expectError( error.InvalidDigest, - decode(.ready, invalid_digest, &invalid_digest_source), + decode(.ready, &invalid_digest_stream), ); var finished: std.Io.Writer.Allocating = .init(testing.allocator); defer finished.deinit(); - try finished.writer.writeAll("prefix"); - const finish_offset = finished.written().len; - try encode(.finish, &finished); + var finished_stream: record.Writer = .init( + testing.allocator, + &finished.writer, + ); + defer finished_stream.deinit(); + try finished_stream.writer().writeAll("prefix"); + try encode(.finish, &finished_stream); try finished.writer.writeByte(0); - var trailing_source: std.Io.Reader = .fixed( - finished.written()[finish_offset..], + var trailing_source: std.Io.Reader = .fixed(finished.written()); + var trailing_stream: record.StreamReader = .init( + &trailing_source, ); + try trailing_stream.reader().discardAll("prefix".len); try testing.expectError( error.TrailingData, - decode(.finish, digest, &trailing_source), + decode(.finish, &trailing_stream), ); } diff --git a/src/terminal/snapshot/history.zig b/src/terminal/snapshot/history.zig index 9de7b4c3f..b49c05081 100644 --- a/src/terminal/snapshot/history.zig +++ b/src/terminal/snapshot/history.zig @@ -162,16 +162,13 @@ pub const EncodeError = Allocator.Error || /// Encode one screen's HISTORY and its complete historical pages. /// /// Pages are encoded newest-to-oldest. Compressed source pages are inspected -/// without changing their storage state. On failure, the entire sequence is -/// removed while earlier destination bytes remain. +/// without changing their storage state. Completed records may already be +/// emitted if a later page fails. pub fn encode( terminal_screen: *const TerminalScreen, key: TerminalScreenKey, - destination: *std.Io.Writer.Allocating, + destination: *record.Writer, ) EncodeError!void { - const sequence_start = destination.written().len; - errdefer destination.shrinkRetainingCapacity(sequence_start); - // SCREEN begins at the page containing the active area's first row. Its // leading rows are already resident; every previous complete page belongs // to this HISTORY sequence. @@ -205,10 +202,10 @@ pub fn encode( // HISTORY declares exactly how many PAGE records follow and how their rows // combine with the overlap already carried by SCREEN. { - var record_writer = try record.Writer.init(destination, .history); - errdefer record_writer.cancel(); - try header.encode(record_writer.payloadWriter()); - try record_writer.finish(); + const payload = destination.begin(.history); + errdefer destination.cancel(); + try header.encode(payload); + try destination.finish(); } // Walk backward so each page can be prepended by the decoder immediately. @@ -432,13 +429,18 @@ test "HISTORY encodes newest first and restores complete history" { std.testing.allocator, ); defer destination.deinit(); - try screen.encode(&source_screen, .primary, &destination); + var stream: record.Writer = .init( + std.testing.allocator, + &destination.writer, + ); + defer stream.deinit(); + try screen.encode(&source_screen, .primary, &stream); const history_offset = destination.written().len; _ = source_screen.pages.compress(.full); const oldest_storage = oldest_history.storage(); const newest_storage = newest_history.storage(); - try encode(&source_screen, .primary, &destination); + try encode(&source_screen, .primary, &stream); try std.testing.expectEqual(oldest_storage, oldest_history.storage()); try std.testing.expectEqual(newest_storage, newest_history.storage()); @@ -645,7 +647,12 @@ test "HISTORY encodes and restores an empty sequence" { std.testing.allocator, ); defer destination.deinit(); - try encode(&terminal_screen, .primary, &destination); + var stream: record.Writer = .init( + std.testing.allocator, + &destination.writer, + ); + defer stream.deinit(); + try encode(&terminal_screen, .primary, &stream); var inspect_source: std.Io.Reader = .fixed(destination.written()); var history_record: record.Reader = undefined; @@ -705,9 +712,15 @@ test "HISTORY accepts row metadata mismatches and rejects structural ones" { std.testing.allocator, ); defer destination.deinit(); - var record_writer = try record.Writer.init(&destination, .history); - try header.encode(record_writer.payloadWriter()); - try record_writer.finish(); + var stream: record.Writer = .init( + std.testing.allocator, + &destination.writer, + ); + defer stream.deinit(); + const payload = stream.begin(.history); + errdefer stream.cancel(); + try header.encode(payload); + try stream.finish(); var source: std.Io.Reader = .fixed(destination.written()); try decode( @@ -752,9 +765,15 @@ test "HISTORY accepts row metadata mismatches and rejects structural ones" { std.testing.allocator, ); defer destination.deinit(); - var record_writer = try record.Writer.init(&destination, .history); - try case.header.encode(record_writer.payloadWriter()); - try record_writer.finish(); + var stream: record.Writer = .init( + std.testing.allocator, + &destination.writer, + ); + defer stream.deinit(); + const payload = stream.begin(.history); + errdefer stream.cancel(); + try case.header.encode(payload); + try stream.finish(); var source: std.Io.Reader = .fixed(destination.written()); try std.testing.expectError( diff --git a/src/terminal/snapshot/main.zig b/src/terminal/snapshot/main.zig index 2c7d58335..188cab684 100644 --- a/src/terminal/snapshot/main.zig +++ b/src/terminal/snapshot/main.zig @@ -77,20 +77,26 @@ //! //! ## Encoding //! -//! Encode a complete snapshot into an empty allocating writer: +//! Encode a complete snapshot into any writer: //! //! ```zig //! var output: std.Io.Writer.Allocating = .init(alloc); //! defer output.deinit(); //! -//! try snapshot.encode(&terminal, &output); +//! try snapshot.encode(alloc, &output.writer, &terminal); //! //! const bytes = output.written(); //! ``` //! -//! We have to use an allocating writer because record formats require -//! encoding the length and CRC in the header, so we need a seekable -//! format. +//! Encoding begins at the writer's current position, so unrelated bytes may +//! precede the snapshot. The encoder buffers only the current record payload +//! to calculate its length and CRC32C; completed records stream immediately +//! and BLAKE3 checkpoint coverage is updated incrementally. Buffering is an +//! encoder implementation detail, not a requirement of the wire format. +//! +//! A failure may leave prior complete records, or a partial record if the +//! destination itself fails. Such a prefix has no valid FINISH checkpoint and +//! cannot be restored as a complete snapshot. //! //! Each record type usually exposes an `encode` function that encodes //! a complete record, such as `screen.encode`. diff --git a/src/terminal/snapshot/page.zig b/src/terminal/snapshot/page.zig index 335d74127..e71bc53c4 100644 --- a/src/terminal/snapshot/page.zig +++ b/src/terminal/snapshot/page.zig @@ -135,17 +135,16 @@ pub const EncodeError = PayloadEncodeError || record.Writer.FinishError; /// Encode one complete PAGE record from a native page. /// -/// The record is appended to `destination`. Its header is reserved before the -/// payload is encoded, then backpatched with the payload length and CRC32C. If -/// encoding fails, the partial record is removed while earlier bytes remain. +/// The payload is built in the stream's reusable record buffer. If payload +/// encoding fails, no bytes from this record are emitted. pub fn encode( page: *const TerminalPage, - destination: *std.Io.Writer.Allocating, + destination: *record.Writer, ) EncodeError!void { - var record_writer = try record.Writer.init(destination, .page); - errdefer record_writer.cancel(); - try encodePayload(page, record_writer.payloadWriter()); - try record_writer.finish(); + const payload = destination.begin(.page); + errdefer destination.cancel(); + try encodePayload(page, payload); + try destination.finish(); } /// Errors possible while decoding and validating a complete PAGE record. @@ -755,7 +754,12 @@ test "framed PAGE golden empty record" { var destination: std.Io.Writer.Allocating = .init(std.testing.allocator); defer destination.deinit(); - try encode(&page, &destination); + var stream: record.Writer = .init( + std.testing.allocator, + &destination.writer, + ); + defer stream.deinit(); + try encode(&page, &stream); try test_fixture.expectEqual( .bytes, "src/terminal/snapshot/testdata/page-empty-record-v1.hex", @@ -794,9 +798,14 @@ test "framed PAGE rejects incomplete wide cells transactionally" { var destination: std.Io.Writer.Allocating = .init(testing.allocator); defer destination.deinit(); try destination.writer.writeAll("prefix"); + var stream: record.Writer = .init( + testing.allocator, + &destination.writer, + ); + defer stream.deinit(); try testing.expectError( error.InvalidWideCell, - encode(&page, &destination), + encode(&page, &stream), ); try testing.expectEqualStrings("prefix", destination.written()); } diff --git a/src/terminal/snapshot/record.zig b/src/terminal/snapshot/record.zig index 19c9285a5..0126c2bcd 100644 --- a/src/terminal/snapshot/record.zig +++ b/src/terminal/snapshot/record.zig @@ -19,9 +19,16 @@ //! Supported tags are in `Tag`. const std = @import("std"); +const assert = std.debug.assert; +const Allocator = std.mem.Allocator; const test_fixture = @import("fixture.zig"); const io = @import("io.zig"); +const Blake3 = std.crypto.hash.Blake3; + +/// The running digest shared by snapshot stream codecs and checkpoints. +pub const PrefixDigest = [Blake3.digest_length]u8; + /// CRC32C as specified by the snapshot format. Zig names this standard /// parameter set after its iSCSI use. pub const Crc32c = std.hash.crc.Crc32Iscsi; @@ -57,7 +64,7 @@ pub const Header = struct { comptime { // This size is part of the wire format. If it changes, the snapshot // version and golden fixtures must also change. - std.debug.assert(len == 10); + assert(len == 10); } /// Determines how the payload is decoded. @@ -142,71 +149,140 @@ pub const Checksum = struct { } }; -/// Builds one complete record. +/// Streams complete records while retaining only one payload at a time. /// -/// This Writer requires an Allocating std.Io.Writer because the record -/// format requires reading the full payload and rewinding in order to -/// write the length + CRC without encoding twice. +/// All emitted bytes pass through one unbuffered BLAKE3 writer. The scratch +/// allocation is retained between records so a stream's peak memory is the +/// largest record payload rather than the complete snapshot. pub const Writer = struct { - destination: *std.Io.Writer.Allocating, - tag: Tag, - record_start: usize, + hashing: std.Io.Writer.Hashed(Blake3), + scratch: std.Io.Writer.Allocating, + active_tag: ?Tag, - /// Reserve space for a record at the current end of `destination`. - /// Once this is called, callers MUST NOT write anything else to - /// the writer until `finish` or `cancel` is called. pub fn init( - destination: *std.Io.Writer.Allocating, - tag: Tag, - ) std.Io.Writer.Error!Writer { - const record_start = destination.written().len; - errdefer destination.shrinkRetainingCapacity(record_start); - try destination.writer.splatByteAll(0, Header.len); + alloc: Allocator, + destination: *std.Io.Writer, + ) Writer { return .{ - .destination = destination, - .tag = tag, - .record_start = record_start, + .hashing = destination.hashed(Blake3.init(.{}), &.{}), + .scratch = .init(alloc), + .active_tag = null, }; } - /// Return the writer through which the payload is encoded exactly once. - pub fn payloadWriter(self: *Writer) *std.Io.Writer { - return &self.destination.writer; + pub fn deinit(self: *Writer) void { + assert(self.active_tag == null); + assert(self.scratch.writer.end == 0); + self.scratch.deinit(); + self.* = undefined; } - pub const FinishError = error{ + /// Return the digest-updating writer for unframed snapshot bytes. + pub fn writer(self: *Writer) *std.Io.Writer { + assert(self.active_tag == null); + return &self.hashing.writer; + } + + /// Begin one record and return its reusable payload writer. + pub fn begin(self: *Writer, tag: Tag) *std.Io.Writer { + assert(self.active_tag == null); + assert(self.scratch.writer.end == 0); + self.active_tag = tag; + return &self.scratch.writer; + } + + pub const FinishError = std.Io.Writer.Error || error{ /// The payload cannot be represented by the record's `u32` length. PayloadTooLarge, }; - /// Marked the completed record with the payload length and CRC32C. + /// Finish and emit the active record's header followed by its payload. + /// + /// Payload validation failures occur before this function and emit no part + /// of the record. A destination failure here may have emitted a prefix. pub fn finish(self: *Writer) FinishError!void { - const bytes = self.destination.written(); - const payload_start = self.record_start + Header.len; + assert(self.active_tag != null); + const tag = self.active_tag.?; + defer self.cancel(); + + const payload = self.scratch.written(); const payload_len = std.math.cast( u32, - bytes.len - payload_start, + payload.len, ) orelse return error.PayloadTooLarge; - const payload = bytes[payload_start..]; - // Calculate our CRC - var checksum: Checksum = .init(self.tag, payload_len); + // The CRC prefix contains the payload length, so checksum the retained + // payload only after its final length is known. + var checksum: Checksum = .init(tag, payload_len); checksum.writer().writeAll(payload) catch unreachable; - // Build the header and encode it directly into the header const header: Header = .{ - .tag = self.tag, + .tag = tag, .payload_len = payload_len, .crc32c = checksum.final(), }; - var header_writer: std.Io.Writer = .fixed(bytes[self.record_start..payload_start]); + var header_bytes: [Header.len]u8 = undefined; + var header_writer: std.Io.Writer = .fixed(&header_bytes); header.encode(&header_writer) catch unreachable; + + try self.hashing.writer.writeAll(&header_bytes); + try self.hashing.writer.writeAll(payload); } - /// Discard this record. This makes it safe to use the underlying - /// alloating writer again as if nothing happened. + /// Discard the active record without emitting any bytes. + /// + /// This is idempotent so an outer error cleanup may call it after `finish` + /// has already cleared state following a destination failure. pub fn cancel(self: *Writer) void { - self.destination.shrinkRetainingCapacity(self.record_start); + self.scratch.shrinkRetainingCapacity(0); + self.active_tag = null; + } + + /// Finalize the prefix written so far without consuming the hasher. + pub fn prefixDigest(self: *const Writer) PrefixDigest { + // Checkpoints require an exact byte boundary. Writer owns this + // adapter and always constructs it without a buffer. + assert(self.active_tag == null); + assert(self.scratch.writer.end == 0); + assert(self.hashing.writer.buffer.len == 0); + assert(self.hashing.writer.buffered().len == 0); + + var result: PrefixDigest = undefined; + self.hashing.hasher.final(&result); + return result; + } +}; + +/// Hashes snapshot bytes as they are consumed without reading ahead. +pub const StreamReader = struct { + hashing: std.Io.Reader.Hashed(Blake3), + + pub fn init(input: *std.Io.Reader) StreamReader { + return .{ + .hashing = input.hashed(Blake3.init(.{}), &.{}), + }; + } + + /// Return the digest-updating reader used before and through READY. + pub fn reader(self: *StreamReader) *std.Io.Reader { + return &self.hashing.reader; + } + + /// Return the source used to consume FINISH without hashing it. + pub fn source(self: *StreamReader) *std.Io.Reader { + return self.hashing.in; + } + + /// Finalize the prefix consumed so far without consuming the hasher. + pub fn prefixDigest(self: *const StreamReader) PrefixDigest { + // A nonempty adapter buffer could contain bytes beyond a checkpoint. + // Construction is private to this type and fixes its capacity at zero. + assert(self.hashing.reader.buffer.len == 0); + assert(self.hashing.reader.bufferedLen() == 0); + + var result: PrefixDigest = undefined; + self.hashing.hasher.final(&result); + return result; } }; @@ -437,14 +513,20 @@ test "payload limit does not consume the next record" { try std.testing.expectEqualStrings("next", try source.take(4)); } -test "record writer appends and backpatches framing" { +test "record writer buffers a payload and appends framing" { var destination: std.Io.Writer.Allocating = .init(std.testing.allocator); defer destination.deinit(); try destination.writer.writeAll("prefix"); + var stream: Writer = .init( + std.testing.allocator, + &destination.writer, + ); + defer stream.deinit(); - var record_writer = try Writer.init(&destination, .page); - try record_writer.payloadWriter().writeAll("payload"); - try record_writer.finish(); + const payload = stream.begin(.page); + errdefer stream.cancel(); + try payload.writeAll("payload"); + try stream.finish(); const encoded = destination.written(); try std.testing.expectEqualStrings("prefix", encoded[0..6]); @@ -468,10 +550,50 @@ test "record writer cancel preserves preceding bytes" { var destination: std.Io.Writer.Allocating = .init(std.testing.allocator); defer destination.deinit(); try destination.writer.writeAll("prefix"); + var stream: Writer = .init( + std.testing.allocator, + &destination.writer, + ); + defer stream.deinit(); - var record_writer = try Writer.init(&destination, .page); - try record_writer.payloadWriter().writeAll("partial"); - record_writer.cancel(); + const partial = stream.begin(.page); + errdefer stream.cancel(); + try partial.writeAll("partial"); + const scratch_capacity = stream.scratch.writer.buffer.len; + stream.cancel(); try std.testing.expectEqualStrings("prefix", destination.written()); + try std.testing.expectEqual( + scratch_capacity, + stream.scratch.writer.buffer.len, + ); + + // Cancel retains the allocation but leaves it ready for the next record. + const next = stream.begin(.page); + try next.writeAll("next"); + try stream.finish(); + + var source: std.Io.Reader = .fixed(destination.written()[6..]); + var record_reader: Reader = undefined; + try record_reader.init(&source); + try std.testing.expectEqualStrings( + "next", + try record_reader.payloadReader().take(4), + ); + try record_reader.finish(); +} + +test "record writer clears active state after destination failure" { + // The header fits but the payload does not, leaving a permitted partial + // record in the destination while the reusable writer becomes idle again. + var destination_bytes: [Header.len]u8 = undefined; + var destination: std.Io.Writer = .fixed(&destination_bytes); + var stream: Writer = .init(std.testing.allocator, &destination); + defer stream.deinit(); + + const payload = stream.begin(.page); + try payload.writeByte(0); + try std.testing.expectError(error.WriteFailed, stream.finish()); + try std.testing.expectEqual(@as(?Tag, null), stream.active_tag); + try std.testing.expectEqual(@as(usize, 0), stream.scratch.writer.end); } diff --git a/src/terminal/snapshot/screen.zig b/src/terminal/snapshot/screen.zig index c2ef53f1c..2fc3d0d8b 100644 --- a/src/terminal/snapshot/screen.zig +++ b/src/terminal/snapshot/screen.zig @@ -238,16 +238,13 @@ pub const EncodeError = PayloadEncodeError || page.EncodeError || error{ /// Encode one SCREEN and its minimal suffix of complete native pages. /// /// The suffix begins with the page containing the active area's first row and -/// ends with the newest page. If encoding any record fails, the entire sequence -/// is removed while earlier destination bytes remain. +/// ends with the newest page. Completed records may already be emitted if a +/// later record fails; the missing READY checkpoint makes that prefix invalid. pub fn encode( screen: *const TerminalScreen, key: TerminalScreenKey, - destination: *std.Io.Writer.Allocating, + destination: *record.Writer, ) EncodeError!void { - const sequence_start = destination.written().len; - errdefer destination.shrinkRetainingCapacity(sequence_start); - // The active top may fall inside this page, leaving an incidental history // prefix. Every earlier complete page is history and is omitted. const first = screen.pages.getTopLeft(.active).node; @@ -262,15 +259,15 @@ pub fn encode( // SCREEN declares exactly how many immediately following PAGE records // belong to it. { - var record_writer = try record.Writer.init(destination, .screen); - errdefer record_writer.cancel(); + const payload = destination.begin(.screen); + errdefer destination.cancel(); try encodePayload( screen, key, encoded_page_count, - record_writer.payloadWriter(), + payload, ); - try record_writer.finish(); + try destination.finish(); } // PageList never compresses the active-boundary page or any later page. @@ -1705,7 +1702,12 @@ test "framed native SCREEN and PAGE sequence" { ); defer destination.deinit(); try destination.writer.writeAll("prefix"); - try encode(&screen, .alternate, &destination); + var stream: record.Writer = .init( + std.testing.allocator, + &destination.writer, + ); + defer stream.deinit(); + try encode(&screen, .alternate, &stream); try std.testing.expectEqualStrings( "prefix", @@ -1849,7 +1851,12 @@ test "SCREEN encodes the minimal complete-page active suffix" { std.testing.allocator, ); defer destination.deinit(); - try encode(&screen, .primary, &destination); + var stream: record.Writer = .init( + std.testing.allocator, + &destination.writer, + ); + defer stream.deinit(); + try encode(&screen, .primary, &stream); var source: std.Io.Reader = .fixed(destination.written()); var screen_record: record.Reader = undefined; @@ -1964,6 +1971,11 @@ test "SCREEN restoration normalizes invalid cursor positions" { std.testing.allocator, ); defer destination.deinit(); + var stream: record.Writer = .init( + std.testing.allocator, + &destination.writer, + ); + defer stream.deinit(); var header = Header.init(&screen, .primary, 1); header.cursor_x = case.x; @@ -1971,14 +1983,15 @@ test "SCREEN restoration normalizes invalid cursor positions" { header.cursor_flags.pending_wrap = case.pending_wrap; header.saved_cursor_present = false; - var screen_writer = try record.Writer.init(&destination, .screen); - try header.encode(screen_writer.payloadWriter()); - try screen_writer.payloadWriter().writeByte(0); - try screen_writer.finish(); + const screen_payload = stream.begin(.screen); + errdefer stream.cancel(); + try header.encode(screen_payload); + try screen_payload.writeByte(0); + try stream.finish(); const native_page = screen.pages.getTopLeft(.active).node .pageAssumeResident(); - try page.encode(native_page, &destination); + try page.encode(native_page, &stream); var source: std.Io.Reader = .fixed(destination.written()); var decoded = try decode( @@ -2011,7 +2024,12 @@ test "SCREEN restoration rejects invalid and incomplete sequences" { std.testing.allocator, ); defer destination.deinit(); - try encode(&screen, .primary, &destination); + var stream: record.Writer = .init( + std.testing.allocator, + &destination.writer, + ); + defer stream.deinit(); + try encode(&screen, .primary, &stream); // Every truncation must fail without leaking a partially restored // PageList, cursor pin, or cursor-owned state. @@ -2052,16 +2070,19 @@ test "SCREEN restoration rejects invalid and incomplete sequences" { std.testing.allocator, ); defer empty_sequence.deinit(); + var empty_stream: record.Writer = .init( + std.testing.allocator, + &empty_sequence.writer, + ); + defer empty_stream.deinit(); { - var record_writer = try record.Writer.init( - &empty_sequence, - .screen, - ); + const payload = empty_stream.begin(.screen); + errdefer empty_stream.cancel(); try Header.init(&screen, .primary, 0).encode( - record_writer.payloadWriter(), + payload, ); - try record_writer.payloadWriter().writeByte(0); - try record_writer.finish(); + try payload.writeByte(0); + try empty_stream.finish(); } var empty_source: std.Io.Reader = .fixed(empty_sequence.written()); try std.testing.expectError( @@ -2083,21 +2104,25 @@ test "SCREEN sequence failure preserves preceding bytes" { ); defer screen.deinit(); - var failing = std.testing.FailingAllocator.init( + var destination: std.Io.Writer.Allocating = .init( std.testing.allocator, - .{}, - ); - var destination = try std.Io.Writer.Allocating.initCapacity( - failing.allocator(), - 6 + record.Header.len + Header.len + 1, ); defer destination.deinit(); try destination.writer.writeAll("prefix"); - failing.fail_index = failing.alloc_index; + var failing = std.testing.FailingAllocator.init( + std.testing.allocator, + .{ .fail_index = 0 }, + ); + var stream: record.Writer = .init( + failing.allocator(), + &destination.writer, + ); + defer stream.deinit(); + try std.testing.expectError( error.WriteFailed, - encode(&screen, .primary, &destination), + encode(&screen, .primary, &stream), ); try std.testing.expectEqualStrings("prefix", destination.written()); } @@ -2114,6 +2139,11 @@ test "SCREEN decode ignores an invalid cursor hyperlink" { std.testing.allocator, ); defer destination.deinit(); + var stream: record.Writer = .init( + std.testing.allocator, + &destination.writer, + ); + defer stream.deinit(); var header = testHeader(); header.key = .primary; @@ -2123,16 +2153,17 @@ test "SCREEN decode ignores an invalid cursor hyperlink" { header.cursor_flags.pending_wrap = false; header.saved_cursor_present = false; - var record_writer = try record.Writer.init(&destination, .screen); - try header.encode(record_writer.payloadWriter()); - try record_writer.payloadWriter().writeByte(1); // Implicit hyperlink. - try io.writeInt(record_writer.payloadWriter(), u32, 1); - try io.writeInt(record_writer.payloadWriter(), u32, 0); // Empty URI. - try record_writer.finish(); + const screen_payload = stream.begin(.screen); + errdefer stream.cancel(); + try header.encode(screen_payload); + try screen_payload.writeByte(1); // Implicit hyperlink. + try io.writeInt(screen_payload, u32, 1); + try io.writeInt(screen_payload, u32, 0); // Empty URI. + try stream.finish(); try page.encode( screen.pages.getTopLeft(.active).node.pageAssumeResident(), - &destination, + &stream, ); var source: std.Io.Reader = .fixed(destination.written()); @@ -2151,6 +2182,11 @@ test "SCREEN decode ignores a PAGE with an empty hyperlink URI" { std.testing.allocator, ); defer destination.deinit(); + var stream: record.Writer = .init( + std.testing.allocator, + &destination.writer, + ); + defer stream.deinit(); var header = testHeader(); header.key = .primary; @@ -2160,10 +2196,11 @@ test "SCREEN decode ignores a PAGE with an empty hyperlink URI" { header.cursor_flags.pending_wrap = false; header.saved_cursor_present = false; - var screen_writer = try record.Writer.init(&destination, .screen); - try header.encode(screen_writer.payloadWriter()); - try screen_writer.payloadWriter().writeByte(0); - try screen_writer.finish(); + const screen_payload = stream.begin(.screen); + errdefer stream.cancel(); + try header.encode(screen_payload); + try screen_payload.writeByte(0); + try stream.finish(); const page_header: page.Header = .{ .columns = 1, @@ -2175,34 +2212,35 @@ test "SCREEN decode ignores a PAGE with an empty hyperlink URI" { .grapheme_capacity_bytes = 0, .string_capacity_bytes = 0, }; - var page_writer = try record.Writer.init(&destination, .page); - try page_header.encode(page_writer.payloadWriter()); + const page_payload = stream.begin(.page); + errdefer stream.cancel(); + try page_header.encode(page_payload); try io.writeInt( - page_writer.payloadWriter(), + page_payload, terminal_hyperlink.Id, 1, ); - try page_writer.payloadWriter().writeByte(1); // Implicit hyperlink. - try io.writeInt(page_writer.payloadWriter(), u32, 1); - try io.writeInt(page_writer.payloadWriter(), u32, 0); // Empty URI. + try page_payload.writeByte(1); // Implicit hyperlink. + try io.writeInt(page_payload, u32, 1); + try io.writeInt(page_payload, u32, 0); // Empty URI. // One narrow codepoint cell refers to the hyperlink table entry above. // Since that entry is ignored, the cell must restore without a hyperlink. - try page_writer.payloadWriter().writeByte(0); - try page_writer.payloadWriter().writeAll(&.{ 0, 0, 0, 0 }); + try page_payload.writeByte(0); + try page_payload.writeAll(&.{ 0, 0, 0, 0 }); try io.writeInt( - page_writer.payloadWriter(), + page_payload, terminal_style.Id, 0, ); try io.writeInt( - page_writer.payloadWriter(), + page_payload, terminal_hyperlink.Id, 1, ); - try io.writeInt(page_writer.payloadWriter(), u32, 'A'); - try io.writeInt(page_writer.payloadWriter(), u32, 0); - try page_writer.finish(); + try io.writeInt(page_payload, u32, 'A'); + try io.writeInt(page_payload, u32, 0); + try stream.finish(); var source: std.Io.Reader = .fixed(destination.written()); var decoded = try decode( diff --git a/src/terminal/snapshot/snapshot.zig b/src/terminal/snapshot/snapshot.zig index 0f3635487..d083b54d5 100644 --- a/src/terminal/snapshot/snapshot.zig +++ b/src/terminal/snapshot/snapshot.zig @@ -24,61 +24,58 @@ const test_complete_fixture = test_fixture.parse( pub const EncodeError = terminal.EncodeError || screen.EncodeError || history.EncodeError || - checkpoint.EncodeError || - error{ - /// A snapshot envelope must begin at byte zero. - DestinationNotEmpty, - }; + checkpoint.EncodeError; /// Encode one complete terminal snapshot. /// -/// `destination` must be empty because checkpoint digests cover every byte -/// from the snapshot envelope onward. The operation is transactional: any -/// failure restores the destination to empty. +/// Encoding starts at the destination's current position. Only one record +/// payload is buffered at a time; completed records stream immediately. On +/// failure, the destination may contain a snapshot prefix without its required +/// checkpoints, and an output failure may have written part of a record. pub fn encode( + alloc: Allocator, + destination: *std.Io.Writer, t: *const Terminal, - destination: *std.Io.Writer.Allocating, ) EncodeError!void { - // We require empty for checkpoint digests - if (destination.written().len != 0) return error.DestinationNotEmpty; - errdefer destination.shrinkRetainingCapacity(0); + var stream: record.Writer = .init(alloc, destination); + defer stream.deinit(); // 1. Envelope - try envelope.encode(&destination.writer); + try envelope.encode(stream.writer()); // 2. Terminal - try terminal.encode(t, destination); + try terminal.encode(t, &stream); // 3. Primary and alt screen try screen.encode( t.screens.get(.primary).?, .primary, - destination, + &stream, ); if (t.screens.get(.alternate)) |alternate| try screen.encode( alternate, .alternate, - destination, + &stream, ); // 4. Ready checkpoint. In the future we'll put our continuation // state before this so pty bytes can also flow. - try checkpoint.encode(.ready, destination); + try checkpoint.encode(.ready, &stream); // 5. History try history.encode( t.screens.get(.primary).?, .primary, - destination, + &stream, ); if (t.screens.get(.alternate)) |alternate| try history.encode( alternate, .alternate, - destination, + &stream, ); // 6. Finish - try checkpoint.encode(.finish, destination); + try checkpoint.encode(.finish, &stream); } /// Errors possible while restoring one complete terminal snapshot. @@ -113,11 +110,10 @@ pub fn decode( io_: std.Io, alloc: Allocator, ) DecodeError!Terminal { - // Keep this reader unbuffered so the hasher never reads past a checkpoint - // boundary. Record readers provide their own bounded buffers. - const Blake3 = std.crypto.hash.Blake3; - var hashing = source.hashed(Blake3.init(.{}), &.{}); - const reader = &hashing.reader; + // StreamReader owns a zero-buffer hashing adapter, making checkpoint + // boundaries part of its API rather than a caller-maintained invariant. + var stream: record.StreamReader = .init(source); + const reader = stream.reader(); // Read the envelope, which is currently just a verification step. try envelope.decode(reader); @@ -179,9 +175,7 @@ pub fn decode( // READY covers the exact envelope-through-SCREEN prefix. Finalizing does // not consume the hasher, so the same stream continues toward FINISH. - var digest: checkpoint.Digest = undefined; - hashing.hasher.final(&digest); - try checkpoint.decode(.ready, digest, reader); + try checkpoint.decode(.ready, &stream); // HISTORY keys make this sequence order-independent just like SCREEN. // Although a decoder may publish recent pages as they validate, any later @@ -206,8 +200,7 @@ pub fn decode( // FINISH authenticates READY and all history. Decode it directly from the // underlying reader so the digest does not include FINISH itself. - hashing.hasher.final(&digest); - try checkpoint.decode(.finish, digest, source); + try checkpoint.decode(.finish, &stream); const keys = [_]TerminalScreenKey{ .primary, .alternate }; if (comptime build_options.slow_runtime_safety) { @@ -311,7 +304,7 @@ test "complete snapshot round trip with history and alternate screen" { var encoded: std.Io.Writer.Allocating = .init(testing.allocator); defer encoded.deinit(); - try encode(&t, &encoded); + try encode(testing.allocator, &encoded.writer, &t); try testing.expectEqualDeep(source_memory, primary.pages.memoryStats()); try test_fixture.expectEqual( .snapshot, @@ -321,6 +314,29 @@ test "complete snapshot round trip with history and alternate screen" { encoded.written(), ); + // A complete snapshot can stream through a non-allocating destination. + // Independently hash that output so both its length and complete byte + // sequence are checked without retaining a second snapshot copy. + var discard: std.Io.Writer.Discarding = .init(&.{}); + var hashing = discard.writer.hashed( + std.crypto.hash.Blake3.init(.{}), + &.{}, + ); + try encode(testing.allocator, &hashing.writer, &t); + try testing.expectEqual( + @as(u64, test_complete_fixture.len), + discard.fullCount(), + ); + var expected_digest: checkpoint.Digest = undefined; + std.crypto.hash.Blake3.hash( + &test_complete_fixture, + &expected_digest, + .{}, + ); + var actual_digest: checkpoint.Digest = undefined; + hashing.hasher.final(&actual_digest); + try testing.expectEqual(expected_digest, actual_digest); + // Restore the checked-in reference rather than the just-generated bytes. var encoded_source: std.Io.Reader = .fixed(&test_complete_fixture); var source_buffer: [1]u8 = undefined; @@ -354,7 +370,7 @@ test "complete snapshot round trip with history and alternate screen" { // SCREEN, PAGE, and HISTORY fields and both checkpoint boundaries. var reencoded: std.Io.Writer.Allocating = .init(testing.allocator); defer reencoded.deinit(); - try encode(&restored, &reencoded); + try encode(testing.allocator, &reencoded.writer, &restored); try testing.expectEqualStrings( &test_complete_fixture, reencoded.written(), @@ -363,18 +379,27 @@ test "complete snapshot round trip with history and alternate screen" { // SCREEN and HISTORY keys make both sequence groups order independent. var reversed: std.Io.Writer.Allocating = .init(testing.allocator); defer reversed.deinit(); - try envelope.encode(&reversed.writer); - try terminal.encode(&t, &reversed); - try screen.encode(t.screens.get(.alternate).?, .alternate, &reversed); - try screen.encode(primary, .primary, &reversed); - try checkpoint.encode(.ready, &reversed); + var reversed_stream: record.Writer = .init( + testing.allocator, + &reversed.writer, + ); + defer reversed_stream.deinit(); + try envelope.encode(reversed_stream.writer()); + try terminal.encode(&t, &reversed_stream); + try screen.encode( + t.screens.get(.alternate).?, + .alternate, + &reversed_stream, + ); + try screen.encode(primary, .primary, &reversed_stream); + try checkpoint.encode(.ready, &reversed_stream); try history.encode( t.screens.get(.alternate).?, .alternate, - &reversed, + &reversed_stream, ); - try history.encode(primary, .primary, &reversed); - try checkpoint.encode(.finish, &reversed); + try history.encode(primary, .primary, &reversed_stream); + try checkpoint.encode(.finish, &reversed_stream); var reversed_source: std.Io.Reader = .fixed(reversed.written()); var reversed_restored = try decode( @@ -440,7 +465,7 @@ test "complete snapshot preserves Kitty virtual placeholders" { var encoded: std.Io.Writer.Allocating = .init(testing.allocator); defer encoded.deinit(); - try encode(&t, &encoded); + try encode(testing.allocator, &encoded.writer, &t); var encoded_source: std.Io.Reader = .fixed(encoded.written()); var restored = try decode( @@ -469,7 +494,7 @@ test "complete snapshot preserves Kitty virtual placeholders" { ); } -test "complete snapshot encoding is transactional" { +test "complete snapshot encoding streams from the current writer position" { const testing = std.testing; var t = try Terminal.init(testing.io, testing.allocator, .{ @@ -478,27 +503,44 @@ test "complete snapshot encoding is transactional" { }); defer t.deinit(testing.allocator); - // A complete snapshot cannot be appended after unrelated bytes because - // its envelope and checkpoint coverage both begin at byte zero. + // Prefix hashing begins with this call's envelope, independent of bytes + // that were already present in the destination. var nonempty: std.Io.Writer.Allocating = .init(testing.allocator); defer nonempty.deinit(); try nonempty.writer.writeAll("prefix"); - try testing.expectError( - error.DestinationNotEmpty, - encode(&t, &nonempty), + const snapshot_offset = nonempty.written().len; + try encode(testing.allocator, &nonempty.writer, &t); + try testing.expectEqualStrings( + "prefix", + nonempty.written()[0..snapshot_offset], ); - try testing.expectEqualStrings("prefix", nonempty.written()); + var appended_source: std.Io.Reader = .fixed( + nonempty.written()[snapshot_offset..], + ); + var appended = try decode( + &appended_source, + testing.io, + testing.allocator, + ); + appended.deinit(testing.allocator); - // A failure after validation enters a record codec still rolls back every - // preceding record in this complete snapshot operation. + // Payload validation happens in the record-local scratch allocation. The + // already-streamed envelope remains, but no partial TERMINAL is emitted. t.colors.palette.current[7] = .{ .r = 1, .g = 2, .b = 3 }; var destination: std.Io.Writer.Allocating = .init(testing.allocator); defer destination.deinit(); + try destination.writer.writeAll("prefix"); try testing.expectError( error.InvalidPalette, - encode(&t, &destination), + encode(testing.allocator, &destination.writer, &t), + ); + var expected_envelope: [envelope.encoded_len]u8 = undefined; + var envelope_writer: std.Io.Writer = .fixed(&expected_envelope); + try envelope.encode(&envelope_writer); + try testing.expectEqualStrings( + &expected_envelope, + destination.written()["prefix".len..], ); - try testing.expectEqual(@as(usize, 0), destination.written().len); } test "complete snapshot rejects ordering checkpoints and trailing data" { @@ -515,9 +557,14 @@ test "complete snapshot rejects ordering checkpoints and trailing data" { // primary SCREEN before READY. var reordered: std.Io.Writer.Allocating = .init(testing.allocator); defer reordered.deinit(); - try envelope.encode(&reordered.writer); - try terminal.encode(&t, &reordered); - try history.encode(primary, .primary, &reordered); + var reordered_stream: record.Writer = .init( + testing.allocator, + &reordered.writer, + ); + defer reordered_stream.deinit(); + try envelope.encode(reordered_stream.writer()); + try terminal.encode(&t, &reordered_stream); + try history.encode(primary, .primary, &reordered_stream); var reordered_source: std.Io.Reader = .fixed(reordered.written()); try testing.expectError( error.UnexpectedRecordTag, @@ -528,15 +575,21 @@ test "complete snapshot rejects ordering checkpoints and trailing data" { // digest so the full driver, rather than record CRC validation, rejects it. var invalid_ready: std.Io.Writer.Allocating = .init(testing.allocator); defer invalid_ready.deinit(); - try envelope.encode(&invalid_ready.writer); - try terminal.encode(&t, &invalid_ready); - try screen.encode(primary, .primary, &invalid_ready); - var record_writer = try record.Writer.init(&invalid_ready, .ready); - try record_writer.payloadWriter().splatByteAll( + var invalid_ready_stream: record.Writer = .init( + testing.allocator, + &invalid_ready.writer, + ); + defer invalid_ready_stream.deinit(); + try envelope.encode(invalid_ready_stream.writer()); + try terminal.encode(&t, &invalid_ready_stream); + try screen.encode(primary, .primary, &invalid_ready_stream); + const ready_payload = invalid_ready_stream.begin(.ready); + errdefer invalid_ready_stream.cancel(); + try ready_payload.splatByteAll( 0, @sizeOf(checkpoint.Digest), ); - try record_writer.finish(); + try invalid_ready_stream.finish(); var invalid_ready_source: std.Io.Reader = .fixed( invalid_ready.written(), ); @@ -548,7 +601,7 @@ test "complete snapshot rejects ordering checkpoints and trailing data" { // FINISH is the exact end of a complete snapshot. var trailing: std.Io.Writer.Allocating = .init(testing.allocator); defer trailing.deinit(); - try encode(&t, &trailing); + try encode(testing.allocator, &trailing.writer, &t); try trailing.writer.writeByte(0); var trailing_source: std.Io.Reader = .fixed(trailing.written()); try testing.expectError( @@ -559,9 +612,14 @@ test "complete snapshot rejects ordering checkpoints and trailing data" { // A SCREEN key must name one of the slots declared by TERMINAL. var undeclared: std.Io.Writer.Allocating = .init(testing.allocator); defer undeclared.deinit(); - try envelope.encode(&undeclared.writer); - try terminal.encode(&t, &undeclared); - try screen.encode(primary, .alternate, &undeclared); + var undeclared_stream: record.Writer = .init( + testing.allocator, + &undeclared.writer, + ); + defer undeclared_stream.deinit(); + try envelope.encode(undeclared_stream.writer()); + try terminal.encode(&t, &undeclared_stream); + try screen.encode(primary, .alternate, &undeclared_stream); var undeclared_source: std.Io.Reader = .fixed(undeclared.written()); try testing.expectError( error.UnexpectedScreenKey, @@ -572,11 +630,16 @@ test "complete snapshot rejects ordering checkpoints and trailing data" { // screen even when the sequence contains no PAGE records. var undeclared_history: std.Io.Writer.Allocating = .init(testing.allocator); defer undeclared_history.deinit(); - try envelope.encode(&undeclared_history.writer); - try terminal.encode(&t, &undeclared_history); - try screen.encode(primary, .primary, &undeclared_history); - try checkpoint.encode(.ready, &undeclared_history); - try history.encode(primary, .alternate, &undeclared_history); + var undeclared_history_stream: record.Writer = .init( + testing.allocator, + &undeclared_history.writer, + ); + defer undeclared_history_stream.deinit(); + try envelope.encode(undeclared_history_stream.writer()); + try terminal.encode(&t, &undeclared_history_stream); + try screen.encode(primary, .primary, &undeclared_history_stream); + try checkpoint.encode(.ready, &undeclared_history_stream); + try history.encode(primary, .alternate, &undeclared_history_stream); var undeclared_history_source: std.Io.Reader = .fixed( undeclared_history.written(), ); @@ -589,10 +652,15 @@ test "complete snapshot rejects ordering checkpoints and trailing data" { _ = try t.switchScreen(.alternate); var duplicate: std.Io.Writer.Allocating = .init(testing.allocator); defer duplicate.deinit(); - try envelope.encode(&duplicate.writer); - try terminal.encode(&t, &duplicate); - try screen.encode(primary, .primary, &duplicate); - try screen.encode(primary, .primary, &duplicate); + var duplicate_stream: record.Writer = .init( + testing.allocator, + &duplicate.writer, + ); + defer duplicate_stream.deinit(); + try envelope.encode(duplicate_stream.writer()); + try terminal.encode(&t, &duplicate_stream); + try screen.encode(primary, .primary, &duplicate_stream); + try screen.encode(primary, .primary, &duplicate_stream); var duplicate_source: std.Io.Reader = .fixed(duplicate.written()); try testing.expectError( error.DuplicateScreen, @@ -602,17 +670,22 @@ test "complete snapshot rejects ordering checkpoints and trailing data" { // The declared count cannot be satisfied by repeating one HISTORY key. var duplicate_history: std.Io.Writer.Allocating = .init(testing.allocator); defer duplicate_history.deinit(); - try envelope.encode(&duplicate_history.writer); - try terminal.encode(&t, &duplicate_history); - try screen.encode(primary, .primary, &duplicate_history); + var duplicate_history_stream: record.Writer = .init( + testing.allocator, + &duplicate_history.writer, + ); + defer duplicate_history_stream.deinit(); + try envelope.encode(duplicate_history_stream.writer()); + try terminal.encode(&t, &duplicate_history_stream); + try screen.encode(primary, .primary, &duplicate_history_stream); try screen.encode( t.screens.get(.alternate).?, .alternate, - &duplicate_history, + &duplicate_history_stream, ); - try checkpoint.encode(.ready, &duplicate_history); - try history.encode(primary, .primary, &duplicate_history); - try history.encode(primary, .primary, &duplicate_history); + try checkpoint.encode(.ready, &duplicate_history_stream); + try history.encode(primary, .primary, &duplicate_history_stream); + try history.encode(primary, .primary, &duplicate_history_stream); var duplicate_history_source: std.Io.Reader = .fixed( duplicate_history.written(), ); diff --git a/src/terminal/snapshot/terminal.zig b/src/terminal/snapshot/terminal.zig index 343392fa4..f20b97707 100644 --- a/src/terminal/snapshot/terminal.zig +++ b/src/terminal/snapshot/terminal.zig @@ -909,26 +909,25 @@ pub const EncodeError = HeaderInitError || /// Encode terminal-wide native state as one framed TERMINAL record. /// -/// State not represented by this snapshot version is ignored. Any failure -/// removes the incomplete record while preserving bytes that were already -/// present in `destination`. +/// State not represented by this snapshot version is ignored. Payload failures +/// emit no part of the TERMINAL record. pub fn encode( terminal: *const Terminal, - destination: *std.Io.Writer.Allocating, + destination: *record.Writer, ) EncodeError!void { const header = try Header.init(terminal); - var record_writer = try record.Writer.init(destination, .terminal); - errdefer record_writer.cancel(); + const payload = destination.begin(.terminal); + errdefer destination.cancel(); try encodePayload( header, &terminal.tabstops, &terminal.colors.palette, terminal.getPwd() orelse "", terminal.getTitle() orelse "", - record_writer.payloadWriter(), + payload, ); - try record_writer.finish(); + try destination.finish(); } /// Errors possible while decoding one native TERMINAL record. @@ -1638,7 +1637,12 @@ test "TERMINAL record encodes native terminal state" { // Encode and restore through the public record codec. var destination: std.Io.Writer.Allocating = .init(testing.allocator); defer destination.deinit(); - try encode(&terminal, &destination); + var stream: record.Writer = .init( + testing.allocator, + &destination.writer, + ); + defer stream.deinit(); + try encode(&terminal, &stream); var source: std.Io.Reader = .fixed(destination.written()); var restored = try decode( @@ -1714,7 +1718,12 @@ test "TERMINAL record declares an initialized alternate screen" { var destination: std.Io.Writer.Allocating = .init(testing.allocator); defer destination.deinit(); - try encode(&terminal, &destination); + var stream: record.Writer = .init( + testing.allocator, + &destination.writer, + ); + defer stream.deinit(); + try encode(&terminal, &stream); var source: std.Io.Reader = .fixed(destination.written()); var restored = try decode( @@ -1737,7 +1746,7 @@ test "TERMINAL record declares an initialized alternate screen" { ); } -test "TERMINAL record encoding is transactional" { +test "TERMINAL payload failure emits no incomplete record" { const testing = std.testing; var terminal = try Terminal.init(testing.io, testing.allocator, .{ @@ -1749,13 +1758,18 @@ test "TERMINAL record encoding is transactional" { var destination: std.Io.Writer.Allocating = .init(testing.allocator); defer destination.deinit(); try destination.writer.writeAll("prefix"); + var stream: record.Writer = .init( + testing.allocator, + &destination.writer, + ); + defer stream.deinit(); - // Payload validation occurs after framing is reserved, so this covers the - // record writer's rollback path. + // Payload validation occurs in the reusable record scratch buffer. The + // invalid TERMINAL is never emitted to the destination. terminal.colors.palette.current[7] = .{ .r = 1, .g = 2, .b = 3 }; try testing.expectError( error.InvalidPalette, - encode(&terminal, &destination), + encode(&terminal, &stream), ); try testing.expectEqualStrings("prefix", destination.written()); } @@ -1766,9 +1780,15 @@ test "TERMINAL record decoding rejects malformed input transactionally" { // A valid record of another type is not accepted as TERMINAL state. var wrong_tag: std.Io.Writer.Allocating = .init(testing.allocator); defer wrong_tag.deinit(); + var wrong_tag_stream: record.Writer = .init( + testing.allocator, + &wrong_tag.writer, + ); + defer wrong_tag_stream.deinit(); { - var record_writer = try record.Writer.init(&wrong_tag, .screen); - try record_writer.finish(); + _ = wrong_tag_stream.begin(.screen); + errdefer wrong_tag_stream.cancel(); + try wrong_tag_stream.finish(); } var wrong_tag_source: std.Io.Reader = .fixed(wrong_tag.written()); try testing.expectError( @@ -1788,7 +1808,12 @@ test "TERMINAL record decoding rejects malformed input transactionally" { var encoded: std.Io.Writer.Allocating = .init(testing.allocator); defer encoded.deinit(); - try encode(&terminal, &encoded); + var encoded_stream: record.Writer = .init( + testing.allocator, + &encoded.writer, + ); + defer encoded_stream.deinit(); + try encode(&terminal, &encoded_stream); for (0..encoded.written().len) |fixture_len| { var source: std.Io.Reader = .fixed(