Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
81 changes: 65 additions & 16 deletions internal/engine/phase.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
})
Expand All @@ -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,
Expand All @@ -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,
Expand All @@ -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
}

Expand Down Expand Up @@ -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 {
Expand All @@ -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,
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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
Expand All @@ -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, ""
}
Expand Down Expand Up @@ -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.
Expand Down
Loading
Loading