| 1 | package turnevent |
| 2 | |
| 3 | import ( |
| 4 | "bytes" |
| 5 | "encoding/json" |
| 6 | "errors" |
| 7 | "fmt" |
| 8 | "os" |
| 9 | "path/filepath" |
| 10 | "testing" |
| 11 | |
| 12 | "reasonix/internal/event" |
| 13 | "reasonix/internal/eventwire" |
| 14 | "reasonix/internal/store" |
| 15 | ) |
| 16 | |
| 17 | func testSessionPath(t *testing.T) string { |
| 18 | t.Helper() |
| 19 | return filepath.Join(t.TempDir(), "session.jsonl") |
| 20 | } |
| 21 | |
| 22 | func openTestLedger(t *testing.T, sessionPath, sessionID string) *Ledger { |
| 23 | t.Helper() |
| 24 | ledger, err := Open(sessionPath, sessionID) |
| 25 | if err != nil { |
| 26 | t.Fatalf("Open: %v", err) |
| 27 | } |
| 28 | t.Cleanup(func() { |
| 29 | if err := ledger.Close(); err != nil { |
| 30 | t.Errorf("Close test ledger: %v", err) |
| 31 | } |
| 32 | }) |
| 33 | return ledger |
| 34 | } |
| 35 | |
| 36 | func TestLedgerPersistsMonotonicEventsAndExactlyOneTerminal(t *testing.T) { |
| 37 | path := testSessionPath(t) |
| 38 | l, err := Open(path, "session") |
| 39 | if err != nil { |
| 40 | t.Fatalf("Open: %v", err) |
| 41 | } |
| 42 | l.SetRoutingMetadata("epoch-1", "submission-1") |
| 43 | turnID, err := l.Begin() |
| 44 | if err != nil { |
| 45 | t.Fatalf("Begin: %v", err) |
| 46 | } |
| 47 | l.SetTranscriptSnapshot(7, "digest-1") |
| 48 | started, ok, err := l.Append(event.Event{Kind: event.TurnStarted}, event.TurnInProgress) |
| 49 | if err != nil || !ok { |
| 50 | t.Fatalf("append started: ok=%v err=%v", ok, err) |
| 51 | } |
| 52 | done, ok, err := l.Append(event.Event{Kind: event.TurnDone}, event.TurnCompleted) |
| 53 | if err != nil || !ok { |
| 54 | t.Fatalf("append done: ok=%v err=%v", ok, err) |
| 55 | } |
| 56 | if started.TurnID != turnID || done.TurnID != turnID || started.Sequence != 1 || done.Sequence != 2 { |
| 57 | t.Fatalf("stamps = started(%q,%d) done(%q,%d), want turn %q seq 1,2", started.TurnID, started.Sequence, done.TurnID, done.Sequence, turnID) |
| 58 | } |
| 59 | if _, ok, err := l.Append(event.Event{Kind: event.TurnDone}, event.TurnFailed); err != nil || ok { |
| 60 | t.Fatalf("second terminal: ok=%v err=%v, want compare-and-append rejection", ok, err) |
| 61 | } |
| 62 | recs, err := l.EventsAfter(0) |
| 63 | if err != nil { |
| 64 | t.Fatalf("EventsAfter: %v", err) |
| 65 | } |
| 66 | if len(recs) != 2 || recs[0].Sequence != 1 || recs[1].Sequence != 2 { |
| 67 | t.Fatalf("records = %#v, want two monotonic events", recs) |
| 68 | } |
| 69 | if recs[1].RuntimeEpoch != "epoch-1" || recs[1].SubmissionID != "submission-1" || recs[1].TranscriptRevision != 7 || recs[1].TranscriptDigest != "digest-1" { |
| 70 | t.Fatalf("terminal metadata = %+v, want routing and transcript identity", recs[1]) |
| 71 | } |
| 72 | if latest, replayAfter := l.ProjectionCursor(); latest != 2 || replayAfter != 2 { |
| 73 | t.Fatalf("terminal cursor = (%d,%d), want (2,2)", latest, replayAfter) |
| 74 | } |
| 75 | } |
| 76 | |
| 77 | func TestLedgerProjectionCursorReplaysOnlyActiveTurn(t *testing.T) { |
| 78 | path := testSessionPath(t) |
| 79 | l := openTestLedger(t, path, "session") |
| 80 | if _, err := l.Begin(); err != nil { |
| 81 | t.Fatalf("Begin first: %v", err) |
| 82 | } |
| 83 | if _, ok, err := l.Append(event.Event{Kind: event.TurnDone}, event.TurnCompleted); err != nil || !ok { |
| 84 | t.Fatalf("complete first: ok=%v err=%v", ok, err) |
| 85 | } |
| 86 | if _, err := l.Begin(); err != nil { |
| 87 | t.Fatalf("Begin second: %v", err) |
| 88 | } |
| 89 | if _, ok, err := l.Append(event.Event{Kind: event.TurnStatusChanged}, event.TurnQueued); err != nil || !ok { |
| 90 | t.Fatalf("queue second: ok=%v err=%v", ok, err) |
| 91 | } |
| 92 | if _, ok, err := l.Append(event.Event{Kind: event.TurnStarted}, event.TurnInProgress); err != nil || !ok { |
| 93 | t.Fatalf("start second: ok=%v err=%v", ok, err) |
| 94 | } |
| 95 | if latest, replayAfter := l.ProjectionCursor(); latest != 3 || replayAfter != 1 { |
| 96 | t.Fatalf("active cursor = (%d,%d), want (3,1)", latest, replayAfter) |
| 97 | } |
| 98 | } |
| 99 | |
| 100 | func TestLedgerSubmissionReceiptSurvivesCompletionAndIsOneShot(t *testing.T) { |
| 101 | path := testSessionPath(t) |
| 102 | l := openTestLedger(t, path, "session") |
| 103 | l.SetRoutingMetadata("epoch-1", "submission-1") |
| 104 | first, err := l.Begin() |
| 105 | if err != nil { |
| 106 | t.Fatalf("Begin first: %v", err) |
| 107 | } |
| 108 | if _, ok, err := l.Append(event.Event{Kind: event.TurnDone}, event.TurnCompleted); err != nil || !ok { |
| 109 | t.Fatalf("complete first: ok=%v err=%v", ok, err) |
| 110 | } |
| 111 | if got := l.TurnIDForSubmission("submission-1"); got != first { |
| 112 | t.Fatalf("receipt after completion = %q, want %q", got, first) |
| 113 | } |
| 114 | second, err := l.Begin() |
| 115 | if err != nil { |
| 116 | t.Fatalf("Begin automatic follow-up: %v", err) |
| 117 | } |
| 118 | if second == first { |
| 119 | t.Fatal("follow-up reused the completed turn id") |
| 120 | } |
| 121 | if _, ok, err := l.Append(event.Event{Kind: event.TurnStarted}, event.TurnInProgress); err != nil || !ok { |
| 122 | t.Fatalf("start follow-up: ok=%v err=%v", ok, err) |
| 123 | } |
| 124 | records, err := l.EventsAfter(0) |
| 125 | if err != nil { |
| 126 | t.Fatalf("EventsAfter: %v", err) |
| 127 | } |
| 128 | last := records[len(records)-1] |
| 129 | if last.SubmissionID != "" || last.Event.SubmissionID != "" { |
| 130 | t.Fatalf("automatic follow-up inherited the previous submission: %+v", last) |
| 131 | } |
| 132 | if last.RuntimeEpoch != "epoch-1" || last.Event.RuntimeEpoch != "epoch-1" { |
| 133 | t.Fatalf("automatic follow-up lost its owning runtime epoch: %+v", last) |
| 134 | } |
| 135 | if got := l.TurnIDForSubmission("submission-1"); got != first { |
| 136 | t.Fatalf("receipt after replacement = %q, want original %q", got, first) |
| 137 | } |
| 138 | } |
| 139 | |
| 140 | func TestLedgerRejectsStatusRegressionAndKeepsCancellationSticky(t *testing.T) { |
| 141 | l := openTestLedger(t, testSessionPath(t), "session") |
| 142 | if _, err := l.Begin(); err != nil { |
| 143 | t.Fatalf("Begin: %v", err) |
| 144 | } |
| 145 | if _, ok, err := l.Append(event.Event{Kind: event.TurnStarted}, event.TurnInProgress); err != nil || !ok { |
| 146 | t.Fatalf("start: ok=%v err=%v", ok, err) |
| 147 | } |
| 148 | if _, ok, err := l.Append(event.Event{Kind: event.TurnStatusChanged}, event.TurnQueued); err == nil || ok { |
| 149 | t.Fatalf("status regression: ok=%v err=%v, want rejection", ok, err) |
| 150 | } |
| 151 | if _, ok, err := l.Append(event.Event{Kind: event.TurnStatusChanged}, event.TurnCancelling); err != nil || !ok { |
| 152 | t.Fatalf("cancel: ok=%v err=%v", ok, err) |
| 153 | } |
| 154 | stamped, ok, err := l.Append(event.Event{Kind: event.PromptAnswered}, event.TurnInProgress) |
| 155 | if err != nil || !ok || stamped.Status != event.TurnCancelling || l.CurrentStatus() != event.TurnCancelling { |
| 156 | t.Fatalf("late prompt answer = (%+v,%v,%v), current=%q, want sticky cancelling", stamped, ok, err, l.CurrentStatus()) |
| 157 | } |
| 158 | } |
| 159 | |
| 160 | func TestLedgerRepairsTornTailAndKeepsValidPrefix(t *testing.T) { |
| 161 | path := testSessionPath(t) |
| 162 | l := openTestLedger(t, path, "session") |
| 163 | if _, err := l.Begin(); err != nil { |
| 164 | t.Fatalf("Begin: %v", err) |
| 165 | } |
| 166 | if _, ok, err := l.Append(event.Event{Kind: event.TurnStarted}, event.TurnInProgress); err != nil || !ok { |
| 167 | t.Fatalf("append: ok=%v err=%v", ok, err) |
| 168 | } |
| 169 | ledgerPath := store.SessionTurnEventLog(path) |
| 170 | f, err := os.OpenFile(ledgerPath, os.O_WRONLY|os.O_APPEND, 0o600) |
| 171 | if err != nil { |
| 172 | t.Fatalf("open ledger tail: %v", err) |
| 173 | } |
| 174 | if _, err := f.WriteString(`{"schemaVersion":1,"seq":2`); err != nil { |
| 175 | t.Fatalf("write torn tail: %v", err) |
| 176 | } |
| 177 | _ = f.Close() |
| 178 | |
| 179 | reopened := openTestLedger(t, path, "session") |
| 180 | recs, err := reopened.EventsAfter(0) |
| 181 | if err != nil { |
| 182 | t.Fatalf("EventsAfter: %v", err) |
| 183 | } |
| 184 | if len(recs) != 2 || recs[0].Status != event.TurnInProgress || recs[1].Status != event.TurnInterrupted { |
| 185 | t.Fatalf("recovered records = %#v, want valid start plus interrupted terminal", recs) |
| 186 | } |
| 187 | damaged, err := os.ReadFile(store.SessionTurnEventLogDamaged(path)) |
| 188 | if err != nil { |
| 189 | t.Fatalf("read damaged tail: %v", err) |
| 190 | } |
| 191 | if len(damaged) == 0 { |
| 192 | t.Fatal("damaged tail was not isolated") |
| 193 | } |
| 194 | } |
| 195 | |
| 196 | func TestLedgerIsolatesNonMonotonicTail(t *testing.T) { |
| 197 | path := testSessionPath(t) |
| 198 | l := openTestLedger(t, path, "session") |
| 199 | if _, err := l.Begin(); err != nil { |
| 200 | t.Fatalf("Begin: %v", err) |
| 201 | } |
| 202 | if _, ok, err := l.Append(event.Event{Kind: event.TurnStarted}, event.TurnInProgress); err != nil || !ok { |
| 203 | t.Fatalf("append: ok=%v err=%v", ok, err) |
| 204 | } |
| 205 | ledgerPath := store.SessionTurnEventLog(path) |
| 206 | f, err := os.OpenFile(ledgerPath, os.O_WRONLY|os.O_APPEND, 0o600) |
| 207 | if err != nil { |
| 208 | t.Fatalf("open ledger: %v", err) |
| 209 | } |
| 210 | if _, err := f.WriteString(`{"schemaVersion":1,"sessionId":"session","turnId":"duplicate","seq":1,"kind":"turn_started","status":"in_progress","createdAt":1,"event":{"kind":"turn_started"}}` + "\n"); err != nil { |
| 211 | t.Fatalf("append duplicate sequence: %v", err) |
| 212 | } |
| 213 | _ = f.Close() |
| 214 | |
| 215 | reopened := openTestLedger(t, path, "session") |
| 216 | recs, err := reopened.EventsAfter(0) |
| 217 | if err != nil { |
| 218 | t.Fatalf("EventsAfter: %v", err) |
| 219 | } |
| 220 | if len(recs) != 2 || recs[0].Sequence != 1 || recs[1].Status != event.TurnInterrupted { |
| 221 | t.Fatalf("records = %#v, want valid seq 1 plus interrupted recovery", recs) |
| 222 | } |
| 223 | damaged, err := os.ReadFile(store.SessionTurnEventLogDamaged(path)) |
| 224 | if err != nil || len(damaged) == 0 { |
| 225 | t.Fatalf("non-monotonic tail was not isolated: bytes=%d err=%v", len(damaged), err) |
| 226 | } |
| 227 | } |
| 228 | |
| 229 | func TestLedgerRecoveryClosesRunningToolsWithoutReplay(t *testing.T) { |
| 230 | path := testSessionPath(t) |
| 231 | l := openTestLedger(t, path, "session") |
| 232 | if _, err := l.Begin(); err != nil { |
| 233 | t.Fatalf("Begin: %v", err) |
| 234 | } |
| 235 | if _, ok, err := l.Append(event.Event{Kind: event.TurnStarted}, event.TurnInProgress); err != nil || !ok { |
| 236 | t.Fatalf("append start: ok=%v err=%v", ok, err) |
| 237 | } |
| 238 | if _, ok, err := l.Append(event.Event{Kind: event.ToolDispatch, Tool: event.Tool{ |
| 239 | ID: "write-1", Name: "write_file", Args: `{"path":"important.txt"}`, |
| 240 | }}, event.TurnInProgress); err != nil || !ok { |
| 241 | t.Fatalf("append dispatch: ok=%v err=%v", ok, err) |
| 242 | } |
| 243 | |
| 244 | reopened := openTestLedger(t, path, "session") |
| 245 | recs, err := reopened.EventsAfter(0) |
| 246 | if err != nil { |
| 247 | t.Fatalf("EventsAfter: %v", err) |
| 248 | } |
| 249 | if len(recs) != 4 { |
| 250 | t.Fatalf("records = %#v, want start, dispatch, synthetic result, terminal", recs) |
| 251 | } |
| 252 | result := recs[2] |
| 253 | if result.Kind != "tool_result" || result.Event.Tool == nil || result.Event.Tool.ID != "write-1" || result.Event.Tool.Err == "" { |
| 254 | t.Fatalf("synthetic result = %#v, want interrupted write-1", result) |
| 255 | } |
| 256 | if result.Event.Tool.Args != "" || result.Event.Tool.Output != "" { |
| 257 | t.Fatalf("synthetic result must not replay tool input/output: %#v", result.Event.Tool) |
| 258 | } |
| 259 | if recs[3].Kind != "turn_done" || recs[3].Status != event.TurnInterrupted || recs[3].Event.Recovery == nil || recs[3].Event.Recovery.State != "unknown" || recs[3].Event.Recovery.RequiresUserDecision || result.Event.Tool.RunState != "unknown" { |
| 260 | t.Fatalf("terminal = %#v, want an interrupted turn with a fact-only unknown write", recs[3]) |
| 261 | } |
| 262 | } |
| 263 | |
| 264 | func TestLedgerBootstrapsLegacyTranscriptWithoutRewritingIt(t *testing.T) { |
| 265 | path := testSessionPath(t) |
| 266 | original := []byte("legacy provider transcript\n") |
| 267 | if err := os.WriteFile(path, original, 0o600); err != nil { |
| 268 | t.Fatalf("write transcript: %v", err) |
| 269 | } |
| 270 | l, err := Open(path, "legacy") |
| 271 | if err != nil { |
| 272 | t.Fatalf("Open: %v", err) |
| 273 | } |
| 274 | recs, err := l.EventsAfter(0) |
| 275 | if err != nil { |
| 276 | t.Fatalf("EventsAfter: %v", err) |
| 277 | } |
| 278 | if len(recs) != 1 || recs[0].Kind != "turn_status" || recs[0].Status != event.TurnCompleted { |
| 279 | t.Fatalf("bootstrap = %#v, want one completed turn_status", recs) |
| 280 | } |
| 281 | got, err := os.ReadFile(path) |
| 282 | if err != nil { |
| 283 | t.Fatalf("read transcript: %v", err) |
| 284 | } |
| 285 | if string(got) != string(original) { |
| 286 | t.Fatalf("legacy transcript changed: %q", got) |
| 287 | } |
| 288 | if _, err := l.Begin(); err != nil { |
| 289 | t.Fatalf("Begin after bootstrap: %v", err) |
| 290 | } |
| 291 | } |
| 292 | |
| 293 | func TestLedgerReplayEmptyArrayAndUnknownFields(t *testing.T) { |
| 294 | path := testSessionPath(t) |
| 295 | l := openTestLedger(t, path, "session") |
| 296 | recs, err := l.EventsAfter(99) |
| 297 | if err != nil { |
| 298 | t.Fatalf("EventsAfter empty: %v", err) |
| 299 | } |
| 300 | if recs == nil || len(recs) != 0 { |
| 301 | t.Fatalf("empty replay = %#v, want non-nil empty slice", recs) |
| 302 | } |
| 303 | if _, err := l.Begin(); err != nil { |
| 304 | t.Fatalf("Begin: %v", err) |
| 305 | } |
| 306 | if _, ok, err := l.Append(event.Event{Kind: event.TurnStarted}, event.TurnInProgress); err != nil || !ok { |
| 307 | t.Fatalf("append: ok=%v err=%v", ok, err) |
| 308 | } |
| 309 | ledgerPath := store.SessionTurnEventLog(path) |
| 310 | data, err := os.ReadFile(ledgerPath) |
| 311 | if err != nil { |
| 312 | t.Fatalf("read ledger: %v", err) |
| 313 | } |
| 314 | data = append(data[:len(data)-2], []byte(`,"futureField":{"nested":true}}\n`)...) |
| 315 | if err := os.WriteFile(ledgerPath, data, 0o600); err != nil { |
| 316 | t.Fatalf("rewrite with unknown field: %v", err) |
| 317 | } |
| 318 | openTestLedger(t, path, "session") |
| 319 | } |
| 320 | |
| 321 | func TestLedgerPersistentHandleTerminalSyncAndProjectionAckOpenCounts(t *testing.T) { |
| 322 | l, err := Open(testSessionPath(t), "session") |
| 323 | if err != nil { |
| 324 | t.Fatalf("Open: %v", err) |
| 325 | } |
| 326 | l.RequireProjectionAck(true) |
| 327 | turnID, err := l.Begin() |
| 328 | if err != nil { |
| 329 | t.Fatalf("Begin: %v", err) |
| 330 | } |
| 331 | for _, e := range []event.Event{ |
| 332 | {Kind: event.TurnStarted}, |
| 333 | {Kind: event.Text, Text: "one"}, |
| 334 | {Kind: event.Text, Text: "two"}, |
| 335 | } { |
| 336 | if _, ok, appendErr := l.Append(e, event.TurnInProgress); appendErr != nil || !ok { |
| 337 | t.Fatalf("Append(%v): ok=%v err=%v", e.Kind, ok, appendErr) |
| 338 | } |
| 339 | } |
| 340 | if _, ok, err := l.Append(event.Event{Kind: event.TurnDone}, event.TurnCompleted); err != nil || !ok { |
| 341 | t.Fatalf("terminal: ok=%v err=%v", ok, err) |
| 342 | } |
| 343 | beforeAck := l.MetricsSnapshot() |
| 344 | if beforeAck.OpenCount != 1 || beforeAck.SyncCount != 1 || beforeAck.CloseCount != 1 { |
| 345 | t.Fatalf("WAL lifecycle = open:%d sync:%d close:%d, want 1/1/1", beforeAck.OpenCount, beforeAck.SyncCount, beforeAck.CloseCount) |
| 346 | } |
| 347 | if err := l.AcknowledgeProjection(turnID); err != nil { |
| 348 | t.Fatalf("AcknowledgeProjection: %v", err) |
| 349 | } |
| 350 | afterAck := l.MetricsSnapshot() |
| 351 | if afterAck.OpenCount != 2 || afterAck.SyncCount != 1 || afterAck.CloseCount != 2 { |
| 352 | t.Fatalf("WAL+ack lifecycle = open:%d sync:%d close:%d, want 2/1/2", afterAck.OpenCount, afterAck.SyncCount, afterAck.CloseCount) |
| 353 | } |
| 354 | } |
| 355 | |
| 356 | func TestLedgerCompactionWaitsForProjectionAckAndReplayResetsAtCheckpoint(t *testing.T) { |
| 357 | path := testSessionPath(t) |
| 358 | l, err := Open(path, "session") |
| 359 | if err != nil { |
| 360 | t.Fatalf("Open: %v", err) |
| 361 | } |
| 362 | l.RequireProjectionAck(true) |
| 363 | l.SetRoutingMetadata("epoch", "submission") |
| 364 | turnID, err := l.Begin() |
| 365 | if err != nil { |
| 366 | t.Fatalf("Begin: %v", err) |
| 367 | } |
| 368 | l.SetTranscriptSnapshot(9, "digest") |
| 369 | if _, ok, err := l.Append(event.Event{Kind: event.Text, Text: "display-only"}, event.TurnInProgress); err != nil || !ok { |
| 370 | t.Fatalf("text: ok=%v err=%v", ok, err) |
| 371 | } |
| 372 | done, ok, err := l.Append(event.Event{Kind: event.TurnDone, Outcome: "partial"}, event.TurnInterrupted) |
| 373 | if err != nil || !ok { |
| 374 | t.Fatalf("terminal: ok=%v err=%v", ok, err) |
| 375 | } |
| 376 | if err := l.Compact(); err != nil { |
| 377 | t.Fatalf("Compact before ack: %v", err) |
| 378 | } |
| 379 | if pending := l.PendingProjections(); len(pending) != 1 || pending[0].TurnID != turnID || len(pending[0].Events) != 2 { |
| 380 | t.Fatalf("pending projection = %+v, want retained terminal turn", pending) |
| 381 | } |
| 382 | if err := l.AcknowledgeProjection(turnID); err != nil { |
| 383 | t.Fatalf("AcknowledgeProjection: %v", err) |
| 384 | } |
| 385 | if err := l.Compact(); err != nil { |
| 386 | t.Fatalf("Compact after ack: %v", err) |
| 387 | } |
| 388 | if pending := l.PendingProjections(); len(pending) != 0 { |
| 389 | t.Fatalf("pending after checkpoint = %+v", pending) |
| 390 | } |
| 391 | reset, err := l.Replay(done.Sequence - 1) |
| 392 | if err != nil { |
| 393 | t.Fatalf("Replay reset: %v", err) |
| 394 | } |
| 395 | if !reset.ResetRequired || reset.FloorSequence != done.Sequence+1 || reset.LatestSequence != done.Sequence || reset.NextAfterSequence != done.Sequence || reset.Events == nil { |
| 396 | t.Fatalf("reset replay = %+v", reset) |
| 397 | } |
| 398 | current, err := l.Replay(done.Sequence) |
| 399 | if err != nil || current.ResetRequired || current.Events == nil || len(current.Events) != 0 { |
| 400 | t.Fatalf("checkpoint replay = %+v err=%v", current, err) |
| 401 | } |
| 402 | if got := l.TurnIDForSubmission("submission"); got != turnID { |
| 403 | t.Fatalf("submission receipt = %q, want %q", got, turnID) |
| 404 | } |
| 405 | data, err := os.ReadFile(store.SessionTurnEventLog(path)) |
| 406 | if err != nil { |
| 407 | t.Fatalf("read compacted ledger: %v", err) |
| 408 | } |
| 409 | if len(data) >= 64<<10 || !bytes.Contains(data, []byte(`"recordType":"checkpoint"`)) || bytes.Contains(data, []byte("display-only")) { |
| 410 | t.Fatalf("unsafe or oversized checkpoint: bytes=%d data=%s", len(data), data) |
| 411 | } |
| 412 | } |
| 413 | |
| 414 | func TestLedgerReplayPaginationAndOutOfRangeCursor(t *testing.T) { |
| 415 | l := openTestLedger(t, testSessionPath(t), "session") |
| 416 | if _, err := l.Begin(); err != nil { |
| 417 | t.Fatalf("Begin: %v", err) |
| 418 | } |
| 419 | for i := range replayMaxEvents + 25 { |
| 420 | if _, ok, err := l.Append(event.Event{Kind: event.ToolProgress, Tool: event.Tool{ID: "tool", Name: "bash", Output: "tick"}}, event.TurnInProgress); err != nil || !ok { |
| 421 | t.Fatalf("Append %d: ok=%v err=%v", i, ok, err) |
| 422 | } |
| 423 | } |
| 424 | first, err := l.Replay(0) |
| 425 | if err != nil || len(first.Events) != replayMaxEvents || !first.HasMore || first.NextAfterSequence != replayMaxEvents { |
| 426 | t.Fatalf("first page = events:%d next:%d more:%v err=%v", len(first.Events), first.NextAfterSequence, first.HasMore, err) |
| 427 | } |
| 428 | second, err := l.Replay(first.NextAfterSequence) |
| 429 | if err != nil || len(second.Events) != 25 || second.HasMore { |
| 430 | t.Fatalf("second page = events:%d more:%v err=%v", len(second.Events), second.HasMore, err) |
| 431 | } |
| 432 | outOfRange, err := l.Replay(second.LatestSequence + 1) |
| 433 | if err != nil || !outOfRange.ResetRequired { |
| 434 | t.Fatalf("out-of-range replay = %+v err=%v", outOfRange, err) |
| 435 | } |
| 436 | } |
| 437 | |
| 438 | func TestLedgerUnsupportedSchemaLeavesOriginalUntouched(t *testing.T) { |
| 439 | path := testSessionPath(t) |
| 440 | ledgerPath := store.SessionTurnEventLog(path) |
| 441 | if err := os.MkdirAll(filepath.Dir(ledgerPath), 0o700); err != nil { |
| 442 | t.Fatalf("MkdirAll: %v", err) |
| 443 | } |
| 444 | original := []byte(`{"schemaVersion":99,"recordType":"checkpoint"}` + "\n") |
| 445 | if err := os.WriteFile(ledgerPath, original, 0o600); err != nil { |
| 446 | t.Fatalf("write future ledger: %v", err) |
| 447 | } |
| 448 | _, err := Open(path, "session") |
| 449 | var unsupported *UnsupportedSchemaError |
| 450 | if !errors.As(err, &unsupported) || unsupported.Version != 99 { |
| 451 | t.Fatalf("Open error = %v, want unsupported schema 99", err) |
| 452 | } |
| 453 | got, readErr := os.ReadFile(ledgerPath) |
| 454 | if readErr != nil || !bytes.Equal(got, original) { |
| 455 | t.Fatalf("future ledger changed: bytes=%q err=%v", got, readErr) |
| 456 | } |
| 457 | if _, statErr := os.Stat(store.SessionTurnEventLogDamaged(path)); !errors.Is(statErr, os.ErrNotExist) { |
| 458 | t.Fatalf("future schema was isolated as damage: %v", statErr) |
| 459 | } |
| 460 | } |
| 461 | |
| 462 | func TestLedgerContinuesActiveV1TurnBeforeUpgradingToV2(t *testing.T) { |
| 463 | path := testSessionPath(t) |
| 464 | ledgerPath := store.SessionTurnEventLog(path) |
| 465 | v1 := Envelope{ |
| 466 | SchemaVersion: legacySchemaVersion, SessionID: "session", TurnID: "legacy-turn", |
| 467 | Sequence: 1, Kind: "turn_started", Status: event.TurnInProgress, CreatedAt: 1, |
| 468 | Event: eventwire.Event{Kind: "turn_started"}, |
| 469 | } |
| 470 | line, err := json.Marshal(v1) |
| 471 | if err != nil { |
| 472 | t.Fatalf("marshal v1: %v", err) |
| 473 | } |
| 474 | if err := os.WriteFile(ledgerPath, append(line, '\n'), 0o600); err != nil { |
| 475 | t.Fatalf("write v1: %v", err) |
| 476 | } |
| 477 | |
| 478 | l, err := Open(path, "session") |
| 479 | if err != nil { |
| 480 | t.Fatalf("Open: %v", err) |
| 481 | } |
| 482 | recovered, err := l.EventsAfter(0) |
| 483 | if err != nil || len(recovered) != 2 || recovered[1].Status != event.TurnInterrupted { |
| 484 | t.Fatalf("v1 recovery = %+v err=%v", recovered, err) |
| 485 | } |
| 486 | if _, err := l.Begin(); err != nil { |
| 487 | t.Fatalf("Begin v2 turn: %v", err) |
| 488 | } |
| 489 | if _, ok, err := l.Append(event.Event{Kind: event.TurnDone}, event.TurnCompleted); err != nil || !ok { |
| 490 | t.Fatalf("complete v2 turn: ok=%v err=%v", ok, err) |
| 491 | } |
| 492 | |
| 493 | data, err := os.ReadFile(ledgerPath) |
| 494 | if err != nil { |
| 495 | t.Fatalf("read upgraded ledger: %v", err) |
| 496 | } |
| 497 | lines := bytes.Split(bytes.TrimSpace(data), []byte{'\n'}) |
| 498 | if len(lines) != 3 || bytes.Contains(lines[0], []byte(`"recordType"`)) || bytes.Contains(lines[1], []byte(`"recordType"`)) || !bytes.Contains(lines[2], []byte(`"recordType":"event"`)) { |
| 499 | t.Fatalf("v1/v2 append boundary changed: %s", data) |
| 500 | } |
| 501 | } |
| 502 | |
| 503 | func TestLedgerCheckpointKeepsOnlyRecentTerminalSummariesAndReceipts(t *testing.T) { |
| 504 | path := testSessionPath(t) |
| 505 | l, err := Open(path, "session") |
| 506 | if err != nil { |
| 507 | t.Fatalf("Open: %v", err) |
| 508 | } |
| 509 | l.RequireProjectionAck(true) |
| 510 | turns := make([]string, 20) |
| 511 | for i := range turns { |
| 512 | submission := fmt.Sprintf("submission-%02d", i) |
| 513 | l.SetRoutingMetadata("epoch", submission) |
| 514 | turns[i], err = l.Begin() |
| 515 | if err != nil { |
| 516 | t.Fatalf("Begin %d: %v", i, err) |
| 517 | } |
| 518 | if _, ok, appendErr := l.Append(event.Event{Kind: event.TurnDone, Outcome: "completed"}, event.TurnCompleted); appendErr != nil || !ok { |
| 519 | t.Fatalf("complete %d: ok=%v err=%v", i, ok, appendErr) |
| 520 | } |
| 521 | if err := l.AcknowledgeProjection(turns[i]); err != nil { |
| 522 | t.Fatalf("ack %d: %v", i, err) |
| 523 | } |
| 524 | } |
| 525 | if err := l.Compact(); err != nil { |
| 526 | t.Fatalf("Compact: %v", err) |
| 527 | } |
| 528 | |
| 529 | reopened, err := Open(path, "session") |
| 530 | if err != nil { |
| 531 | t.Fatalf("reopen checkpoint: %v", err) |
| 532 | } |
| 533 | if got := len(reopened.summaries); got != terminalSummaryLimit { |
| 534 | t.Fatalf("terminal summaries = %d, want %d", got, terminalSummaryLimit) |
| 535 | } |
| 536 | if reopened.summaries[0].TurnID != turns[4] || reopened.summaries[len(reopened.summaries)-1].TurnID != turns[19] { |
| 537 | t.Fatalf("summary window = first:%q last:%q", reopened.summaries[0].TurnID, reopened.summaries[len(reopened.summaries)-1].TurnID) |
| 538 | } |
| 539 | if got := reopened.TurnIDForSubmission("submission-03"); got != "" { |
| 540 | t.Fatalf("expired submission receipt = %q, want bounded eviction", got) |
| 541 | } |
| 542 | if got := reopened.TurnIDForSubmission("submission-19"); got != turns[19] { |
| 543 | t.Fatalf("recent submission receipt = %q, want %q", got, turns[19]) |
| 544 | } |
| 545 | data, err := os.ReadFile(store.SessionTurnEventLog(path)) |
| 546 | if err != nil { |
| 547 | t.Fatalf("read checkpoint: %v", err) |
| 548 | } |
| 549 | if bytes.Contains(data, []byte("submission-03")) || !bytes.Contains(data, []byte("submission-19")) { |
| 550 | t.Fatalf("checkpoint did not enforce the bounded receipt window: %s", data) |
| 551 | } |
| 552 | } |
| 553 | |
| 554 | func TestLedgerRebuildsCompleteTerminalSummaryBeforeCheckpoint(t *testing.T) { |
| 555 | path := testSessionPath(t) |
| 556 | l, err := Open(path, "session") |
| 557 | if err != nil { |
| 558 | t.Fatalf("Open: %v", err) |
| 559 | } |
| 560 | l.RequireProjectionAck(true) |
| 561 | l.SetRoutingMetadata("epoch", "submission") |
| 562 | turnID, err := l.Begin() |
| 563 | if err != nil { |
| 564 | t.Fatalf("Begin: %v", err) |
| 565 | } |
| 566 | l.SetTranscriptSnapshot(7, "digest") |
| 567 | if _, ok, appendErr := l.Append(event.Event{Kind: event.TurnStarted}, event.TurnInProgress); appendErr != nil || !ok { |
| 568 | t.Fatalf("start: ok=%v err=%v", ok, appendErr) |
| 569 | } |
| 570 | if _, ok, appendErr := l.Append(event.Event{Kind: event.TurnDone, Outcome: "partial"}, event.TurnInterrupted); appendErr != nil || !ok { |
| 571 | t.Fatalf("complete: ok=%v err=%v", ok, appendErr) |
| 572 | } |
| 573 | if err := l.AcknowledgeProjection(turnID); err != nil { |
| 574 | t.Fatalf("ack: %v", err) |
| 575 | } |
| 576 | if err := l.Close(); err != nil { |
| 577 | t.Fatalf("Close: %v", err) |
| 578 | } |
| 579 | |
| 580 | reopened, err := Open(path, "session") |
| 581 | if err != nil { |
| 582 | t.Fatalf("reopen: %v", err) |
| 583 | } |
| 584 | if len(reopened.summaries) != 1 { |
| 585 | t.Fatalf("summaries = %+v, want one reconstructed terminal", reopened.summaries) |
| 586 | } |
| 587 | summary := reopened.summaries[0] |
| 588 | if summary.TurnID != turnID || summary.Outcome != "partial" || summary.Status != event.TurnInterrupted { |
| 589 | t.Fatalf("summary identity = %+v", summary) |
| 590 | } |
| 591 | if summary.StartedAt <= 0 || summary.FinishedAt < summary.StartedAt || summary.DurationMs != summary.FinishedAt-summary.StartedAt { |
| 592 | t.Fatalf("summary timing = %+v", summary) |
| 593 | } |
| 594 | if summary.RuntimeEpoch != "epoch" || summary.SubmissionID != "submission" || summary.TranscriptRevision != 7 || summary.TranscriptDigest != "digest" { |
| 595 | t.Fatalf("summary metadata = %+v", summary) |
| 596 | } |
| 597 | } |
| 598 | |
| 599 | func TestLedgerCompactionFailurePreservesOriginalSidecar(t *testing.T) { |
| 600 | path := testSessionPath(t) |
| 601 | l, err := Open(path, "session") |
| 602 | if err != nil { |
| 603 | t.Fatalf("Open: %v", err) |
| 604 | } |
| 605 | l.RequireProjectionAck(true) |
| 606 | turnID, err := l.Begin() |
| 607 | if err != nil { |
| 608 | t.Fatalf("Begin: %v", err) |
| 609 | } |
| 610 | if _, ok, err := l.Append(event.Event{Kind: event.Text, Text: "retained until atomic replace"}, event.TurnInProgress); err != nil || !ok { |
| 611 | t.Fatalf("text: ok=%v err=%v", ok, err) |
| 612 | } |
| 613 | if _, ok, err := l.Append(event.Event{Kind: event.TurnDone}, event.TurnCompleted); err != nil || !ok { |
| 614 | t.Fatalf("terminal: ok=%v err=%v", ok, err) |
| 615 | } |
| 616 | if err := l.AcknowledgeProjection(turnID); err != nil { |
| 617 | t.Fatalf("ack: %v", err) |
| 618 | } |
| 619 | ledgerPath := store.SessionTurnEventLog(path) |
| 620 | before, err := os.ReadFile(ledgerPath) |
| 621 | if err != nil { |
| 622 | t.Fatalf("read original sidecar: %v", err) |
| 623 | } |
| 624 | originalWriter := atomicWriteLedgerFile |
| 625 | atomicWriteLedgerFile = func(string, []byte, os.FileMode) error { return errors.New("simulated atomic replace failure") } |
| 626 | t.Cleanup(func() { atomicWriteLedgerFile = originalWriter }) |
| 627 | |
| 628 | if err := l.Compact(); err == nil { |
| 629 | t.Fatal("Compact succeeded through an injected atomic replacement failure") |
| 630 | } |
| 631 | after, err := os.ReadFile(ledgerPath) |
| 632 | if err != nil { |
| 633 | t.Fatalf("read sidecar after failed compact: %v", err) |
| 634 | } |
| 635 | if !bytes.Equal(after, before) { |
| 636 | t.Fatalf("failed compaction changed the original sidecar\nbefore=%s\nafter=%s", before, after) |
| 637 | } |
| 638 | if metrics := l.MetricsSnapshot(); metrics.CompactionFailures != 1 { |
| 639 | t.Fatalf("compaction failures = %d, want 1", metrics.CompactionFailures) |
| 640 | } |
| 641 | } |
| 642 | |
| 643 | func TestProjectionAckDoesNotFailWhenBestEffortCompactionFails(t *testing.T) { |
| 644 | path := testSessionPath(t) |
| 645 | l, err := Open(path, "session") |
| 646 | if err != nil { |
| 647 | t.Fatalf("Open: %v", err) |
| 648 | } |
| 649 | l.RequireProjectionAck(true) |
| 650 | l.compactBytes = 1 |
| 651 | turnID, err := l.Begin() |
| 652 | if err != nil { |
| 653 | t.Fatalf("Begin: %v", err) |
| 654 | } |
| 655 | if _, ok, err := l.Append(event.Event{Kind: event.Text, Text: "keep the WAL"}, event.TurnInProgress); err != nil || !ok { |
| 656 | t.Fatalf("text: ok=%v err=%v", ok, err) |
| 657 | } |
| 658 | if _, ok, err := l.Append(event.Event{Kind: event.TurnDone}, event.TurnCompleted); err != nil || !ok { |
| 659 | t.Fatalf("terminal: ok=%v err=%v", ok, err) |
| 660 | } |
| 661 | |
| 662 | originalWriter := atomicWriteLedgerFile |
| 663 | atomicWriteLedgerFile = func(string, []byte, os.FileMode) error { return errors.New("simulated atomic replace failure") } |
| 664 | if err := l.AcknowledgeProjection(turnID); err != nil { |
| 665 | t.Fatalf("durable projection ack inherited best-effort compaction failure: %v", err) |
| 666 | } |
| 667 | if metrics := l.MetricsSnapshot(); metrics.CompactionFailures != 1 { |
| 668 | t.Fatalf("compaction failures = %d, want 1", metrics.CompactionFailures) |
| 669 | } |
| 670 | if pending := l.PendingProjections(); len(pending) != 0 { |
| 671 | t.Fatalf("durably acknowledged projection remained pending: %+v", pending) |
| 672 | } |
| 673 | |
| 674 | atomicWriteLedgerFile = originalWriter |
| 675 | t.Cleanup(func() { atomicWriteLedgerFile = originalWriter }) |
| 676 | if err := l.AcknowledgeProjection(turnID); err != nil { |
| 677 | t.Fatalf("retry acknowledged projection: %v", err) |
| 678 | } |
| 679 | if l.compactedThrough == 0 { |
| 680 | t.Fatal("idempotent acknowledgement did not retry checkpoint compaction") |
| 681 | } |
| 682 | } |
| 683 | |
| 684 | func BenchmarkLedgerReplayMillionEventCheckpointTail(b *testing.B) { |
| 685 | dir := b.TempDir() |
| 686 | path := filepath.Join(dir, "session.jsonl") |
| 687 | ledgerPath := store.SessionTurnEventLog(path) |
| 688 | cp := checkpointRecord{SchemaVersion: schemaVersion, RecordType: "checkpoint", SessionID: "session", CompactedThroughSequence: 1_000_000, ProjectionCommittedThrough: 1_000_000, TerminalSummaries: []TerminalSummary{}} |
| 689 | data, _ := json.Marshal(cp) |
| 690 | data = append(data, '\n') |
| 691 | for seq := uint64(1_000_001); seq <= 1_000_100; seq++ { |
| 692 | kind := "tool_progress" |
| 693 | status := event.TurnInProgress |
| 694 | if seq == 1_000_100 { |
| 695 | kind = "turn_done" |
| 696 | status = event.TurnCompleted |
| 697 | } |
| 698 | rec := diskEventRecord{RecordType: "event", Envelope: Envelope{SchemaVersion: schemaVersion, SessionID: "session", TurnID: "tail", Sequence: seq, Kind: kind, Status: status, CreatedAt: 1, Event: eventwire.Event{Kind: kind}}} |
| 699 | line, _ := json.Marshal(rec) |
| 700 | data = append(data, line...) |
| 701 | data = append(data, '\n') |
| 702 | } |
| 703 | if err := os.WriteFile(ledgerPath, data, 0o600); err != nil { |
| 704 | b.Fatal(err) |
| 705 | } |
| 706 | l, err := Open(path, "session") |
| 707 | if err != nil { |
| 708 | b.Fatal(err) |
| 709 | } |
| 710 | b.ResetTimer() |
| 711 | for range b.N { |
| 712 | view, err := l.Replay(1_000_000) |
| 713 | if err != nil || len(view.Events) != 100 { |
| 714 | b.Fatalf("Replay: events=%d err=%v", len(view.Events), err) |
| 715 | } |
| 716 | } |
| 717 | } |
| 718 |