返回 DeepSeek-Reasonix
transcript_api_test.go
根目录 / internal / control / transcript_api_test.go
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
259 lines GO