lib/sql/src/history/validate.zig
daab053ee43316e1809a84551d573ddd1e5bf3d2
1 const std = @import("std");
2 const sql = @import("../root.zig");
3 const conflict_mod = @import("conflict.zig");
4 const materialize_mod = @import("materialize.zig");
5 const record_mod = @import("record.zig");
6 const version = sql.version;
7
8 const Allocator = std.mem.Allocator;
9
10 pub const Limits = struct {
11 metadata_bytes_max: usize,
12 dependencies_max: usize,
13 conflict_bytes_max: usize,
14
15 fn assertValid(self: Limits) void {
16 std.debug.assert(self.metadata_bytes_max > 0);
17 std.debug.assert(self.dependencies_max > 0);
18 std.debug.assert(self.conflict_bytes_max > 0);
19 std.debug.assert(self.conflict_bytes_max <= std.math.maxInt(u32));
20 }
21 };
22
23 const stream_bytes: usize = 64 * 1024;
24 const prefix_read_bytes: usize = 1024 * 1024;
25 const refs_max: usize = 65_536;
26 const ref_name_bytes_each_max: usize = 4_096;
27 const all_record_kinds: u32 = std.math.maxInt(u32);
28
29 pub const BaselineRef = struct {
30 name: []const u8,
31 head: version.Hash,
32 };
33
34 const Record = struct {
35 offset: usize,
36 payload_offset: usize,
37 end: usize,
38 kind: record_mod.RecordKind,
39 expected: version.Hash,
40 payload: []const u8,
41 row_digest: ?version.Hash,
42 };
43
44 const Scanner = struct {
45 allocator: Allocator,
46 io: std.Io,
47 file: std.Io.File,
48 limit: usize,
49 offset: usize,
50 metadata_bytes_max: usize,
51 control: sql.wal.Control,
52 payload: std.ArrayList(u8) = .empty,
53 read_buffer: []u8 = &.{},
54 read_start: usize = 0,
55 read_filled: usize = 0,
56 read_refills: usize = 0,
57 read_bytes: usize = 0,
58
59 fn init(
60 allocator: Allocator,
61 io: std.Io,
62 file: std.Io.File,
63 start: usize,
64 limit: usize,
65 metadata_bytes_max: usize,
66 control: sql.wal.Control,
67 ) Scanner {
68 std.debug.assert(start <= limit);
69 return .{
70 .allocator = allocator,
71 .io = io,
72 .file = file,
73 .limit = limit,
74 .offset = start,
75 .metadata_bytes_max = metadata_bytes_max,
76 .control = control,
77 };
78 }
79
80 fn deinit(self: *Scanner) void {
81 if (self.read_buffer.len != 0) {
82 self.allocator.free(self.read_buffer);
83 }
84 self.payload.deinit(self.allocator);
85 self.* = undefined;
86 }
87
88 fn enableReadBuffer(self: *Scanner) !void {
89 std.debug.assert(self.read_buffer.len == 0);
90 self.read_buffer = try self.allocator.alloc(u8, prefix_read_bytes);
91 }
92
93 fn next(self: *Scanner, materialize: u32, verify_hash: bool) !?Record {
94 try self.control.check();
95 if (self.offset == self.limit) return null;
96 if (self.limit - self.offset < record_mod.record_header_size) {
97 return error.InvalidHistory;
98 }
99 var header: [record_mod.record_header_size]u8 = undefined;
100 try self.readExact(&header, self.offset);
101 const decoded = try record_mod.Header.decode(&header);
102 const kind = decoded.kind;
103 const payload_len: usize = decoded.payload_len;
104 const payload_offset = std.math.add(
105 usize,
106 self.offset,
107 record_mod.record_header_size,
108 ) catch return error.InvalidHistory;
109 const payload_end = std.math.add(usize, payload_offset, payload_len) catch
110 return error.InvalidHistory;
111 if (payload_end > self.limit) return error.InvalidHistory;
112 const expected = decoded.envelope_hash;
113 self.payload.clearRetainingCapacity();
114 const wanted = materialize & kindMask(kind) != 0;
115 const row_digest = try self.readPayload(
116 payload_offset,
117 payload_end,
118 @backingInt(kind),
119 kind,
120 expected,
121 wanted,
122 verify_hash,
123 );
124 const record_offset = self.offset;
125 self.offset = payload_end;
126 return .{
127 .offset = record_offset,
128 .payload_offset = payload_offset,
129 .end = payload_end,
130 .kind = kind,
131 .expected = expected,
132 .payload = self.payload.items,
133 .row_digest = row_digest,
134 };
135 }
136
137 fn readPayload(
138 self: *Scanner,
139 payload_offset: usize,
140 payload_end: usize,
141 kind_value: u32,
142 kind: record_mod.RecordKind,
143 expected: version.Hash,
144 wanted: bool,
145 verify_hash: bool,
146 ) !?version.Hash {
147 const payload_len = payload_end - payload_offset;
148 if (wanted and kind != .row_chunk) {
149 if (payload_len > self.metadata_bytes_max) return error.StreamTooLong;
150 try self.payload.resize(self.allocator, payload_len);
151 try self.readExact(self.payload.items, payload_offset);
152 if (verify_hash) {
153 const actual = record_mod.recordHash(kind_value, self.payload.items);
154 if (!version.same(actual, expected)) return error.InvalidHistory;
155 }
156 return null;
157 }
158 if (verify_hash) {
159 return try self.streamHash(
160 payload_offset,
161 payload_end,
162 kind_value,
163 kind,
164 expected,
165 );
166 }
167 if (wanted and kind == .row_chunk) {
168 var digest: version.Hash = undefined;
169 try self.readExact(&digest, payload_offset);
170 return digest;
171 }
172 return null;
173 }
174
175 fn streamHash(
176 self: *Scanner,
177 payload_offset: usize,
178 payload_end: usize,
179 kind_value: u32,
180 kind: record_mod.RecordKind,
181 expected: version.Hash,
182 ) !?version.Hash {
183 var hasher = record_mod.envelopeHasher(
184 kind_value,
185 @intCast(payload_end - payload_offset),
186 );
187 var scratch: [stream_bytes]u8 = undefined;
188 var row_digest: ?version.Hash = null;
189 var cursor = payload_offset;
190 var chunks: usize = 0;
191 const chunks_max = std.math.divCeil(
192 usize,
193 payload_end - payload_offset,
194 scratch.len,
195 ) catch unreachable;
196 while (cursor < payload_end) : (chunks += 1) {
197 std.debug.assert(chunks < chunks_max);
198 try self.control.check();
199 const count = @min(scratch.len, payload_end - cursor);
200 try self.readExact(scratch[0..count], cursor);
201 if (kind == .row_chunk and cursor == payload_offset) {
202 row_digest = scratch[0..version.hash_bytes].*;
203 }
204 hasher.update(scratch[0..count]);
205 cursor += count;
206 }
207 var actual: version.Hash = undefined;
208 hasher.final(&actual);
209 if (!version.same(actual, expected)) return error.InvalidHistory;
210 return row_digest;
211 }
212
213 fn readExact(self: *Scanner, target: []u8, start: usize) !void {
214 var filled: usize = 0;
215 while (filled < target.len) {
216 try self.control.check();
217 const cursor = start + filled;
218 if (self.read_buffer.len == 0) {
219 const count = try self.file.readPositionalAll(
220 self.io,
221 target[filled..],
222 cursor,
223 );
224 if (count == 0) return error.InvalidHistory;
225 filled += count;
226 continue;
227 }
228 if (cursor < self.read_start or
229 cursor >= self.read_start + self.read_filled)
230 {
231 self.read_start = cursor;
232 self.read_filled = try self.file.readPositionalAll(
233 self.io,
234 self.read_buffer[0..@min(
235 self.read_buffer.len,
236 self.limit - cursor,
237 )],
238 cursor,
239 );
240 if (self.read_filled == 0) return error.InvalidHistory;
241 self.read_refills = std.math.add(
242 usize,
243 self.read_refills,
244 1,
245 ) catch return error.StreamTooLong;
246 self.read_bytes = std.math.add(
247 usize,
248 self.read_bytes,
249 self.read_filled,
250 ) catch return error.StreamTooLong;
251 }
252 const buffered_offset = cursor - self.read_start;
253 const count = @min(
254 target.len - filled,
255 self.read_filled - buffered_offset,
256 );
257 @memcpy(
258 target[filled..][0..count],
259 self.read_buffer[buffered_offset..][0..count],
260 );
261 filled += count;
262 }
263 }
264 };
265
266 const RefEntry = struct {
267 name: []const u8,
268 head: version.Hash,
269 owned: bool,
270 };
271
272 const RefState = struct {
273 allocator: Allocator,
274 entries: std.ArrayList(RefEntry) = .empty,
275 by_name: std.StringHashMapUnmanaged(usize) = .empty,
276 name_bytes: usize = 0,
277 name_bytes_max: usize,
278
279 fn init(
280 allocator: Allocator,
281 baseline: anytype,
282 name_bytes_max: usize,
283 control: sql.wal.Control,
284 ) !RefState {
285 var state = RefState{
286 .allocator = allocator,
287 .name_bytes_max = name_bytes_max,
288 };
289 errdefer state.deinit();
290 if (baseline.len > refs_max) return error.StreamTooLong;
291 try state.entries.ensureTotalCapacity(allocator, baseline.len);
292 try state.by_name.ensureTotalCapacity(
293 allocator,
294 @intCast(baseline.len),
295 );
296 for (baseline, 0..) |entry, index| {
297 if (index % 256 == 0) try control.check();
298 if (entry.name.len == 0 or entry.name.len > ref_name_bytes_each_max) {
299 return error.InvalidHistory;
300 }
301 if (state.by_name.contains(entry.name)) return error.InvalidHistory;
302 state.name_bytes = std.math.add(
303 usize,
304 state.name_bytes,
305 entry.name.len,
306 ) catch return error.StreamTooLong;
307 if (state.name_bytes > state.name_bytes_max) return error.StreamTooLong;
308 state.entries.appendAssumeCapacity(.{
309 .name = entry.name,
310 .head = entry.head,
311 .owned = false,
312 });
313 state.by_name.putAssumeCapacity(entry.name, index);
314 }
315 return state;
316 }
317
318 fn deinit(self: *RefState) void {
319 self.by_name.deinit(self.allocator);
320 for (self.entries.items) |entry| {
321 if (entry.owned) self.allocator.free(entry.name);
322 }
323 self.entries.deinit(self.allocator);
324 self.* = undefined;
325 }
326
327 fn find(self: *const RefState, name: []const u8) ?usize {
328 return self.by_name.get(name);
329 }
330
331 fn set(self: *RefState, name: []const u8, head: version.Hash) !void {
332 if (name.len == 0 or name.len > ref_name_bytes_each_max) {
333 return error.InvalidHistory;
334 }
335 if (self.find(name)) |index| {
336 self.entries.items[index].head = head;
337 return;
338 }
339 if (self.entries.items.len >= refs_max) return error.StreamTooLong;
340 const next_bytes = std.math.add(usize, self.name_bytes, name.len) catch
341 return error.StreamTooLong;
342 if (next_bytes > self.name_bytes_max) return error.StreamTooLong;
343 try self.entries.ensureUnusedCapacity(self.allocator, 1);
344 try self.by_name.ensureUnusedCapacity(self.allocator, 1);
345 const owned = try self.allocator.dupe(u8, name);
346 errdefer self.allocator.free(owned);
347 const index = self.entries.items.len;
348 self.entries.appendAssumeCapacity(.{
349 .name = owned,
350 .head = head,
351 .owned = true,
352 });
353 self.by_name.putAssumeCapacity(owned, index);
354 self.name_bytes = next_bytes;
355 }
356
357 fn delete(self: *RefState, name: []const u8) void {
358 const removed_lookup = self.by_name.fetchRemove(name) orelse return;
359 const index = removed_lookup.value;
360 const last_index = self.entries.items.len - 1;
361 const removed = self.entries.swapRemove(index);
362 if (index != last_index) {
363 self.by_name.getPtr(self.entries.items[index].name).?.* = index;
364 }
365 self.name_bytes -= removed.name.len;
366 if (removed.owned) self.allocator.free(removed.name);
367 }
368 };
369
370 const ObjectKind = enum(u8) {
371 commit,
372 conflict,
373 index_page,
374 row_chunk,
375 tree_node,
376 };
377
378 const Slot = struct {
379 kind: version.ConflictKind,
380 relation_offset: u32,
381 relation_len: u32,
382 rowid: i64,
383 };
384
385 const Producer = struct {
386 kind: ObjectKind,
387 hash: version.Hash,
388 order: usize,
389 slot: ?Slot = null,
390
391 fn lessThan(_: void, left: Producer, right: Producer) bool {
392 if (left.kind != right.kind) {
393 return @backingInt(left.kind) < @backingInt(right.kind);
394 }
395 const hash_order = std.mem.order(u8, left.hash[0..], right.hash[0..]);
396 if (hash_order != .eq) return hash_order == .lt;
397 return left.order < right.order;
398 }
399 };
400
401 const Need = struct {
402 kind: ObjectKind,
403 hash: version.Hash,
404 required_before: usize,
405 source_offset: usize,
406 source_kind: record_mod.RecordKind,
407 slot: ?Slot = null,
408 resolved: bool = false,
409 prefix_seen: bool = false,
410
411 fn lessThan(_: void, left: Need, right: Need) bool {
412 if (left.kind != right.kind) {
413 return @backingInt(left.kind) < @backingInt(right.kind);
414 }
415 const hash_order = std.mem.order(u8, left.hash[0..], right.hash[0..]);
416 if (hash_order != .eq) return hash_order == .lt;
417 return left.required_before < right.required_before;
418 }
419 };
420
421 const Objects = struct {
422 allocator: Allocator,
423 producers: std.ArrayList(Producer) = .empty,
424 needs: std.ArrayList(Need) = .empty,
425 relations: std.ArrayList(u8) = .empty,
426 unresolved: usize = 0,
427 dependencies_max: usize,
428 conflict_bytes_max: usize,
429 source_offset: usize = 0,
430 source_kind: record_mod.RecordKind = .commit,
431
432 fn init(allocator: Allocator, limits: Limits) Objects {
433 return .{
434 .allocator = allocator,
435 .dependencies_max = limits.dependencies_max,
436 .conflict_bytes_max = limits.conflict_bytes_max,
437 };
438 }
439
440 fn deinit(self: *Objects) void {
441 self.producers.deinit(self.allocator);
442 self.needs.deinit(self.allocator);
443 self.relations.deinit(self.allocator);
444 self.* = undefined;
445 }
446
447 fn addProducer(
448 self: *Objects,
449 kind: ObjectKind,
450 hash: version.Hash,
451 order: usize,
452 ) !void {
453 if (self.producers.items.len >= self.dependencies_max) {
454 return error.StreamTooLong;
455 }
456 try self.producers.append(self.allocator, .{
457 .kind = kind,
458 .hash = hash,
459 .order = order,
460 });
461 }
462
463 fn addConflictProducer(
464 self: *Objects,
465 entry: version.ConflictEntry,
466 order: usize,
467 ) !void {
468 if (self.producers.items.len >= self.dependencies_max) {
469 return error.StreamTooLong;
470 }
471 try self.producers.append(self.allocator, .{
472 .kind = .conflict,
473 .hash = entry.hash,
474 .order = order,
475 .slot = try self.copySlot(entry),
476 });
477 }
478
479 fn addNeed(
480 self: *Objects,
481 kind: ObjectKind,
482 hash: version.Hash,
483 required_before: usize,
484 ) !void {
485 if (self.needs.items.len >= self.dependencies_max) return error.StreamTooLong;
486 try self.needs.append(self.allocator, .{
487 .kind = kind,
488 .hash = hash,
489 .required_before = required_before,
490 .source_offset = self.source_offset,
491 .source_kind = self.source_kind,
492 });
493 self.unresolved += 1;
494 }
495
496 fn addConflictNeed(
497 self: *Objects,
498 entry: version.ConflictEntry,
499 required_before: usize,
500 ) !void {
501 if (self.needs.items.len >= self.dependencies_max) return error.StreamTooLong;
502 try self.needs.append(self.allocator, .{
503 .kind = .conflict,
504 .hash = entry.hash,
505 .required_before = required_before,
506 .source_offset = self.source_offset,
507 .source_kind = self.source_kind,
508 .slot = try self.copySlot(entry),
509 });
510 self.unresolved += 1;
511 }
512
513 fn copySlot(self: *Objects, entry: version.ConflictEntry) !Slot {
514 const next_len = std.math.add(
515 usize,
516 self.relations.items.len,
517 entry.relation.len,
518 ) catch return error.StreamTooLong;
519 if (next_len > self.conflict_bytes_max) return error.StreamTooLong;
520 const offset = self.relations.items.len;
521 try self.relations.appendSlice(self.allocator, entry.relation);
522 return .{
523 .kind = entry.kind,
524 .relation_offset = @intCast(offset),
525 .relation_len = @intCast(entry.relation.len),
526 .rowid = entry.rowid,
527 };
528 }
529
530 fn resolveSuffix(self: *Objects, control: sql.wal.Control) !void {
531 std.mem.sort(Producer, self.producers.items, {}, Producer.lessThan);
532 std.mem.sort(Need, self.needs.items, {}, Need.lessThan);
533 var producer_index: usize = 0;
534 var need_index: usize = 0;
535 while (need_index < self.needs.items.len) {
536 try control.check();
537 const first_need = self.needs.items[need_index];
538 while (producer_index < self.producers.items.len and keyBefore(
539 self.producers.items[producer_index].kind,
540 self.producers.items[producer_index].hash,
541 first_need.kind,
542 first_need.hash,
543 )) : (producer_index += 1) {
544 if (producer_index % 256 == 0) try control.check();
545 }
546 const producer = if (producer_index < self.producers.items.len and
547 sameKey(self.producers.items[producer_index], first_need))
548 self.producers.items[producer_index]
549 else
550 null;
551 const need_end = try self.needGroupEnd(need_index, control);
552 for (self.needs.items[need_index..need_end], 0..) |*need, index| {
553 if (index % 256 == 0) try control.check();
554 const available = producer orelse continue;
555 if (available.order >= need.required_before or
556 !self.slotsEqual(need.slot, available.slot)) continue;
557 need.resolved = true;
558 self.unresolved -= 1;
559 }
560 need_index = need_end;
561 }
562 }
563
564 fn prefixMask(
565 self: *const Objects,
566 control: sql.wal.Control,
567 ) !u32 {
568 var mask: u32 = 0;
569 for (self.needs.items, 0..) |need, index| {
570 if (index % 256 == 0) try control.check();
571 if (need.resolved) continue;
572 mask |= kindMask(switch (need.kind) {
573 .commit => .commit,
574 .conflict => .conflict,
575 .index_page => .chunk_index_page,
576 .row_chunk => .row_chunk,
577 .tree_node => .tree_nodes,
578 });
579 }
580 return mask;
581 }
582
583 fn markPrefix(
584 self: *Objects,
585 kind: ObjectKind,
586 hash: version.Hash,
587 control: sql.wal.Control,
588 ) !void {
589 const first = self.needLowerBound(kind, hash);
590 if (first == self.needs.items.len or
591 !sameNeedKey(self.needs.items[first], kind, hash) or
592 self.needs.items[first].prefix_seen) return;
593 const end = try self.needGroupEnd(first, control);
594 for (self.needs.items[first..end], 0..) |*need, index| {
595 if (index % 256 == 0) try control.check();
596 need.prefix_seen = true;
597 if (need.resolved or need.slot != null) continue;
598 need.resolved = true;
599 self.unresolved -= 1;
600 }
601 }
602
603 fn markPrefixConflict(
604 self: *Objects,
605 entry: version.ConflictEntry,
606 control: sql.wal.Control,
607 ) !void {
608 const first = self.needLowerBound(.conflict, entry.hash);
609 if (first == self.needs.items.len or
610 !sameNeedKey(self.needs.items[first], .conflict, entry.hash) or
611 self.needs.items[first].prefix_seen) return;
612 const end = try self.needGroupEnd(first, control);
613 for (self.needs.items[first..end], 0..) |*need, index| {
614 if (index % 256 == 0) try control.check();
615 need.prefix_seen = true;
616 if (need.resolved or !self.slotMatches(need.slot.?, entry)) continue;
617 need.resolved = true;
618 self.unresolved -= 1;
619 }
620 }
621
622 fn needGroupEnd(
623 self: *const Objects,
624 start: usize,
625 control: sql.wal.Control,
626 ) !usize {
627 std.debug.assert(start < self.needs.items.len);
628 const first = self.needs.items[start];
629 var end = start + 1;
630 while (end < self.needs.items.len and
631 sameNeedKey(self.needs.items[end], first.kind, first.hash)) : (end += 1)
632 {
633 if (end % 256 == 0) try control.check();
634 }
635 return end;
636 }
637
638 fn needLowerBound(
639 self: *const Objects,
640 kind: ObjectKind,
641 hash: version.Hash,
642 ) usize {
643 var low: usize = 0;
644 var high = self.needs.items.len;
645 while (low < high) {
646 const middle = low + (high - low) / 2;
647 const candidate = self.needs.items[middle];
648 if (keyBefore(candidate.kind, candidate.hash, kind, hash)) {
649 low = middle + 1;
650 } else {
651 high = middle;
652 }
653 }
654 return low;
655 }
656
657 fn slotsEqual(self: *const Objects, left: ?Slot, right: ?Slot) bool {
658 if (left == null or right == null) return left == null and right == null;
659 const a = left.?;
660 const b = right.?;
661 return a.kind == b.kind and
662 (a.kind == .relation or a.rowid == b.rowid) and
663 std.mem.eql(u8, self.relation(a), self.relation(b));
664 }
665
666 fn slotMatches(
667 self: *const Objects,
668 slot: Slot,
669 entry: version.ConflictEntry,
670 ) bool {
671 return slot.kind == entry.kind and
672 (slot.kind == .relation or slot.rowid == entry.rowid) and
673 std.mem.eql(u8, self.relation(slot), entry.relation);
674 }
675
676 fn relation(self: *const Objects, slot: Slot) []const u8 {
677 const start: usize = slot.relation_offset;
678 return self.relations.items[start..][0..slot.relation_len];
679 }
680 };
681
682 const Decision = enum {
683 pending,
684 baseline,
685 target,
686 };
687
688 const FastForward = struct {
689 id: version.Hash,
690 ref_index: usize,
691 expected: version.Hash,
692 target: version.Hash,
693 decision: Decision = .pending,
694 };
695
696 const Validation = struct {
697 allocator: Allocator,
698 io: std.Io,
699 file: std.Io.File,
700 control: sql.wal.Control,
701 refs: RefState,
702 objects: Objects,
703 fast_forward: ?FastForward = null,
704
705 fn init(
706 allocator: Allocator,
707 io: std.Io,
708 file: std.Io.File,
709 baseline: anytype,
710 limits: Limits,
711 control: sql.wal.Control,
712 ) !Validation {
713 return .{
714 .allocator = allocator,
715 .io = io,
716 .file = file,
717 .control = control,
718 .refs = try RefState.init(
719 allocator,
720 baseline,
721 limits.metadata_bytes_max,
722 control,
723 ),
724 .objects = Objects.init(allocator, limits),
725 };
726 }
727
728 fn deinit(self: *Validation) void {
729 self.refs.deinit();
730 self.objects.deinit();
731 self.* = undefined;
732 }
733
734 fn apply(self: *Validation, record: Record) !void {
735 self.objects.source_offset = record.offset;
736 self.objects.source_kind = record.kind;
737 if (self.fast_forward != null and !isFastForward(record.kind)) {
738 return error.InvalidHistory;
739 }
740 switch (record.kind) {
741 .commit => try self.commit(record),
742 .ref => try self.refUpdate(record),
743 .conflict => try self.conflict(record),
744 .database_root => try self.databaseRoot(record),
745 .relation_root => try self.relationRoot(record),
746 .relation_rows => try self.relationRows(record),
747 .conflict_root => try self.conflictRoot(record),
748 .ref_delete => try self.refDelete(record),
749 .row_chunk => try self.objects.addProducer(
750 .row_chunk,
751 record.row_digest.?,
752 record.end,
753 ),
754 .chunk_index_page => try self.indexPage(record),
755 .tree_nodes => try self.treeNodes(record),
756 .relation_spans => try self.relationSpans(record),
757 .fast_forward_prepare => try self.fastForwardPrepare(record),
758 .fast_forward_commit => try self.fastForwardDecision(record, .target),
759 .fast_forward_abort => try self.fastForwardDecision(record, .baseline),
760 .fast_forward_complete => try self.fastForwardComplete(record),
761 }
762 }
763
764 fn commit(self: *Validation, record: Record) !void {
765 const identity = try commitIdentity(
766 self.allocator,
767 record.payload,
768 self.control,
769 );
770 try self.objects.addProducer(.commit, identity, record.end);
771 }
772
773 fn refUpdate(self: *Validation, record: Record) !void {
774 var reader = record_mod.PayloadReader.init(record.payload);
775 const target = try reader.hash();
776 const name = try reader.readBytes();
777 try reader.finish();
778 try self.objects.addNeed(.commit, target, record.payload_offset);
779 try self.refs.set(name, target);
780 }
781
782 fn refDelete(self: *Validation, record: Record) !void {
783 var reader = record_mod.PayloadReader.init(record.payload);
784 const name = try reader.readBytes();
785 try reader.finish();
786 self.refs.delete(name);
787 }
788
789 fn conflict(self: *Validation, record: Record) !void {
790 var reader = record_mod.PayloadReader.init(record.payload);
791 const artifact = try conflict_mod.decodeConflictArtifactPayload(&reader);
792 try reader.finish();
793 try self.objects.addConflictProducer(artifact.entry(), record.end);
794 }
795
796 fn databaseRoot(self: *Validation, record: Record) !void {
797 try preflightDatabaseRoot(record.payload, self.control);
798 var reader = record_mod.PayloadReader.init(record.payload);
799 var root = try record_mod.readDatabaseRootValue(self.allocator, &reader);
800 defer root.deinit();
801 try reader.finish();
802 }
803
804 fn relationRoot(self: *Validation, record: Record) !void {
805 try preflightRelationRoot(record.payload, self.control);
806 var reader = record_mod.PayloadReader.init(record.payload);
807 var decoded = try materialize_mod.readRelationRootShallow(
808 self.allocator,
809 &reader,
810 );
811 defer decoded.root.deinit();
812 defer self.allocator.free(decoded.index_keys);
813 try reader.finish();
814 if (decoded.table_key) |key| {
815 try self.objects.addNeed(.tree_node, key, record.payload_offset);
816 }
817 for (decoded.index_keys) |index_key| {
818 if (index_key) |key| {
819 try self.objects.addNeed(.tree_node, key, record.payload_offset);
820 }
821 }
822 }
823
824 fn relationRows(self: *Validation, record: Record) !void {
825 var reader = record_mod.PayloadReader.init(record.payload);
826 _ = try reader.hash();
827 const count = try reader.readU32();
828 var index: u32 = 0;
829 while (index < count) : (index += 1) {
830 try self.control.check();
831 try self.objects.addNeed(
832 .index_page,
833 try reader.hash(),
834 record.payload_offset,
835 );
836 }
837 try reader.finish();
838 }
839
840 fn indexPage(self: *Validation, record: Record) !void {
841 var reader = record_mod.PayloadReader.init(record.payload);
842 const digest = try reader.hash();
843 const count = try reader.readU32();
844 var index: u32 = 0;
845 while (index < count) : (index += 1) {
846 try self.control.check();
847 try self.objects.addNeed(
848 .row_chunk,
849 try reader.hash(),
850 record.payload_offset,
851 );
852 }
853 try reader.finish();
854 try self.objects.addProducer(.index_page, digest, record.end);
855 }
856
857 fn treeNodes(self: *Validation, record: Record) !void {
858 var reader = record_mod.PayloadReader.init(record.payload);
859 const count = try reader.readU32();
860 var index: u32 = 0;
861 while (index < count) : (index += 1) {
862 try self.control.check();
863 const key = try reader.hash();
864 try skipTreeNodeHeader(&reader);
865 const child_count = try reader.readU32();
866 var child_index: u32 = 0;
867 while (child_index < child_count) : (child_index += 1) {
868 const child = try reader.hash();
869 try self.objects.addNeed(
870 .tree_node,
871 child,
872 record.payload_offset + reader.cursor,
873 );
874 }
875 try self.objects.addProducer(
876 .tree_node,
877 key,
878 record.payload_offset + reader.cursor,
879 );
880 }
881 try reader.finish();
882 }
883
884 fn relationSpans(self: *Validation, record: Record) !void {
885 _ = self;
886 var reader = record_mod.PayloadReader.init(record.payload);
887 _ = try reader.hash();
888 const count = try reader.readU32();
889 const bytes = std.math.mul(usize, count, 2 * @sizeOf(u64)) catch
890 return error.InvalidHistory;
891 if (bytes != reader.remaining()) return error.InvalidHistory;
892 reader.cursor += bytes;
893 try reader.finish();
894 }
895
896 fn conflictRoot(self: *Validation, record: Record) !void {
897 var reader = record_mod.PayloadReader.init(record.payload);
898 const entries = try conflict_mod.readConflictEntries(self.allocator, &reader);
899 defer conflict_mod.deinitConflictEntries(self.allocator, entries);
900 try reader.finish();
901 var previous: ?version.ConflictEntry = null;
902 for (entries) |entry| {
903 try self.control.check();
904 if (previous) |prior| {
905 if (prior.sameSlot(entry) or
906 !version.ConflictEntry.lessThan({}, prior, entry))
907 {
908 return error.InvalidHistory;
909 }
910 }
911 try self.objects.addConflictNeed(entry, record.payload_offset);
912 previous = entry;
913 }
914 }
915
916 fn fastForwardPrepare(self: *Validation, record: Record) !void {
917 if (self.fast_forward != null) return error.InvalidHistory;
918 var reader = record_mod.PayloadReader.init(record.payload);
919 const name = try reader.readBytes();
920 const expected = try reader.hash();
921 const target = try reader.hash();
922 try reader.finish();
923 const ref_index = self.refs.find(name) orelse return error.InvalidHistory;
924 if (!version.same(self.refs.entries.items[ref_index].head, expected)) {
925 return error.InvalidHistory;
926 }
927 try self.objects.addNeed(.commit, expected, record.payload_offset);
928 try self.objects.addNeed(.commit, target, record.payload_offset);
929 self.fast_forward = .{
930 .id = record.expected,
931 .ref_index = ref_index,
932 .expected = expected,
933 .target = target,
934 };
935 }
936
937 fn fastForwardDecision(
938 self: *Validation,
939 record: Record,
940 decision: Decision,
941 ) !void {
942 var reader = record_mod.PayloadReader.init(record.payload);
943 const id = try reader.hash();
944 try reader.finish();
945 const active = if (self.fast_forward) |*value| value else return error.InvalidHistory;
946 if (!version.same(active.id, id) or active.decision != .pending) {
947 return error.InvalidHistory;
948 }
949 const ref_entry = &self.refs.entries.items[active.ref_index];
950 if (!version.same(ref_entry.head, active.expected)) {
951 return error.InvalidHistory;
952 }
953 if (decision == .target) ref_entry.head = active.target;
954 active.decision = decision;
955 }
956
957 fn fastForwardComplete(self: *Validation, record: Record) !void {
958 var reader = record_mod.PayloadReader.init(record.payload);
959 const id = try reader.hash();
960 try reader.finish();
961 const active = self.fast_forward orelse return error.InvalidHistory;
962 if (!version.same(active.id, id) or active.decision == .pending) {
963 return error.InvalidHistory;
964 }
965 const selected = if (active.decision == .target)
966 active.target
967 else
968 active.expected;
969 if (!version.same(self.refs.entries.items[active.ref_index].head, selected)) {
970 return error.InvalidHistory;
971 }
972 self.fast_forward = null;
973 }
974 };
975
976 pub fn validateSuffix(
977 allocator: Allocator,
978 io: std.Io,
979 file: std.Io.File,
980 start: usize,
981 end: usize,
982 baseline_refs: anytype,
983 limits: Limits,
984 control: sql.wal.Control,
985 ) !void {
986 std.debug.assert(start <= end);
987 limits.assertValid();
988 try control.check();
989 var validation = try Validation.init(
990 allocator,
991 io,
992 file,
993 baseline_refs,
994 limits,
995 control,
996 );
997 defer validation.deinit();
998 var scanner = Scanner.init(
999 allocator,
1000 io,
1001 file,
1002 start,
1003 end,
1004 limits.metadata_bytes_max,
1005 control,
1006 );
1007 defer scanner.deinit();
1008 var records: usize = 0;
1009 const records_max = (end - start) / record_mod.record_header_size + 1;
1010 while (try scanner.next(all_record_kinds, true)) |record| : (records += 1) {
1011 std.debug.assert(records < records_max);
1012 try validation.apply(record);
1013 }
1014 if (validation.fast_forward != null) return error.InvalidHistory;
1015 try validation.objects.resolveSuffix(control);
1016 if (validation.objects.unresolved == 0) return;
1017 try resolvePrefix(
1018 &validation.objects,
1019 allocator,
1020 io,
1021 file,
1022 start,
1023 limits.metadata_bytes_max,
1024 control,
1025 );
1026 if (validation.objects.unresolved != 0) return error.InvalidHistory;
1027 }
1028
1029 /// A diagnostic for the first record that a full read-only scan rejects.
1030 pub const BadRecord = struct {
1031 offset: usize,
1032 kind: ?record_mod.RecordKind = null,
1033 key: ?version.Hash = null,
1034 dependency: ?Dependency = null,
1035 };
1036
1037 pub const Dependency = struct {
1038 kind: record_mod.RecordKind,
1039 key: version.Hash,
1040 };
1041
1042 pub const Inspection = struct {
1043 last_valid_boundary: usize,
1044 bad: ?BadRecord = null,
1045 };
1046
1047 /// Verifies every envelope and the dependencies covered by validateSuffix.
1048 pub fn inspectWhole(
1049 allocator: Allocator,
1050 io: std.Io,
1051 file: std.Io.File,
1052 end: usize,
1053 limits: Limits,
1054 control: sql.wal.Control,
1055 ) !Inspection {
1056 limits.assertValid();
1057 var validation = try Validation.init(
1058 allocator,
1059 io,
1060 file,
1061 &[_]BaselineRef{},
1062 limits,
1063 control,
1064 );
1065 defer validation.deinit();
1066 var scanner = Scanner.init(
1067 allocator,
1068 io,
1069 file,
1070 0,
1071 end,
1072 limits.metadata_bytes_max,
1073 control,
1074 );
1075 defer scanner.deinit();
1076 var boundary: usize = 0;
1077 var records: usize = 0;
1078 const records_max = end / record_mod.record_header_size + 1;
1079 while (true) {
1080 const next = scanner.next(all_record_kinds, true) catch |err| switch (err) {
1081 error.InvalidHistory => return .{
1082 .last_valid_boundary = boundary,
1083 .bad = try badAt(allocator, io, file, scanner.offset, end),
1084 },
1085 else => return err,
1086 };
1087 const record = next orelse break;
1088 std.debug.assert(records < records_max);
1089 records += 1;
1090 validation.apply(record) catch |err| switch (err) {
1091 error.InvalidHistory => return .{
1092 .last_valid_boundary = boundary,
1093 .bad = .{
1094 .offset = record.offset,
1095 .kind = record.kind,
1096 .key = try keyAt(allocator, record.kind, record.payload),
1097 },
1098 },
1099 else => return err,
1100 };
1101 boundary = record.end;
1102 }
1103 if (validation.fast_forward != null) {
1104 return .{ .last_valid_boundary = boundary, .bad = .{ .offset = boundary } };
1105 }
1106 try validation.objects.resolveSuffix(control);
1107 var earliest: ?Need = null;
1108 for (validation.objects.needs.items) |need| {
1109 if (need.resolved) continue;
1110 if (earliest == null or need.source_offset < earliest.?.source_offset) {
1111 earliest = need;
1112 }
1113 }
1114 if (earliest) |need| {
1115 var bad = try badAt(allocator, io, file, need.source_offset, end);
1116 bad.dependency = .{
1117 .kind = dependencyRecordKind(need.kind),
1118 .key = need.hash,
1119 };
1120 return .{ .last_valid_boundary = need.source_offset, .bad = bad };
1121 }
1122 return .{ .last_valid_boundary = end };
1123 }
1124
1125 fn dependencyRecordKind(kind: ObjectKind) record_mod.RecordKind {
1126 return switch (kind) {
1127 .commit => .commit,
1128 .conflict => .conflict,
1129 .index_page => .chunk_index_page,
1130 .row_chunk => .row_chunk,
1131 .tree_node => .tree_nodes,
1132 };
1133 }
1134
1135 fn badAt(
1136 allocator: Allocator,
1137 io: std.Io,
1138 file: std.Io.File,
1139 offset: usize,
1140 end: usize,
1141 ) !BadRecord {
1142 var bad = BadRecord{ .offset = offset };
1143 if (end - offset < record_mod.record_header_size) return bad;
1144 var header_bytes: [record_mod.record_header_size]u8 = undefined;
1145 if (try file.readPositionalAll(io, &header_bytes, offset) != header_bytes.len) return bad;
1146 const header = record_mod.Header.decode(&header_bytes) catch return bad;
1147 bad.kind = header.kind;
1148 const payload_start = offset + record_mod.record_header_size;
1149 if (header.payload_len > end - payload_start) return bad;
1150 if (header.kind == .row_chunk) {
1151 var key: version.Hash = undefined;
1152 if (try file.readPositionalAll(io, &key, payload_start) == key.len) bad.key = key;
1153 return bad;
1154 }
1155 const payload = try allocator.alloc(u8, header.payload_len);
1156 defer allocator.free(payload);
1157 if (try file.readPositionalAll(io, payload, payload_start) != payload.len) return bad;
1158 bad.key = try keyAt(allocator, header.kind, payload);
1159 return bad;
1160 }
1161
1162 fn keyAt(
1163 allocator: Allocator,
1164 kind: record_mod.RecordKind,
1165 payload: []const u8,
1166 ) !?version.Hash {
1167 const identity = @import("identity.zig");
1168 const entries_max = payload.len / version.hash_bytes + 1;
1169 const relations = try allocator.alloc(version.RelationEntry, entries_max);
1170 defer allocator.free(relations);
1171 const conflicts = try allocator.alloc(version.ConflictEntry, entries_max);
1172 defer allocator.free(conflicts);
1173 var key: ?version.Hash = null;
1174 const Capture = struct {
1175 fn accept(target: *?version.Hash, entry: identity.Entry) !void {
1176 if (target.* == null) target.* = entry.key.bytes;
1177 }
1178 };
1179 identity.each(
1180 kind,
1181 payload,
1182 .{ .record_offset = 0, .payload_len = @intCast(payload.len), .envelope_hash = @splat(0) },
1183 .{ .relations = relations, .conflicts = conflicts },
1184 &key,
1185 Capture.accept,
1186 ) catch return null;
1187 return key;
1188 }
1189
1190 fn resolvePrefix(
1191 objects: *Objects,
1192 allocator: Allocator,
1193 io: std.Io,
1194 file: std.Io.File,
1195 end: usize,
1196 metadata_bytes_max: usize,
1197 control: sql.wal.Control,
1198 ) !void {
1199 if (objects.unresolved == 0 or end == 0) return;
1200 const mask = try objects.prefixMask(control);
1201 var scanner = Scanner.init(
1202 allocator,
1203 io,
1204 file,
1205 0,
1206 end,
1207 metadata_bytes_max,
1208 control,
1209 );
1210 defer scanner.deinit();
1211 try scanner.enableReadBuffer();
1212 var records: usize = 0;
1213 const records_max = end / record_mod.record_header_size + 1;
1214 while (try scanner.next(mask, false)) |record| : (records += 1) {
1215 std.debug.assert(records < records_max);
1216 if (mask & kindMask(record.kind) != 0) {
1217 try applyPrefixProducer(objects, allocator, record, control);
1218 if (objects.unresolved == 0) break;
1219 }
1220 }
1221 }
1222
1223 fn applyPrefixProducer(
1224 objects: *Objects,
1225 allocator: Allocator,
1226 record: Record,
1227 control: sql.wal.Control,
1228 ) !void {
1229 switch (record.kind) {
1230 .commit => try objects.markPrefix(
1231 .commit,
1232 try commitIdentity(allocator, record.payload, control),
1233 control,
1234 ),
1235 .conflict => {
1236 var reader = record_mod.PayloadReader.init(record.payload);
1237 const artifact = try conflict_mod.decodeConflictArtifactPayload(&reader);
1238 try reader.finish();
1239 try objects.markPrefixConflict(artifact.entry(), control);
1240 },
1241 .chunk_index_page => {
1242 var reader = record_mod.PayloadReader.init(record.payload);
1243 const digest = try reader.hash();
1244 try skipHashArray(&reader, try reader.readU32(), control);
1245 try reader.finish();
1246 try objects.markPrefix(.index_page, digest, control);
1247 },
1248 .row_chunk => try objects.markPrefix(
1249 .row_chunk,
1250 record.row_digest.?,
1251 control,
1252 ),
1253 .tree_nodes => try markPrefixTreeNodes(objects, record.payload, control),
1254 else => {},
1255 }
1256 }
1257
1258 fn markPrefixTreeNodes(
1259 objects: *Objects,
1260 payload: []const u8,
1261 control: sql.wal.Control,
1262 ) !void {
1263 var reader = record_mod.PayloadReader.init(payload);
1264 const count = try reader.readU32();
1265 var index: u32 = 0;
1266 while (index < count) : (index += 1) {
1267 try control.check();
1268 try objects.markPrefix(.tree_node, try reader.hash(), control);
1269 try skipTreeNodeHeader(&reader);
1270 try skipHashArray(&reader, try reader.readU32(), control);
1271 }
1272 try reader.finish();
1273 }
1274
1275 fn preflightDatabaseRoot(
1276 payload: []const u8,
1277 control: sql.wal.Control,
1278 ) !void {
1279 var reader = record_mod.PayloadReader.init(payload);
1280 _ = try reader.hash();
1281 const count = try reader.readU32();
1282 const entry_bytes_min = @sizeOf(u32) + version.hash_bytes;
1283 if (@as(usize, count) > reader.remaining() / entry_bytes_min) {
1284 return error.InvalidHistory;
1285 }
1286 var index: u32 = 0;
1287 while (index < count) : (index += 1) {
1288 if (index % 256 == 0) try control.check();
1289 _ = try reader.readBytes();
1290 _ = try reader.hash();
1291 }
1292 try reader.finish();
1293 }
1294
1295 fn preflightRelationRoot(
1296 payload: []const u8,
1297 control: sql.wal.Control,
1298 ) !void {
1299 var reader = record_mod.PayloadReader.init(payload);
1300 _ = try reader.readU32();
1301 _ = try reader.readBytes();
1302 _ = try reader.readU32();
1303 _ = try reader.readU64();
1304 _ = try reader.hash();
1305 try preflightRelationSchema(&reader, control);
1306 _ = try materialize_mod.readMapRootHeader(&reader);
1307 _ = try record_mod.readStatsRoot(&reader);
1308
1309 const count = try reader.readU32();
1310 const map_header_bytes_min = 10 * @sizeOf(u64) +
1311 2 * version.hash_bytes + @sizeOf(u8);
1312 const index_bytes_min = version.hash_bytes + map_header_bytes_min +
1313 2 * version.hash_bytes;
1314 if (@as(usize, count) > reader.remaining() / index_bytes_min) {
1315 return error.InvalidHistory;
1316 }
1317 var index: u32 = 0;
1318 while (index < count) : (index += 1) {
1319 if (index % 256 == 0) try control.check();
1320 _ = try reader.hash();
1321 _ = try materialize_mod.readMapRootHeader(&reader);
1322 _ = try reader.hash();
1323 _ = try reader.hash();
1324 }
1325 _ = try reader.hash();
1326 try reader.finish();
1327 }
1328
1329 fn preflightRelationSchema(
1330 reader: *record_mod.PayloadReader,
1331 control: sql.wal.Control,
1332 ) !void {
1333 const column_count = try reader.readU32();
1334 const column_bytes_min = @sizeOf(u32) + 2 * @sizeOf(u8);
1335 if (@as(usize, column_count) > reader.remaining() / column_bytes_min) {
1336 return error.InvalidHistory;
1337 }
1338 var column_index: u32 = 0;
1339 while (column_index < column_count) : (column_index += 1) {
1340 if (column_index % 256 == 0) try control.check();
1341 _ = try reader.readBytes();
1342 _ = try record_mod.collation(try reader.readU8());
1343 try preflightValue(reader);
1344 }
1345
1346 const index_count = try reader.readU32();
1347 const schema_index_bytes_min = 3 * @sizeOf(u32);
1348 if (@as(usize, index_count) >
1349 reader.remaining() / schema_index_bytes_min)
1350 {
1351 return error.InvalidHistory;
1352 }
1353 var index: u32 = 0;
1354 while (index < index_count) : (index += 1) {
1355 try control.check();
1356 _ = try reader.readBytes();
1357 const field_count = try reader.readU32();
1358 if (@as(usize, field_count) > reader.remaining() / @sizeOf(u64)) {
1359 return error.InvalidHistory;
1360 }
1361 var field_index: u32 = 0;
1362 while (field_index < field_count) : (field_index += 1) {
1363 if (field_index % 256 == 0) try control.check();
1364 _ = try record_mod.readUsize(reader);
1365 }
1366 const projection_count = try reader.readU32();
1367 if (@as(usize, projection_count) > reader.remaining()) {
1368 return error.InvalidHistory;
1369 }
1370 var projection_index: u32 = 0;
1371 while (projection_index < projection_count) : (projection_index += 1) {
1372 if (projection_index % 256 == 0) try control.check();
1373 _ = try record_mod.collation(try reader.readU8());
1374 }
1375 }
1376 }
1377
1378 fn preflightValue(reader: *record_mod.PayloadReader) !void {
1379 switch (try reader.readU8()) {
1380 @backingInt(sql.row.Storage.nil) => {},
1381 @backingInt(sql.row.Storage.integer) => _ = try reader.readI64(),
1382 @backingInt(sql.row.Storage.text),
1383 @backingInt(sql.row.Storage.blob),
1384 => _ = try reader.readBytes(),
1385 else => return error.InvalidHistory,
1386 }
1387 }
1388
1389 fn commitIdentity(
1390 allocator: Allocator,
1391 payload: []const u8,
1392 control: sql.wal.Control,
1393 ) !version.Hash {
1394 var reader = record_mod.PayloadReader.init(payload);
1395 const root = try reader.hash();
1396 const parent_count = try reader.readU32();
1397 const parent_bytes = std.math.mul(usize, parent_count, version.hash_bytes) catch
1398 return error.InvalidHistory;
1399 if (parent_bytes != reader.remaining()) return error.InvalidHistory;
1400 const parents = try allocator.alloc(version.Hash, parent_count);
1401 defer allocator.free(parents);
1402 for (parents, 0..) |*parent, index| {
1403 if (index % 256 == 0) try control.check();
1404 parent.* = try reader.hash();
1405 }
1406 try reader.finish();
1407 return version.Commit.init(root, parents).hash;
1408 }
1409
1410 fn skipTreeNodeHeader(reader: *record_mod.PayloadReader) !void {
1411 _ = try record_mod.nodeKind(try reader.readU8());
1412 _ = try reader.readBytes();
1413 _ = try reader.optionalBytes();
1414 _ = try record_mod.readUsize(reader);
1415 _ = try record_mod.readSummary(reader);
1416 _ = try reader.hash();
1417 }
1418
1419 fn skipHashArray(
1420 reader: *record_mod.PayloadReader,
1421 count: u32,
1422 control: sql.wal.Control,
1423 ) !void {
1424 const bytes = std.math.mul(usize, count, version.hash_bytes) catch
1425 return error.InvalidHistory;
1426 if (bytes > reader.remaining()) return error.InvalidHistory;
1427 try control.check();
1428 reader.cursor += bytes;
1429 }
1430
1431 fn keyBefore(
1432 left_kind: ObjectKind,
1433 left_hash: version.Hash,
1434 right_kind: ObjectKind,
1435 right_hash: version.Hash,
1436 ) bool {
1437 if (left_kind != right_kind) {
1438 return @backingInt(left_kind) < @backingInt(right_kind);
1439 }
1440 return std.mem.lessThan(u8, left_hash[0..], right_hash[0..]);
1441 }
1442
1443 fn sameKey(producer: Producer, need: Need) bool {
1444 return producer.kind == need.kind and version.same(producer.hash, need.hash);
1445 }
1446
1447 fn sameNeedKey(need: Need, kind: ObjectKind, hash: version.Hash) bool {
1448 return need.kind == kind and version.same(need.hash, hash);
1449 }
1450
1451 fn isFastForward(kind: record_mod.RecordKind) bool {
1452 return switch (kind) {
1453 .fast_forward_prepare,
1454 .fast_forward_commit,
1455 .fast_forward_abort,
1456 .fast_forward_complete,
1457 => true,
1458 else => false,
1459 };
1460 }
1461
1462 fn kindMask(kind: record_mod.RecordKind) u32 {
1463 const shift: u5 = @intCast(@backingInt(kind) - 1);
1464 return @as(u32, 1) << shift;
1465 }
1466
1467 const testing_io = std.Options.debug_io;
1468 const testing_limits = Limits{
1469 .metadata_bytes_max = 4 * 1024 * 1024,
1470 .dependencies_max = (4 * 1024 * 1024) / version.hash_bytes,
1471 .conflict_bytes_max = 4 * 1024 * 1024,
1472 };
1473
1474 fn appendTestingRecord(
1475 file: std.Io.File,
1476 offset: usize,
1477 kind: record_mod.RecordKind,
1478 payload: []const u8,
1479 ) !usize {
1480 var bytes: std.ArrayList(u8) = .empty;
1481 defer bytes.deinit(std.testing.allocator);
1482 try record_mod.appendU32(std.testing.allocator, &bytes, record_mod.magic);
1483 try record_mod.appendU32(std.testing.allocator, &bytes, record_mod.format_version);
1484 try record_mod.appendU32(std.testing.allocator, &bytes, @backingInt(kind));
1485 try record_mod.appendU32(std.testing.allocator, &bytes, @intCast(payload.len));
1486 try record_mod.appendHash(
1487 std.testing.allocator,
1488 &bytes,
1489 record_mod.recordHash(@backingInt(kind), payload),
1490 );
1491 try bytes.appendSlice(std.testing.allocator, payload);
1492 try file.writePositionalAll(testing_io, bytes.items, offset);
1493 return offset + bytes.items.len;
1494 }
1495
1496 test "attested prefix scanner skips an unrelated sparse payload" {
1497 var tmp = std.testing.tmpDir(.{});
1498 defer tmp.cleanup();
1499 var file = try tmp.dir.createFile(testing_io, "sparse.history", .{ .read = true });
1500 defer file.close(testing_io);
1501
1502 const sparse_payload_bytes: usize = 64 * 1024 * 1024;
1503 const sparse_expected = version.Hash{
1504 0xfc, 0x66, 0xb7, 0x1e, 0x01, 0xf4, 0x22, 0x01,
1505 0xf9, 0x79, 0x8c, 0x2f, 0x32, 0xd4, 0x96, 0x7d,
1506 0x20, 0x86, 0x43, 0x6f, 0x49, 0x49, 0x6e, 0xc4,
1507 0x1e, 0xb1, 0xee, 0x94, 0xbd, 0x4d, 0x0f, 0x8a,
1508 };
1509 var header: std.ArrayList(u8) = .empty;
1510 defer header.deinit(std.testing.allocator);
1511 try record_mod.appendU32(std.testing.allocator, &header, record_mod.magic);
1512 try record_mod.appendU32(
1513 std.testing.allocator,
1514 &header,
1515 record_mod.format_version,
1516 );
1517 try record_mod.appendU32(
1518 std.testing.allocator,
1519 &header,
1520 @backingInt(record_mod.RecordKind.row_chunk),
1521 );
1522 try record_mod.appendU32(
1523 std.testing.allocator,
1524 &header,
1525 @intCast(sparse_payload_bytes),
1526 );
1527 try record_mod.appendHash(std.testing.allocator, &header, sparse_expected);
1528 try std.testing.expectEqual(record_mod.record_header_size, header.items.len);
1529 try file.writePositionalAll(testing_io, header.items, 0);
1530 const commit_offset = header.items.len + sparse_payload_bytes;
1531 try file.setLength(testing_io, @intCast(commit_offset));
1532
1533 const root = version.emptyHash("validate.sparse.root");
1534 var commit_payload: std.ArrayList(u8) = .empty;
1535 defer commit_payload.deinit(std.testing.allocator);
1536 try record_mod.appendHash(std.testing.allocator, &commit_payload, root);
1537 try record_mod.appendU32(std.testing.allocator, &commit_payload, 0);
1538 const end = try appendTestingRecord(
1539 file,
1540 commit_offset,
1541 .commit,
1542 commit_payload.items,
1543 );
1544
1545 var scanner = Scanner.init(
1546 std.testing.allocator,
1547 testing_io,
1548 file,
1549 0,
1550 end,
1551 testing_limits.metadata_bytes_max,
1552 .{},
1553 );
1554 defer scanner.deinit();
1555 try scanner.enableReadBuffer();
1556 const sparse = (try scanner.next(kindMask(.commit), false)).?;
1557 try std.testing.expectEqual(record_mod.RecordKind.row_chunk, sparse.kind);
1558 try std.testing.expectEqual(@as(usize, 0), sparse.payload.len);
1559 const commit = (try scanner.next(kindMask(.commit), false)).?;
1560 try std.testing.expectEqual(record_mod.RecordKind.commit, commit.kind);
1561 try std.testing.expectEqual(commit_payload.items.len, commit.payload.len);
1562 try std.testing.expectEqual(@as(usize, 2), scanner.read_refills);
1563 try std.testing.expect(scanner.read_bytes <= prefix_read_bytes + 128);
1564 }
1565
1566 test "history suffix validation resolves a commit in the prefix" {
1567 var tmp = std.testing.tmpDir(.{});
1568 defer tmp.cleanup();
1569 var file = try tmp.dir.createFile(testing_io, "valid.history", .{ .read = true });
1570 defer file.close(testing_io);
1571 const root = version.emptyHash("validate.prefix.root");
1572 const commit = version.Commit.init(root, &.{});
1573 var commit_payload: std.ArrayList(u8) = .empty;
1574 defer commit_payload.deinit(std.testing.allocator);
1575 try record_mod.appendHash(std.testing.allocator, &commit_payload, root);
1576 try record_mod.appendU32(std.testing.allocator, &commit_payload, 0);
1577 const suffix_start = try appendTestingRecord(file, 0, .commit, commit_payload.items);
1578 var ref_payload: std.ArrayList(u8) = .empty;
1579 defer ref_payload.deinit(std.testing.allocator);
1580 try record_mod.appendHash(std.testing.allocator, &ref_payload, commit.hash);
1581 try record_mod.appendBytes(std.testing.allocator, &ref_payload, "main");
1582 const end = try appendTestingRecord(file, suffix_start, .ref, ref_payload.items);
1583 try validateSuffix(
1584 std.testing.allocator,
1585 testing_io,
1586 file,
1587 suffix_start,
1588 end,
1589 &[_]BaselineRef{},
1590 testing_limits,
1591 .{},
1592 );
1593 }
1594
1595 test "history suffix validation rejects a ref before its future commit" {
1596 var tmp = std.testing.tmpDir(.{});
1597 defer tmp.cleanup();
1598 var file = try tmp.dir.createFile(testing_io, "future.history", .{ .read = true });
1599 defer file.close(testing_io);
1600 const root = version.emptyHash("validate.future.root");
1601 const commit = version.Commit.init(root, &.{});
1602 var ref_payload: std.ArrayList(u8) = .empty;
1603 defer ref_payload.deinit(std.testing.allocator);
1604 try record_mod.appendHash(std.testing.allocator, &ref_payload, commit.hash);
1605 try record_mod.appendBytes(std.testing.allocator, &ref_payload, "main");
1606 var end = try appendTestingRecord(file, 0, .ref, ref_payload.items);
1607 var commit_payload: std.ArrayList(u8) = .empty;
1608 defer commit_payload.deinit(std.testing.allocator);
1609 try record_mod.appendHash(std.testing.allocator, &commit_payload, root);
1610 try record_mod.appendU32(std.testing.allocator, &commit_payload, 0);
1611 end = try appendTestingRecord(file, end, .commit, commit_payload.items);
1612 try std.testing.expectError(
1613 error.InvalidHistory,
1614 validateSuffix(
1615 std.testing.allocator,
1616 testing_io,
1617 file,
1618 0,
1619 end,
1620 &[_]BaselineRef{},
1621 testing_limits,
1622 .{},
1623 ),
1624 );
1625 }
1626
1627 test "history suffix validation agrees with replay on malformed conflict" {
1628 var tmp = std.testing.tmpDir(.{});
1629 defer tmp.cleanup();
1630 const path = "conflict.history";
1631 const malformed = [_]u8{0xff};
1632 const end = blk: {
1633 var file = try tmp.dir.createFile(testing_io, path, .{ .read = true });
1634 defer file.close(testing_io);
1635 const written = try appendTestingRecord(file, 0, .conflict, &malformed);
1636 try std.testing.expectError(
1637 error.InvalidHistory,
1638 validateSuffix(
1639 std.testing.allocator,
1640 testing_io,
1641 file,
1642 0,
1643 written,
1644 &[_]BaselineRef{},
1645 testing_limits,
1646 .{},
1647 ),
1648 );
1649 break :blk written;
1650 };
1651 try std.testing.expect(end > record_mod.record_header_size);
1652 try std.testing.expectError(
1653 error.InvalidHistory,
1654 sql.History.open(std.testing.allocator, tmp.dir, .{
1655 .path = path,
1656 .recovery = .reject,
1657 }),
1658 );
1659 }
1660
1661 test "history suffix validation rejects a conflict root before its artifact" {
1662 var tmp = std.testing.tmpDir(.{});
1663 defer tmp.cleanup();
1664 var file = try tmp.dir.createFile(testing_io, "root.history", .{ .read = true });
1665 defer file.close(testing_io);
1666 const artifact = version.ConflictArtifact.init("issues", 7, null, null, null);
1667 const entry = artifact.entry();
1668 var payload: std.ArrayList(u8) = .empty;
1669 defer payload.deinit(std.testing.allocator);
1670 try conflict_mod.appendConflictEntries(std.testing.allocator, &payload, &.{entry});
1671 const end = try appendTestingRecord(file, 0, .conflict_root, payload.items);
1672 try std.testing.expectError(
1673 error.InvalidHistory,
1674 validateSuffix(
1675 std.testing.allocator,
1676 testing_io,
1677 file,
1678 0,
1679 end,
1680 &[_]BaselineRef{},
1681 testing_limits,
1682 .{},
1683 ),
1684 );
1685 }
1686
1687 test "history suffix validation accepts an adjacent conflict artifact and root" {
1688 var tmp = std.testing.tmpDir(.{});
1689 defer tmp.cleanup();
1690 const path = "conflict-root.history";
1691 {
1692 var file = try tmp.dir.createFile(testing_io, path, .{ .read = true });
1693 defer file.close(testing_io);
1694 const artifact = version.ConflictArtifact.init("issues", 7, null, "ours", null);
1695 var artifact_payload: std.ArrayList(u8) = .empty;
1696 defer artifact_payload.deinit(std.testing.allocator);
1697 try conflict_mod.appendConflictArtifactRecordPayload(
1698 std.testing.allocator,
1699 &artifact_payload,
1700 artifact,
1701 );
1702 var end = try appendTestingRecord(file, 0, .conflict, artifact_payload.items);
1703 var root_payload: std.ArrayList(u8) = .empty;
1704 defer root_payload.deinit(std.testing.allocator);
1705 try conflict_mod.appendConflictEntries(
1706 std.testing.allocator,
1707 &root_payload,
1708 &.{artifact.entry()},
1709 );
1710 end = try appendTestingRecord(file, end, .conflict_root, root_payload.items);
1711 try validateSuffix(
1712 std.testing.allocator,
1713 testing_io,
1714 file,
1715 0,
1716 end,
1717 &[_]BaselineRef{},
1718 testing_limits,
1719 .{},
1720 );
1721 }
1722 var recovered = try sql.History.open(std.testing.allocator, tmp.dir, .{
1723 .path = path,
1724 .recovery = .reject,
1725 });
1726 defer recovered.deinit();
1727 }
1728
1729 test "history suffix validation rejects a first tree node that names itself" {
1730 var tmp = std.testing.tmpDir(.{});
1731 defer tmp.cleanup();
1732 const path = "self-tree.history";
1733 const key = version.emptyHash("validate.self.tree");
1734 const end = blk: {
1735 var file = try tmp.dir.createFile(testing_io, path, .{ .read = true });
1736 defer file.close(testing_io);
1737 var payload: std.ArrayList(u8) = .empty;
1738 defer payload.deinit(std.testing.allocator);
1739 try record_mod.appendU32(std.testing.allocator, &payload, 1);
1740 try record_mod.appendHash(std.testing.allocator, &payload, key);
1741 try record_mod.appendU8(
1742 std.testing.allocator,
1743 &payload,
1744 @backingInt(sql.tree.NodeKind.leaf),
1745 );
1746 try record_mod.appendBytes(std.testing.allocator, &payload, "");
1747 try record_mod.appendOptionalBytes(std.testing.allocator, &payload, null);
1748 try record_mod.appendU64(std.testing.allocator, &payload, 0);
1749 try record_mod.appendSummary(std.testing.allocator, &payload, .{});
1750 try record_mod.appendHash(
1751 std.testing.allocator,
1752 &payload,
1753 version.emptyHash("validate.self.tree.node"),
1754 );
1755 try record_mod.appendU32(std.testing.allocator, &payload, 1);
1756 try record_mod.appendHash(std.testing.allocator, &payload, key);
1757 const written = try appendTestingRecord(file, 0, .tree_nodes, payload.items);
1758 try std.testing.expectError(
1759 error.InvalidHistory,
1760 validateSuffix(
1761 std.testing.allocator,
1762 testing_io,
1763 file,
1764 0,
1765 written,
1766 &[_]BaselineRef{},
1767 testing_limits,
1768 .{},
1769 ),
1770 );
1771 break :blk written;
1772 };
1773 try std.testing.expect(end > record_mod.record_header_size);
1774 try std.testing.expectError(
1775 error.InvalidHistory,
1776 sql.History.open(std.testing.allocator, tmp.dir, .{
1777 .path = path,
1778 .recovery = .reject,
1779 }),
1780 );
1781 }
1782
1783 test "history suffix validation rejects a byte-corrupted dependency prefix" {
1784 var tmp = std.testing.tmpDir(.{});
1785 defer tmp.cleanup();
1786 var file = try tmp.dir.createFile(testing_io, "corrupt.history", .{ .read = true });
1787 defer file.close(testing_io);
1788 const root = version.emptyHash("validate.corrupt.root");
1789 const commit = version.Commit.init(root, &.{});
1790 var commit_payload: std.ArrayList(u8) = .empty;
1791 defer commit_payload.deinit(std.testing.allocator);
1792 try record_mod.appendHash(std.testing.allocator, &commit_payload, root);
1793 try record_mod.appendU32(std.testing.allocator, &commit_payload, 0);
1794 const suffix_start = try appendTestingRecord(file, 0, .commit, commit_payload.items);
1795 var ref_payload: std.ArrayList(u8) = .empty;
1796 defer ref_payload.deinit(std.testing.allocator);
1797 try record_mod.appendHash(std.testing.allocator, &ref_payload, commit.hash);
1798 try record_mod.appendBytes(std.testing.allocator, &ref_payload, "main");
1799 const end = try appendTestingRecord(file, suffix_start, .ref, ref_payload.items);
1800 var byte: [1]u8 = undefined;
1801 try std.testing.expectEqual(
1802 @as(usize, 1),
1803 try file.readPositionalAll(
1804 testing_io,
1805 &byte,
1806 record_mod.record_header_size,
1807 ),
1808 );
1809 byte[0] ^= 0xff;
1810 try file.writePositionalAll(testing_io, &byte, record_mod.record_header_size);
1811 try std.testing.expectError(
1812 error.InvalidHistory,
1813 validateSuffix(
1814 std.testing.allocator,
1815 testing_io,
1816 file,
1817 suffix_start,
1818 end,
1819 &[_]BaselineRef{},
1820 testing_limits,
1821 .{},
1822 ),
1823 );
1824 }
1825
1826 test "history suffix validation applies complete fast forward and rejects its crash cut" {
1827 var tmp = std.testing.tmpDir(.{});
1828 defer tmp.cleanup();
1829 var file = try tmp.dir.createFile(testing_io, "forward.history", .{ .read = true });
1830 defer file.close(testing_io);
1831 const first = version.Commit.init(version.emptyHash("forward.first"), &.{});
1832 const second = version.Commit.init(version.emptyHash("forward.second"), &.{});
1833 var offset: usize = 0;
1834 for ([_]version.Commit{ first, second }) |commit| {
1835 var payload: std.ArrayList(u8) = .empty;
1836 defer payload.deinit(std.testing.allocator);
1837 try record_mod.appendHash(std.testing.allocator, &payload, commit.root);
1838 try record_mod.appendU32(std.testing.allocator, &payload, 0);
1839 offset = try appendTestingRecord(file, offset, .commit, payload.items);
1840 }
1841 var ref_payload: std.ArrayList(u8) = .empty;
1842 defer ref_payload.deinit(std.testing.allocator);
1843 try record_mod.appendHash(std.testing.allocator, &ref_payload, first.hash);
1844 try record_mod.appendBytes(std.testing.allocator, &ref_payload, "main");
1845 offset = try appendTestingRecord(file, offset, .ref, ref_payload.items);
1846 const suffix_start = offset;
1847 var prepare: std.ArrayList(u8) = .empty;
1848 defer prepare.deinit(std.testing.allocator);
1849 try record_mod.appendBytes(std.testing.allocator, &prepare, "main");
1850 try record_mod.appendHash(std.testing.allocator, &prepare, first.hash);
1851 try record_mod.appendHash(std.testing.allocator, &prepare, second.hash);
1852 const id = record_mod.recordHash(
1853 @backingInt(record_mod.RecordKind.fast_forward_prepare),
1854 prepare.items,
1855 );
1856 const crash_end = try appendTestingRecord(file, offset, .fast_forward_prepare, prepare.items);
1857 var decision: std.ArrayList(u8) = .empty;
1858 defer decision.deinit(std.testing.allocator);
1859 try record_mod.appendHash(std.testing.allocator, &decision, id);
1860 offset = try appendTestingRecord(file, crash_end, .fast_forward_commit, decision.items);
1861 const end = try appendTestingRecord(file, offset, .fast_forward_complete, decision.items);
1862 const baseline = [_]BaselineRef{.{ .name = "main", .head = first.hash }};
1863 try validateSuffix(
1864 std.testing.allocator,
1865 testing_io,
1866 file,
1867 suffix_start,
1868 end,
1869 &baseline,
1870 testing_limits,
1871 .{},
1872 );
1873 try std.testing.expectError(
1874 error.InvalidHistory,
1875 validateSuffix(
1876 std.testing.allocator,
1877 testing_io,
1878 file,
1879 suffix_start,
1880 crash_end,
1881 &baseline,
1882 testing_limits,
1883 .{},
1884 ),
1885 );
1886 }
1887
1888 test "history suffix validation streams an oversized row chunk" {
1889 var tmp = std.testing.tmpDir(.{});
1890 defer tmp.cleanup();
1891 var file = try tmp.dir.createFile(testing_io, "row.history", .{ .read = true });
1892 defer file.close(testing_io);
1893 const payload = try std.testing.allocator.alloc(
1894 u8,
1895 testing_limits.metadata_bytes_max + 1,
1896 );
1897 defer std.testing.allocator.free(payload);
1898 @memset(payload, 0xa5);
1899 const end = try appendTestingRecord(file, 0, .row_chunk, payload);
1900 try validateSuffix(
1901 std.testing.allocator,
1902 testing_io,
1903 file,
1904 0,
1905 end,
1906 &[_]BaselineRef{},
1907 testing_limits,
1908 .{},
1909 );
1910 }
1911
1912 test "history suffix validation preflights allocation-driving counts" {
1913 var tmp = std.testing.tmpDir(.{});
1914 defer tmp.cleanup();
1915
1916 {
1917 var file = try tmp.dir.createFile(
1918 testing_io,
1919 "database-count.history",
1920 .{ .read = true },
1921 );
1922 defer file.close(testing_io);
1923 var payload: std.ArrayList(u8) = .empty;
1924 defer payload.deinit(std.testing.allocator);
1925 try record_mod.appendHash(
1926 std.testing.allocator,
1927 &payload,
1928 version.emptyHash("validate.database.count"),
1929 );
1930 try record_mod.appendU32(
1931 std.testing.allocator,
1932 &payload,
1933 std.math.maxInt(u32),
1934 );
1935 const end = try appendTestingRecord(file, 0, .database_root, payload.items);
1936 try std.testing.expectError(
1937 error.InvalidHistory,
1938 validateSuffix(
1939 std.testing.allocator,
1940 testing_io,
1941 file,
1942 0,
1943 end,
1944 &[_]BaselineRef{},
1945 testing_limits,
1946 .{},
1947 ),
1948 );
1949 }
1950
1951 {
1952 var file = try tmp.dir.createFile(
1953 testing_io,
1954 "relation-count.history",
1955 .{ .read = true },
1956 );
1957 defer file.close(testing_io);
1958 var payload: std.ArrayList(u8) = .empty;
1959 defer payload.deinit(std.testing.allocator);
1960 try record_mod.appendU32(std.testing.allocator, &payload, 1);
1961 try record_mod.appendBytes(std.testing.allocator, &payload, "issues");
1962 try record_mod.appendU32(std.testing.allocator, &payload, 1);
1963 try record_mod.appendU64(std.testing.allocator, &payload, 1);
1964 try record_mod.appendHash(
1965 std.testing.allocator,
1966 &payload,
1967 version.emptyHash("validate.relation.schema"),
1968 );
1969 try record_mod.appendU32(
1970 std.testing.allocator,
1971 &payload,
1972 std.math.maxInt(u32),
1973 );
1974 const end = try appendTestingRecord(file, 0, .relation_root, payload.items);
1975 try std.testing.expectError(
1976 error.InvalidHistory,
1977 validateSuffix(
1978 std.testing.allocator,
1979 testing_io,
1980 file,
1981 0,
1982 end,
1983 &[_]BaselineRef{},
1984 testing_limits,
1985 .{},
1986 ),
1987 );
1988 }
1989 }