返回 DeepSeek-Reasonix
business_test.go
根目录 / internal / transcript / business_test.go
1 package transcript
2
3 import (
4 "fmt"
5 "testing"
6
7 "reasonix/internal/event"
8 "reasonix/internal/eventwire"
9 "reasonix/internal/provider"
10 "reasonix/internal/turnevent"
11 )
12
13 func businessFrame(t *testing.T, p *Projection, covered uint64, e event.Event) {
14 t.Helper()
15 wire := eventwire.ToWire(e)
16 status := event.TurnInProgress
17 if e.Kind == event.TurnDone {
18 status = event.TurnCompleted
19 }
20 if err := p.ApplyFrame(turnevent.Envelope{SessionID: testIdentity.SessionID, RuntimeEpoch: testIdentity.RuntimeEpoch, TurnID: "turn", Kind: wire.Kind, Status: status, Event: wire}, covered); err != nil {
21 t.Fatal(err)
22 }
23 }
24
25 func TestAcceptBusinessPublishesTheTurnIDItAssigns(t *testing.T) {
26 p, initial := newFollowProjection(t)
27 p.AcceptBusiness([]Message{
28 {RecordID: "m:user", MessageID: "user", Role: "user", Content: "question"},
29 {RecordID: "m:answer", MessageID: "answer", Role: "assistant", Content: "answer"},
30 {RecordID: "m:kept", MessageID: "kept", Role: "user", Content: "other", TurnID: "kept-turn"},
31 }, 1, "turn-7", false)
32 suffix := followChanges(t, p, FollowRequest{Subscription: initial.Subscription, AfterRevision: initial.Snapshot.ProjectionRevision})
33 if len(suffix.Changes) != 1 || len(suffix.Changes[0].Records) != 3 {
34 t.Fatalf("published records = %+v", suffix.Changes)
35 }
36 records := suffix.Changes[0].Records
37 if records[0].TurnID != "turn-7" || records[1].TurnID != "turn-7" || records[2].TurnID != "kept-turn" {
38 t.Fatalf("published turn identity = %q %q %q", records[0].TurnID, records[1].TurnID, records[2].TurnID)
39 }
40 cut := snapshot(t, p)
41 if len(cut.Records) != 3 || cut.Records[0].Message.TurnID != "turn-7" || cut.Records[1].Message.TurnID != "turn-7" || cut.Records[2].Message.TurnID != "kept-turn" {
42 t.Fatalf("buffer turn identity diverged from the published change: %+v", cut.Records)
43 }
44 }
45
46 func TestBusinessSettlementUpdatesStreamingRowWithoutDuplicateOrPendingState(t *testing.T) {
47 p, err := NewProjection(testIdentity, nil, 0)
48 if err != nil {
49 t.Fatal(err)
50 }
51 businessFrame(t, p, 0, event.Event{Kind: event.StreamAttempt, MessageID: "answer", AttemptID: "answer", StreamAttempt: event.StreamAttemptInfo{ID: "answer", Action: event.StreamAttemptBegin}})
52 businessFrame(t, p, 0, event.Event{Kind: event.Text, MessageID: "answer", AttemptID: "answer", Text: "partial"})
53 p.AcceptBusiness([]Message{{RecordID: "m:answer", MessageID: "answer", Role: "assistant", Content: "complete answer", Reasoning: "complete thinking"}}, 1, "turn", false)
54 businessFrame(t, p, 1, event.Event{Kind: event.Message, MessageID: "answer", AttemptID: "answer", Text: "complete answer", Reasoning: "complete thinking"})
55 businessFrame(t, p, 1, event.Event{Kind: event.StreamAttempt, MessageID: "answer", AttemptID: "answer", StreamAttempt: event.StreamAttemptInfo{ID: "answer", Action: event.StreamAttemptCommit}})
56 businessFrame(t, p, 1, event.Event{Kind: event.TurnDone})
57 cut := snapshot(t, p)
58 if len(cut.Records) != 1 {
59 t.Fatalf("settlement duplicated streaming node: records=%d", len(cut.Records))
60 }
61 message := cut.Records[0].Message
62 if message.Content != "complete answer" || message.Reasoning != "complete thinking" || message.Pending {
63 t.Fatalf("settled content or pending state is wrong: %+v", message)
64 }
65 if len(cut.ActiveAttempts) != 0 || len(cut.ActiveRecords) != 0 {
66 t.Fatalf("settlement retains active state: attempts=%d records=%d", len(cut.ActiveAttempts), len(cut.ActiveRecords))
67 }
68 }
69
70 func TestToolResultCompletionPreservesCanonicalIdentityAndMetadata(t *testing.T) {
71 p, err := NewProjection(testIdentity, nil, 0)
72 if err != nil {
73 t.Fatal(err)
74 }
75 execution := &provider.ToolExecution{Kind: "shell", State: "completed", DurationMs: 23}
76 readCompletion := &provider.ReadCompletion{ID: "read-1", Omitted: 2}
77 p.AcceptBusiness([]Message{{
78 RecordID: "tool:call-1", MessageID: "result-1", Role: "tool", ToolCallID: "call-1", ToolName: "PowerShell",
79 Content: "persisted", Source: "history", TurnID: "turn-1", HistoryTurn: 4, CreatedAt: 123,
80 Execution: execution, ToolResultArchived: true, ReadCompletion: readCompletion,
81 }}, 1, "turn-1", false)
82 businessFrame(t, p, 1, event.Event{Kind: event.ToolResult, MessageID: "assistant-owner", Tool: event.Tool{
83 ID: "call-1", Name: "PowerShell", Output: "completed", PresentedFiles: []provider.PresentedFile{{Path: "report.txt"}},
84 }})
85
86 cut := snapshot(t, p)
87 if len(cut.Records) != 1 {
88 t.Fatalf("tool completion duplicated canonical row: records=%d", len(cut.Records))
89 }
90 got := cut.Records[0].Message
91 if got.RecordID != "tool:call-1" || got.MessageID != "result-1" || got.Source != "history" || got.TurnID != "turn-1" || got.HistoryTurn != 4 || got.CreatedAt != 123 {
92 t.Fatalf("tool completion lost canonical identity/location: %+v", got)
93 }
94 if got.Execution == nil || got.Execution.Kind != execution.Kind || got.Execution.State != execution.State || got.Execution.DurationMs != execution.DurationMs ||
95 !got.ToolResultArchived || got.ReadCompletion == nil || got.ReadCompletion.ID != readCompletion.ID || got.ReadCompletion.Omitted != readCompletion.Omitted {
96 t.Fatalf("tool completion lost canonical metadata: %+v", got)
97 }
98 if got.Content != "completed" || len(got.PresentedFiles) != 1 || got.PresentedFiles[0].Path != "report.txt" {
99 t.Fatalf("tool completion did not update event-owned fields: %+v", got)
100 }
101 }
102
103 func TestToolResultEventThenCanonicalRecordKeepsOneStableRow(t *testing.T) {
104 p, err := NewProjection(testIdentity, nil, 0)
105 if err != nil {
106 t.Fatal(err)
107 }
108 businessFrame(t, p, 0, event.Event{Kind: event.ToolResult, Tool: event.Tool{ID: "call-1", Name: "edit_file", Output: "event result"}})
109 p.AcceptBusiness([]Message{{RecordID: "tool:call-1", MessageID: "result-1", Role: "tool", ToolCallID: "call-1",
110 ToolName: "edit_file", Content: "canonical result", HistoryTurn: 2, Execution: &provider.ToolExecution{State: "completed"}}}, 1, "turn-1", false)
111 businessFrame(t, p, 1, event.Event{Kind: event.ToolResult, Tool: event.Tool{ID: "call-1", Name: "edit_file", Output: "event result"}})
112
113 cut := snapshot(t, p)
114 if len(cut.Records) != 1 {
115 t.Fatalf("event/formal handoff duplicated row: records=%d", len(cut.Records))
116 }
117 got := cut.Records[0].Message
118 if got.RecordID != "tool:call-1" || got.MessageID != "result-1" || got.HistoryTurn != 2 || got.Execution == nil {
119 t.Fatalf("event/formal handoff degraded canonical row: %+v", got)
120 }
121 }
122
123 func TestUnappliedSteerNoticeAndCanonicalRecordShareOneRow(t *testing.T) {
124 for _, eventFirst := range []bool{false, true} {
125 p, err := NewProjection(testIdentity, nil, 0)
126 if err != nil {
127 t.Fatal(err)
128 }
129 const recordID = "m:queued:notice:0"
130 warning := event.Event{Kind: event.Notice, Code: event.NoticeCodeUnappliedSteer,
131 MessageID: "queued", Level: event.LevelWarn, Text: "Guidance was not applied:\nUse plan B"}
132 formal := []Message{{RecordID: recordID, MessageID: "queued", Role: "notice", Code: event.NoticeCodeUnappliedSteer,
133 Level: "warn", Content: warning.Text}}
134 if eventFirst {
135 businessFrame(t, p, 0, warning)
136 p.AcceptBusiness(formal, 1, "turn", false)
137 } else {
138 p.AcceptBusiness(formal, 1, "turn", false)
139 businessFrame(t, p, 1, warning)
140 }
141 businessFrame(t, p, 1, warning)
142 cut := snapshot(t, p)
143 if len(cut.Records) != 1 || cut.Records[0].ID != recordID || cut.Records[0].Message.MessageID != "queued" {
144 t.Fatalf("eventFirst=%v duplicated unapplied steer: %+v", eventFirst, cut.Records)
145 }
146 }
147 }
148
149 func TestBusinessRowsWithoutCanonicalIdentityReceiveDistinctViewIdentity(t *testing.T) {
150 p, err := NewProjection(testIdentity, nil, 0)
151 if err != nil {
152 t.Fatal(err)
153 }
154 initial, err := p.Follow(t.Context(), FollowRequest{})
155 if err != nil {
156 t.Fatal(err)
157 }
158 p.AcceptBusiness([]Message{{Role: "notice", Content: "first"}, {Role: "notice", Content: "second"}}, 1, "turn", false)
159 cut := snapshot(t, p)
160 if len(cut.Records) != 2 || cut.Records[0].ID == "" || cut.Records[0].ID == cut.Records[1].ID {
161 t.Fatalf("business identity allocation = %+v", cut.Records)
162 }
163 suffix := followChanges(t, p, FollowRequest{Subscription: initial.Subscription, AfterRevision: initial.Snapshot.ProjectionRevision})
164 if len(suffix.Changes) != 1 || len(suffix.Changes[0].Records) != 2 || suffix.Changes[0].Records[0].RecordID == "" ||
165 suffix.Changes[0].Records[0].RecordID == suffix.Changes[0].Records[1].RecordID {
166 t.Fatalf("published business identities = %+v", suffix.Changes)
167 }
168 }
169
170 func TestBusinessPublisherReleasesCompletedActiveRowsWithinHistoryBudget(t *testing.T) {
171 p, err := NewProjection(testIdentity, nil, 0)
172 if err != nil {
173 t.Fatal(err)
174 }
175 businessFrame(t, p, 0, event.Event{Kind: event.StreamAttempt, MessageID: "old-answer", AttemptID: "old-answer", StreamAttempt: event.StreamAttemptInfo{ID: "old-answer", Action: event.StreamAttemptBegin}})
176 businessFrame(t, p, 0, event.Event{Kind: event.Text, MessageID: "old-answer", Text: "old partial answer"})
177 for i := range 120 {
178 id := fmt.Sprintf("new-answer-%03d", i)
179 p.AcceptBusiness([]Message{{RecordID: "m:" + id, MessageID: id, Role: "assistant", Content: id}}, uint64(i+1), fmt.Sprintf("turn-%d", i), false)
180 }
181 active := snapshot(t, p)
182 retained := false
183 for _, records := range [][]Record{active.Records, active.ActiveRecords} {
184 for _, record := range records {
185 retained = retained || record.Message.MessageID == "old-answer" && record.Message.Content == "old partial answer"
186 }
187 }
188 if len(active.ActiveAttempts) != 1 || !retained {
189 t.Fatal("appended history evicted the active assistant prefix")
190 }
191 // A terminal turn must release pending state even if its attempt end frame
192 // was unavailable. No later user submission should be needed to reclaim
193 // rows that were exempted from the resident budget while still active.
194 businessFrame(t, p, 120, event.Event{Kind: event.TurnDone})
195 cut := snapshot(t, p)
196 if cut.TotalRecords > 96 || len(cut.ActiveRecords) != 0 || len(cut.ActiveAttempts) != 0 {
197 t.Fatalf("completed rows escaped resident budget: total=%d activeRows=%d activeAttempts=%d", cut.TotalRecords, len(cut.ActiveRecords), len(cut.ActiveAttempts))
198 }
199 for _, record := range cut.Records {
200 if record.Message.MessageID == "old-answer" || record.Message.Pending {
201 t.Fatalf("completed active row was retained indefinitely: %+v", record.Message)
202 }
203 }
204 }
205
205 lines GO