返回 DeepSeek-Reasonix
recorder_test.go
根目录 / internal / trajectory / recorder_test.go
1 package trajectory
2
3 import (
4 "encoding/json"
5 "os"
6 "path/filepath"
7 "strings"
8 "testing"
9 "time"
10
11 "reasonix/internal/event"
12 "reasonix/internal/evidence"
13 )
14
15 type capabilitySink struct {
16 events []event.Event
17 readiness []evidence.ReadinessAudit
18 anchorSafety []event.AnchorSafetyAudit
19 recoveries []event.ProtocolRecoveryAudit
20 outcomes []evidence.OutcomeSample
21 reports []event.CompletionReportAudit
22 workspace []event.WorkspaceMutation
23 runBudgets []event.RunBudgetSample
24 turns int
25 }
26
27 func (s *capabilitySink) Emit(e event.Event) { s.events = append(s.events, e) }
28 func (s *capabilitySink) RecordReadinessAudit(a evidence.ReadinessAudit) {
29 s.readiness = append(s.readiness, a)
30 }
31 func (s *capabilitySink) RecordAnchorSafetyAudit(a event.AnchorSafetyAudit) {
32 s.anchorSafety = append(s.anchorSafety, a)
33 }
34 func (s *capabilitySink) RecordProtocolRecovery(a event.ProtocolRecoveryAudit) {
35 s.recoveries = append(s.recoveries, a)
36 }
37 func (s *capabilitySink) RecordTurnCompletion() { s.turns++ }
38 func (s *capabilitySink) RecordOutcomeProgress(sample evidence.OutcomeSample) {
39 s.outcomes = append(s.outcomes, sample)
40 }
41 func (s *capabilitySink) RecordCompletionReport(a event.CompletionReportAudit) {
42 s.reports = append(s.reports, a)
43 }
44 func (s *capabilitySink) RecordWorkspaceMutation(m event.WorkspaceMutation) {
45 s.workspace = append(s.workspace, m)
46 }
47 func (s *capabilitySink) RecordRunBudget(sample event.RunBudgetSample) {
48 s.runBudgets = append(s.runBudgets, sample)
49 }
50
51 func readRecords(t *testing.T, path string) []Record {
52 t.Helper()
53 data, err := os.ReadFile(path)
54 if err != nil {
55 t.Fatalf("read trajectory: %v", err)
56 }
57 var out []Record
58 for line := range strings.SplitSeq(strings.TrimSpace(string(data)), "\n") {
59 var r Record
60 if err := json.Unmarshal([]byte(line), &r); err != nil {
61 t.Fatalf("bad record %q: %v", line, err)
62 }
63 out = append(out, r)
64 }
65 return out
66 }
67
68 func TestRecorderAppendsOrderedTimestampedRecords(t *testing.T) {
69 path := filepath.Join(t.TempDir(), "run.trajectory.jsonl")
70 inner := &capabilitySink{}
71 now := time.UnixMilli(1754500000000)
72 r, err := New(inner, path, func() time.Time { return now })
73 if err != nil {
74 t.Fatalf("New: %v", err)
75 }
76
77 r.Emit(event.Event{Kind: event.ToolDispatch, Tool: event.Tool{ID: "c1", Name: "bash", Args: `{"command":"ls"}`}})
78 now = now.Add(120 * time.Millisecond)
79 r.Emit(event.Event{Kind: event.ToolResult, Tool: event.Tool{
80 ID: "c1", Name: "bash", Output: "ok", DurationMs: 100,
81 StartedAt: 1754500000010, EndedAt: 1754500000110,
82 }})
83 r.Emit(event.Event{Kind: event.Reasoning, Text: "thinking about the next step"})
84 if err := r.Close(); err != nil {
85 t.Fatalf("Close: %v", err)
86 }
87
88 recs := readRecords(t, path)
89 if len(recs) != 3 {
90 t.Fatalf("got %d records, want 3", len(recs))
91 }
92 for i, rec := range recs {
93 if rec.Seq != uint64(i+1) {
94 t.Errorf("record %d seq = %d, want %d", i, rec.Seq, i+1)
95 }
96 if rec.SchemaVersion != SchemaVersion {
97 t.Errorf("record %d schema = %d, want %d", i, rec.SchemaVersion, SchemaVersion)
98 }
99 if rec.Event == nil {
100 t.Fatalf("record %d has no event payload", i)
101 }
102 }
103 if recs[0].TS != 1754500000000 || recs[1].TS != 1754500000120 {
104 t.Errorf("timestamps = %d, %d, want recorder-clock values", recs[0].TS, recs[1].TS)
105 }
106 if recs[1].Event.Tool == nil || recs[1].Event.Tool.StartedAt != 1754500000010 || recs[1].Event.Tool.EndedAt != 1754500000110 {
107 t.Errorf("tool result record lost execution bounds: %+v", recs[1].Event.Tool)
108 }
109 if recs[2].Event.Kind != "reasoning" || recs[2].Event.Text != "thinking about the next step" {
110 t.Errorf("reasoning record = %+v", recs[2].Event)
111 }
112 if len(inner.events) != 3 {
113 t.Errorf("inner sink saw %d events, want 3", len(inner.events))
114 }
115 }
116
117 func TestRecorderRecordsAndForwardsOptionalCapabilities(t *testing.T) {
118 path := filepath.Join(t.TempDir(), "run.trajectory.jsonl")
119 inner := &capabilitySink{}
120 r, err := New(inner, path, nil)
121 if err != nil {
122 t.Fatalf("New: %v", err)
123 }
124
125 r.RecordReadinessAudit(evidence.ReadinessAudit{Result: evidence.ReadinessBlocked, MissingVerification: 2})
126 r.RecordProtocolRecovery(event.ProtocolRecoveryAudit{Kind: event.ProtocolRecoveryMissingReasoningDetected})
127 r.RecordTurnCompletion()
128 r.RecordOutcomeProgress(evidence.OutcomeSample{
129 Round: 3, Exploration: 2, Objective: 1, LegacyGain: 4,
130 Runway: 0, RunwayDry: 6, RunwayIdle: 6, RunwaySpent: true,
131 })
132 r.RecordDelegationAdmission(event.DelegationAdmissionAudit{Tool: "research", Verdict: "deny", Reason: "local_fix_no_external_need", Intent: "mutation"})
133 r.RecordCompletionReport(event.CompletionReportAudit{
134 Verdict: "partial", Changes: 2, ChangesUnreviewed: 1, Gaps: 1, GapKinds: []string{"unreviewed_change"},
135 })
136 r.RecordWorkspaceMutation(event.WorkspaceMutation{ToolName: "write_file"})
137 r.RecordRunBudget(event.RunBudgetSample{Currency: "USD"})
138 if err := r.Close(); err != nil {
139 t.Fatalf("Close: %v", err)
140 }
141
142 recs := readRecords(t, path)
143 if len(recs) != 6 {
144 t.Fatalf("got %d records, want 6", len(recs))
145 }
146 if recs[0].ReadinessAudit == nil || recs[0].ReadinessAudit.Result != "blocked" || recs[0].ReadinessAudit.MissingVerification != 2 {
147 t.Errorf("readiness record = %+v", recs[0].ReadinessAudit)
148 }
149 if recs[1].ProtocolRecovery != string(event.ProtocolRecoveryMissingReasoningDetected) {
150 t.Errorf("protocol recovery record = %q", recs[1].ProtocolRecovery)
151 }
152 if !recs[2].TurnCompletion {
153 t.Errorf("turn completion record = %+v", recs[2])
154 }
155 if recs[3].OutcomeProgress == nil || recs[3].OutcomeProgress.Round != 3 || recs[3].OutcomeProgress.Objective != 1 || recs[3].OutcomeProgress.LegacyGain != 4 {
156 t.Errorf("outcome progress record = %+v", recs[3].OutcomeProgress)
157 }
158 if rec := recs[3].OutcomeProgress; rec.Runway == nil || *rec.Runway != 0 || rec.RunwayDry != 6 || rec.RunwayIdle != 6 || !rec.RunwaySpent {
159 t.Errorf("runway shadow record = %+v, want explicit spent balance", rec)
160 }
161 if recs[4].DelegationAdmission == nil || recs[4].DelegationAdmission.Verdict != "deny" || recs[4].DelegationAdmission.Tool != "research" {
162 t.Errorf("delegation admission record = %+v", recs[4].DelegationAdmission)
163 }
164 if rec := recs[5].CompletionReport; rec == nil || rec.Verdict != "partial" || rec.ChangesUnreviewed != 1 || len(rec.GapKinds) != 1 {
165 t.Errorf("completion report record = %+v", recs[5].CompletionReport)
166 }
167 if len(inner.workspace) != 1 || len(inner.runBudgets) != 1 {
168 t.Errorf("host capabilities not forwarded: workspace=%d run_budget=%d", len(inner.workspace), len(inner.runBudgets))
169 }
170 if len(inner.readiness) != 1 || len(inner.recoveries) != 1 || inner.turns != 1 || len(inner.outcomes) != 1 || len(inner.reports) != 1 {
171 t.Errorf("inner capabilities = %d/%d/%d/%d/%d, want 1/1/1/1/1", len(inner.readiness), len(inner.recoveries), inner.turns, len(inner.outcomes), len(inner.reports))
172 }
173 }
174
175 func TestOutcomeProgressRunwayIsAdditiveAndPresenceAware(t *testing.T) {
176 var old Record
177 if err := json.Unmarshal([]byte(`{"schema_version":1,"outcome_progress":{"round":1}}`), &old); err != nil {
178 t.Fatalf("decode old record: %v", err)
179 }
180 if old.OutcomeProgress == nil || old.OutcomeProgress.Runway != nil {
181 t.Fatalf("old record runway = %+v, want unobserved nil", old.OutcomeProgress)
182 }
183
184 zero := 0
185 data, err := json.Marshal(Record{
186 SchemaVersion: SchemaVersion,
187 OutcomeProgress: &OutcomeProgress{Round: 1, Runway: &zero, RunwaySpent: true},
188 })
189 if err != nil {
190 t.Fatalf("encode new record: %v", err)
191 }
192 if !strings.Contains(string(data), `"runway":0`) {
193 t.Fatalf("zero balance was omitted: %s", data)
194 }
195 // A previous reader ignores the additive fields and keeps its known data.
196 var legacy struct {
197 OutcomeProgress *struct {
198 Round int `json:"round"`
199 } `json:"outcome_progress"`
200 }
201 if err := json.Unmarshal(data, &legacy); err != nil || legacy.OutcomeProgress == nil || legacy.OutcomeProgress.Round != 1 {
202 t.Fatalf("legacy decode = %+v, %v", legacy, err)
203 }
204 }
205
206 func TestRecorderForwardsAfterCloseWithoutRecording(t *testing.T) {
207 path := filepath.Join(t.TempDir(), "run.trajectory.jsonl")
208 inner := &capabilitySink{}
209 r, err := New(inner, path, nil)
210 if err != nil {
211 t.Fatalf("New: %v", err)
212 }
213 r.Emit(event.Event{Kind: event.Text, Text: "before"})
214 if err := r.Close(); err != nil {
215 t.Fatalf("Close: %v", err)
216 }
217 r.Emit(event.Event{Kind: event.Text, Text: "after"})
218
219 if len(readRecords(t, path)) != 1 {
220 t.Fatalf("post-close event must not be recorded")
221 }
222 if len(inner.events) != 2 {
223 t.Fatalf("inner sink saw %d events, want 2 (forwarding survives Close)", len(inner.events))
224 }
225 }
226
227 func TestNewFailsOnUnwritablePath(t *testing.T) {
228 if _, err := New(event.Discard, filepath.Join(t.TempDir(), "missing", "run.jsonl"), nil); err == nil {
229 t.Fatal("New must fail when the parent directory does not exist")
230 }
231 }
232
232 lines GO