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