terminal/snapshot: less buffering, better stream writing

Stream complete snapshot records to any std.Io.Writer while retaining one reusable payload buffer for length and CRC calculation. Update BLAKE3 incrementally so checkpoints no longer require rehashing an allocating destination.

Wrap decode hashing in StreamReader to enforce exact checkpoint boundaries. Preserve v1 bytes while allowing snapshots to begin at the current writer position and retaining only valid prefixes on failures.
This commit is contained in:
Mitchell Hashimoto
2026-07-31 11:17:56 -07:00
parent 58e92098a2
commit f0fe788fcc
8 changed files with 622 additions and 285 deletions

View File

@@ -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),
);
}

View File

@@ -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(

View File

@@ -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`.

View File

@@ -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());
}

View File

@@ -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);
}

View File

@@ -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(

View File

@@ -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(),
);

View File

@@ -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(