From 0a74d801a4e609bf339eb0a2f097b7a74853fcc4 Mon Sep 17 00:00:00 2001 From: Victor Garcia Date: Mon, 28 Sep 2026 00:27:12 -0600 Subject: [PATCH 1/3] feat(engine): record agent spend on every phase exit (U1) agent_end was emitted only when an agent phase entry succeeded, so any cost figure undercounted every failed, dead, or attempt-ending entry. Every entry that emits agent_start now closes it with agent_end carrying an outcome (passed, failed, died, aborted, cancelled, ceiling, send_budget) and an unmetered_sends count, before the event that closes the phase or the attempt. A runtime-error send's reported usage is now added to spend. A send killed in flight (watchdog, phase clock, ceiling, cancel) or ending without a usable result is counted as unmetered, not as free. Claude Code's total_cost_usd is the session's running total across --resume, not the invocation's cost: a real haiku resume reported 0.0585011 after a first send of 0.0405199. The adapter now records the session's reported total on runtime.Session and returns each send's own share, so the engine's per-send sum is correct. Co-Authored-By: Claude Opus 5.5 --- internal/engine/phase.go | 81 +++- internal/engine/phase_test.go | 387 +++++++++++++++++++- internal/protocol/types.go | 13 + internal/runtime/claudecode/adapter.go | 16 +- internal/runtime/claudecode/adapter_test.go | 16 +- internal/runtime/runtime.go | 18 +- 6 files changed, 501 insertions(+), 30 deletions(-) diff --git a/internal/engine/phase.go b/internal/engine/phase.go index 4160be7..716b43e 100644 --- a/internal/engine/phase.go +++ b/internal/engine/phase.go @@ -688,6 +688,24 @@ func terminalSend(kind sendEnd, detail string) *phaseRun { return nil } +// agentOutcome names how an agent phase entry ended, for its agent_end. +// A death is named at the death itself: it is a phaseFailed run here. +func (run phaseRun) agentOutcome() string { + switch run.outcome { + case phasePassed: + return protocol.AgentPassed + case phaseAborted: + return protocol.AgentAborted + case phaseCancelled: + return protocol.AgentCancelled + case phaseCeiling: + return protocol.AgentCeiling + case phaseSendBudget: + return protocol.AgentSendBudget + } + return protocol.AgentFailed +} + func (run phaseRun) envelopeRef() *parsedEnvelope { if !run.hasEnvelope { return nil @@ -1000,10 +1018,26 @@ func (e *execution) runAgentPhaseAttempt( "model": role.Model, "session": session.Key, "can_resume": e.capability.CanResume, "effort": role.Effort, "budget_usd": role.BudgetUSD, "tools": role.Tools, "phase_attempt": entry, }) + // agentEnd closes the agent_start above on every exit, so an entry's + // spend is recorded whether it passed, failed, died, or ended the attempt + // (R1). Spend accumulated across every send that returned a result; + // unmetered sends were killed first, and their cost is unknown, not zero. + // Occupancy is the LAST metered send's number — the one whose context is + // current (sssf tracer discipline). + agentEnd := func(outcome string) { + e.emit.emit(protocol.EventAgentEnd, phase.Name, phase.Owner, map[string]any{ + "outcome": outcome, + "tokens": sender.spend.TotalTokens, "cost": sender.spend.CostUSD, + "context_tokens": sender.last.Usage.ContextTokens, + "sends": sender.sendCount, + "unmetered_sends": sender.unmetered, + }) + } // death wraps up one dead entry: roll the worktree back to the pre-phase // snapshot and emit the envelope-less terminal trace event (R11). death := func(detail string) (phaseRun, bool) { + agentEnd(protocol.AgentDied) e.emit.emit(protocol.EventPhaseDeath, phase.Name, phase.Owner, map[string]any{ "phase_attempt": entry, "error": detail, }) @@ -1026,8 +1060,10 @@ func (e *execution) runAgentPhaseAttempt( if _, terminal := e.enforceWriteBoundary( ctx, phase, entry, started, before, role.Writes); terminal != nil { terminal.failure = detail + " — and " + terminal.failure + agentEnd(terminal.agentOutcome()) return *terminal, false } + agentEnd(protocol.AgentFailed) e.recordResult(protocol.PhaseResult{ Phase: phase.Name, Kind: phase.Kind, Status: protocol.EnvelopeFail, PhaseAttempt: entry, Error: detail, StartedAt: &started, @@ -1052,6 +1088,7 @@ func (e *execution) runAgentPhaseAttempt( if sender.sendCount == 0 { // The budget check refuses before the subprocess starts: no agent // ran in this entry, so there is nothing of its to enforce. + agentEnd(run.agentOutcome()) return run, false } // Detached on purpose: cancellation is one of the exits this guards, @@ -1061,8 +1098,10 @@ func (e *execution) runAgentPhaseAttempt( if _, breach := e.enforceWriteBoundary( context.WithoutCancel(ctx), phase, entry, started, before, role.Writes); breach != nil { breach.failure = detail + " — and " + breach.failure + agentEnd(breach.agentOutcome()) return *breach, false } + agentEnd(run.agentOutcome()) return run, false } @@ -1127,6 +1166,7 @@ func (e *execution) runAgentPhaseAttempt( // which it wrote somewhere it was not allowed to (R10). touched, terminal := e.enforceWriteBoundary(ctx, phase, entry, started, before, role.Writes) if terminal != nil { + agentEnd(terminal.agentOutcome()) return *terminal, false } for _, path := range touched { @@ -1141,18 +1181,14 @@ func (e *execution) runAgentPhaseAttempt( e.emit.emit(protocol.EventHandoff, phase.Name, phase.Owner, map[string]any{ "artifacts": envelope.Base.Artifacts, "summary": envelope.Base.Summary, }) - // Spend accumulated across every send; occupancy is the LAST send's - // number — the one whose context is current (sssf tracer discipline). - e.emit.emit(protocol.EventAgentEnd, phase.Name, phase.Owner, map[string]any{ - "tokens": sender.spend.TotalTokens, "cost": sender.spend.CostUSD, - "context_tokens": sender.last.Usage.ContextTokens, - "sends": sender.sendCount, - }) status := protocol.EnvelopeFail + outcome := protocol.AgentFailed if envelope.Base.Status == protocol.EnvelopeSuccess { status = protocol.EnvelopeSuccess + outcome = protocol.AgentPassed } + agentEnd(outcome) e.recordResult(protocol.PhaseResult{ Phase: phase.Name, Kind: phase.Kind, Status: status, PhaseAttempt: entry, Envelope: envelope.Raw, Gates: e.lastMergedReport(phase), StartedAt: &started, @@ -1237,6 +1273,9 @@ type agentSender struct { spend runtime.Usage last runtime.Result sendCount int + // unmetered counts started sends whose result was never read — killed, + // or ended without a usable result — so their cost is unknown (R1). + unmetered int } func (s *agentSender) send(ctx context.Context, prompt string) (runtime.Result, sendEnd, string) { @@ -1286,25 +1325,33 @@ func (s *agentSender) send(ctx context.Context, prompt string) (runtime.Result, resetTimer(silence, s.timeouts.silence) s.emitRuntimeEvent(event) case <-silence.C: - killAndDrain(handle) + s.abandon(handle) return runtime.Result{}, sendDeath, fmt.Sprintf( "no output for %s (watchdog)", s.timeouts.silence) case <-phaseClock.C: - killAndDrain(handle) + s.abandon(handle) return runtime.Result{}, sendDeath, fmt.Sprintf( "phase wall clock (%s) exceeded", s.timeouts.phase) case <-ceilingClock.C: - killAndDrain(handle) + s.abandon(handle) return runtime.Result{}, sendCeiling, "attempt wall-clock ceiling exceeded" case <-s.attempt.Cancelled: - killAndDrain(handle) + s.abandon(handle) return runtime.Result{}, sendCancelled, "cancelled during phase " + s.phase } } result, err := handle.Result() if err != nil { + s.unmetered++ return runtime.Result{}, sendDeath, "agent subprocess: " + err.Error() } + // Metered before the error check: a runtime error still reports what it + // cost (KTD2). + s.spend.InputTokens += result.Usage.InputTokens + s.spend.OutputTokens += result.Usage.OutputTokens + s.spend.TotalTokens += result.Usage.TotalTokens + s.spend.CostUSD += result.Usage.CostUSD + s.last = result if result.IsError { // The runtime's own terminal-error flag, checked before anything tries // to read an envelope out of the text. What comes back here is the @@ -1315,11 +1362,6 @@ func (s *agentSender) send(ctx context.Context, prompt string) (runtime.Result, result.ExitCode, truncateText(strings.TrimSpace(result.Text), protocol.MaxCommandOutputTailBytes)) } - s.spend.InputTokens += result.Usage.InputTokens - s.spend.OutputTokens += result.Usage.OutputTokens - s.spend.TotalTokens += result.Usage.TotalTokens - s.spend.CostUSD += result.Usage.CostUSD - s.last = result s.transcripts[s.role] = append(s.transcripts[s.role], exchange{Prompt: prompt, Response: result.Text}) return result, sendOK, "" } @@ -1347,6 +1389,13 @@ func (s *agentSender) recordProcess(groupID int64, active bool) { } } +// abandon kills a send in flight. Its result is never read, so it is counted +// as unmetered rather than as free (R1). +func (s *agentSender) abandon(handle runtime.Handle) { + killAndDrain(handle) + s.unmetered++ +} + // killAndDrain stops the process group and drains the event stream so the // adapter's reader goroutine can finish; the result is discarded — a killed // send never yields an envelope. diff --git a/internal/engine/phase_test.go b/internal/engine/phase_test.go index 3832d6f..50f0bf3 100644 --- a/internal/engine/phase_test.go +++ b/internal/engine/phase_test.go @@ -781,7 +781,6 @@ acceptance: [all_phases_passed] t.Fatalf("state = %q (%s), want accepted_unpublished", outcome.State, outcome.Error) } - type triple struct{ Type, Phase, Name string } want := []triple{ {protocol.EventLog, "", "attempt_start"}, {protocol.EventPhaseStart, "plan", ""}, @@ -813,6 +812,15 @@ acceptance: [all_phases_passed] {protocol.EventPhaseEnd, "report", ""}, {protocol.EventLog, "", "acceptance"}, } + requireSequence(t, sink, want) + passed := agentEndPayload{Outcome: protocol.AgentPassed, Sends: 1} + requireAgentEnds(t, sink, passed, passed, passed) +} + +type triple struct{ Type, Phase, Name string } + +func requireSequence(t *testing.T, sink *recordingSink, want []triple) { + t.Helper() events := sink.all() if len(events) != len(want) { var got []string @@ -832,6 +840,52 @@ acceptance: [all_phases_passed] } } +// The exits the three-phase chain never takes: a dead entry closes its +// agent_start before its phase_death, and an entry refused its first send by +// the attempt's send budget closes its agent_start before the attempt's +// terminal error. +func TestDeathAndSendBudgetExitsEmitTheExactExpectedEventSequence(t *testing.T) { + repo := initRepo(t) + fake := enginetest.New( + enginetest.Step{Crash: true}, + enginetest.Step{Text: envelope(map[string]any{"status": "success", "summary": "built"})}, + ) + sink := &recordingSink{} + runner := newTestRunner(t, fake, sink, func(config *engine.Config) { + config.MaxAttemptSends = 2 + }) + + snapshot := "name: t\n" + repairRoster + ` +phases: + - {name: build, kind: agent, owner: builder} + - {name: review, kind: agent, owner: builder} +` + if outcome := runner.Execute(context.Background(), testAttempt(snapshot, nil, repo)); outcome.State != protocol.AttemptFailed { + t.Fatalf("state = %q, want failed", outcome.State) + } + requireSequence(t, sink, []triple{ + {protocol.EventLog, "", "attempt_start"}, + {protocol.EventPhaseStart, "build", ""}, + {protocol.EventAgentStart, "build", "builder"}, + {protocol.EventAgentEnd, "build", "builder"}, + {protocol.EventPhaseDeath, "build", "builder"}, + {protocol.EventPhaseStart, "build", ""}, + {protocol.EventAgentStart, "build", "builder"}, + {protocol.EventHandoff, "build", "builder"}, + {protocol.EventAgentEnd, "build", "builder"}, + {protocol.EventPhaseEnd, "build", ""}, + {protocol.EventPhaseStart, "review", ""}, + {protocol.EventAgentStart, "review", "builder"}, + {protocol.EventAgentEnd, "review", "builder"}, + {protocol.EventError, "", "attempt_send_budget_exhausted"}, + }) + requireAgentEnds(t, sink, + agentEndPayload{Outcome: protocol.AgentDied, Sends: 1, UnmeteredSends: 1}, + agentEndPayload{Outcome: protocol.AgentPassed, Sends: 1}, + agentEndPayload{Outcome: protocol.AgentSendBudget}, + ) +} + // ---- cancellation ---------------------------------------------------------- func TestCancellationDuringAnAgentPhaseKillsTheSendAndReportsCancelled(t *testing.T) { @@ -1705,3 +1759,334 @@ publish: t.Fatalf("outcome = %+v, want work with no risk report held", outcome) } } + +// ---- scenario: spend on every exit (R1, KTD2) ------------------------------ + +type agentEndPayload struct { + Outcome string `json:"outcome"` + Cost float64 `json:"cost"` + Tokens int `json:"tokens"` + Sends int `json:"sends"` + UnmeteredSends int `json:"unmetered_sends"` +} + +func (s *recordingSink) agentEnds(t *testing.T) []agentEndPayload { + t.Helper() + var ends []agentEndPayload + for _, event := range s.all() { + if event.Type != protocol.EventAgentEnd { + continue + } + var payload agentEndPayload + if err := json.Unmarshal(event.Payload, &payload); err != nil { + t.Fatalf("agent_end payload %s: %v", event.Payload, err) + } + ends = append(ends, payload) + } + return ends +} + +// index is the position of the first event of the type and name, or -1. +func (s *recordingSink) index(eventType, name string) int { + for i, event := range s.all() { + if event.Type == eventType && (name == "" || event.Name == name) { + return i + } + } + return -1 +} + +func requireAgentEnds(t *testing.T, sink *recordingSink, want ...agentEndPayload) { + t.Helper() + got := sink.agentEnds(t) + if len(got) != len(want) { + t.Fatalf("agent_end events = %+v, want %+v", got, want) + } + for i := range want { + if got[i] != want[i] { + t.Fatalf("agent_end %d = %+v, want %+v", i, got[i], want[i]) + } + } +} + +// requireBefore pins that the entry's agent_end precedes the event that +// closes the phase or the attempt, so spend is never recorded after the trace +// says it is over. +func requireBefore(t *testing.T, sink *recordingSink, eventType, name string) { + t.Helper() + end, closing := sink.index(protocol.EventAgentEnd, ""), sink.index(eventType, name) + if end < 0 || closing < 0 || end > closing { + t.Fatalf("agent_end at %d, %s/%s at %d: agent_end must come first", end, eventType, name, closing) + } +} + +func TestAPassingAgentPhaseRecordsItsSpend(t *testing.T) { + repo := initRepo(t) + fake := enginetest.New(enginetest.Step{ + Text: envelope(map[string]any{"status": "success", "summary": "done"}), + Usage: runtime.Usage{TotalTokens: 120, CostUSD: 0.25}, + }) + sink := &recordingSink{} + runner := newTestRunner(t, fake, sink, nil) + + snapshot := "name: t\n" + repairRoster + ` +phases: + - {name: build, kind: agent, owner: builder} +` + outcome := runner.Execute(context.Background(), testAttempt(snapshot, nil, repo)) + if outcome.State != protocol.AttemptAcceptedUnpublished { + t.Fatalf("state = %q (%s), want accepted_unpublished", outcome.State, outcome.Error) + } + requireAgentEnds(t, sink, agentEndPayload{ + Outcome: protocol.AgentPassed, Cost: 0.25, Tokens: 120, Sends: 1, + }) +} + +func TestAnAgentReportingFailStillRecordsItsSpend(t *testing.T) { + repo := initRepo(t) + fake := enginetest.New(enginetest.Step{ + Text: envelope(map[string]any{"status": "fail", "summary": "I could not do it"}), + Usage: runtime.Usage{TotalTokens: 80, CostUSD: 0.5}, + }) + sink := &recordingSink{} + runner := newTestRunner(t, fake, sink, nil) + + snapshot := "name: t\n" + repairRoster + ` +phases: + - {name: build, kind: agent, owner: builder} +` + if outcome := runner.Execute(context.Background(), testAttempt(snapshot, nil, repo)); outcome.State != protocol.AttemptFailed { + t.Fatalf("state = %q, want failed", outcome.State) + } + requireAgentEnds(t, sink, agentEndPayload{ + Outcome: protocol.AgentFailed, Cost: 0.5, Tokens: 80, Sends: 1, + }) + requireBefore(t, sink, protocol.EventPhaseEnd, "") +} + +// A runtime error ends the phase before anything is parsed, but the CLI still +// reported what the send cost, and that money was spent. +func TestARuntimeErrorSendContributesItsReportedUsage(t *testing.T) { + repo := initRepo(t) + fake := enginetest.New(enginetest.Step{ + IsError: true, + Text: "API Error: 529 overloaded", + Usage: runtime.Usage{TotalTokens: 30, CostUSD: 0.125}, + }) + sink := &recordingSink{} + runner := newTestRunner(t, fake, sink, nil) + + snapshot := "name: t\n" + repairRoster + ` +phases: + - {name: build, kind: agent, owner: builder} +` + if outcome := runner.Execute(context.Background(), testAttempt(snapshot, nil, repo)); outcome.State != protocol.AttemptFailed { + t.Fatalf("state = %q, want failed", outcome.State) + } + requireAgentEnds(t, sink, agentEndPayload{ + Outcome: protocol.AgentFailed, Cost: 0.125, Tokens: 30, Sends: 1, + }) + requireBefore(t, sink, protocol.EventPhaseEnd, "") +} + +// Each entry of a phase closes its own agent_start. The killed send never +// returned a result, so its scripted usage is never read: it counts as +// unmetered, not as free and not as its would-be cost. +func TestADeadEntryReportsItsKilledSendAsUnmeteredAndTheReentryItsOwnSpend(t *testing.T) { + repo := initRepo(t) + fake := enginetest.New( + enginetest.Step{Hang: true, Usage: runtime.Usage{TotalTokens: 999, CostUSD: 9}}, + enginetest.Step{ + Text: envelope(map[string]any{"status": "success", "summary": "recovered"}), + Usage: runtime.Usage{TotalTokens: 40, CostUSD: 0.25}, + }, + ) + sink := &recordingSink{} + runner := newTestRunner(t, fake, sink, func(config *engine.Config) { + config.NoOutputTimeout = 100 * time.Millisecond + }) + + snapshot := "name: t\n" + repairRoster + ` +phases: + - {name: build, kind: agent, owner: builder} +` + outcome := runner.Execute(context.Background(), testAttempt(snapshot, nil, repo)) + if outcome.State != protocol.AttemptAcceptedUnpublished { + t.Fatalf("state = %q (%s), want accepted_unpublished", outcome.State, outcome.Error) + } + requireAgentEnds(t, sink, + agentEndPayload{Outcome: protocol.AgentDied, Sends: 1, UnmeteredSends: 1}, + agentEndPayload{Outcome: protocol.AgentPassed, Cost: 0.25, Tokens: 40, Sends: 1}, + ) + requireBefore(t, sink, protocol.EventPhaseDeath, "") +} + +func TestASendBudgetExhaustionRecordsSpendBeforeItsTerminalEvent(t *testing.T) { + repo := initRepo(t) + const cap = 3 + steps := make([]enginetest.Step, 0, cap) + for i := 0; i < cap; i++ { + steps = append(steps, enginetest.Step{ + Text: "never parseable " + fmt.Sprint(i), + Usage: runtime.Usage{TotalTokens: 10, CostUSD: 0.125}, + }) + } + fake := enginetest.New(steps...) + sink := &recordingSink{} + runner := newTestRunner(t, fake, sink, func(config *engine.Config) { + config.MaxAttemptSends = cap + }) + + snapshot := "name: t\n" + repairRoster + ` +phases: + - {name: build, kind: agent, owner: builder} +` + if outcome := runner.Execute(context.Background(), testAttempt(snapshot, nil, repo)); outcome.State != protocol.AttemptFailed { + t.Fatalf("state = %q, want failed", outcome.State) + } + requireAgentEnds(t, sink, agentEndPayload{ + Outcome: protocol.AgentSendBudget, Cost: 0.375, Tokens: 30, Sends: cap, + }) + requireBefore(t, sink, protocol.EventError, "attempt_send_budget_exhausted") +} + +// The send the budget refuses never starts, so an entry that opens on an +// exhausted budget still closes its agent_start, with nothing sent. +func TestAnEntryRefusedItsFirstSendStillClosesItsAgentStart(t *testing.T) { + repo := initRepo(t) + fake := enginetest.New(enginetest.Step{ + Text: envelope(map[string]any{"status": "success", "summary": "planned"}), + Usage: runtime.Usage{TotalTokens: 10, CostUSD: 0.5}, + }) + sink := &recordingSink{} + runner := newTestRunner(t, fake, sink, func(config *engine.Config) { + config.MaxAttemptSends = 1 + }) + + snapshot := "name: t\n" + repairRoster + ` +phases: + - {name: plan, kind: agent, owner: builder} + - {name: build, kind: agent, owner: builder} +` + if outcome := runner.Execute(context.Background(), testAttempt(snapshot, nil, repo)); outcome.State != protocol.AttemptFailed { + t.Fatalf("state = %q, want failed", outcome.State) + } + requireAgentEnds(t, sink, + agentEndPayload{Outcome: protocol.AgentPassed, Cost: 0.5, Tokens: 10, Sends: 1}, + agentEndPayload{Outcome: protocol.AgentSendBudget}, + ) + if got := sink.count(protocol.EventAgentStart, ""); got != 2 { + t.Fatalf("agent_start events = %d, want 2", got) + } +} + +func TestACeilingHitRecordsTheKilledSendBeforeItsTerminalEvent(t *testing.T) { + repo := initRepo(t) + fake := enginetest.New(enginetest.Step{Hang: true, Usage: runtime.Usage{CostUSD: 9}}) + sink := &recordingSink{} + runner := newTestRunner(t, fake, sink, func(config *engine.Config) { + config.AttemptCeiling = 100 * time.Millisecond + config.NoOutputTimeout = time.Minute // the ceiling must fire first + }) + + snapshot := "name: t\n" + repairRoster + ` +phases: + - {name: build, kind: agent, owner: builder} +` + if outcome := runner.Execute(context.Background(), testAttempt(snapshot, nil, repo)); outcome.State != protocol.AttemptFailed { + t.Fatalf("state = %q, want failed", outcome.State) + } + requireAgentEnds(t, sink, agentEndPayload{ + Outcome: protocol.AgentCeiling, Sends: 1, UnmeteredSends: 1, + }) + requireBefore(t, sink, protocol.EventError, "attempt_ceiling_exceeded") +} + +func TestACancelledEntryRecordsTheKilledSend(t *testing.T) { + repo := initRepo(t) + fake := enginetest.New(enginetest.Step{Hang: true, Usage: runtime.Usage{CostUSD: 9}}) + cancelled := make(chan struct{}) + go func() { + time.Sleep(100 * time.Millisecond) + close(cancelled) + }() + sink := &recordingSink{} + runner := newTestRunner(t, fake, sink, nil) + + snapshot := "name: t\n" + repairRoster + ` +phases: + - {name: build, kind: agent, owner: builder} +` + attempt := testAttempt(snapshot, nil, repo) + attempt.Cancelled = cancelled + if outcome := runner.Execute(context.Background(), attempt); outcome.State != protocol.AttemptCancelled { + t.Fatalf("state = %q, want cancelled", outcome.State) + } + requireAgentEnds(t, sink, agentEndPayload{ + Outcome: protocol.AgentCancelled, Sends: 1, UnmeteredSends: 1, + }) +} + +// An entry that breaches the write boundary is aborted, and what it spent +// getting there is recorded all the same. +func TestAnAbortedEntryRecordsItsSpend(t *testing.T) { + repo := initRepo(t) + fake := enginetest.New(enginetest.Step{ + Files: map[string]string{"stray.txt": "outside the allowlist"}, + Text: envelope(map[string]any{"status": "success", "summary": "done"}), + Usage: runtime.Usage{TotalTokens: 10, CostUSD: 0.5}, + }) + sink := &recordingSink{} + runner := newTestRunner(t, fake, sink, nil) + + snapshot := "name: t\n" + repairRoster + ` +phases: + - {name: build, kind: agent, owner: builder} +` + if outcome := runner.Execute(context.Background(), testAttempt(snapshot, nil, repo)); outcome.State != protocol.AttemptFailed { + t.Fatalf("state = %q, want failed", outcome.State) + } + requireAgentEnds(t, sink, agentEndPayload{ + Outcome: protocol.AgentAborted, Cost: 0.5, Tokens: 10, Sends: 1, + }) +} + +// The unit's verification: over a scripted attempt, agent_end cost sums to +// the usage of every send that returned a result, and unmetered_sends to the +// sends that did not — here a crash, an unparseable emission, its parse +// correction, and a runtime error in a later phase. +func TestAgentEndSpendSumsToEveryMeteredSendOverAnAttempt(t *testing.T) { + repo := initRepo(t) + fake := enginetest.New( + enginetest.Step{Crash: true, Usage: runtime.Usage{CostUSD: 9}}, + enginetest.Step{Text: "not an envelope", Usage: runtime.Usage{TotalTokens: 1, CostUSD: 0.5}}, + enginetest.Step{ + Text: envelope(map[string]any{"status": "success", "summary": "built"}), + Usage: runtime.Usage{TotalTokens: 2, CostUSD: 0.25}, + }, + enginetest.Step{IsError: true, Text: "rate limited", Usage: runtime.Usage{TotalTokens: 4, CostUSD: 0.125}}, + ) + sink := &recordingSink{} + runner := newTestRunner(t, fake, sink, nil) + + snapshot := "name: t\n" + repairRoster + ` +phases: + - {name: build, kind: agent, owner: builder} + - {name: review, kind: agent, owner: builder} +` + if outcome := runner.Execute(context.Background(), testAttempt(snapshot, nil, repo)); outcome.State != protocol.AttemptFailed { + t.Fatalf("state = %q, want failed", outcome.State) + } + var cost float64 + var tokens, sends, unmetered int + for _, end := range sink.agentEnds(t) { + cost += end.Cost + tokens += end.Tokens + sends += end.Sends + unmetered += end.UnmeteredSends + } + if cost != 0.875 || tokens != 7 || unmetered != 1 || sends != len(fake.Calls()) { + t.Fatalf("summed agent_end: cost %v tokens %d unmetered %d sends %d; want 0.875, 7, 1, %d", + cost, tokens, unmetered, sends, len(fake.Calls())) + } +} diff --git a/internal/protocol/types.go b/internal/protocol/types.go index 4467d20..8be8c79 100644 --- a/internal/protocol/types.go +++ b/internal/protocol/types.go @@ -80,6 +80,19 @@ const ( EventPhaseDeath = "phase_death" ) +// Agent outcomes: how an agent phase entry ended, carried on its agent_end +// event. Every entry that emitted agent_start emits exactly one agent_end, so +// summing agent_end cost over an attempt is its whole metered spend (R1). +const ( + AgentPassed = "passed" + AgentFailed = "failed" + AgentDied = "died" + AgentAborted = "aborted" + AgentCancelled = "cancelled" + AgentCeiling = "ceiling" + AgentSendBudget = "send_budget" +) + // Envelope statuses. Success must be earned: everything that constructs a // phase or envelope result starts from fail. const ( diff --git a/internal/runtime/claudecode/adapter.go b/internal/runtime/claudecode/adapter.go index d7d1e36..ebe18da 100644 --- a/internal/runtime/claudecode/adapter.go +++ b/internal/runtime/claudecode/adapter.go @@ -319,6 +319,7 @@ func (a *Adapter) StartOrContinue( watchdog: watchdog, groupID: groupID, resumedID: resumedID, + session: session, events: make(chan runtime.Event, eventChannelDepth), done: make(chan struct{}), stderrDone: make(chan struct{}), @@ -360,8 +361,11 @@ type claudeHandle struct { // resumedID is the session this send was told to continue, empty on a // creating send. Result compares it against what the CLI reports back. resumedID string - events chan runtime.Event - done chan struct{} + // session is where Result records the CLI's running cost total, so the + // next send on the conversation can report only its own share. + session *runtime.Session + events chan runtime.Event + done chan struct{} // stderrDone closes when the stderr capture goroutine has finished. stderrDone chan struct{} stopped chan struct{} @@ -586,6 +590,14 @@ func (h *claudeHandle) Result() (runtime.Result, error) { h.finalErr = fmt.Errorf("%w: asked to resume %s, the CLI answered under %s", ErrSessionDiscontinuity, h.resumedID, h.result.SessionID) } + if h.resultSeen && !errors.Is(h.finalErr, ErrSessionDiscontinuity) { + // total_cost_usd is the session's running total across --resume, not + // this invocation's cost (verified against claude 2.1.283: a resume + // reported the first send's cost plus its own). Usage is per send. + total := h.result.Usage.CostUSD + h.result.Usage.CostUSD = total - h.session.ReportedCostUSD + h.session.ReportedCostUSD = total + } return h.result, h.finalErr } diff --git a/internal/runtime/claudecode/adapter_test.go b/internal/runtime/claudecode/adapter_test.go index 59a4a32..0aea97f 100644 --- a/internal/runtime/claudecode/adapter_test.go +++ b/internal/runtime/claudecode/adapter_test.go @@ -23,7 +23,8 @@ import ( // stubScript stands in for the CLI. It echoes back the session id it was // given (--session-id or --resume), the way the real CLI does, so session // continuity is observable; STUB_MODE=coldstart makes it answer under a -// different session instead. +// different session instead. STUB_TOTAL_COST is the session's running cost +// total the result reports, the way the real CLI reports total_cost_usd. const stubScript = `#!/bin/sh case "$1" in --version) echo "9.9.9 (Claude Code)"; exit 0 ;; @@ -65,7 +66,7 @@ case "$STUB_MODE" in printf '{"type":"result","result":"survived the flood","is_error":false,"session_id":"%s"}\n' "$SESSION" ;; *) printf '%s\n' '{"type":"assistant","message":{"content":[{"type":"tool_use","name":"Write","input":{"path":"x.txt"}},{"type":"text","text":"working on it"}]}}' - printf '{"type":"result","result":"hello from stub","is_error":false,"total_cost_usd":0.5,"session_id":"%s","usage":{"input_tokens":5,"output_tokens":7,"cache_read_input_tokens":2}}\n' "$SESSION" + printf '{"type":"result","result":"hello from stub","is_error":false,"total_cost_usd":%s,"session_id":"%s","usage":{"input_tokens":5,"output_tokens":7,"cache_read_input_tokens":2}}\n' "${STUB_TOTAL_COST:-0.5}" "$SESSION" ;; esac ` @@ -185,15 +186,22 @@ func TestFirstSendCreatesTheSessionAndTheSecondResumesIt(t *testing.T) { t.Errorf("tool_call event missing: %s", joined) } - // Second send: same session continues via --resume. + // Second send: same session continues via --resume. The CLI reports the + // session's running total, so the send's own cost is the difference. + opts.Env = append(opts.Env, "STUB_TOTAL_COST=1.25") handle, err = adapter.StartOrContinue(context.Background(), session, "second prompt", opts) if err != nil { t.Fatal(err) } drain(t, handle) - if _, err := handle.Result(); err != nil { + result, err = handle.Result() + if err != nil { t.Fatal(err) } + if result.Usage.CostUSD != 0.75 || session.ReportedCostUSD != 1.25 { + t.Fatalf("resumed send cost = %v (session total %v), want 0.75 of 1.25", + result.Usage.CostUSD, session.ReportedCostUSD) + } if session.Sends != 2 { t.Fatalf("session sends = %d, want 2", session.Sends) } diff --git a/internal/runtime/runtime.go b/internal/runtime/runtime.go index 02a54ca..5ea2cfb 100644 --- a/internal/runtime/runtime.go +++ b/internal/runtime/runtime.go @@ -41,9 +41,10 @@ type Event struct { } // Usage is one send's cost accounting. Spend accumulates across retries at -// the engine (every send costs); ContextTokens is occupancy — only the last -// send's value is current. Runtimes that do not report cost leave zeros -// (capability flag ReportsCost=false). +// the engine (every send costs), so CostUSD is this send's cost alone, never +// a session running total; ContextTokens is occupancy — only the last send's +// value is current. Runtimes that do not report cost leave zeros (capability +// flag ReportsCost=false). type Usage struct { InputTokens int OutputTokens int @@ -71,11 +72,14 @@ type Options struct { // engine-chosen stable name (attempt-scoped by declaration, KTD5); NativeID // is the per-CLI identity the adapter mints on the first send and reuses to // continue. Sends counts completed sends so create-or-continue is a property -// of the session, not a separate method. +// of the session, not a separate method. ReportedCostUSD is the running total +// the CLI last reported for this conversation, kept by adapters whose CLI +// reports cost that way so they can hand back a per-send Usage.CostUSD. type Session struct { - Key string - NativeID string - Sends int + Key string + NativeID string + Sends int + ReportedCostUSD float64 } // Result is one send's terminal outcome: the final response text plus From 62c41a9808603b067f016fba3b072315ad11d9e7 Mon Sep 17 00:00:00 2001 From: Victor Garcia Date: Mon, 28 Sep 2026 00:27:49 -0600 Subject: [PATCH 2/3] refactor(runtime): drop the unobservable discontinuity guard on cost A discontinuity is a send death and the engine drops that session, so the recorded total is never read again. Co-Authored-By: Claude Opus 5.5 --- internal/runtime/claudecode/adapter.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/internal/runtime/claudecode/adapter.go b/internal/runtime/claudecode/adapter.go index ebe18da..d8c287d 100644 --- a/internal/runtime/claudecode/adapter.go +++ b/internal/runtime/claudecode/adapter.go @@ -590,7 +590,7 @@ func (h *claudeHandle) Result() (runtime.Result, error) { h.finalErr = fmt.Errorf("%w: asked to resume %s, the CLI answered under %s", ErrSessionDiscontinuity, h.resumedID, h.result.SessionID) } - if h.resultSeen && !errors.Is(h.finalErr, ErrSessionDiscontinuity) { + if h.resultSeen { // total_cost_usd is the session's running total across --resume, not // this invocation's cost (verified against claude 2.1.283: a resume // reported the first send's cost plus its own). Usage is per send. From 3589b380035a0cd593ffb6275c0dc5891d57a7e3 Mon Sep 17 00:00:00 2001 From: Victor Garcia Date: Mon, 28 Sep 2026 00:32:09 -0600 Subject: [PATCH 3/3] test(engine): pin agent_end on the breach exits of failed and terminal entries Mutation testing showed the abort paths inside fail and terminalExit could drop agent_end without a test noticing. Co-Authored-By: Claude Opus 5.5 --- internal/engine/phase_test.go | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/internal/engine/phase_test.go b/internal/engine/phase_test.go index 50f0bf3..f1a043c 100644 --- a/internal/engine/phase_test.go +++ b/internal/engine/phase_test.go @@ -1162,6 +1162,11 @@ phases: if fake.Remaining() != 0 { t.Fatalf("unconsumed scripted steps: %d", fake.Remaining()) } + // The breach turns the failure into an abort, and the entry's sends are + // still recorded. + requireAgentEnds(t, sink, agentEndPayload{ + Outcome: protocol.AgentAborted, Sends: protocol.ParseBudgetPerEmission + 1, + }) } // ---- scenario: the send ladder is really bounded --------------------------- @@ -1345,6 +1350,9 @@ phases: if _, err := os.Stat(filepath.Join(repo, "src", "ok.txt")); err != nil { t.Fatalf("the allowed write was discarded on the way out: %v", err) } + requireAgentEnds(t, sink, agentEndPayload{ + Outcome: protocol.AgentAborted, Sends: 1, UnmeteredSends: 1, + }) } func TestACancelledAttemptStillEnforcesTheWriteBoundary(t *testing.T) {