| 1 | package control |
| 2 | |
| 3 | import ( |
| 4 | "encoding/json" |
| 5 | "os" |
| 6 | "path/filepath" |
| 7 | "strings" |
| 8 | "testing" |
| 9 | |
| 10 | "reasonix/internal/agent" |
| 11 | "reasonix/internal/event" |
| 12 | "reasonix/internal/provider" |
| 13 | "reasonix/internal/store" |
| 14 | "reasonix/internal/tool" |
| 15 | "reasonix/internal/transcript" |
| 16 | ) |
| 17 | |
| 18 | func TestTranscriptReplayResetsOversizedWirePage(t *testing.T) { |
| 19 | for _, body := range []string{strings.Repeat("x", 2<<20), strings.Repeat("<", 400000)} { |
| 20 | c := newOwnedTestController(t, Options{SessionPath: filepath.Join(t.TempDir(), "session.jsonl"), Sink: event.Discard}) |
| 21 | t.Cleanup(c.Close) |
| 22 | before, err := c.TranscriptSnapshot(transcript.PageRequest{}) |
| 23 | if err != nil { |
| 24 | t.Fatal(err) |
| 25 | } |
| 26 | if _, err := c.turnEventLedger().Begin(); err != nil { |
| 27 | t.Fatal(err) |
| 28 | } |
| 29 | if err := c.emitTurnEventChecked(event.Event{Kind: event.ToolResult, Tool: event.Tool{ID: "read", Name: "read_file", Output: body}}); err != nil { |
| 30 | t.Fatal(err) |
| 31 | } |
| 32 | replay, err := c.TranscriptReplay(TranscriptReplayRequest{Identity: before.Identity}) |
| 33 | if err != nil { |
| 34 | t.Fatal(err) |
| 35 | } |
| 36 | encoded, err := json.Marshal(replay) |
| 37 | if err != nil { |
| 38 | t.Fatal(err) |
| 39 | } |
| 40 | if len(encoded)+1 > 2<<20 || !replay.ResetRequired || len(replay.Events) != 0 { |
| 41 | t.Fatalf("oversized replay did not reset: bytes=%d reset=%v", len(encoded), replay.ResetRequired) |
| 42 | } |
| 43 | cut, err := c.TranscriptSnapshot(transcript.PageRequest{}) |
| 44 | if err != nil || cut.CoveredThroughSeq != replay.LatestSequence || len(cut.Records) != 1 { |
| 45 | t.Fatalf("replacement cut: %v", err) |
| 46 | } |
| 47 | ref := cut.Records[0].Refs[0] |
| 48 | var full strings.Builder |
| 49 | for offset := 0; ; { |
| 50 | chunk, err := c.TranscriptContent(transcript.ContentRequest{ContentRef: ref, Offset: offset}) |
| 51 | if err != nil || chunk.Stale { |
| 52 | t.Fatalf("content: %v", err) |
| 53 | } |
| 54 | full.WriteString(chunk.Data) |
| 55 | offset = chunk.NextOffset |
| 56 | if chunk.Done { |
| 57 | break |
| 58 | } |
| 59 | } |
| 60 | if full.String() != body { |
| 61 | t.Fatal("fallback snapshot lost the oversized event body") |
| 62 | } |
| 63 | } |
| 64 | } |
| 65 | |
| 66 | func TestTranscriptProjectionCommitsBeforePublicationAndAllowsReentry(t *testing.T) { |
| 67 | var c *Controller |
| 68 | publications := 0 |
| 69 | c = newOwnedTestController(t, Options{SessionPath: filepath.Join(t.TempDir(), "session.jsonl"), Sink: event.FuncSink(func(e event.Event) { |
| 70 | if e.Sequence == 0 { |
| 71 | return |
| 72 | } |
| 73 | publications++ |
| 74 | snap, err := c.TranscriptSnapshot(transcript.PageRequest{}) |
| 75 | if err != nil { |
| 76 | t.Errorf("snapshot during publish: %v", err) |
| 77 | return |
| 78 | } |
| 79 | if snap.CoveredThroughSeq < e.Sequence { |
| 80 | t.Errorf("published %d before snapshot coverage %d", e.Sequence, snap.CoveredThroughSeq) |
| 81 | } |
| 82 | if e.Kind == event.AskRequest { |
| 83 | if err := c.emitTurnEventChecked(event.Event{Kind: event.PromptAnswered, ItemID: "prompt"}); err != nil { |
| 84 | t.Error(err) |
| 85 | } |
| 86 | } |
| 87 | })}) |
| 88 | t.Cleanup(c.Close) |
| 89 | c.SetTurnEventRoutingMetadata("runtime", "submission") |
| 90 | if _, err := c.turnEventLedger().Begin(); err != nil { |
| 91 | t.Fatal(err) |
| 92 | } |
| 93 | for _, e := range []event.Event{ |
| 94 | {Kind: event.TurnStarted}, |
| 95 | {Kind: event.UserMessage, MessageID: "u", Text: "question"}, |
| 96 | {Kind: event.Reasoning, MessageID: "a", Text: "think"}, |
| 97 | {Kind: event.AskRequest, ItemID: "prompt"}, |
| 98 | } { |
| 99 | if err := c.emitTurnEventChecked(e); err != nil { |
| 100 | t.Fatal(err) |
| 101 | } |
| 102 | } |
| 103 | snap, err := c.TranscriptSnapshot(transcript.PageRequest{}) |
| 104 | if err != nil { |
| 105 | t.Fatal(err) |
| 106 | } |
| 107 | if publications != 5 || snap.CoveredThroughSeq != 5 || len(snap.Records) != 2 || snap.Records[1].Message.Reasoning != "think" { |
| 108 | t.Fatalf("publications=%d snapshot=%+v", publications, snap) |
| 109 | } |
| 110 | if len(snap.Runtime.PendingEvents) != 0 { |
| 111 | t.Fatal("answered prompt survived snapshot") |
| 112 | } |
| 113 | if snap.Records[0].Message.SubmissionID != "submission" { |
| 114 | t.Fatal("lost submission identity") |
| 115 | } |
| 116 | view, err := c.TranscriptReplay(TranscriptReplayRequest{Identity: snap.Identity, After: 2}) |
| 117 | if err != nil { |
| 118 | t.Fatal(err) |
| 119 | } |
| 120 | if view.CoveredThroughSeq != 5 || len(view.Events) != 3 { |
| 121 | t.Fatalf("replay=%+v", view) |
| 122 | } |
| 123 | } |
| 124 | |
| 125 | func TestTranscriptCheckpointFailureRetainsWALWithoutFailingCompletedTurn(t *testing.T) { |
| 126 | path := filepath.Join(t.TempDir(), "session.jsonl") |
| 127 | c := newOwnedTestController(t, Options{SessionPath: path, Sink: event.Discard}) |
| 128 | t.Cleanup(c.Close) |
| 129 | if err := os.Mkdir(store.SessionTranscriptProjection(path), 0o700); err != nil { |
| 130 | t.Fatal(err) |
| 131 | } |
| 132 | turn, err := c.turnEventLedger().Begin() |
| 133 | if err != nil { |
| 134 | t.Fatal(err) |
| 135 | } |
| 136 | for _, e := range []event.Event{{Kind: event.TurnStarted}, {Kind: event.Text, MessageID: "a", Text: "kept"}, {Kind: event.TurnDone}} { |
| 137 | if err := c.emitTurnEventChecked(e); err != nil { |
| 138 | t.Fatal(err) |
| 139 | } |
| 140 | } |
| 141 | if c.turnEventLedgerError() != nil { |
| 142 | t.Fatal("display checkpoint failure poisoned the runtime WAL") |
| 143 | } |
| 144 | if len(c.PendingTurnProjections()) == 0 { |
| 145 | t.Fatal("failed checkpoint allowed WAL compaction") |
| 146 | } |
| 147 | if err := os.Remove(store.SessionTranscriptProjection(path)); err != nil { |
| 148 | t.Fatal(err) |
| 149 | } |
| 150 | if err := c.AcknowledgeTurnProjection(turn); err != nil { |
| 151 | t.Fatal(err) |
| 152 | } |
| 153 | if len(c.PendingTurnProjections()) != 0 { |
| 154 | t.Fatal("successful retry did not acknowledge projection") |
| 155 | } |
| 156 | state, ok, err := transcript.LoadCheckpoint(path) |
| 157 | if err != nil || !ok || len(state.Records) != 1 || state.Records[0].Content != "kept" { |
| 158 | t.Fatalf("retry checkpoint=%+v %v", state, err) |
| 159 | } |
| 160 | } |
| 161 | |
| 162 | func TestTranscriptCheckpointRestoreDoesNotReplayAutosavedTextTwice(t *testing.T) { |
| 163 | path := filepath.Join(t.TempDir(), "session.jsonl") |
| 164 | session := agent.NewSession("system") |
| 165 | if err := session.Save(path); err != nil { |
| 166 | t.Fatal(err) |
| 167 | } |
| 168 | newController := func(session *agent.Session) *Controller { |
| 169 | executor := agent.New(nil, tool.NewRegistry(), session, agent.Options{}, event.Discard) |
| 170 | return newOwnedTestController(t, Options{Executor: executor, SessionPath: path, Sink: event.Discard}) |
| 171 | } |
| 172 | c := newController(session) |
| 173 | emit := func(c *Controller, e event.Event) { |
| 174 | t.Helper() |
| 175 | if err := c.emitTurnEventChecked(e); err != nil { |
| 176 | t.Fatal(err) |
| 177 | } |
| 178 | } |
| 179 | start := func(c *Controller, userID, assistantID, body string) { |
| 180 | t.Helper() |
| 181 | if _, err := c.turnEventLedger().Begin(); err != nil { |
| 182 | t.Fatal(err) |
| 183 | } |
| 184 | emit(c, event.Event{Kind: event.TurnStarted}) |
| 185 | c.executor.Session().Add(provider.Message{ID: userID, Role: provider.RoleUser, Content: "question", Origin: provider.MessageOriginUser}) |
| 186 | emit(c, event.Event{Kind: event.UserMessage, MessageID: userID, Text: "question"}) |
| 187 | emit(c, event.Event{Kind: event.StreamAttempt, MessageID: assistantID, AttemptID: assistantID, StreamAttempt: event.StreamAttemptInfo{ID: assistantID, Action: event.StreamAttemptBegin}}) |
| 188 | emit(c, event.Event{Kind: event.Text, MessageID: assistantID, Text: body}) |
| 189 | emit(c, event.Event{Kind: event.Message, MessageID: assistantID, Text: body}) |
| 190 | emit(c, event.Event{Kind: event.StreamAttempt, MessageID: assistantID, AttemptID: assistantID, StreamAttempt: event.StreamAttemptInfo{ID: assistantID, Action: event.StreamAttemptCommit}}) |
| 191 | c.executor.Session().Add(provider.Message{ID: assistantID, Role: provider.RoleAssistant, Content: body}) |
| 192 | if err := c.executor.Session().Save(path); err != nil { |
| 193 | t.Fatal(err) |
| 194 | } |
| 195 | } |
| 196 | start(c, "u1", "a1", "first") |
| 197 | emit(c, event.Event{Kind: event.TurnDone}) |
| 198 | before, err := c.TranscriptSnapshot(transcript.PageRequest{}) |
| 199 | if err != nil { |
| 200 | t.Fatal(err) |
| 201 | } |
| 202 | c.Close() |
| 203 | loaded, err := agent.LoadSession(path) |
| 204 | if err != nil { |
| 205 | t.Fatal(err) |
| 206 | } |
| 207 | c = newController(loaded) |
| 208 | after, err := c.TranscriptSnapshot(transcript.PageRequest{}) |
| 209 | if err != nil { |
| 210 | t.Fatal(err) |
| 211 | } |
| 212 | if len(before.Records) != len(after.Records) { |
| 213 | t.Fatalf("rebuild changed record count: before=%d after=%d", len(before.Records), len(after.Records)) |
| 214 | } |
| 215 | for i := range before.Records { |
| 216 | want, got := before.Records[i].Message, after.Records[i].Message |
| 217 | if want.MessageID != got.MessageID || want.Role != got.Role || want.Content != got.Content { |
| 218 | t.Fatalf("rebuild changed stable message %d: before=%+v after=%+v", i, want, got) |
| 219 | } |
| 220 | } |
| 221 | start(c, "u2", "a2", "autosaved tail") |
| 222 | c.Close() // No terminal event: simulates a process ending after autosave. |
| 223 | loaded, err = agent.LoadSession(path) |
| 224 | if err != nil { |
| 225 | t.Fatal(err) |
| 226 | } |
| 227 | c = newController(loaded) |
| 228 | t.Cleanup(c.Close) |
| 229 | recovered, err := c.TranscriptSnapshot(transcript.PageRequest{}) |
| 230 | if err != nil { |
| 231 | t.Fatal(err) |
| 232 | } |
| 233 | // The legacy helper above never emits a v3 turn/start for its second tail. |
| 234 | // Cold history therefore preserves the four durable messages without |
| 235 | // manufacturing an interruption fact from transcript wording alone. |
| 236 | if len(recovered.Records) != 4 || recovered.Records[3].Message.Content != "autosaved tail" { |
| 237 | t.Fatalf("recovered suffix duplicated or lost: %+v", recovered) |
| 238 | } |
| 239 | } |
| 240 |