lib/trace/src/session.zig
daab053ee43316e1809a84551d573ddd1e5bf3d2
1 const std = @import("std");
2 const event = @import("event.zig");
3 const checkpoint = @import("checkpoint.zig");
4
5 const Allocator = std.mem.Allocator;
6
7 pub const ReplayProgress = struct {
8 cursor: u64,
9 event_count: u64,
10 remaining_count: u64,
11
12 pub fn complete(self: ReplayProgress) bool {
13 return self.remaining_count == 0;
14 }
15 };
16
17 pub const Session = struct {
18 allocator: Allocator,
19 mode: event.Mode,
20 sink: ?event.Sink = null,
21 source: ?event.Source = null,
22 epoch: u64 = 0,
23 next_seq: u64 = 1,
24 sequence_exhausted: bool = false,
25 replay_cursor: u64 = 0,
26 last_checkpoint: ?event.Event = null,
27
28 pub fn initOff(allocator: Allocator) Session {
29 return .{ .allocator = allocator, .mode = .off };
30 }
31
32 pub fn initRecord(allocator: Allocator, sink: event.Sink) Session {
33 return .{ .allocator = allocator, .mode = .record, .sink = sink };
34 }
35
36 pub fn initReplay(allocator: Allocator, source: event.Source) Session {
37 return .{ .allocator = allocator, .mode = .replay, .source = source };
38 }
39
40 pub fn deinit(self: *Session) void {
41 if (self.last_checkpoint) |*item| item.deinit(self.allocator);
42 self.* = undefined;
43 }
44
45 pub fn replayProgress(self: *const Session) ReplayProgress {
46 const event_count = if (self.source) |source| source.count() else 0;
47 std.debug.assert(self.replay_cursor <= event_count);
48 return .{
49 .cursor = self.replay_cursor,
50 .event_count = event_count,
51 .remaining_count = event_count - self.replay_cursor,
52 };
53 }
54
55 pub fn sessionStart(
56 self: *Session,
57 thread_id: event.ThreadId,
58 label: []const u8,
59 ) !event.Timepoint {
60 const timepoint = try self.prepareTimepoint(thread_id);
61 return try self.handleEvent(event.Event.sessionStart(timepoint, label));
62 }
63
64 pub fn sessionEnd(
65 self: *Session,
66 thread_id: event.ThreadId,
67 status: i64,
68 ) !event.Timepoint {
69 const timepoint = try self.prepareTimepoint(thread_id);
70 return try self.handleEvent(event.Event.sessionEnd(timepoint, status));
71 }
72
73 pub fn functionEnter(
74 self: *Session,
75 thread_id: event.ThreadId,
76 site: event.Safepoint,
77 ) !event.Timepoint {
78 const timepoint = try self.prepareTimepoint(thread_id);
79 return try self.handleEvent(event.Event.functionEnter(timepoint, site));
80 }
81
82 pub fn functionExit(
83 self: *Session,
84 thread_id: event.ThreadId,
85 site: event.Safepoint,
86 ) !event.Timepoint {
87 const timepoint = try self.prepareTimepoint(thread_id);
88 return try self.handleEvent(event.Event.functionExit(timepoint, site));
89 }
90
91 pub fn safepoint(
92 self: *Session,
93 thread_id: event.ThreadId,
94 site: event.Safepoint,
95 ) !event.Timepoint {
96 const timepoint = try self.prepareTimepoint(thread_id);
97 return try self.handleEvent(event.Event.safepointReached(timepoint, site));
98 }
99
100 pub fn allocation(
101 self: *Session,
102 thread_id: event.ThreadId,
103 object_id: u64,
104 size: u64,
105 alignment: u32,
106 label: ?[]const u8,
107 ) !event.Timepoint {
108 const timepoint = try self.prepareTimepoint(thread_id);
109 return try self.handleEvent(event.Event.allocation(
110 timepoint,
111 object_id,
112 size,
113 alignment,
114 label,
115 ));
116 }
117
118 pub fn free(
119 self: *Session,
120 thread_id: event.ThreadId,
121 object_id: u64,
122 ) !event.Timepoint {
123 const timepoint = try self.prepareTimepoint(thread_id);
124 return try self.handleEvent(event.Event.free(timepoint, object_id));
125 }
126
127 pub fn boundaryBytesAlloc(
128 self: *Session,
129 thread_id: event.ThreadId,
130 operation: []const u8,
131 recorded_value: ?[]const u8,
132 ) ![]u8 {
133 const timepoint = try self.prepareTimepoint(thread_id);
134 switch (self.mode) {
135 .off => {
136 const source_value = recorded_value orelse return error.MissingBoundaryValue;
137 const value = try self.allocator.dupe(u8, source_value);
138 self.commitTimepoint(timepoint);
139 return value;
140 },
141 .record => {
142 const value = recorded_value orelse return error.MissingBoundaryValue;
143 const owned = try self.allocator.dupe(u8, value);
144 errdefer self.allocator.free(owned);
145 try self.sink.?.append(event.Event.boundaryBytes(timepoint, operation, value));
146 self.commitTimepoint(timepoint);
147 return owned;
148 },
149 .replay => {
150 const source = self.source.?;
151 const actual = (try source.peek()) orelse return error.MissingReplayEvent;
152 const recorded = actual.data orelse return error.ReplayEventMismatch;
153 const expected = event.Event.boundaryBytes(timepoint, operation, recorded);
154 if (!expected.eqlForReplay(actual.*)) return error.ReplayEventMismatch;
155 const next_replay_cursor = try self.prepareReplayCursor();
156 const value = try self.allocator.dupe(u8, recorded);
157 source.advance();
158 self.replay_cursor = next_replay_cursor;
159 self.commitTimepoint(timepoint);
160 return value;
161 },
162 }
163 }
164
165 pub fn checkpointCapture(
166 self: *Session,
167 thread_id: event.ThreadId,
168 label: []const u8,
169 provider: checkpoint.Provider,
170 ) !event.Timepoint {
171 const timepoint = try self.prepareTimepoint(thread_id);
172 const bytes = try provider.captureAlloc(self.allocator);
173 defer self.allocator.free(bytes);
174 return try self.handleEvent(event.Event.checkpoint(timepoint, label, bytes));
175 }
176
177 pub fn checkpointRestore(
178 self: *Session,
179 thread_id: event.ThreadId,
180 label: []const u8,
181 provider: checkpoint.Provider,
182 ) !event.Timepoint {
183 const checkpoint_event = self.last_checkpoint orelse return error.MissingCheckpoint;
184 const checkpoint_bytes = checkpoint_event.data orelse return error.MissingCheckpoint;
185 const timepoint = try self.prepareTimepoint(thread_id);
186 const expected = event.Event.checkpointRestore(timepoint, label);
187 const prepared_replay = if (self.mode == .replay)
188 try self.prepareReplay(expected)
189 else
190 null;
191 const replay_timepoint = if (prepared_replay) |prepared|
192 prepared.actual.timepoint
193 else
194 timepoint;
195
196 provider.prepareRestoreBytes(checkpoint_bytes) catch |err| {
197 provider.cancelRestore();
198 return err;
199 };
200 errdefer provider.cancelRestore();
201
202 if (self.mode == .record) try self.sink.?.append(expected);
203
204 provider.commitRestore();
205 if (prepared_replay) |prepared| self.commitReplay(prepared);
206 self.commitTimepoint(timepoint);
207 return replay_timepoint;
208 }
209
210 pub fn userEvent(
211 self: *Session,
212 thread_id: event.ThreadId,
213 label: []const u8,
214 bytes: []const u8,
215 ) !event.Timepoint {
216 const timepoint = try self.prepareTimepoint(thread_id);
217 return try self.handleEvent(event.Event.user(timepoint, label, bytes));
218 }
219
220 pub fn verifyReplayComplete(self: *Session) !void {
221 if (self.mode != .replay) return;
222 const source = self.source.?;
223 if (self.replay_cursor != source.count()) return error.ReplayNotComplete;
224 if (try source.peek() != null) return error.ReplayNotComplete;
225 }
226
227 fn prepareTimepoint(self: *const Session, thread_id: event.ThreadId) !event.Timepoint {
228 if (self.sequence_exhausted) return error.SequenceExhausted;
229 return event.Timepoint{
230 .epoch = self.epoch,
231 .thread_id = thread_id,
232 .seq = self.next_seq,
233 };
234 }
235
236 fn commitTimepoint(self: *Session, timepoint: event.Timepoint) void {
237 std.debug.assert(!self.sequence_exhausted);
238 std.debug.assert(timepoint.epoch == self.epoch);
239 std.debug.assert(timepoint.seq == self.next_seq);
240 if (self.next_seq == std.math.maxInt(u64)) {
241 self.sequence_exhausted = true;
242 } else {
243 self.next_seq += 1;
244 }
245 }
246
247 fn handleEvent(self: *Session, expected: event.Event) !event.Timepoint {
248 var owned_checkpoint: ?event.Event = null;
249 errdefer if (owned_checkpoint) |*owned| owned.deinit(self.allocator);
250
251 var replay_timepoint: ?event.Timepoint = null;
252 switch (self.mode) {
253 .off => {
254 if (expected.kind == .checkpoint) {
255 owned_checkpoint = try expected.cloneAlloc(self.allocator);
256 }
257 },
258 .record => {
259 if (expected.kind == .checkpoint) {
260 owned_checkpoint = try expected.cloneAlloc(self.allocator);
261 }
262 try self.sink.?.append(expected);
263 },
264 .replay => {
265 const prepared = try self.prepareReplay(expected);
266 replay_timepoint = prepared.actual.timepoint;
267 if (prepared.actual.kind == .checkpoint) {
268 owned_checkpoint = try prepared.actual.cloneAlloc(self.allocator);
269 }
270 self.commitReplay(prepared);
271 },
272 }
273
274 self.commitTimepoint(expected.timepoint);
275 if (owned_checkpoint) |owned| {
276 self.replaceCheckpoint(owned);
277 owned_checkpoint = null;
278 }
279 return replay_timepoint orelse expected.timepoint;
280 }
281
282 const PreparedReplay = struct {
283 source: event.Source,
284 actual: *const event.Event,
285 next_cursor: u64,
286 };
287
288 fn prepareReplay(self: *Session, expected: event.Event) !PreparedReplay {
289 const source = self.source.?;
290 const actual = (try source.peek()) orelse return error.MissingReplayEvent;
291 if (!expected.eqlForReplay(actual.*)) return error.ReplayEventMismatch;
292 return .{
293 .source = source,
294 .actual = actual,
295 .next_cursor = try self.prepareReplayCursor(),
296 };
297 }
298
299 fn prepareReplayCursor(self: *const Session) !u64 {
300 return std.math.add(u64, self.replay_cursor, 1) catch
301 return error.ReplayCursorExhausted;
302 }
303
304 fn commitReplay(self: *Session, prepared: PreparedReplay) void {
305 prepared.source.advance();
306 self.replay_cursor = prepared.next_cursor;
307 }
308
309 fn replaceCheckpoint(self: *Session, owned: event.Event) void {
310 if (self.last_checkpoint) |*previous| previous.deinit(self.allocator);
311 self.last_checkpoint = owned;
312 }
313 };
314
315 const SliceStream = struct {
316 events: []const event.Event,
317 cursor: usize = 0,
318
319 fn source(self: *SliceStream) event.Source {
320 return .{ .context = self, .peekFn = peek, .advanceFn = advance, .countFn = count };
321 }
322
323 fn peek(context: *anyopaque) !?*const event.Event {
324 const self: *SliceStream = @ptrCast(@alignCast(context));
325 if (self.cursor == self.events.len) return null;
326 return &self.events[self.cursor];
327 }
328
329 fn advance(context: *anyopaque) void {
330 const self: *SliceStream = @ptrCast(@alignCast(context));
331 self.cursor += 1;
332 }
333
334 fn count(context: *anyopaque) u64 {
335 const self: *SliceStream = @ptrCast(@alignCast(context));
336 return self.events.len;
337 }
338 };
339
340 const CaptureSink = struct {
341 timepoints: [8]event.Timepoint = undefined,
342 len: usize = 0,
343
344 fn sink(self: *CaptureSink) event.Sink {
345 return .{ .context = self, .appendFn = append };
346 }
347
348 fn append(context: *anyopaque, item: event.Event) !void {
349 const self: *CaptureSink = @ptrCast(@alignCast(context));
350 if (self.len == self.timepoints.len) return error.TestSinkFull;
351 self.timepoints[self.len] = item.timepoint;
352 self.len += 1;
353 }
354 };
355
356 const FailOnceSink = struct {
357 timepoints: [8]event.Timepoint = undefined,
358 len: usize = 0,
359 fail_next: bool = true,
360
361 fn sink(self: *FailOnceSink) event.Sink {
362 return .{ .context = self, .appendFn = append };
363 }
364
365 fn append(context: *anyopaque, item: event.Event) !void {
366 const self: *FailOnceSink = @ptrCast(@alignCast(context));
367 if (self.fail_next) {
368 self.fail_next = false;
369 return error.InjectedSinkFailure;
370 }
371 if (self.len == self.timepoints.len) return error.TestSinkFull;
372 self.timepoints[self.len] = item.timepoint;
373 self.len += 1;
374 }
375 };
376
377 const FailOnceStream = struct {
378 events: []const event.Event,
379 cursor: usize = 0,
380 fail_next: bool = true,
381
382 fn source(self: *FailOnceStream) event.Source {
383 return .{ .context = self, .peekFn = peek, .advanceFn = advance, .countFn = count };
384 }
385
386 fn peek(context: *anyopaque) !?*const event.Event {
387 const self: *FailOnceStream = @ptrCast(@alignCast(context));
388 if (self.fail_next) {
389 self.fail_next = false;
390 return error.InjectedSourceFailure;
391 }
392 if (self.cursor == self.events.len) return null;
393 return &self.events[self.cursor];
394 }
395
396 fn advance(context: *anyopaque) void {
397 const self: *FailOnceStream = @ptrCast(@alignCast(context));
398 self.cursor += 1;
399 }
400
401 fn count(context: *anyopaque) u64 {
402 const self: *FailOnceStream = @ptrCast(@alignCast(context));
403 return self.events.len;
404 }
405 };
406
407 const ProviderProbe = struct {
408 value: [32]u8 = undefined,
409 value_len: usize = 0,
410 staged: [32]u8 = undefined,
411 staged_len: usize = 0,
412 prepared: bool = false,
413 prepare_calls: u64 = 0,
414 commit_calls: u64 = 0,
415 cancel_calls: u64 = 0,
416 fail_prepare: bool = false,
417
418 fn init(bytes: []const u8) ProviderProbe {
419 var self: ProviderProbe = .{};
420 self.setValue(bytes);
421 return self;
422 }
423
424 fn setValue(self: *ProviderProbe, bytes: []const u8) void {
425 std.debug.assert(bytes.len <= self.value.len);
426 @memcpy(self.value[0..bytes.len], bytes);
427 self.value_len = bytes.len;
428 }
429
430 fn current(self: *const ProviderProbe) []const u8 {
431 return self.value[0..self.value_len];
432 }
433
434 fn provider(self: *ProviderProbe) checkpoint.Provider {
435 return .{
436 .context = self,
437 .capture = capture,
438 .prepare_restore = prepareRestore,
439 .commit_restore = commitRestore,
440 .cancel_restore = cancelRestore,
441 };
442 }
443
444 fn capture(context: *anyopaque, allocator: Allocator) ![]u8 {
445 const self: *ProviderProbe = @ptrCast(@alignCast(context));
446 return try allocator.dupe(u8, self.current());
447 }
448
449 fn prepareRestore(context: *anyopaque, bytes: []const u8) !void {
450 const self: *ProviderProbe = @ptrCast(@alignCast(context));
451 self.prepare_calls += 1;
452 if (self.fail_prepare) return error.InjectedProviderFailure;
453 if (bytes.len > self.staged.len) return error.TestProviderCapacityExceeded;
454 @memcpy(self.staged[0..bytes.len], bytes);
455 self.staged_len = bytes.len;
456 self.prepared = true;
457 }
458
459 fn commitRestore(context: *anyopaque) void {
460 const self: *ProviderProbe = @ptrCast(@alignCast(context));
461 std.debug.assert(self.prepared);
462 @memcpy(self.value[0..self.staged_len], self.staged[0..self.staged_len]);
463 self.value_len = self.staged_len;
464 self.prepared = false;
465 self.commit_calls += 1;
466 }
467
468 fn cancelRestore(context: *anyopaque) void {
469 const self: *ProviderProbe = @ptrCast(@alignCast(context));
470 self.prepared = false;
471 self.cancel_calls += 1;
472 }
473 };
474
475 test "recording streams events without retaining copies" {
476 var capture: CaptureSink = .{};
477 var session = Session.initRecord(std.testing.allocator, capture.sink());
478 defer session.deinit();
479 const start = try session.sessionStart(1, "test");
480 const enter = try session.functionEnter(1, .{ .function_id = 10, .site_id = 1 });
481 const exit = try session.functionExit(1, .{ .function_id = 10, .site_id = 2 });
482 try std.testing.expectEqual(@as(u64, 1), start.seq);
483 try std.testing.expectEqual(@as(u64, 2), enter.seq);
484 try std.testing.expectEqual(@as(u64, 3), exit.seq);
485 try std.testing.expectEqual(@as(usize, 3), capture.len);
486 }
487
488 test "session failure atomicity retries a failed record append without a sequence gap" {
489 var capture: FailOnceSink = .{};
490 var session = Session.initRecord(std.testing.allocator, capture.sink());
491 defer session.deinit();
492
493 try std.testing.expectError(
494 error.InjectedSinkFailure,
495 session.sessionStart(1, "test"),
496 );
497 try std.testing.expectEqual(@as(usize, 0), capture.len);
498
499 const retried = try session.sessionStart(1, "test");
500 try std.testing.expectEqual(@as(u64, 1), retried.seq);
501 try std.testing.expectEqual(@as(usize, 1), capture.len);
502 try std.testing.expectEqual(@as(u64, 1), capture.timepoints[0].seq);
503 }
504
505 test "session failure atomicity retries a failed source peek" {
506 const events = [_]event.Event{
507 event.Event.sessionStart(.{ .thread_id = 1, .seq = 1 }, "test"),
508 };
509 var stream = FailOnceStream{ .events = &events };
510 var session = Session.initReplay(std.testing.allocator, stream.source());
511 defer session.deinit();
512
513 try std.testing.expectError(
514 error.InjectedSourceFailure,
515 session.sessionStart(1, "test"),
516 );
517 try std.testing.expectEqual(@as(usize, 0), stream.cursor);
518 try std.testing.expectEqual(@as(u64, 0), session.replayProgress().cursor);
519
520 const retried = try session.sessionStart(1, "test");
521 try std.testing.expectEqual(@as(u64, 1), retried.seq);
522 try session.verifyReplayComplete();
523 }
524
525 test "replay validates a borrowed stream and exposes numeric progress" {
526 const events = [_]event.Event{
527 event.Event.sessionStart(.{ .thread_id = 1, .seq = 1 }, "test"),
528 event.Event.safepointReached(.{ .thread_id = 1, .seq = 2 }, .{ .function_id = 20, .site_id = 5 }),
529 };
530 var stream = SliceStream{ .events = &events };
531 var session = Session.initReplay(std.testing.allocator, stream.source());
532 defer session.deinit();
533 var progress = session.replayProgress();
534 try std.testing.expectEqual(@as(u64, 0), progress.cursor);
535 try std.testing.expectEqual(@as(u64, 2), progress.remaining_count);
536 _ = try session.sessionStart(1, "test");
537 progress = session.replayProgress();
538 try std.testing.expectEqual(@as(u64, 1), progress.cursor);
539 try std.testing.expectEqual(@as(u64, 1), progress.remaining_count);
540 _ = try session.safepoint(1, .{ .function_id = 20, .site_id = 5 });
541 try session.verifyReplayComplete();
542 try std.testing.expect(session.replayProgress().complete());
543 }
544
545 test "session failure atomicity allocates a record boundary before append" {
546 var failing = std.testing.FailingAllocator.init(std.testing.allocator, .{
547 .fail_index = 0,
548 });
549 var capture: CaptureSink = .{};
550 var session = Session.initRecord(failing.allocator(), capture.sink());
551 defer session.deinit();
552
553 try std.testing.expectError(
554 error.OutOfMemory,
555 session.boundaryBytesAlloc(1, "env.TEST", "value"),
556 );
557 try std.testing.expectEqual(@as(usize, 0), capture.len);
558
559 failing.fail_index = std.math.maxInt(usize);
560 failing.resize_fail_index = std.math.maxInt(usize);
561 const retried = try session.boundaryBytesAlloc(1, "env.TEST", "value");
562 defer failing.allocator().free(retried);
563 try std.testing.expectEqualStrings("value", retried);
564 try std.testing.expectEqual(@as(usize, 1), capture.len);
565 try std.testing.expectEqual(@as(u64, 1), capture.timepoints[0].seq);
566 }
567
568 test "session failure atomicity retries boundary replay after allocation failure" {
569 const events = [_]event.Event{
570 event.Event.boundaryBytes(.{ .thread_id = 1, .seq = 1 }, "env.TEST", "value"),
571 };
572 var stream = SliceStream{ .events = &events };
573 var failing = std.testing.FailingAllocator.init(std.testing.allocator, .{
574 .fail_index = 0,
575 });
576 var session = Session.initReplay(failing.allocator(), stream.source());
577 defer session.deinit();
578
579 try std.testing.expectError(
580 error.OutOfMemory,
581 session.boundaryBytesAlloc(1, "env.TEST", null),
582 );
583 try std.testing.expectEqual(@as(usize, 0), stream.cursor);
584 try std.testing.expectEqual(@as(u64, 0), session.replayProgress().cursor);
585
586 failing.fail_index = std.math.maxInt(usize);
587 failing.resize_fail_index = std.math.maxInt(usize);
588 const retried = try session.boundaryBytesAlloc(1, "env.TEST", null);
589 defer failing.allocator().free(retried);
590 try std.testing.expectEqualStrings("value", retried);
591 try session.verifyReplayComplete();
592 }
593
594 test "session failure atomicity rejects noncanonical boundary fields" {
595 var events = [_]event.Event{
596 event.Event.boundaryBytes(.{ .thread_id = 1, .seq = 1 }, "env.TEST", "value"),
597 };
598 events[0].label = "unexpected";
599 var stream = SliceStream{ .events = &events };
600 var session = Session.initReplay(std.testing.allocator, stream.source());
601 defer session.deinit();
602
603 try std.testing.expectError(
604 error.ReplayEventMismatch,
605 session.boundaryBytesAlloc(1, "env.TEST", null),
606 );
607 try std.testing.expectEqual(@as(usize, 0), stream.cursor);
608 try std.testing.expectEqual(@as(u64, 0), session.replayProgress().cursor);
609
610 events[0].label = null;
611 const retried = try session.boundaryBytesAlloc(1, "env.TEST", null);
612 defer std.testing.allocator.free(retried);
613 try std.testing.expectEqualStrings("value", retried);
614 try session.verifyReplayComplete();
615 }
616
617 test "session failure atomicity requires a live boundary value while off" {
618 var session = Session.initOff(std.testing.allocator);
619 defer session.deinit();
620
621 try std.testing.expectError(
622 error.MissingBoundaryValue,
623 session.boundaryBytesAlloc(1, "env.TEST", null),
624 );
625 const retried = try session.boundaryBytesAlloc(1, "env.TEST", "value");
626 defer std.testing.allocator.free(retried);
627 try std.testing.expectEqualStrings("value", retried);
628 const next = try session.sessionEnd(1, 0);
629 try std.testing.expectEqual(@as(u64, 2), next.seq);
630 }
631
632 test "replay mismatch leaves the session retryable" {
633 const events = [_]event.Event{
634 event.Event.safepointReached(.{ .thread_id = 1, .seq = 1 }, .{ .function_id = 20, .site_id = 5 }),
635 };
636 var stream = SliceStream{ .events = &events };
637 var session = Session.initReplay(std.testing.allocator, stream.source());
638 defer session.deinit();
639 try std.testing.expectError(
640 error.ReplayEventMismatch,
641 session.safepoint(1, .{ .function_id = 21, .site_id = 5 }),
642 );
643 try std.testing.expectEqual(@as(usize, 0), stream.cursor);
644 try std.testing.expectEqual(@as(u64, 0), session.replayProgress().cursor);
645 const retried = try session.safepoint(1, .{ .function_id = 20, .site_id = 5 });
646 try std.testing.expectEqual(@as(u64, 1), retried.seq);
647 try session.verifyReplayComplete();
648 }
649
650 test "boundary replay copies borrowed bytes before advancing" {
651 const events = [_]event.Event{
652 event.Event.boundaryBytes(.{ .thread_id = 1, .seq = 1 }, "env.TEST", "hello"),
653 };
654 var stream = SliceStream{ .events = &events };
655 var session = Session.initReplay(std.testing.allocator, stream.source());
656 defer session.deinit();
657 const replayed = try session.boundaryBytesAlloc(1, "env.TEST", null);
658 defer std.testing.allocator.free(replayed);
659 try std.testing.expectEqualStrings("hello", replayed);
660 try session.verifyReplayComplete();
661 }
662
663 test "session failure atomicity validates restore before staging provider state" {
664 const events = [_]event.Event{
665 event.Event.checkpoint(.{ .thread_id = 1, .seq = 1 }, "saved", "checkpoint"),
666 event.Event.checkpointRestore(.{ .thread_id = 1, .seq = 2 }, "restore"),
667 };
668 var stream = SliceStream{ .events = &events };
669 var provider_probe = ProviderProbe.init("checkpoint");
670 var session = Session.initReplay(std.testing.allocator, stream.source());
671 defer session.deinit();
672
673 _ = try session.checkpointCapture(1, "saved", provider_probe.provider());
674 provider_probe.setValue("live");
675
676 try std.testing.expectError(
677 error.ReplayEventMismatch,
678 session.checkpointRestore(1, "wrong", provider_probe.provider()),
679 );
680 try std.testing.expectEqualStrings("live", provider_probe.current());
681 try std.testing.expectEqual(@as(u64, 0), provider_probe.prepare_calls);
682 try std.testing.expectEqual(@as(usize, 1), stream.cursor);
683 try std.testing.expectEqual(@as(u64, 1), session.replayProgress().cursor);
684
685 provider_probe.fail_prepare = true;
686 try std.testing.expectError(
687 error.InjectedProviderFailure,
688 session.checkpointRestore(1, "restore", provider_probe.provider()),
689 );
690 try std.testing.expectEqualStrings("live", provider_probe.current());
691 try std.testing.expect(!provider_probe.prepared);
692 try std.testing.expectEqual(@as(u64, 1), provider_probe.cancel_calls);
693 try std.testing.expectEqual(@as(usize, 1), stream.cursor);
694 try std.testing.expectEqual(@as(u64, 1), session.replayProgress().cursor);
695
696 provider_probe.fail_prepare = false;
697 const restored = try session.checkpointRestore(1, "restore", provider_probe.provider());
698 try std.testing.expectEqual(@as(u64, 2), restored.seq);
699 try std.testing.expectEqualStrings("checkpoint", provider_probe.current());
700 try std.testing.expectEqual(@as(u64, 2), provider_probe.prepare_calls);
701 try std.testing.expectEqual(@as(u64, 1), provider_probe.commit_calls);
702 try session.verifyReplayComplete();
703 }
704
705 test "session failure atomicity preserves a checkpoint when replacement allocation fails" {
706 const events = [_]event.Event{
707 event.Event.checkpoint(.{ .thread_id = 1, .seq = 1 }, "first", "one"),
708 event.Event.checkpoint(.{ .thread_id = 1, .seq = 2 }, "second", "two"),
709 };
710 var stream = SliceStream{ .events = &events };
711 var failing = std.testing.FailingAllocator.init(std.testing.allocator, .{});
712 var provider_probe = ProviderProbe.init("one");
713 var session = Session.initReplay(failing.allocator(), stream.source());
714 defer session.deinit();
715
716 _ = try session.checkpointCapture(1, "first", provider_probe.provider());
717 try std.testing.expectEqualStrings("one", session.last_checkpoint.?.data.?);
718 provider_probe.setValue("two");
719 failing.fail_index = failing.alloc_index + 1;
720
721 try std.testing.expectError(
722 error.OutOfMemory,
723 session.checkpointCapture(1, "second", provider_probe.provider()),
724 );
725 try std.testing.expectEqualStrings("one", session.last_checkpoint.?.data.?);
726 try std.testing.expectEqual(@as(usize, 1), stream.cursor);
727 try std.testing.expectEqual(@as(u64, 1), session.replayProgress().cursor);
728
729 failing.fail_index = std.math.maxInt(usize);
730 failing.resize_fail_index = std.math.maxInt(usize);
731 const retried = try session.checkpointCapture(1, "second", provider_probe.provider());
732 try std.testing.expectEqual(@as(u64, 2), retried.seq);
733 try std.testing.expectEqualStrings("two", session.last_checkpoint.?.data.?);
734 try session.verifyReplayComplete();
735 }
736
737 test "session failure atomicity cancels a staged restore when recording fails" {
738 var capture: FailOnceSink = .{ .fail_next = false };
739 var provider_probe = ProviderProbe.init("checkpoint");
740 var session = Session.initRecord(std.testing.allocator, capture.sink());
741 defer session.deinit();
742
743 _ = try session.checkpointCapture(1, "saved", provider_probe.provider());
744 provider_probe.setValue("live");
745 capture.fail_next = true;
746
747 try std.testing.expectError(
748 error.InjectedSinkFailure,
749 session.checkpointRestore(1, "restore", provider_probe.provider()),
750 );
751 try std.testing.expectEqualStrings("live", provider_probe.current());
752 try std.testing.expect(!provider_probe.prepared);
753 try std.testing.expectEqual(@as(u64, 1), provider_probe.cancel_calls);
754 try std.testing.expectEqual(@as(usize, 1), capture.len);
755
756 const restored = try session.checkpointRestore(1, "restore", provider_probe.provider());
757 try std.testing.expectEqual(@as(u64, 2), restored.seq);
758 try std.testing.expectEqualStrings("checkpoint", provider_probe.current());
759 try std.testing.expectEqual(@as(usize, 2), capture.len);
760 try std.testing.expectEqual(@as(u64, 2), capture.timepoints[1].seq);
761 }
762
763 test "checkpoint provider owns committed restore after checkpoint replacement" {
764 var provider_probe = ProviderProbe.init("one");
765 {
766 var session = Session.initOff(std.testing.allocator);
767 defer session.deinit();
768
769 _ = try session.checkpointCapture(1, "first", provider_probe.provider());
770 provider_probe.setValue("live");
771 _ = try session.checkpointRestore(1, "restore", provider_probe.provider());
772 try std.testing.expectEqualStrings("one", provider_probe.current());
773 _ = try session.checkpointCapture(1, "second", provider_probe.provider());
774 try std.testing.expectEqualStrings("one", provider_probe.current());
775 }
776 try std.testing.expectEqualStrings("one", provider_probe.current());
777 }
778
779 test "session sequence capacity admits the final timepoint exactly once" {
780 var session = Session.initOff(std.testing.allocator);
781 defer session.deinit();
782 session.next_seq = std.math.maxInt(u64);
783
784 const final = try session.sessionStart(1, "last");
785 try std.testing.expectEqual(std.math.maxInt(u64), final.seq);
786 try std.testing.expectError(
787 error.SequenceExhausted,
788 session.sessionEnd(1, 0),
789 );
790 }