返回 DeepSeek-Reasonix
transcript_recovery_test.go
根目录 / internal / session / transcript_recovery_test.go
1 package session
2
3 import (
4 "context"
5 "encoding/json"
6 "fmt"
7 "path/filepath"
8 "testing"
9
10 "reasonix/internal/event"
11 "reasonix/internal/provider"
12 "reasonix/internal/transcript"
13 )
14
15 // The shape reproduces the reported long tool-driven turn without any copied
16 // user text: 147 messages, 46 assistants with visible content, 72 tool results.
17 func syntheticTranscriptRecoveryMessages() []provider.Message {
18 messages := []provider.Message{
19 {ID: "user", Role: provider.RoleUser, Content: "create a synthetic illustration", Origin: provider.MessageOriginUser},
20 {ID: "introduction", Role: provider.RoleAssistant, Content: "Starting the illustration."},
21 }
22 for i := range 72 {
23 callID := fmt.Sprintf("call-%02d", i)
24 assistant := provider.Message{ID: fmt.Sprintf("assistant-%02d", i), Role: provider.RoleAssistant,
25 ToolCalls: []provider.ToolCall{{ID: callID, Name: "synthetic_tool", Arguments: `{}`}}}
26 if i < 44 {
27 assistant.Content = fmt.Sprintf("Iteration %02d is ready for inspection.", i)
28 assistant.ReasoningContent = fmt.Sprintf("Synthetic reasoning for iteration %02d.", i)
29 }
30 messages = append(messages, assistant, provider.Message{ID: fmt.Sprintf("tool-%02d", i), Role: provider.RoleTool, Name: "synthetic_tool", ToolCallID: callID, Content: fmt.Sprintf("synthetic result %02d", i)})
31 }
32 return append(messages, provider.Message{ID: "final-answer", Role: provider.RoleAssistant, Content: "The synthetic illustration is complete.", ReasoningContent: "All requested checks completed.", WorkDurationMs: 933524})
33 }
34
35 func TestTranscriptLongTurnRemainsReachableAfterBoundedEvictionAndRestart(t *testing.T) {
36 service, err := NewService("local", NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions")))
37 if err != nil {
38 t.Fatal(err)
39 }
40 t.Cleanup(func() { _ = service.Shutdown(context.Background()) })
41 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "synthetic-long-turn"})
42 if err != nil {
43 t.Fatal(err)
44 }
45 messages := syntheticTranscriptRecoveryMessages()
46 if _, err := runtime.Session().Append(t.Context(), Batch{OperationID: "turn-start", TurnID: "turn", Events: []Event{{Kind: "turn/start"}}}); err != nil {
47 t.Fatal(err)
48 }
49 for i, message := range messages {
50 if len(message.ToolCalls) > 0 {
51 attempt, _ := json.Marshal(map[string]any{"id": message.ID, "messageId": message.ID, "action": "begin"})
52 call, _ := json.Marshal(map[string]any{"id": message.ToolCalls[0].ID, "name": "synthetic_tool"})
53 // Repeated records for one stable identity must not inflate counts.
54 if _, err := runtime.Session().Append(t.Context(), Batch{OperationID: fmt.Sprintf("sampling-%03d", i), TurnID: "turn", Events: []Event{
55 {Kind: "assistant/attempt", Payload: attempt}, {Kind: "assistant/attempt", Payload: attempt},
56 {Kind: "tool/call", Payload: call}, {Kind: "tool/call", Payload: call},
57 }}); err != nil {
58 t.Fatal(err)
59 }
60 }
61 payload, err := json.Marshal(map[string]any{"message": message})
62 if err != nil {
63 t.Fatal(err)
64 }
65 if _, err := runtime.Session().Append(t.Context(), Batch{OperationID: fmt.Sprintf("message-%03d", i), TurnID: "turn", Events: []Event{{Kind: "message/complete", Payload: payload}}}); err != nil {
66 t.Fatal(err)
67 }
68 }
69 if _, err := runtime.Session().Append(t.Context(), Batch{OperationID: "turn-end", TurnID: "turn", Events: []Event{{Kind: "turn/end", Payload: []byte(`{"status":"completed"}`)}}}); err != nil {
70 t.Fatal(err)
71 }
72 if _, err := runtime.Session().Flush(t.Context()); err != nil {
73 t.Fatal(err)
74 }
75 before, err := runtime.Transcript().Snapshot(transcript.PageRequest{Records: 32})
76 if err != nil || before.TotalRecords > 96 || len(before.Records) > 32 {
77 t.Fatalf("publisher exceeded resident budget: total=%d page=%d error=%v", before.TotalRecords, len(before.Records), err)
78 }
79 ref, query := runtime.Ref(), service.Query()
80 // Test reachability under the package timeout, not a disk-speed cutoff.
81 // Materialize synchronously; history_preparation_test independently covers
82 // nonblocking readers while an index rebuild is pending.
83 if _, _, err := query.prepareHistoryIndex(t.Context(), ref); err != nil {
84 t.Fatal(err)
85 }
86 seen := make(map[string]provider.Message)
87 newest := windowReady(t, query, ref, HistoryWindowRequest{Anchor: "newest", Limit: 32})
88 final := newest.Messages[len(newest.Messages)-1]
89 if !final.TurnFinal || final.MessageID != "final-answer" || final.TurnDurationMs != 933524 {
90 t.Fatalf("history lost the authoritative final identity or complete turn duration: %+v", final)
91 }
92 if final.SamplingCount == nil || *final.SamplingCount != 72 || final.ToolCount == nil || *final.ToolCount != 72 {
93 t.Fatalf("history lost distinct attempt/tool counts: %+v", final)
94 }
95 page := newest
96 for {
97 if page.Status != "ready" || len(page.Messages) > 32 || page.SnapshotSequence != newest.SnapshotSequence {
98 t.Fatalf("invalid bounded history page: %+v", page)
99 }
100 for _, stored := range page.Messages {
101 body := []byte(stored.Inline)
102 if stored.ContentRef != nil {
103 body, err = query.ReadContent(t.Context(), ref, *stored.ContentRef, 0, stored.ContentRef.Bytes)
104 if err != nil {
105 t.Fatal(err)
106 }
107 }
108 var message provider.Message
109 if err := json.Unmarshal(body, &message); err != nil {
110 t.Fatalf("decode %s: %v", stored.MessageID, err)
111 }
112 if _, duplicate := seen[message.ID]; duplicate {
113 t.Fatalf("history repeated message %q", message.ID)
114 }
115 seen[message.ID] = message
116 }
117 if !page.HasOlder {
118 break
119 }
120 if page.OlderCursor == "" {
121 t.Fatal("evicted older history has no continuation")
122 }
123 page = windowReady(t, query, ref, HistoryWindowRequest{Anchor: "cursor", Cursor: page.OlderCursor, Limit: 32})
124 }
125 contentAssistants, toolResults := 0, 0
126 for _, expected := range messages {
127 actual, found := seen[expected.ID]
128 if !found || actual.Content != expected.Content || actual.ReasoningContent != expected.ReasoningContent || actual.ToolCallID != expected.ToolCallID {
129 t.Fatalf("message %q is missing or its body changed after eviction", expected.ID)
130 }
131 if actual.Role == provider.RoleAssistant && (actual.Content != "" || actual.ReasoningContent != "") {
132 contentAssistants++
133 }
134 if actual.Role == provider.RoleTool {
135 toolResults++
136 }
137 }
138 if len(seen) != 147 || contentAssistants != 46 || toolResults != 72 {
139 t.Fatalf("recovered shape: messages=%d visibleAssistants=%d tools=%d", len(seen), contentAssistants, toolResults)
140 }
141 // The oldest window must lead back to the newest; reclaiming one edge is
142 // not deletion and must never strand the reader at an unloaded boundary.
143 newerSeen := make(map[string]bool)
144 for {
145 for _, message := range page.Messages {
146 newerSeen[message.MessageID] = true
147 }
148 if !page.HasNewer {
149 break
150 }
151 if page.NewerCursor == "" {
152 t.Fatal("older history has no forward continuation")
153 }
154 page = windowReady(t, query, ref, HistoryWindowRequest{Anchor: "cursor", Cursor: page.NewerCursor, Limit: 32})
155 }
156 if len(newerSeen) != len(messages) || !newerSeen["final-answer"] {
157 t.Fatal("bidirectional paging did not reach all messages and final answer")
158 }
159 for _, id := range []string{"introduction", "assistant-00", "final-answer"} {
160 location, err := query.LocateMessage(t.Context(), ref, id, newest.SnapshotSequence)
161 if err != nil || location.Status != "ready" || location.MessageID != id {
162 t.Fatalf("locate %q: status=%q error=%v", id, location.Status, err)
163 }
164 }
165 if err := service.Close(t.Context(), ref); err != nil {
166 t.Fatal(err)
167 }
168 binding, err := service.Open(t.Context(), ref)
169 if err != nil {
170 t.Fatal(err)
171 }
172 t.Cleanup(func() { _ = binding.Release(context.Background()) })
173 reopened := binding.Runtime()
174 // This path constructs only session/query objects: reopening and following
175 // must not need an Agent, provider, or a model request to reconstruct UI.
176 initial, err := reopened.FollowTranscript(t.Context(), transcript.FollowRequest{})
177 if err != nil || initial.Snapshot == nil {
178 t.Fatalf("follow reopened session: %v", err)
179 }
180 t.Cleanup(func() {
181 _, _ = reopened.FollowTranscript(context.Background(), transcript.FollowRequest{Subscription: initial.Subscription, Close: true})
182 })
183 cut := initial.Snapshot
184 if cut.Runtime.SamplingCount != 72 || cut.Runtime.ToolCount != 72 {
185 t.Fatalf("restart lost distinct attempt/tool counts: %+v", cut.Runtime)
186 }
187 if cut.TotalRecords > 96 || cut.Runtime.Status != event.TurnCompleted || cut.Runtime.FinalMessageID != "final-answer" || cut.Runtime.DurationMs != 933524 {
188 t.Errorf("reopened transcript lost terminal state or budget: records=%d runtime=%+v", cut.TotalRecords, cut.Runtime)
189 }
190 found := false
191 for _, row := range cut.Records {
192 found = found || row.Message.MessageID == "final-answer" && row.Message.Content == "The synthetic illustration is complete."
193 }
194 if !found || len(cut.ActiveAttempts) != 0 || len(cut.ActiveRecords) != 0 {
195 t.Fatal("reopened completed follow omitted the final answer or manufactured an active turn")
196 }
197 }
198
198 lines GO