| 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 |