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