tiny.sql.wal
Defined in tiny.sql.
API (44)
Actions
Public operations.
Control.checkExtent.classifyFrame.committedScanner.Capacity.deriveScanner.consumeFrameScanner.deinitScanner.extendScanner.findCommittedScanner.initScanner.nextFrameScanner.outcomeScanner.progressScanner.resetScanner.selectedFrameScanner.verifySelectedendMarkendMarkControlledframeOffsetpageAt
Types and contracts
Public types and contracts.
ControlControlledReadErrorErrorExtentExtent.AdmittedExtent.OverCapacityFrameFrameRefFrameRequestHeaderReadErrorReaderSaltScannerScanner.CapacityScanner.FrameOutcomeScanner.MalformedFrameScanner.OutcomeScanner.ProgressScanner.StateScanner.WorkspaceWriter
Values and defaults
Public values and defaults.
Source
Source: lib/sql/src/root.zig:47
zig
pub const wal = @import("wal.zig");Source: lib/sql/src/wal.zig:8
zig
const std = @import("std");const alloc_phase = @import("alloc_phase");const page = @import("page.zig");const trace = @import("trace.zig");const Allocator = std.mem.Allocator;pub const header_size: usize = 32;pub const frame_header_size: usize = 24;pub const frame_size: usize = frame_header_size + page.size;const controlled_copy_chunk_bytes: usize = 1 * 1024 * 1024;pub const ReadError = error{ InvalidChecksum, InvalidWal, UnsupportedPageSize,};pub const ControlledReadError = ReadError || error{Interrupted};pub const Error = Allocator.Error || ReadError || error{ CapacityOverflow, WalFull,};pub const Salt = struct { first: u32, second: u32,};pub const Header = struct { sequence: u32, salt: Salt,};pub const Frame = struct { index: usize, page_id: u32, db_page_count: u32, image: []const u8, pub fn committed(self: Frame) bool { return self.db_page_count != 0; }};pub const FrameRef = extern struct { page_id: u32, frame: u32, checksum_before: [2]u32, checksum_after: [2]u32,};comptime { std.debug.assert(@sizeOf(FrameRef) == 24);}pub const Control = struct { context: ?*anyopaque = null, interrupted_fn: *const fn (context: ?*anyopaque) bool = neverInterrupted, pub fn check(self: Control) error{Interrupted}!void { if (self.interrupted_fn(self.context)) return error.Interrupted; } fn neverInterrupted(_: ?*anyopaque) bool { return false; }};pub const Extent = union(enum) { partial_header: u8, admitted: Admitted, over_capacity: OverCapacity, pub const Admitted = struct { complete_frames: u32, partial_tail_bytes: u16, }; pub const OverCapacity = struct { complete_frames: u64, partial_tail_bytes: u16, }; pub fn classify( file_bytes: u64, frame_capacity: usize, control: Control, ) error{ CapacityOverflow, Interrupted }!Extent { try control.check(); if (frame_capacity > std.math.maxInt(u32)) return error.CapacityOverflow; if (file_bytes < header_size) { return .{ .partial_header = @intCast(file_bytes) }; } const body_bytes = file_bytes - header_size; const complete_frames = body_bytes / frame_size; const partial_tail_bytes: u16 = @intCast(body_bytes % frame_size); const capacity: u64 = @intCast(frame_capacity); if (complete_frames > capacity or (complete_frames == capacity and partial_tail_bytes > 0)) { return .{ .over_capacity = .{ .complete_frames = complete_frames, .partial_tail_bytes = partial_tail_bytes, } }; } return .{ .admitted = .{ .complete_frames = @intCast(complete_frames), .partial_tail_bytes = partial_tail_bytes, } }; }};pub const Scanner = struct { pub const Workspace = struct { scratch: *[frame_size]u8, refs: []FrameRef, }; pub const Capacity = struct { frames: u32, storage_bytes: usize, pub fn derive(frame_capacity: usize) error{CapacityOverflow}!Capacity { if (frame_capacity > std.math.maxInt(u32)) return error.CapacityOverflow; return .{ .frames = @intCast(frame_capacity), .storage_bytes = try storageBytes( frame_capacity, frame_size, @sizeOf(FrameRef), ), }; } }; pub const State = enum { uninitialized, scanning, ready, invalidated, teardown, }; pub const MalformedFrame = enum { invalid_page_id, salt_mismatch, checksum_mismatch, database_page_count_regressed, page_outside_database, }; pub const FrameOutcome = union(enum) { staged, committed, malformed_full_frame: MalformedFrame, }; pub const Outcome = union(enum) { complete, partial_uncommitted_tail: u16, }; pub const Progress = struct { state: State, header: [header_size]u8, tail_checksum: [2]u32, target_frames: u32, covered_frames: u32, committed_end_mark: u32, logical_database_page_count: u32, committed_ref_count: u32, staged_transaction_start: u32, staged_ref_count: u32, partial_tail_bytes: u16, }; const ValidatedFrame = struct { ref: FrameRef, checksum: Checksum, database_page_count: u32, }; const Validation = union(enum) { valid: ValidatedFrame, malformed: MalformedFrame, }; workspace: Workspace, capacity: Capacity, state: State = .uninitialized, header: [header_size]u8 = @splat(0), salt: Salt = .{ .first = 0, .second = 0 }, tail_checksum: Checksum = .{}, target_frames: u32 = 0, covered_frames: u32 = 0, committed_end_mark: u32 = 0, logical_database_page_count: u32 = 0, committed_ref_count: u32 = 0, staged_transaction_start: u32 = 0, staged_ref_count: u32 = 0, partial_tail_bytes: u16 = 0, pub fn init(workspace: Workspace) error{CapacityOverflow}!Scanner { return .{ .workspace = workspace, .capacity = try Capacity.derive(workspace.refs.len), }; } pub fn deinit(self: *Scanner) Workspace { std.debug.assert(self.state != .teardown); const workspace = self.workspace; self.state = .teardown; self.workspace.refs = &.{}; return workspace; } pub fn reset( self: *Scanner, header: *const [header_size]u8, extent: Extent.Admitted, database_page_floor: u32, control: Control, ) (ReadError || error{ CapacityExceeded, Interrupted, InvalidExtent, ScannerNotReady })!void { if (self.state == .teardown) return error.ScannerNotReady; if (extent.partial_tail_bytes >= frame_size) return error.InvalidExtent; if (extent.complete_frames > self.capacity.frames) { self.state = .invalidated; return error.CapacityExceeded; } const decoded = decodeHeaderControlled(header, control) catch |err| switch (err) { error.Interrupted => return error.Interrupted, else => { self.state = .invalidated; return err; }, }; try control.check(); self.header = header.*; self.salt = decoded.salt; self.tail_checksum = decoded.checksum; self.target_frames = extent.complete_frames; self.covered_frames = 0; self.committed_end_mark = 0; self.logical_database_page_count = database_page_floor; self.committed_ref_count = 0; self.staged_transaction_start = 0; self.staged_ref_count = 0; self.partial_tail_bytes = extent.partial_tail_bytes; self.state = if (extent.complete_frames == 0) .ready else .scanning; self.assertBounds(); } pub fn extend( self: *Scanner, extent: Extent.Admitted, control: Control, ) error{ CapacityExceeded, Interrupted, InvalidExtent, RewindRequired, ScannerNotReady }!void { if (self.state != .ready) return error.ScannerNotReady; if (extent.partial_tail_bytes >= frame_size) return error.InvalidExtent; try control.check(); if (extent.complete_frames > self.capacity.frames) { self.state = .invalidated; return error.CapacityExceeded; } const rewound = extent.complete_frames < self.covered_frames or (extent.complete_frames == self.covered_frames and extent.partial_tail_bytes < self.partial_tail_bytes); if (rewound) { self.state = .invalidated; return error.RewindRequired; } self.target_frames = extent.complete_frames; self.partial_tail_bytes = extent.partial_tail_bytes; if (self.covered_frames < self.target_frames) self.state = .scanning; self.assertBounds(); } pub fn nextFrame(self: *Scanner) error{ScannerNotReady}!?FrameRequest { if (self.state != .scanning and self.state != .ready) { return error.ScannerNotReady; } self.assertBounds(); if (self.state == .ready) return null; std.debug.assert(self.covered_frames < self.target_frames); const frame = self.covered_frames + 1; return .{ .frame = frame, .offset = frameOffset(frame), .bytes = self.workspace.scratch, }; } pub fn consumeFrame( self: *Scanner, control: Control, ) error{ Interrupted, ScannerNotReady }!FrameOutcome { if (self.state != .scanning) return error.ScannerNotReady; self.assertBounds(); const validation = try self.validateFrame(control); const validated = switch (validation) { .valid => |value| value, .malformed => |reason| return self.invalidate(reason), }; if (validated.database_page_count != 0) { if (try self.invalidMarker(validated, control)) |reason| { return self.invalidate(reason); } } try control.check(); self.appendValidated(validated); if (validated.database_page_count == 0) { self.finishPhysicalFrame(); return .staged; } self.foldStaged(control) catch |err| { self.state = .invalidated; return err; }; self.committed_end_mark = validated.ref.frame; self.logical_database_page_count = validated.database_page_count; self.staged_transaction_start = 0; self.staged_ref_count = 0; self.finishPhysicalFrame(); return .committed; } pub fn outcome(self: *const Scanner) error{ScannerNotReady}!Outcome { if (self.state != .ready) return error.ScannerNotReady; self.assertBounds(); if (self.partial_tail_bytes == 0) return .complete; return .{ .partial_uncommitted_tail = self.partial_tail_bytes }; } pub fn progress(self: *const Scanner) Progress { return .{ .state = self.state, .header = self.header, .tail_checksum = self.tail_checksum.pair(), .target_frames = self.target_frames, .covered_frames = self.covered_frames, .committed_end_mark = self.committed_end_mark, .logical_database_page_count = self.logical_database_page_count, .committed_ref_count = self.committed_ref_count, .staged_transaction_start = self.staged_transaction_start, .staged_ref_count = self.staged_ref_count, .partial_tail_bytes = self.partial_tail_bytes, }; } pub fn findCommitted( self: *const Scanner, page_id: u32, control: Control, ) error{ Interrupted, ScannerNotReady }!?FrameRef { if (self.state != .ready) return error.ScannerNotReady; self.assertBounds(); try control.check(); if (page_id == 0 or page_id > self.logical_database_page_count) return null; var left: u32 = 0; var right = self.committed_ref_count; while (left < right) { try control.check(); const middle = left + (right - left) / 2; const ref = self.workspace.refs[middle]; if (ref.page_id < page_id) { left = middle + 1; } else { right = middle; } } if (left == self.committed_ref_count) return null; const ref = self.workspace.refs[left]; return if (ref.page_id == page_id) ref else null; } pub fn selectedFrame( self: *Scanner, ref: FrameRef, control: Control, ) error{ Interrupted, ScannerNotReady, StaleFrameRef }!FrameRequest { const current = (try self.findCommitted(ref.page_id, control)) orelse return error.StaleFrameRef; if (!std.meta.eql(current, ref)) return error.StaleFrameRef; return .{ .frame = ref.frame, .offset = frameOffset(ref.frame), .bytes = self.workspace.scratch, }; } pub fn verifySelected( self: *const Scanner, ref: FrameRef, destination: *[page.size]u8, control: Control, ) error{ Interrupted, ScannerNotReady, SelectedFrameMismatch, StaleFrameRef }!void { if (self.state != .ready) return error.ScannerNotReady; std.debug.assert(disjoint(self.workspace.scratch, destination)); const current = (try self.findCommitted(ref.page_id, control)) orelse return error.StaleFrameRef; if (!std.meta.eql(current, ref)) return error.StaleFrameRef; const frame = self.workspace.scratch; try control.check(); if (readU32(frame[0..4]) != ref.page_id) return error.SelectedFrameMismatch; if (readU32(frame[8..12]) != self.salt.first) return error.SelectedFrameMismatch; if (readU32(frame[12..16]) != self.salt.second) return error.SelectedFrameMismatch; var checksum = Checksum.fromPair(ref.checksum_before); try checksum.updateControlled(frame[0..8], control); try checksum.updateControlled(frame[frame_header_size..], control); const on_disk = [2]u32{ readU32(frame[16..20]), readU32(frame[20..24]), }; if (!std.mem.eql(u32, &on_disk, &ref.checksum_after)) { return error.SelectedFrameMismatch; } const computed = checksum.pair(); if (!std.mem.eql(u32, &on_disk, &computed)) { return error.SelectedFrameMismatch; } try control.check(); @memcpy(destination, frame[frame_header_size..]); } fn validateFrame( self: *const Scanner, control: Control, ) error{Interrupted}!Validation { try control.check(); const frame = self.workspace.scratch; const page_id = readU32(frame[0..4]); const database_page_count = readU32(frame[4..8]); if (page_id == 0) return .{ .malformed = .invalid_page_id }; if (readU32(frame[8..12]) != self.salt.first or readU32(frame[12..16]) != self.salt.second) { return .{ .malformed = .salt_mismatch }; } var checksum = self.tail_checksum; try checksum.updateControlled(frame[0..8], control); try checksum.updateControlled(frame[frame_header_size..], control); if (readU32(frame[16..20]) != checksum.first or readU32(frame[20..24]) != checksum.second) { return .{ .malformed = .checksum_mismatch }; } const frame_index = self.covered_frames + 1; return .{ .valid = .{ .ref = .{ .page_id = page_id, .frame = frame_index, .checksum_before = self.tail_checksum.pair(), .checksum_after = checksum.pair(), }, .checksum = checksum, .database_page_count = database_page_count, } }; } fn invalidMarker( self: *const Scanner, validated: ValidatedFrame, control: Control, ) error{Interrupted}!?MalformedFrame { const page_count = validated.database_page_count; if (page_count < self.logical_database_page_count) { return .database_page_count_regressed; } if (validated.ref.page_id > page_count) return .page_outside_database; var index: u32 = 0; while (index < self.staged_ref_count) : (index += 1) { try control.check(); const staged = self.workspace.refs[self.committed_ref_count + index]; if (staged.page_id > page_count) return .page_outside_database; } return null; } fn appendValidated(self: *Scanner, validated: ValidatedFrame) void { const index = self.committed_ref_count + self.staged_ref_count; std.debug.assert(index < self.capacity.frames); std.debug.assert(validated.ref.frame == self.covered_frames + 1); self.workspace.refs[index] = validated.ref; self.tail_checksum = validated.checksum; self.covered_frames += 1; if (self.staged_ref_count == 0) { self.staged_transaction_start = validated.ref.frame; } self.staged_ref_count += 1; std.debug.assert(self.covered_frames <= self.target_frames); std.debug.assert( @as(u64, self.committed_ref_count) + self.staged_ref_count <= self.capacity.frames, ); std.debug.assert( @as(u64, self.committed_ref_count) + self.staged_ref_count <= self.covered_frames, ); } fn foldStaged(self: *Scanner, control: Control) error{Interrupted}!void { var committed = self.committed_ref_count; var staged = self.staged_ref_count; while (staged > 0) { const candidate = self.workspace.refs[committed]; var insertion: u32 = 0; while (insertion < committed and self.workspace.refs[insertion].page_id < candidate.page_id) { try self.checkFold(control); insertion += 1; } const replaces = insertion < committed and self.workspace.refs[insertion].page_id == candidate.page_id; if (replaces) { try self.replaceCommitted(insertion, committed, staged, candidate, control); staged -= 1; } else { try self.insertCommitted(insertion, committed, candidate, control); committed += 1; staged -= 1; } try self.checkFold(control); } self.committed_ref_count = committed; } fn replaceCommitted( self: *Scanner, insertion: u32, committed: u32, staged: u32, candidate: FrameRef, control: Control, ) error{Interrupted}!void { try self.checkFold(control); self.workspace.refs[insertion] = candidate; var index = committed; while (index + 1 < committed + staged) : (index += 1) { try self.checkFold(control); self.workspace.refs[index] = self.workspace.refs[index + 1]; } } fn insertCommitted( self: *Scanner, insertion: u32, committed: u32, candidate: FrameRef, control: Control, ) error{Interrupted}!void { var index = committed; while (index > insertion) { try self.checkFold(control); self.workspace.refs[index] = self.workspace.refs[index - 1]; index -= 1; } try self.checkFold(control); self.workspace.refs[insertion] = candidate; } fn checkFold(self: *Scanner, control: Control) error{Interrupted}!void { control.check() catch { self.state = .invalidated; return error.Interrupted; }; } fn finishPhysicalFrame(self: *Scanner) void { std.debug.assert(self.covered_frames <= self.target_frames); if (self.covered_frames == self.target_frames) self.state = .ready; self.assertBounds(); } fn invalidate(self: *Scanner, reason: MalformedFrame) FrameOutcome { self.state = .invalidated; return .{ .malformed_full_frame = reason }; } fn storageBytes( frame_capacity: usize, scratch_bytes: usize, ref_bytes: usize, ) error{CapacityOverflow}!usize { const refs_bytes = std.math.mul( usize, frame_capacity, ref_bytes, ) catch return error.CapacityOverflow; return std.math.add( usize, scratch_bytes, refs_bytes, ) catch return error.CapacityOverflow; } fn assertBounds(self: *const Scanner) void { std.debug.assert(self.target_frames <= self.capacity.frames); std.debug.assert(self.covered_frames <= self.target_frames); std.debug.assert(self.committed_end_mark <= self.covered_frames); std.debug.assert(self.committed_ref_count <= self.covered_frames); std.debug.assert(self.staged_ref_count <= self.covered_frames); std.debug.assert( @as(u64, self.committed_ref_count) + self.staged_ref_count <= self.capacity.frames, ); std.debug.assert( @as(u64, self.committed_ref_count) + self.staged_ref_count <= self.covered_frames, ); std.debug.assert((self.staged_ref_count == 0) == (self.staged_transaction_start == 0)); if (self.state == .ready) { std.debug.assert(self.covered_frames == self.target_frames); } if (self.state == .scanning) { std.debug.assert(self.covered_frames < self.target_frames); } }};pub const FrameRequest = struct { frame: u32, offset: u64, bytes: *[frame_size]u8,};pub fn frameOffset(frame: u32) u64 { std.debug.assert(frame > 0); return @as(u64, header_size) + @as(u64, frame - 1) * @as(u64, frame_size);}pub const Writer = struct { pub const storage_alignment: usize = @alignOf(u8); pub const Storage = []align(storage_alignment) u8; pub const Limits = struct { header: Header, frames: usize, }; pub const Capacity = struct { frames: usize, storage_bytes: usize, pub const DeriveError = error{CapacityOverflow}; pub fn derive(limits: Limits) DeriveError!Capacity { const frame_bytes = std.math.mul(usize, limits.frames, frame_size) catch { return error.CapacityOverflow; }; const storage_bytes = std.math.add(usize, header_size, frame_bytes) catch { return error.CapacityOverflow; }; return .{ .frames = limits.frames, .storage_bytes = storage_bytes }; } }; pub const Exhaustion = error{WalFull}; pub const InitError = Capacity.DeriveError || error{StorageTooShort}; pub const LoadError = ReadError || Exhaustion; pub const ControlledLoadError = LoadError || error{Interrupted}; pub const work_limits: alloc_phase.capacity.WorkLimits = .{ .transition_steps_max = 1, .cleanup_steps_per_call_max = 0, .cleanup_calls_at_capacity_max = 0, }; pub const claim: alloc_phase.capacity.Declaration = .{ .source = .{ .id = "sql.wal_writer", .kind = .phase_static, .limit_source = .caller, .storage = .{ .covered = &.{ .{ .id = "exact_header_and_committed_frame_byte_region", .lifetime = .steady, .detail = "mutable borrow of the caller-provisioned exact header and frame byte region", }, .{ .id = "invisible_staged_frame_tail_transferred_into_the_co_0ed8d0c7aca0", .lifetime = .transferred, .detail = "invisible staged frame tail transferred into the committed prefix", }, }, .excluded = &.{ "caller-owned storage outside the exact borrowed prefix", "caller-owned page images and recovery source bytes", "pager frame indexes, page indexes, checkpoint plans, files, and I/O runtime state", "trace instrumentation and allocator implementation state", }, }, .capacity = .{ .inputs = &.{ alloc_phase.capacity.bindInput(Limits, "frames", "frames"), }, .type_selectors = &.{ alloc_phase.capacity.bindType([frame_size]u8, "walframebytes"), alloc_phase.capacity.bindType([header_size]u8, "walheaderbytes"), }, .nodes = &.{ .{ .input = 0 }, .{ .scale = .{ .node = 0, .coefficient = .{ .size_of_concrete_type = 0 } } }, .{ .constant = 1 }, .{ .scale = .{ .node = 2, .coefficient = .{ .size_of_concrete_type = 1 } } }, .{ .add = .{ .left = 1, .right = 3 } }, }, .assertions = &.{.{ .scope = .closure_total, .measure = .retained, .relation = .exact, .expression = 4, }}, }, .overload = .{ .kind = .reject_before_mutation, .detail = "append staging and recovery load return WalFull before bytes length salt or checksum mutate", }, .risks = .{ .transitive = .{ .status = .open, .detail = "pager indexes and trace instrumentation are outside the writer-owned byte region", }, .foreign = .{ .status = .open, .detail = "file reads and writes that consume WAL bytes are owned by the file database", }, }, .work = .{ .equation = "transition_steps <= transition_steps_max" }, .obligations = &.{ .{ .key = "sql_wal_capacity", .role = .capacity_model }, .{ .key = "sql_wal_storage_rejection", .role = .initialization_failure }, .{ .key = "sql_wal_sealed", .role = .overload }, .{ .key = "sql_wal_work_bound", .role = .work_bound }, .{ .key = "sql_wal_semantics", .role = .custom }, .{ .key = "sql_wal_staged_transfer", .role = .custom }, .{ .key = "sql_wal_corruption", .role = .custom }, }, }, .bindings = .{ .owner = @This(), .seal = .{ .family = alloc_phase.capacity.selector(@This().activate), .premise = .{ .class = .checked_semantic_fact, .authority = .checker, }, }, .teardown = .{ .family = alloc_phase.capacity.selector(@This().deinit), .premise = .{ .class = .checked_semantic_fact, .authority = .checker, }, }, }, }; phase: alloc_phase.capacity.Phase, capacity: Capacity, storage: Storage, len: usize, salt: Salt, checksum: Checksum, pub const Position = struct { len: usize, checksum_first: u32, checksum_second: u32, }; pub fn init(storage: Storage, limits: Limits) InitError!Writer { const capacity = try Capacity.derive(limits); if (storage.len < capacity.storage_bytes) return error.StorageTooShort; const borrowed = storage[0..capacity.storage_bytes]; const checksum = encodeHeader(borrowed[0..header_size], limits.header); return .{ .phase = .initialization, .capacity = capacity, .storage = borrowed, .len = header_size, .salt = limits.header.salt, .checksum = checksum, }; } pub fn activate(self: *Writer) void { std.debug.assert(self.phase == .initialization); std.debug.assert(self.storage.len == self.capacity.storage_bytes); std.debug.assert(self.len == header_size); self.phase = .steady; } pub fn deinit(self: *Writer) Storage { std.debug.assert(self.phase != .teardown); std.debug.assert(self.storage.len == self.capacity.storage_bytes); self.phase = .teardown; const storage = self.storage; self.* = undefined; return storage; } pub fn bytes(self: *const Writer) []const u8 { std.debug.assert(self.phase == .steady); std.debug.assert(self.len <= self.storage.len); return self.storage[0..self.len]; } pub fn frameCount(self: *const Writer) usize { std.debug.assert(self.phase == .steady); return (self.len - header_size) / frame_size; } pub fn frameCapacity(self: *const Writer) usize { std.debug.assert(self.phase == .steady); return self.capacity.frames; } pub fn byteCapacity(self: *const Writer) usize { std.debug.assert(self.phase == .steady); return self.capacity.storage_bytes; } pub fn remainingFrames(self: *const Writer) usize { std.debug.assert(self.phase == .steady); return self.capacity.frames - self.frameCount(); } pub fn position(self: *const Writer) Position { std.debug.assert(self.phase == .steady); return .{ .len = self.len, .checksum_first = self.checksum.first, .checksum_second = self.checksum.second, }; } pub fn restore(self: *Writer, position_value: Position) void { std.debug.assert(self.phase == .steady); std.debug.assert(position_value.len >= header_size); std.debug.assert((position_value.len - header_size) % frame_size == 0); std.debug.assert(position_value.len <= self.len); self.len = position_value.len; self.checksum = .{ .first = position_value.checksum_first, .second = position_value.checksum_second, }; } pub fn load(self: *Writer, source: []const u8, committed_len: usize) LoadError!void { return self.loadControlled(source, committed_len, .{}) catch |err| switch (err) { error.Interrupted => unreachable, else => return @errorCast(err), }; } pub fn loadControlled( self: *Writer, source: []const u8, committed_len: usize, control: Control, ) ControlledLoadError!void { std.debug.assert(self.phase == .steady); if (committed_len > source.len) return error.InvalidWal; if (committed_len > self.capacity.storage_bytes) return error.WalFull; var reader = try Reader.initControlled(source[0..committed_len], control); const frames_max = (committed_len - header_size) / frame_size; var frame_index: usize = 0; while (frame_index < frames_max) : (frame_index += 1) { if (try reader.nextControlled(control) == null) break; } if (reader.offset != committed_len) return error.InvalidWal; var copied: usize = 0; var chunks: usize = 0; const chunks_max = std.math.divCeil( usize, committed_len, controlled_copy_chunk_bytes, ) catch unreachable; while (copied < committed_len) : (chunks += 1) { std.debug.assert(chunks < chunks_max); try control.check(); const end = copied + @min( controlled_copy_chunk_bytes, committed_len - copied, ); std.mem.copyForwards(u8, self.storage[copied..end], source[copied..end]); copied = end; } try control.check(); self.len = committed_len; self.salt = reader.salt; self.checksum = reader.checksum; } pub fn append(self: *Writer, page_id: u32, db_page_count: u32, image: *const [page.size]u8) Exhaustion!void { const phase = trace.scope("wal.append"); defer phase.end(); std.debug.assert(self.phase == .steady); if (self.remainingFrames() == 0) return error.WalFull; var frame_header = @as([frame_header_size]u8, @splat(0)); writeU32(frame_header[0..4], page_id); writeU32(frame_header[4..8], db_page_count); writeU32(frame_header[8..12], self.salt.first); writeU32(frame_header[12..16], self.salt.second); var checksum = self.checksum; checksum.update(frame_header[0..8]); checksum.update(image[0..]); writeU32(frame_header[16..20], checksum.first); writeU32(frame_header[20..24], checksum.second); @memcpy(self.storage[self.len..][0..frame_header_size], frame_header[0..]); @memcpy(self.storage[self.len + frame_header_size ..][0..page.size], image[0..]); self.len += frame_size; self.checksum = checksum; trace.progress("wal.append.complete"); } pub fn stagePage(self: *Writer, position_value: Position, index: usize, page_id: u32, image: *const [page.size]u8) Exhaustion!void { std.debug.assert(self.phase == .steady); std.debug.assert(position_value.len == self.len); const offset = try self.stagedFrameOffset(position_value, index); writeU32(self.storage[offset..][0..4], page_id); const staged = self.storage[offset + frame_header_size ..][0..page.size]; if (staged != image) @memcpy(staged, image); } pub fn stagedPage(self: *const Writer, position_value: Position, index: usize) []const u8 { std.debug.assert(self.phase == .steady); std.debug.assert(position_value.len == self.len); const offset = self.stagedFrameOffset(position_value, index) catch unreachable; return self.storage[offset + frame_header_size ..][0..page.size]; } /// Returns a staged image for editing in place. Staging the returned /// image again copies nothing, and the checksum covers the edits /// when the frame commits. pub fn stagedPageMut(self: *Writer, position_value: Position, index: usize) *[page.size]u8 { std.debug.assert(self.phase == .steady); std.debug.assert(position_value.len == self.len); const offset = self.stagedFrameOffset(position_value, index) catch unreachable; return self.storage[offset + frame_header_size ..][0..page.size]; } pub fn stagedPageId(self: *const Writer, position_value: Position, index: usize) u32 { std.debug.assert(self.phase == .steady); const offset = self.stagedFrameOffset(position_value, index) catch unreachable; return readU32(self.storage[offset..][0..4]); } pub fn swapStagedFrames(self: *Writer, position_value: Position, left_index: usize, right_index: usize) void { std.debug.assert(self.phase == .steady); std.debug.assert(position_value.len == self.len); if (left_index == right_index) return; const left_offset = self.stagedFrameOffset(position_value, left_index) catch unreachable; const right_offset = self.stagedFrameOffset(position_value, right_index) catch unreachable; const left: *[frame_size]u8 = self.storage[left_offset..][0..frame_size]; const right: *[frame_size]u8 = self.storage[right_offset..][0..frame_size]; const temporary = left.*; left.* = right.*; right.* = temporary; } pub fn commitStagedFrame(self: *Writer, position_value: Position, index: usize, db_page_count: u32) void { std.debug.assert(self.phase == .steady); const offset = self.stagedFrameOffset(position_value, index) catch unreachable; std.debug.assert(self.len == offset); const frame_header = self.storage[offset..][0..frame_header_size]; const image = self.storage[offset + frame_header_size ..][0..page.size]; std.debug.assert(readU32(frame_header[0..4]) != 0); writeU32(frame_header[4..8], db_page_count); writeU32(frame_header[8..12], self.salt.first); writeU32(frame_header[12..16], self.salt.second); var checksum = self.checksum; checksum.update(frame_header[0..8]); checksum.update(image); writeU32(frame_header[16..20], checksum.first); writeU32(frame_header[20..24], checksum.second); self.len += frame_size; self.checksum = checksum; } pub fn rewriteTail(self: *Writer, checkpoint_mark: usize, header: Header) void { std.debug.assert(self.phase == .steady); std.debug.assert(checkpoint_mark <= self.frameCount()); const source = self.bytes(); var reader = Reader.init(source) catch @panic("invalid wal writer storage"); self.writeHeader(header); while (reader.next() catch @panic("invalid wal writer storage")) |frame| { if (frame.index <= checkpoint_mark) continue; const image = frame.image[0..page.size].*; self.append(frame.page_id, frame.db_page_count, &image) catch unreachable; } if (reader.offset != source.len) @panic("invalid wal writer storage"); } fn writeHeader(self: *Writer, header: Header) void { const checksum = encodeHeader(self.storage[0..header_size], header); self.len = header_size; self.salt = header.salt; self.checksum = checksum; } fn stagedFrameOffset(self: *const Writer, position_value: Position, index: usize) Exhaustion!usize { std.debug.assert(position_value.len >= header_size); std.debug.assert((position_value.len - header_size) % frame_size == 0); const staged_bytes = std.math.mul(usize, index, frame_size) catch return error.WalFull; const offset = std.math.add(usize, position_value.len, staged_bytes) catch return error.WalFull; if (offset > self.storage.len or frame_size > self.storage.len - offset) return error.WalFull; return offset; }};comptime { alloc_phase.capacity.requireProvisionedRejectingOwnerShape(Writer);}pub const Reader = struct { data: []const u8, salt: Salt, checksum: Checksum, offset: usize = header_size, index: usize = 0, stopped: bool = false, pub fn init(data: []const u8) ReadError!Reader { return initControlled(data, .{}) catch |err| switch (err) { error.Interrupted => unreachable, else => return @errorCast(err), }; } pub fn initControlled(data: []const u8, control: Control) ControlledReadError!Reader { try control.check(); if (data.len < header_size) return error.InvalidWal; if (readU32(data[0..4]) != magic) return error.InvalidWal; if (readU32(data[4..8]) != format_version) return error.InvalidWal; if (readU32(data[8..12]) != page.size) return error.UnsupportedPageSize; var checksum: Checksum = .{}; checksum.update(data[0..24]); if (readU32(data[24..28]) != checksum.first) return error.InvalidChecksum; if (readU32(data[28..32]) != checksum.second) return error.InvalidChecksum; try control.check(); return .{ .data = data, .salt = .{ .first = readU32(data[16..20]), .second = readU32(data[20..24]), }, .checksum = checksum, }; } pub fn next(self: *Reader) ReadError!?Frame { return self.nextControlled(.{}) catch |err| switch (err) { error.Interrupted => unreachable, else => return @errorCast(err), }; } pub fn nextControlled(self: *Reader, control: Control) ControlledReadError!?Frame { const phase = trace.scope("wal.next"); defer phase.end(); try control.check(); if (self.stopped) return null; if (self.offset + frame_size > self.data.len) { self.stopped = true; return null; } const frame_header = self.data[self.offset..][0..frame_header_size]; const image = self.data[self.offset + frame_header_size ..][0..page.size]; if (readU32(frame_header[8..12]) != self.salt.first or readU32(frame_header[12..16]) != self.salt.second) { self.stopped = true; return null; } var checksum = self.checksum; checksum.update(frame_header[0..8]); checksum.update(image); if (readU32(frame_header[16..20]) != checksum.first or readU32(frame_header[20..24]) != checksum.second) { self.stopped = true; return null; } try control.check(); self.checksum = checksum; self.index += 1; self.offset += frame_size; return .{ .index = self.index, .page_id = readU32(frame_header[0..4]), .db_page_count = readU32(frame_header[4..8]), .image = image, }; }};pub fn endMark(data: []const u8) ReadError!usize { return endMarkControlled(data, .{}) catch |err| switch (err) { error.Interrupted => unreachable, else => return @errorCast(err), };}pub fn endMarkControlled(data: []const u8, control: Control) ControlledReadError!usize { const phase = trace.scope("wal.endmark"); defer phase.end(); var reader = try Reader.initControlled(data, control); var mark: usize = 0; const frames_max = (data.len - header_size) / frame_size; var frame_index: usize = 0; while (frame_index < frames_max) : (frame_index += 1) { const frame = (try reader.nextControlled(control)) orelse break; if (frame.committed()) mark = frame.index; } try control.check(); return mark;}pub fn pageAt(data: []const u8, page_id: u32, max_frame: usize) ReadError!?[]const u8 { const phase = trace.scope("wal.page_at"); defer phase.end(); if (max_frame == 0) return null; var reader = try Reader.init(data); var latest: ?[]const u8 = null; while (try reader.next()) |frame| { if (frame.index > max_frame) break; if (frame.page_id == page_id) latest = frame.image; } return latest;}const magic: u32 = 0x74737177;const format_version: u32 = 1;const Checksum = struct { first: u32 = 0, second: u32 = 0, fn fromPair(words: [2]u32) Checksum { return .{ .first = words[0], .second = words[1] }; } fn pair(self: Checksum) [2]u32 { return .{ self.first, self.second }; } fn update(self: *Checksum, data: []const u8) void { std.debug.assert(data.len % 8 == 0); var index: usize = 0; while (index < data.len) : (index += 8) { self.updatePair(data[index..][0..8]); } } fn updateControlled( self: *Checksum, data: []const u8, control: Control, ) error{Interrupted}!void { std.debug.assert(data.len % 8 == 0); var index: usize = 0; while (index < data.len) : (index += 8) { try control.check(); self.updatePair(data[index..][0..8]); } } fn updatePair(self: *Checksum, bytes: *const [8]u8) void { const left = readU32(bytes[0..4]); const right = readU32(bytes[4..8]); self.first +%= left +% self.second; self.second +%= right +% self.first; }};const DecodedHeader = struct { salt: Salt, checksum: Checksum,};fn decodeHeaderControlled( data: *const [header_size]u8, control: Control,) (ReadError || error{Interrupted})!DecodedHeader { try control.check(); if (readU32(data[0..4]) != magic) return error.InvalidWal; if (readU32(data[4..8]) != format_version) return error.InvalidWal; if (readU32(data[8..12]) != page.size) return error.UnsupportedPageSize; var checksum: Checksum = .{}; try checksum.updateControlled(data[0..24], control); if (readU32(data[24..28]) != checksum.first) return error.InvalidChecksum; if (readU32(data[28..32]) != checksum.second) return error.InvalidChecksum; return .{ .salt = .{ .first = readU32(data[16..20]), .second = readU32(data[20..24]), }, .checksum = checksum, };}fn disjoint(left: *const [frame_size]u8, right: *const [page.size]u8) bool { const left_start = @intFromPtr(left); const left_end = left_start + frame_size; const right_start = @intFromPtr(right); const right_end = right_start + page.size; return left_end <= right_start or right_end <= left_start;}fn readU32(bytes: []const u8) u32 { return std.mem.readInt(u32, bytes[0..4], .big);}fn writeU32(bytes: []u8, value: u32) void { std.mem.writeInt(u32, bytes[0..4], value, .big);}fn encodeHeader(bytes: *[header_size]u8, header: Header) Checksum { @memset(bytes, 0); writeU32(bytes[0..4], magic); writeU32(bytes[4..8], format_version); writeU32(bytes[8..12], page.size); writeU32(bytes[12..16], header.sequence); writeU32(bytes[16..20], header.salt.first); writeU32(bytes[20..24], header.salt.second); var checksum: Checksum = .{}; checksum.update(bytes[0..24]); writeU32(bytes[24..28], checksum.first); writeU32(bytes[28..32], checksum.second); return checksum;}fn testingHeader() Header { return .{ .sequence = 7, .salt = .{ .first = 0x1111_2222, .second = 0x3333_4444 }, };}fn fillPage(bytes: *[page.size]u8, page_id: u8, value: u8) void { @memset(bytes, 0); bytes[0] = page_id; bytes[1] = value;}fn testingWriter(allocator: Allocator, frames: usize) !Writer { const limits: Writer.Limits = .{ .header = testingHeader(), .frames = frames, }; const capacity = try Writer.Capacity.derive(limits); const storage = try allocator.alloc(u8, capacity.storage_bytes); errdefer allocator.free(storage); var writer = try Writer.init(storage, limits); writer.activate(); return writer;}fn testingScanner( scratch: *[frame_size]u8, refs: []FrameRef,) !Scanner { return try Scanner.init(.{ .scratch = scratch, .refs = refs });}fn admittedExtent(data: []const u8, capacity: usize) !Extent.Admitted { return switch (try Extent.classify(data.len, capacity, .{})) { .admitted => |extent| extent, .partial_header, .over_capacity => error.TestUnexpectedResult, };}fn resetTestingScanner( scanner: *Scanner, data: []const u8, database_page_floor: u32,) !void { const extent = try admittedExtent(data, scanner.workspace.refs.len); const header: *const [header_size]u8 = data[0..header_size]; try scanner.reset(header, extent, database_page_floor, .{});}fn consumeTestingFrame( scanner: *Scanner, data: []const u8, control: Control,) !Scanner.FrameOutcome { const request = (try scanner.nextFrame()) orelse return error.TestUnexpectedResult; const offset: usize = @intCast(request.offset); @memcpy(request.bytes, data[offset..][0..frame_size]); return try scanner.consumeFrame(control);}fn scanTestingData( scanner: *Scanner, data: []const u8, database_page_floor: u32,) !Scanner.Outcome { try resetTestingScanner(scanner, data, database_page_floor); var consumed: u32 = 0; while (scanner.state == .scanning) { std.debug.assert(consumed < scanner.capacity.frames); _ = try consumeTestingFrame(scanner, data, .{}); consumed += 1; } return try scanner.outcome();}fn loadSelectedFrame(scanner: *Scanner, data: []const u8, ref: FrameRef) !void { const request = try scanner.selectedFrame(ref, .{}); const offset: usize = @intCast(request.offset); @memcpy(request.bytes, data[offset..][0..frame_size]);}const TestingInterrupt = struct { remaining: usize, fn control(self: *TestingInterrupt) Control { return .{ .context = self, .interrupted_fn = interrupted }; } fn interrupted(context: ?*anyopaque) bool { const self: *TestingInterrupt = @ptrCast(@alignCast(context.?)); if (self.remaining == 0) return true; self.remaining -= 1; return false; }};fn zeroFrameRef() FrameRef { return .{ .page_id = 0, .frame = 0, .checksum_before = .{ 0, 0 }, .checksum_after = .{ 0, 0 }, };}fn refPagesEqual(refs: []const FrameRef, page_ids: []const u32) bool { if (refs.len != page_ids.len) return false; for (refs, page_ids) |ref, page_id| { if (ref.page_id != page_id) return false; } return true;}fn prepareFinalTestingFrame(scanner: *Scanner, data: []const u8) !void { try resetTestingScanner(scanner, data, 0); std.debug.assert(scanner.target_frames > 0); var consumed: u32 = 0; while (scanner.covered_frames + 1 < scanner.target_frames) { std.debug.assert(consumed < scanner.capacity.frames); _ = try consumeTestingFrame(scanner, data, .{}); consumed += 1; } const request = (try scanner.nextFrame()) orelse return error.TestUnexpectedResult; const offset: usize = @intCast(request.offset); @memcpy(request.bytes, data[offset..][0..frame_size]);}fn expectScannerNotReady(scanner: *Scanner) !void { const ref = FrameRef{ .page_id = 1, .frame = 1, .checksum_before = .{ 0, 0 }, .checksum_after = .{ 0, 0 }, }; var destination: [page.size]u8 = @splat(0xa5); try std.testing.expectError(error.ScannerNotReady, scanner.outcome()); try std.testing.expectError(error.ScannerNotReady, scanner.findCommitted(0, .{})); try std.testing.expectError(error.ScannerNotReady, scanner.selectedFrame(ref, .{})); try std.testing.expectError( error.ScannerNotReady, scanner.verifySelected(ref, &destination, .{}), ); try std.testing.expectError(error.ScannerNotReady, scanner.nextFrame()); try std.testing.expectError(error.ScannerNotReady, scanner.consumeFrame(.{})); try std.testing.expectError( error.ScannerNotReady, scanner.extend(.{ .complete_frames = 0, .partial_tail_bytes = 0 }, .{}), ); const expected: [page.size]u8 = @splat(0xa5); try std.testing.expectEqualSlices(u8, &expected, &destination);}fn modelWriterCapacity(frames: usize) ?Writer.Capacity { const maximum_frames = (std.math.maxInt(usize) - header_size) / frame_size; if (frames > maximum_frames) return null; return .{ .frames = frames, .storage_bytes = header_size + frames * frame_size, };}fn modelScannerCapacity(frames: usize) ?Scanner.Capacity { if (frames > std.math.maxInt(u32)) return null; const refs_bytes = std.math.mul( usize, frames, @sizeOf(FrameRef), ) catch return null; const storage_bytes = std.math.add( usize, frame_size, refs_bytes, ) catch return null; return .{ .frames = @intCast(frames), .storage_bytes = storage_bytes, };}test "wal scanner capacity matches an independent typed byte model" { try std.testing.expectEqual(@as(usize, 24), @sizeOf(FrameRef)); for (0..4_097) |frames| { try std.testing.expectEqual( modelScannerCapacity(frames).?, try Scanner.Capacity.derive(frames), ); } const capacity = modelScannerCapacity(1_024).?; try std.testing.expectEqual(@as(u32, 1_024), capacity.frames); try std.testing.expectEqual(@as(usize, 28_696), capacity.storage_bytes); try std.testing.expectEqual(capacity, try Scanner.Capacity.derive(1_024)); const maximum_frames: usize = std.math.maxInt(u32); if (modelScannerCapacity(maximum_frames)) |expected| { try std.testing.expectEqual(expected, try Scanner.Capacity.derive(maximum_frames)); } else { try std.testing.expectError( error.CapacityOverflow, Scanner.Capacity.derive(maximum_frames), ); } if (std.math.add(usize, maximum_frames, 1)) |over_u32| { try std.testing.expectError( error.CapacityOverflow, Scanner.Capacity.derive(over_u32), ); } else |_| {} try std.testing.expectError( error.CapacityOverflow, Scanner.storageBytes(std.math.maxInt(usize), frame_size, @sizeOf(FrameRef)), ); try std.testing.expectError( error.CapacityOverflow, Scanner.storageBytes(1, std.math.maxInt(usize), 1), ); var scratch: [frame_size]u8 = undefined; var refs: [1_024]FrameRef = undefined; var scanner = try testingScanner(&scratch, &refs); try std.testing.expectEqual(&scratch, scanner.workspace.scratch); try std.testing.expectEqual(refs[0..].ptr, scanner.workspace.refs.ptr); const returned = scanner.deinit(); try std.testing.expectEqual(&scratch, returned.scratch); try std.testing.expectEqual(refs[0..].ptr, returned.refs.ptr);}test "wal scanner extent distinguishes every incomplete boundary" { var header_bytes: u64 = 0; while (header_bytes < header_size) : (header_bytes += 1) { const extent = try Extent.classify(header_bytes, 2, .{}); try std.testing.expectEqual(@as(u8, @intCast(header_bytes)), extent.partial_header); } var tail_bytes: u64 = 0; while (tail_bytes < frame_size) : (tail_bytes += 1) { const file_bytes = header_size + frame_size + tail_bytes; const extent = (try Extent.classify(file_bytes, 2, .{})).admitted; try std.testing.expectEqual(@as(u32, 1), extent.complete_frames); try std.testing.expectEqual(@as(u16, @intCast(tail_bytes)), extent.partial_tail_bytes); }}test "wal scanner rejects capacity plus one before frame payload" { const exact_bytes = header_size + 2 * frame_size; const exact = (try Extent.classify(exact_bytes, 2, .{})).admitted; try std.testing.expectEqual(@as(u32, 2), exact.complete_frames); try std.testing.expectEqual(@as(u16, 0), exact.partial_tail_bytes); const partial = (try Extent.classify(exact_bytes + 1, 2, .{})).over_capacity; try std.testing.expectEqual(@as(u64, 2), partial.complete_frames); try std.testing.expectEqual(@as(u16, 1), partial.partial_tail_bytes); const full = (try Extent.classify(exact_bytes + frame_size, 2, .{})).over_capacity; try std.testing.expectEqual(@as(u64, 3), full.complete_frames); try std.testing.expectEqual(@as(u16, 0), full.partial_tail_bytes); const maximum_admitted_bytes = header_size + 1_024 * frame_size; const maximum = (try Extent.classify(maximum_admitted_bytes, 1_024, .{})).admitted; try std.testing.expectEqual(@as(u32, 1_024), maximum.complete_frames); _ = (try Extent.classify(maximum_admitted_bytes + 1, 1_024, .{})).over_capacity;}test "wal scanner reports a partial uncommitted tail without consuming it" { var writer = try testingWriter(std.testing.allocator, 0); defer std.testing.allocator.free(writer.deinit()); var scratch: [frame_size]u8 = undefined; var refs: [1]FrameRef = undefined; var scanner = try testingScanner(&scratch, &refs); const header: *const [header_size]u8 = writer.bytes()[0..header_size]; try scanner.reset(header, .{ .complete_frames = 0, .partial_tail_bytes = 73, }, 0, .{}); try std.testing.expectEqual( @as(u16, 73), (try scanner.outcome()).partial_uncommitted_tail, ); const progress = scanner.progress(); try std.testing.expectEqual(@as(u32, 0), progress.covered_frames); try std.testing.expectEqualSlices( u8, writer.bytes()[0..header_size], &progress.header, );}test "wal scanner rejects every forged noncanonical partial tail" { var writer = try testingWriter(std.testing.allocator, 0); defer std.testing.allocator.free(writer.deinit()); var scratch: [frame_size]u8 = undefined; var refs: [1]FrameRef = undefined; var scanner = try testingScanner(&scratch, &refs); const header: *const [header_size]u8 = writer.bytes()[0..header_size]; try scanner.reset(header, .{ .complete_frames = 0, .partial_tail_bytes = frame_size - 1, }, 0, .{}); try std.testing.expectEqual( @as(u16, frame_size - 1), (try scanner.outcome()).partial_uncommitted_tail, ); const before = scanner.progress(); try std.testing.expectError(error.InvalidExtent, scanner.reset( header, .{ .complete_frames = 0, .partial_tail_bytes = frame_size }, 0, .{}, )); try std.testing.expectEqualDeep(before, scanner.progress()); try std.testing.expectError(error.InvalidExtent, scanner.extend(.{ .complete_frames = 0, .partial_tail_bytes = std.math.maxInt(u16), }, .{})); try std.testing.expectEqualDeep(before, scanner.progress());}test "wal scanner makes rewind and capacity invalidation unreadable" { var writer = try testingWriter(std.testing.allocator, 0); defer std.testing.allocator.free(writer.deinit()); var scratch: [frame_size]u8 = undefined; var refs: [1]FrameRef = undefined; var scanner = try testingScanner(&scratch, &refs); const header: *const [header_size]u8 = writer.bytes()[0..header_size]; try scanner.reset(header, .{ .complete_frames = 0, .partial_tail_bytes = 10, }, 0, .{}); try std.testing.expectError(error.RewindRequired, scanner.extend(.{ .complete_frames = 0, .partial_tail_bytes = 9, }, .{})); try expectScannerNotReady(&scanner); try scanner.reset(header, .{ .complete_frames = 0, .partial_tail_bytes = 0, }, 0, .{}); try std.testing.expectError(error.CapacityExceeded, scanner.extend(.{ .complete_frames = 2, .partial_tail_bytes = 0, }, .{})); try expectScannerNotReady(&scanner); try scanner.reset(header, .{ .complete_frames = 0, .partial_tail_bytes = 0, }, 0, .{}); try std.testing.expectEqual( std.meta.Tag(Scanner.Outcome).complete, std.meta.activeTag(try scanner.outcome()), );}test "wal scanner reset interruption sweeps every header validation cut" { var writer = try testingWriter(std.testing.allocator, 0); defer std.testing.allocator.free(writer.deinit()); const header: *const [header_size]u8 = writer.bytes()[0..header_size]; var completed = false; for (0..16) |budget| { var scratch: [frame_size]u8 = undefined; var refs: [1]FrameRef = @splat(zeroFrameRef()); var scanner = try testingScanner(&scratch, &refs); const before = scanner.progress(); var interrupt = TestingInterrupt{ .remaining = budget }; if (scanner.reset( header, .{ .complete_frames = 0, .partial_tail_bytes = 0 }, 0, interrupt.control(), )) |_| { try std.testing.expectEqual(Scanner.State.ready, scanner.state); completed = true; break; } else |err| switch (err) { error.Interrupted => { try std.testing.expectEqualDeep(before, scanner.progress()); try std.testing.expectEqualDeep(@as([1]FrameRef, @splat(zeroFrameRef())), refs); }, else => return err, } } try std.testing.expect(completed);}test "wal scanner rejects a flip of every complete header byte" { var writer = try testingWriter(std.testing.allocator, 0); defer std.testing.allocator.free(writer.deinit()); var scratch: [frame_size]u8 = undefined; var refs: [1]FrameRef = undefined; var scanner = try testingScanner(&scratch, &refs); const header: *[header_size]u8 = writer.storage[0..header_size]; for (0..header_size) |index| { header[index] ^= 0xff; if (scanner.reset( header, .{ .complete_frames = 0, .partial_tail_bytes = 0 }, 0, .{}, )) |_| { return error.TestUnexpectedResult; } else |err| switch (err) { error.InvalidChecksum, error.InvalidWal, error.UnsupportedPageSize, => {}, error.CapacityExceeded, error.Interrupted, error.InvalidExtent, error.ScannerNotReady, => return err, } try std.testing.expectEqual(Scanner.State.invalidated, scanner.state); header[index] ^= 0xff; }}test "wal scanner admits a completed prior partial frame" { var writer = try testingWriter(std.testing.allocator, 1); defer std.testing.allocator.free(writer.deinit()); var image: [page.size]u8 = undefined; fillPage(&image, 1, 10); try writer.append(1, 1, &image); var scratch: [frame_size]u8 = undefined; var refs: [1]FrameRef = undefined; var scanner = try testingScanner(&scratch, &refs); const partial = writer.bytes()[0 .. header_size + 101]; try resetTestingScanner(&scanner, partial, 0); try std.testing.expectEqual( @as(u16, 101), (try scanner.outcome()).partial_uncommitted_tail, ); try scanner.extend(try admittedExtent(writer.bytes(), refs.len), .{}); try std.testing.expectEqual( std.meta.Tag(Scanner.FrameOutcome).committed, std.meta.activeTag(try consumeTestingFrame(&scanner, writer.bytes(), .{})), ); try std.testing.expectEqual(Scanner.State.ready, scanner.state);}test "wal scanner terminal-invalidates page zero and malformed full frames" { var writer = try testingWriter(std.testing.allocator, 1); defer std.testing.allocator.free(writer.deinit()); var image: [page.size]u8 = undefined; fillPage(&image, 1, 10); try writer.append(0, 1, &image); var scratch: [frame_size]u8 = undefined; var refs: [1]FrameRef = undefined; var scanner = try testingScanner(&scratch, &refs); try resetTestingScanner(&scanner, writer.bytes(), 0); const page_zero = try consumeTestingFrame(&scanner, writer.bytes(), .{}); try std.testing.expectEqual( Scanner.MalformedFrame.invalid_page_id, page_zero.malformed_full_frame, ); try std.testing.expectEqual(Scanner.State.invalidated, scanner.state); try expectScannerNotReady(&scanner); writer.restore(.{ .len = header_size, .checksum_first = readU32(writer.bytes()[24..28]), .checksum_second = readU32(writer.bytes()[28..32]), }); try writer.append(1, 1, &image); _ = try scanTestingData(&scanner, writer.bytes(), 0); try std.testing.expectEqual(Scanner.State.ready, scanner.state); writer.storage[header_size + 8] ^= 0xff; try resetTestingScanner(&scanner, writer.bytes(), 0); const salt = try consumeTestingFrame(&scanner, writer.bytes(), .{}); try std.testing.expectEqual( Scanner.MalformedFrame.salt_mismatch, salt.malformed_full_frame, ); try expectScannerNotReady(&scanner); writer.storage[header_size + 8] ^= 0xff; _ = try scanTestingData(&scanner, writer.bytes(), 0); writer.storage[header_size + frame_header_size + 9] ^= 0xff; try resetTestingScanner(&scanner, writer.bytes(), 0); const checksum = try consumeTestingFrame(&scanner, writer.bytes(), .{}); try std.testing.expectEqual( Scanner.MalformedFrame.checksum_mismatch, checksum.malformed_full_frame, ); try std.testing.expectEqual(Scanner.State.invalidated, scanner.state); try expectScannerNotReady(&scanner); writer.storage[header_size + frame_header_size + 9] ^= 0xff; _ = try scanTestingData(&scanner, writer.bytes(), 0); try std.testing.expectEqual(Scanner.State.ready, scanner.state);}fn scanSingleMarker(marker: u32, database_page_floor: u32) !Scanner.FrameOutcome { var writer = try testingWriter(std.testing.allocator, 1); defer std.testing.allocator.free(writer.deinit()); var image: [page.size]u8 = undefined; fillPage(&image, 1, 10); try writer.append(1, marker, &image); var scratch: [frame_size]u8 = undefined; var refs: [1]FrameRef = undefined; var scanner = try testingScanner(&scratch, &refs); try resetTestingScanner(&scanner, writer.bytes(), database_page_floor); return try consumeTestingFrame(&scanner, writer.bytes(), .{});}test "wal scanner enforces the monotonic base page count floor" { const lower = try scanSingleMarker(9, 10); try std.testing.expectEqual( Scanner.MalformedFrame.database_page_count_regressed, lower.malformed_full_frame, ); try std.testing.expectEqual( std.meta.Tag(Scanner.FrameOutcome).committed, std.meta.activeTag(try scanSingleMarker(10, 10)), ); try std.testing.expectEqual( std.meta.Tag(Scanner.FrameOutcome).committed, std.meta.activeTag(try scanSingleMarker(12, 10)), );}test "wal scanner rejects a marker below any staged page" { var writer = try testingWriter(std.testing.allocator, 2); defer std.testing.allocator.free(writer.deinit()); var image: [page.size]u8 = undefined; fillPage(&image, 11, 10); try writer.append(11, 0, &image); fillPage(&image, 1, 20); try writer.append(1, 10, &image); var scratch: [frame_size]u8 = undefined; var refs: [2]FrameRef = undefined; var scanner = try testingScanner(&scratch, &refs); try resetTestingScanner(&scanner, writer.bytes(), 0); try std.testing.expectEqual( std.meta.Tag(Scanner.FrameOutcome).staged, std.meta.activeTag(try consumeTestingFrame(&scanner, writer.bytes(), .{})), ); const marker = try consumeTestingFrame(&scanner, writer.bytes(), .{}); try std.testing.expectEqual( Scanner.MalformedFrame.page_outside_database, marker.malformed_full_frame, ); try std.testing.expectEqual(Scanner.State.invalidated, scanner.state); try expectScannerNotReady(&scanner);}test "wal scanner keeps an uncommitted high page outside authority" { var writer = try testingWriter(std.testing.allocator, 1); defer std.testing.allocator.free(writer.deinit()); var image: [page.size]u8 = undefined; fillPage(&image, 100, 10); try writer.append(100, 0, &image); var scratch: [frame_size]u8 = undefined; var refs: [1]FrameRef = undefined; var scanner = try testingScanner(&scratch, &refs); _ = try scanTestingData(&scanner, writer.bytes(), 10); try std.testing.expect((try scanner.findCommitted(100, .{})) == null); try std.testing.expectEqual(@as(u32, 10), scanner.logical_database_page_count); try std.testing.expectEqual(@as(u32, 0), scanner.committed_ref_count); try std.testing.expectEqual(@as(u32, 1), scanner.staged_ref_count);}test "wal scanner folds duplicate pages with the greatest frame winning" { var writer = try testingWriter(std.testing.allocator, 5); defer std.testing.allocator.free(writer.deinit()); var image: [page.size]u8 = undefined; fillPage(&image, 2, 20); try writer.append(2, 0, &image); fillPage(&image, 1, 10); try writer.append(1, 2, &image); fillPage(&image, 1, 11); try writer.append(1, 0, &image); fillPage(&image, 1, 12); try writer.append(1, 0, &image); fillPage(&image, 2, 22); try writer.append(2, 2, &image); var scratch: [frame_size]u8 = undefined; var refs: [5]FrameRef = undefined; var scanner = try testingScanner(&scratch, &refs); _ = try scanTestingData(&scanner, writer.bytes(), 0); const first = (try scanner.findCommitted(1, .{})).?; const second = (try scanner.findCommitted(2, .{})).?; try std.testing.expectEqual(@as(u32, 4), first.frame); try std.testing.expectEqual(@as(u32, 5), second.frame); try std.testing.expectEqual(@as(u32, 2), scanner.committed_ref_count); try std.testing.expectEqual(@as(u32, 5), scanner.committed_end_mark); var destination: [page.size]u8 = undefined; try loadSelectedFrame(&scanner, writer.bytes(), first); try scanner.verifySelected(first, &destination, .{}); try std.testing.expectEqual(@as(u8, 12), destination[1]); try loadSelectedFrame(&scanner, writer.bytes(), second); try scanner.verifySelected(second, &destination, .{}); try std.testing.expectEqual(@as(u8, 22), destination[1]);}test "wal scanner resumes from the physical checksum tail" { var writer = try testingWriter(std.testing.allocator, 4); defer std.testing.allocator.free(writer.deinit()); var image: [page.size]u8 = undefined; fillPage(&image, 1, 10); try writer.append(1, 1, &image); fillPage(&image, 1, 15); try writer.append(1, 1, &image); fillPage(&image, 1, 20); try writer.append(1, 0, &image); var scratch: [frame_size]u8 = undefined; var refs: [4]FrameRef = undefined; var scanner = try testingScanner(&scratch, &refs); _ = try scanTestingData(&scanner, writer.bytes(), 0); try std.testing.expectEqual( @as(u32, 2), (try scanner.findCommitted(1, .{})).?.frame, ); fillPage(&image, 1, 30); try writer.append(1, 1, &image); const extent = try admittedExtent(writer.bytes(), refs.len); try scanner.extend(extent, .{}); _ = try consumeTestingFrame(&scanner, writer.bytes(), .{}); const latest = (try scanner.findCommitted(1, .{})).?; try std.testing.expectEqual(@as(u32, 4), latest.frame); try std.testing.expectEqual(@as(u32, 4), scanner.covered_frames); try std.testing.expectEqual(@as(u32, 4), scanner.committed_end_mark); var destination: [page.size]u8 = undefined; try loadSelectedFrame(&scanner, writer.bytes(), latest); try scanner.verifySelected(latest, &destination, .{}); try std.testing.expectEqual(@as(u8, 30), destination[1]);}const InsertSweep = struct { completed: bool = false, pre_mutation: bool = false, shifted_one: bool = false, shifted_two: bool = false, shifted_three: bool = false, inserted: bool = false,};fn exerciseInsertCut(data: []const u8, budget: usize) !InsertSweep { var scratch: [frame_size]u8 = undefined; var refs: [4]FrameRef = @splat(zeroFrameRef()); var scanner = try testingScanner(&scratch, &refs); try prepareFinalTestingFrame(&scanner, data); const before = scanner.progress(); const refs_before = refs; var interrupt = TestingInterrupt{ .remaining = budget }; if (scanner.consumeFrame(interrupt.control())) |outcome| { try std.testing.expectEqual( std.meta.Tag(Scanner.FrameOutcome).committed, std.meta.activeTag(outcome), ); try std.testing.expect(refPagesEqual(&refs, &.{ 1, 2, 4, 6 })); return .{ .completed = true }; } else |err| switch (err) { error.Interrupted => {}, else => return err, } if (scanner.state == .scanning) { try std.testing.expectEqualDeep(before, scanner.progress()); try std.testing.expectEqualDeep(refs_before, refs); _ = try scanner.consumeFrame(.{}); try std.testing.expectEqual(Scanner.State.ready, scanner.state); return .{ .pre_mutation = true }; } try std.testing.expectEqual(Scanner.State.invalidated, scanner.state); const result = InsertSweep{ .shifted_one = refPagesEqual(&refs, &.{ 2, 4, 6, 6 }), .shifted_two = refPagesEqual(&refs, &.{ 2, 4, 4, 6 }), .shifted_three = refPagesEqual(&refs, &.{ 2, 2, 4, 6 }), .inserted = refPagesEqual(&refs, &.{ 1, 2, 4, 6 }), }; try expectScannerNotReady(&scanner); _ = try scanTestingData(&scanner, data, 0); try std.testing.expect((try scanner.findCommitted(1, .{})) != null); return result;}test "wal scanner interruption sweeps every insert fold cut" { var writer = try testingWriter(std.testing.allocator, 4); defer std.testing.allocator.free(writer.deinit()); var image: [page.size]u8 = undefined; for ([_]u8{ 2, 4, 6 }) |page_id| { fillPage(&image, page_id, page_id); try writer.append(page_id, if (page_id == 6) 6 else 0, &image); } fillPage(&image, 1, 1); try writer.append(1, 6, &image); var observed = InsertSweep{}; for (0..frame_size / 8 + 128) |budget| { const cut = try exerciseInsertCut(writer.bytes(), budget); observed.completed = observed.completed or cut.completed; observed.pre_mutation = observed.pre_mutation or cut.pre_mutation; observed.shifted_one = observed.shifted_one or cut.shifted_one; observed.shifted_two = observed.shifted_two or cut.shifted_two; observed.shifted_three = observed.shifted_three or cut.shifted_three; observed.inserted = observed.inserted or cut.inserted; if (cut.completed) break; } try std.testing.expect(observed.completed); try std.testing.expect(observed.pre_mutation); try std.testing.expect(observed.shifted_one); try std.testing.expect(observed.shifted_two); try std.testing.expect(observed.shifted_three); try std.testing.expect(observed.inserted);}const ReplaceSweep = struct { completed: bool = false, pre_mutation: bool = false, first_replacement: bool = false, tail_shift: bool = false, duplicate_replacement: bool = false, later_candidate: bool = false,};fn exerciseReplaceCut(data: []const u8, budget: usize) !ReplaceSweep { var scratch: [frame_size]u8 = undefined; var refs: [5]FrameRef = @splat(zeroFrameRef()); var scanner = try testingScanner(&scratch, &refs); try prepareFinalTestingFrame(&scanner, data); const before = scanner.progress(); const refs_before = refs; var interrupt = TestingInterrupt{ .remaining = budget }; if (scanner.consumeFrame(interrupt.control())) |outcome| { try std.testing.expectEqual( std.meta.Tag(Scanner.FrameOutcome).committed, std.meta.activeTag(outcome), ); try std.testing.expectEqual(@as(u32, 4), refs[0].frame); try std.testing.expectEqual(@as(u32, 5), refs[1].frame); return .{ .completed = true }; } else |err| switch (err) { error.Interrupted => {}, else => return err, } if (scanner.state == .scanning) { try std.testing.expectEqualDeep(before, scanner.progress()); try std.testing.expectEqualDeep(refs_before, refs); _ = try scanner.consumeFrame(.{}); try std.testing.expectEqual(Scanner.State.ready, scanner.state); return .{ .pre_mutation = true }; } try std.testing.expectEqual(Scanner.State.invalidated, scanner.state); const result = ReplaceSweep{ .first_replacement = refs[0].frame == 3, .tail_shift = refs[2].frame == 4 and refs[3].frame == 5, .duplicate_replacement = refs[0].frame == 4, .later_candidate = refs[1].frame == 5, }; try expectScannerNotReady(&scanner); _ = try scanTestingData(&scanner, data, 0); try std.testing.expectEqual(@as(u32, 4), (try scanner.findCommitted(1, .{})).?.frame); try std.testing.expectEqual(@as(u32, 5), (try scanner.findCommitted(2, .{})).?.frame); return result;}test "wal scanner interruption sweeps replacement and duplicate fold cuts" { var writer = try testingWriter(std.testing.allocator, 5); defer std.testing.allocator.free(writer.deinit()); var image: [page.size]u8 = undefined; const pages = [_]u8{ 1, 2, 1, 1, 2 }; for (pages, 0..) |page_id, index| { fillPage(&image, page_id, @intCast(index + 1)); const marker: u32 = if (index == 1 or index == 4) 2 else 0; try writer.append(page_id, marker, &image); } var observed = ReplaceSweep{}; for (0..frame_size / 8 + 128) |budget| { const cut = try exerciseReplaceCut(writer.bytes(), budget); observed.completed = observed.completed or cut.completed; observed.pre_mutation = observed.pre_mutation or cut.pre_mutation; observed.first_replacement = observed.first_replacement or cut.first_replacement; observed.tail_shift = observed.tail_shift or cut.tail_shift; observed.duplicate_replacement = observed.duplicate_replacement or cut.duplicate_replacement; observed.later_candidate = observed.later_candidate or cut.later_candidate; if (cut.completed) break; } try std.testing.expect(observed.completed); try std.testing.expect(observed.pre_mutation); try std.testing.expect(observed.first_replacement); try std.testing.expect(observed.tail_shift); try std.testing.expect(observed.duplicate_replacement); try std.testing.expect(observed.later_candidate);}fn expectSelectedMutationRejected( scanner: *Scanner, data: []const u8, ref: FrameRef, mutation_index: usize,) !void { try loadSelectedFrame(scanner, data, ref); scanner.workspace.scratch[mutation_index] ^= 0xff; var destination: [page.size]u8 = @splat(0xa5); const before = scanner.progress(); try std.testing.expectError( error.SelectedFrameMismatch, scanner.verifySelected(ref, &destination, .{}), ); try std.testing.expectEqualDeep(before, scanner.progress()); const expected: [page.size]u8 = @splat(0xa5); try std.testing.expectEqualSlices(u8, &expected, &destination);}fn exerciseVerifyCut( scanner: *Scanner, data: []const u8, ref: FrameRef, expected: *const [page.size]u8, budget: usize,) !bool { try loadSelectedFrame(scanner, data, ref); var destination: [page.size]u8 = @splat(0xa5); const before = scanner.progress(); const ref_before = (try scanner.findCommitted(ref.page_id, .{})).?; var interrupt = TestingInterrupt{ .remaining = budget }; if (scanner.verifySelected(ref, &destination, interrupt.control())) |_| { try std.testing.expectEqualSlices(u8, expected, &destination); return true; } else |err| switch (err) { error.Interrupted => {}, else => return err, } try std.testing.expectEqualDeep(before, scanner.progress()); try std.testing.expectEqualDeep( ref_before, (try scanner.findCommitted(ref.page_id, .{})).?, ); const sentinel: [page.size]u8 = @splat(0xa5); try std.testing.expectEqualSlices(u8, &sentinel, &destination); return false;}test "wal scanner selected verification changes destination only after proof" { var writer = try testingWriter(std.testing.allocator, 1); defer std.testing.allocator.free(writer.deinit()); var image: [page.size]u8 = undefined; fillPage(&image, 1, 10); try writer.append(1, 1, &image); var scratch: [frame_size]u8 = undefined; var refs: [1]FrameRef = undefined; var scanner = try testingScanner(&scratch, &refs); _ = try scanTestingData(&scanner, writer.bytes(), 0); const ref = (try scanner.findCommitted(1, .{})).?; try expectSelectedMutationRejected(&scanner, writer.bytes(), ref, 0); try expectSelectedMutationRejected(&scanner, writer.bytes(), ref, 8); try expectSelectedMutationRejected(&scanner, writer.bytes(), ref, 16); try expectSelectedMutationRejected( &scanner, writer.bytes(), ref, frame_header_size + 9, ); var completed = false; for (0..frame_size / 8 + 128) |budget| { if (try exerciseVerifyCut(&scanner, writer.bytes(), ref, &image, budget)) { completed = true; break; } } try std.testing.expect(completed);}test "wal scanner rolling checksum does not claim colliding byte identity" { var writer = try testingWriter(std.testing.allocator, 1); defer std.testing.allocator.free(writer.deinit()); var image: [page.size]u8 = undefined; fillPage(&image, 1, 10); try writer.append(1, 1, &image); var scratch: [frame_size]u8 = undefined; var refs: [1]FrameRef = undefined; var scanner = try testingScanner(&scratch, &refs); _ = try scanTestingData(&scanner, writer.bytes(), 0); const ref = (try scanner.findCommitted(1, .{})).?; try loadSelectedFrame(&scanner, writer.bytes(), ref); const payload = scanner.workspace.scratch[frame_header_size..]; writeU32(payload[0..4], readU32(payload[0..4]) +% 1); writeU32(payload[8..12], readU32(payload[8..12]) -% 2); writeU32(payload[12..16], readU32(payload[12..16]) -% 1); var destination: [page.size]u8 = undefined; const before = scanner.progress(); try scanner.verifySelected(ref, &destination, .{}); try std.testing.expectEqualDeep(before, scanner.progress()); try std.testing.expectEqualSlices(u8, payload, &destination); try std.testing.expect(!std.mem.eql(u8, &image, &destination));}test "wal writer capacity matches an independent additive model" { comptime { @stardustClaim( @import("alloc_phase").capacity.witness(Writer, "sql_wal_capacity"), null, null, null, null, null, null, ); } for (0..4097) |frames| { try std.testing.expectEqual( modelWriterCapacity(frames).?, try Writer.Capacity.derive(.{ .header = testingHeader(), .frames = frames }), ); } const maximum_frames = (std.math.maxInt(usize) - header_size) / frame_size; try std.testing.expectEqual( modelWriterCapacity(maximum_frames).?, try Writer.Capacity.derive(.{ .header = testingHeader(), .frames = maximum_frames }), ); const overflow_frames = maximum_frames + 1; try std.testing.expect(modelWriterCapacity(overflow_frames) == null); try std.testing.expectError( error.CapacityOverflow, Writer.Capacity.derive(.{ .header = testingHeader(), .frames = overflow_frames }), ); var empty_storage: [0]u8 = .{}; try std.testing.expectError( error.CapacityOverflow, Writer.init(&empty_storage, .{ .header = testingHeader(), .frames = overflow_frames }), );}test "wal writer rejects short storage and releases its exact borrow" { comptime { @stardustClaim( @import("alloc_phase").capacity.witness(Writer, "sql_wal_storage_rejection"), null, null, null, null, null, null, ); @stardustClaim( @import("alloc_phase").capacity.witness(Writer, "sql_wal_work_bound"), null, null, null, null, null, null, ); } const limits: Writer.Limits = .{ .header = testingHeader(), .frames = 3 }; const capacity = try Writer.Capacity.derive(limits); const storage = try std.testing.allocator.alloc(u8, capacity.storage_bytes); defer std.testing.allocator.free(storage); try std.testing.expectError(error.StorageTooShort, Writer.init(storage[0 .. storage.len - 1], limits)); var writer = try Writer.init(storage, limits); try std.testing.expectEqual(alloc_phase.capacity.Phase.initialization, writer.phase); writer.activate(); try std.testing.expectEqual(alloc_phase.capacity.Phase.steady, writer.phase); const released = writer.deinit(); try std.testing.expectEqual(storage.ptr, released.ptr); try std.testing.expectEqual(storage.len, released.len);}test "wal writer is sealed before maximum append recovery and rewrite" { comptime { @stardustClaim( @import("alloc_phase").capacity.witness(Writer, "sql_wal_sealed"), null, null, null, null, null, null, ); } var source = try testingWriter(std.testing.allocator, 4); defer std.testing.allocator.free(source.deinit()); var image: [page.size]u8 = undefined; for (0..4) |index| { fillPage(&image, @intCast(index + 1), @intCast(10 + index)); try source.append(@intCast(index + 1), if (index == 0) 4 else 0, &image); } var phase_allocator = try alloc_phase.SealedPhaseAllocator.init(std.testing.allocator); const limits: Writer.Limits = .{ .header = testingHeader(), .frames = 3 }; const storage_capacity = try Writer.Capacity.derive(limits); const storage = phase_allocator.initializationAllocator().alloc(u8, storage_capacity.storage_bytes) catch |err| { phase_allocator.abortInitialization(); phase_allocator.deinit(); return err; }; var writer = Writer.init(storage, limits) catch |err| { phase_allocator.abortInitialization(); phase_allocator.deinit(); return err; }; defer { if (phase_allocator.phase() == .initialization) phase_allocator.abortInitialization(); if (phase_allocator.phase() == .steady) phase_allocator.beginTeardown(); if (writer.phase != .teardown) phase_allocator.teardownAllocator().free(writer.deinit()); phase_allocator.deinit(); } const storage_pointer = writer.storage.ptr; const capacity = writer.capacity; phase_allocator.seal(); writer.activate(); const committed_len = header_size + 3 * frame_size; try writer.load(source.bytes(), committed_len); try std.testing.expectEqual(@as(usize, 3), writer.frameCount()); writer.rewriteTail(1, .{ .sequence = 8, .salt = .{ .first = 0x5555_6666, .second = 0x7777_8888 }, }); try std.testing.expectEqual(@as(usize, 2), writer.frameCount()); fillPage(&image, 9, 99); const staged_position = writer.position(); try writer.stagePage(staged_position, 0, 9, &image); writer.commitStagedFrame(staged_position, 0, 9); try std.testing.expectEqual(@as(usize, 3), writer.frameCount()); const position_before_overload = writer.position(); var bytes_before_overload: [header_size + 3 * frame_size]u8 = undefined; @memcpy(bytes_before_overload[0..], writer.bytes()); try std.testing.expectError(error.WalFull, writer.append(10, 10, &image)); try std.testing.expectError(error.WalFull, writer.load(source.bytes(), source.bytes().len)); try std.testing.expectEqual(position_before_overload, writer.position()); try std.testing.expectEqual(storage_pointer, writer.storage.ptr); try std.testing.expectEqual(capacity, writer.capacity); try std.testing.expectEqualSlices(u8, bytes_before_overload[0..], writer.bytes()); try std.testing.expectEqual(@as(u64, 0), phase_allocator.violations().total());}test "wal header validates and has no committed frames" { var writer = try testingWriter(std.testing.allocator, 0); defer std.testing.allocator.free(writer.deinit()); var reader = try Reader.init(writer.bytes()); try std.testing.expect(try reader.next() == null); try std.testing.expectEqual(@as(usize, 0), try endMark(writer.bytes()));}test "wal commit marker makes preceding frames visible" { comptime { @stardustClaim( @import("alloc_phase").capacity.witness(Writer, "sql_wal_semantics"), null, null, null, null, null, null, ); } var writer = try testingWriter(std.testing.allocator, 2); defer std.testing.allocator.free(writer.deinit()); var first: [page.size]u8 = undefined; var second: [page.size]u8 = undefined; fillPage(&first, 1, 10); fillPage(&second, 2, 20); try writer.append(1, 0, &first); try writer.append(2, 2, &second); const mark = try endMark(writer.bytes()); try std.testing.expectEqual(@as(usize, 2), mark); const recovered_first = (try pageAt(writer.bytes(), 1, mark)).?; const recovered_second = (try pageAt(writer.bytes(), 2, mark)).?; try std.testing.expectEqual(@as(u8, 10), recovered_first[1]); try std.testing.expectEqual(@as(u8, 20), recovered_second[1]);}test "wal staged frames transfer in place with byte exact output" { comptime { @stardustClaim( @import("alloc_phase").capacity.witness(Writer, "sql_wal_staged_transfer"), null, null, null, null, null, null, ); } var staged = try testingWriter(std.testing.allocator, 2); defer std.testing.allocator.free(staged.deinit()); var appended = try testingWriter(std.testing.allocator, 2); defer std.testing.allocator.free(appended.deinit()); var first: [page.size]u8 = undefined; var second: [page.size]u8 = undefined; fillPage(&first, 1, 10); fillPage(&second, 2, 20); const position = staged.position(); try staged.stagePage(position, 0, 2, &second); try staged.stagePage(position, 1, 1, &first); try std.testing.expectEqual(@as(u32, 2), staged.stagedPageId(position, 0)); try std.testing.expectEqual(@as(u8, 20), staged.stagedPage(position, 0)[1]); staged.swapStagedFrames(position, 0, 1); try std.testing.expectEqual(@as(u32, 1), staged.stagedPageId(position, 0)); try std.testing.expectEqual(@as(u32, 2), staged.stagedPageId(position, 1)); staged.commitStagedFrame(position, 0, 0); staged.commitStagedFrame(position, 1, 2); try appended.append(1, 0, &first); try appended.append(2, 2, &second); try std.testing.expectEqualSlices(u8, appended.bytes(), staged.bytes()); const full_position = staged.position(); try std.testing.expectError(error.WalFull, staged.stagePage(full_position, 0, 3, &first)); try std.testing.expectEqual(full_position, staged.position());}test "wal ignores uncommitted tail frames" { var writer = try testingWriter(std.testing.allocator, 2); defer std.testing.allocator.free(writer.deinit()); var first: [page.size]u8 = undefined; var second: [page.size]u8 = undefined; fillPage(&first, 1, 10); fillPage(&second, 1, 99); try writer.append(1, 1, &first); try writer.append(1, 0, &second); const mark = try endMark(writer.bytes()); try std.testing.expectEqual(@as(usize, 1), mark); const recovered = (try pageAt(writer.bytes(), 1, mark)).?; try std.testing.expectEqual(@as(u8, 10), recovered[1]);}test "wal writer restores a prior append position" { var writer = try testingWriter(std.testing.allocator, 1); defer std.testing.allocator.free(writer.deinit()); var first: [page.size]u8 = undefined; var second: [page.size]u8 = undefined; fillPage(&first, 1, 10); fillPage(&second, 1, 20); const position = writer.position(); try writer.append(1, 1, &first); writer.restore(position); try std.testing.expectEqual(@as(usize, 0), writer.frameCount()); try writer.append(1, 1, &second); const mark = try endMark(writer.bytes()); try std.testing.expectEqual(@as(usize, 1), mark); const recovered = (try pageAt(writer.bytes(), 1, mark)).?; try std.testing.expectEqual(@as(u8, 20), recovered[1]);}test "wal recovery copies the committed prefix into fixed storage" { var writer = try testingWriter(std.testing.allocator, 2); defer std.testing.allocator.free(writer.deinit()); var first: [page.size]u8 = undefined; var second: [page.size]u8 = undefined; fillPage(&first, 1, 10); fillPage(&second, 2, 20); try writer.append(1, 1, &first); try writer.append(2, 0, &second); const committed_len = header_size + frame_size; var recovered = try testingWriter(std.testing.allocator, 2); defer std.testing.allocator.free(recovered.deinit()); const recovered_storage = recovered.storage.ptr; try recovered.load(writer.bytes(), committed_len); try std.testing.expectEqual(recovered_storage, recovered.bytes().ptr); try std.testing.expect(writer.bytes().ptr != recovered.bytes().ptr); try std.testing.expectEqual(committed_len, recovered.bytes().len); try std.testing.expectEqual(@as(usize, 1), recovered.frameCount()); try recovered.append(2, 2, &second); try std.testing.expectEqual(@as(usize, 2), try endMark(recovered.bytes()));}test "wal recovery rejects an invalid prefix before mutation" { var writer = try testingWriter(std.testing.allocator, 1); defer std.testing.allocator.free(writer.deinit()); const position = writer.position(); try std.testing.expectError(error.InvalidWal, writer.load("invalid", "invalid".len)); try std.testing.expectEqual(position, writer.position());}test "wal stops recovery before a corrupted frame" { var writer = try testingWriter(std.testing.allocator, 2); defer std.testing.allocator.free(writer.deinit()); var first: [page.size]u8 = undefined; var second: [page.size]u8 = undefined; fillPage(&first, 1, 10); fillPage(&second, 2, 20); try writer.append(1, 1, &first); try writer.append(2, 2, &second); writer.storage[header_size + frame_size + frame_header_size + 9] ^= 0xff; const mark = try endMark(writer.bytes()); try std.testing.expectEqual(@as(usize, 1), mark); try std.testing.expect(try pageAt(writer.bytes(), 2, mark) == null);}test "wal stops recovery at a torn frame" { var writer = try testingWriter(std.testing.allocator, 2); defer std.testing.allocator.free(writer.deinit()); var first: [page.size]u8 = undefined; var second: [page.size]u8 = undefined; fillPage(&first, 1, 10); fillPage(&second, 2, 20); try writer.append(1, 1, &first); try writer.append(2, 2, &second); writer.len -= 17; const mark = try endMark(writer.bytes()); try std.testing.expectEqual(@as(usize, 1), mark);}fn sweepWriter() !Writer { var writer = try testingWriter(std.testing.allocator, 3); errdefer std.testing.allocator.free(writer.deinit()); var image: [page.size]u8 = undefined; fillPage(&image, 1, 10); try writer.append(1, 0, &image); fillPage(&image, 2, 20); try writer.append(2, 2, &image); fillPage(&image, 1, 30); try writer.append(1, 0, &image); return writer;}fn sweepMarkForReadableFrames(readable: usize) usize { return if (readable >= 2) 2 else 0;}fn expectSweepState(data: []const u8, mark: usize) !void { switch (mark) { 0 => { try std.testing.expect(try pageAt(data, 1, mark) == null); try std.testing.expect(try pageAt(data, 2, mark) == null); }, 2 => { try std.testing.expectEqual(@as(u8, 10), (try pageAt(data, 1, mark)).?[1]); try std.testing.expectEqual(@as(u8, 20), (try pageAt(data, 2, mark)).?[1]); }, else => return error.TestUnexpectedResult, }}test "wal recovery survives truncation at every byte offset" { var writer = try sweepWriter(); defer std.testing.allocator.free(writer.deinit()); const data = writer.bytes(); var length: usize = 0; while (length <= data.len) : (length += 1) { const prefix = data[0..length]; if (length < header_size) { try std.testing.expectError(error.InvalidWal, endMark(prefix)); continue; } const readable = (length - header_size) / frame_size; const mark = try endMark(prefix); try std.testing.expectEqual(sweepMarkForReadableFrames(readable), mark); try expectSweepState(prefix, mark); }}test "wal recovery stops before a flip of any byte" { comptime { @stardustClaim( @import("alloc_phase").capacity.witness(Writer, "sql_wal_corruption"), null, null, null, null, null, null, ); } var writer = try sweepWriter(); defer std.testing.allocator.free(writer.deinit()); const data = writer.storage[0..writer.len]; var index: usize = 0; while (index < data.len) : (index += 1) { data[index] ^= 0xff; defer data[index] ^= 0xff; if (index < header_size) { if (endMark(data)) |_| { return error.TestUnexpectedResult; } else |err| switch (err) { error.InvalidWal, error.InvalidChecksum, error.UnsupportedPageSize => {}, } continue; } const flipped_frame = (index - header_size) / frame_size; const mark = try endMark(data); try std.testing.expectEqual(sweepMarkForReadableFrames(flipped_frame), mark); try expectSweepState(data, mark); }}test "wal recovery loads exactly the frame-aligned valid prefixes" { var writer = try sweepWriter(); defer std.testing.allocator.free(writer.deinit()); const data = writer.bytes(); var recovered = try testingWriter(std.testing.allocator, 3); defer std.testing.allocator.free(recovered.deinit()); var length: usize = 0; while (length <= data.len) : (length += 1) { const aligned = length >= header_size and (length - header_size) % frame_size == 0; if (!aligned) { try std.testing.expectError(error.InvalidWal, recovered.load(data, length)); continue; } try recovered.load(data, length); try std.testing.expectEqual((length - header_size) / frame_size, recovered.frameCount()); }}Complete call list for wal.Scanner.consumeFrame
7 direct calls.
lib.sql.src.wal.Scanner.appendValidated[method] — private source atlib/sql/src/wal.zig:482in nearest public ownertiny.sql.wallib.sql.src.wal.Scanner.assertBounds[method] — private source atlib/sql/src/wal.zig:600in nearest public ownertiny.sql.wallib.sql.src.wal.Scanner.finishPhysicalFrame[method] — private source atlib/sql/src/wal.zig:572in nearest public ownertiny.sql.wallib.sql.src.wal.Scanner.foldStaged[method] — private source atlib/sql/src/wal.zig:504in nearest public ownertiny.sql.wallib.sql.src.wal.Scanner.invalidMarker[method] — private source atlib/sql/src/wal.zig:463in nearest public ownertiny.sql.wallib.sql.src.wal.Scanner.invalidate[method] — private source atlib/sql/src/wal.zig:578in nearest public ownertiny.sql.wallib.sql.src.wal.Scanner.validateFrame[method] — private source atlib/sql/src/wal.zig:428in nearest public ownertiny.sql.wal
Complete caller list for wal.endMark
9 direct callers.
lib.sql.src.wal.test_wal_commit_marker_makes_preceding_frames_visible[function] — test source atlib/sql/src/wal.zig:2303in nearest public ownertiny.sql.wallib.sql.src.wal.test_wal_header_validates_and_has_no_committed_frames[function] — test source atlib/sql/src/wal.zig:2294in nearest public ownertiny.sql.wallib.sql.src.wal.test_wal_ignores_uncommitted_tail_frames[function] — test source atlib/sql/src/wal.zig:2376in nearest public ownertiny.sql.wallib.sql.src.wal.test_wal_recovery_copies_the_committed_prefix_into_fixed_storage[function] — test source atlib/sql/src/wal.zig:2414in nearest public ownertiny.sql.wallib.sql.src.wal.test_wal_recovery_stops_before_a_flip_of_any_byte[function] — test source atlib/sql/src/wal.zig:2531in nearest public ownertiny.sql.wallib.sql.src.wal.test_wal_recovery_survives_truncation_at_every_byte_offset[function] — test source atlib/sql/src/wal.zig:2512in nearest public ownertiny.sql.wallib.sql.src.wal.test_wal_stops_recovery_at_a_torn_frame[function] — test source atlib/sql/src/wal.zig:2464in nearest public ownertiny.sql.wallib.sql.src.wal.test_wal_stops_recovery_before_a_corrupted_frame[function] — test source atlib/sql/src/wal.zig:2447in nearest public ownertiny.sql.wallib.sql.src.wal.test_wal_writer_restores_a_prior_append_position[function] — test source atlib/sql/src/wal.zig:2393in nearest public ownertiny.sql.wal
Audit
| Definitions | 43 |
|---|---|
| Public names | 43 |
| Members | 71 |
| Version | 26.7.0 |
| Revision | daab053ee433 |