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