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
83 changes: 83 additions & 0 deletions pkg/agent/errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,33 @@ var (
// }
ErrLLMAuth = errors.New("agent: LLM provider authentication or configuration failed")

// ErrLLMTruncated is returned when the provider stopped generating
// before the response was finished — an output-token cap, typically.
//
// The danger is that a truncated response is a *valid prefix*: the text
// looks fine and only fails when something downstream parses it, so the
// operator reads "unexpected end of JSON input" and inspects a schema
// that was never wrong. Providers detect the cut at the source, so route
// on the sentinel instead of matching decode-error text.
//
// The cut is a function of this response's length rather than of the
// request, so isRetryable leaves it retryable — but a truncation that
// already streamed content is not retried, because re-running it would
// replay the same text into the consumer's stream. The real fix is a
// larger provider token cap or a smaller ask, both caller decisions.
ErrLLMTruncated = errors.New("agent: LLM response truncated before completion")

// ErrLLMContentBlocked is returned when the provider stopped generating
// for a content policy: safety, recitation, a blocklist, or an
// unsupported language.
//
// Deliberately distinct from ErrLLMTruncated because the two demand
// opposite responses. A truncation is length-dependent and may pass on
// the next attempt; a policy stop is deterministic for a given prompt,
// so every retry reproduces it. isRetryable treats it as terminal — the
// caller must change the request or surface the block to the user.
ErrLLMContentBlocked = errors.New("agent: LLM stopped generating for a content policy")

// ErrContextCancelled is returned when the request context is cancelled mid-loop.
ErrContextCancelled = errors.New("agent: operation cancelled")

Expand Down Expand Up @@ -201,3 +228,59 @@ func (e *LLMFailureError) Is(target error) bool {
func (e *LLMFailureError) Unwrap() error {
return e.Cause
}

// IncompleteKind classifies why a provider stopped generating early. Every
// provider has its own stop-reason vocabulary ("MAX_TOKENS", "length",
// "max_tokens"); each translates into these three cases so consumers route
// on one contract instead of three.
type IncompleteKind string

const (
// IncompleteTruncated: an output-length cap cut the response short.
IncompleteTruncated IncompleteKind = "truncated"
// IncompleteBlocked: a content policy stopped the generation.
IncompleteBlocked IncompleteKind = "blocked"
// IncompleteOther: neither a length cap nor a policy — a malformed
// tool call, or a backend stop the provider does not classify. The
// response is still partial, but it matches no sentinel, so retry
// behavior stays at the default.
IncompleteOther IncompleteKind = "other"
)

// IncompleteResponseError reports a generation that ended before the model
// finished. Every provider returns it for the same situations — an output
// cap fired, or a content filter stopped the stream — so an adopter writes
// one branch rather than one per vendor:
//
// if errors.Is(err, agent.ErrLLMTruncated) {
// // raise the provider's token cap, or ask for less
// }
//
// Whatever streamed before the stop still rides on the returned LLMResult
// (content and usage), because it is real output that cost real tokens —
// it is a prefix, though, never a complete answer. Reason carries the raw
// provider stop reason for logs; Kind is what routing should key on, via
// the sentinels above.
type IncompleteResponseError struct {
// Provider names the vendor package that produced the error, matching
// the error-message prefix convention ("anthropic", "openai", "gemini").
Provider string
// Reason is the provider's own stop reason, verbatim.
Reason string
// Kind is the provider-neutral classification.
Kind IncompleteKind
}

func (e *IncompleteResponseError) Error() string {
return fmt.Sprintf("%s: generation stopped early (%s): the response is a partial prefix, not a complete answer", e.Provider, e.Reason)
}

func (e *IncompleteResponseError) Is(target error) bool {
switch e.Kind {
case IncompleteTruncated:
return target == ErrLLMTruncated
case IncompleteBlocked:
return target == ErrLLMContentBlocked
}
return false
}
9 changes: 7 additions & 2 deletions pkg/agent/llm_call.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,13 @@ import (
// streaming (consumer has already seen partial output). Returns the
// final content, result, and (possibly retry-exhausted) error.
func (al *AgentLoop) callLLMWithRetry(ctx context.Context, st *iterationState, msgs []history.Message) (string, LLMResult, error) {
finalContent, result, _, err := al.callLLM(ctx, st, msgs)
if err == nil || al.Retry == nil || !isRetryable(err) {
finalContent, result, contentEmitted, err := al.callLLM(ctx, st, msgs)
// The content gate applies to the first attempt too, not just to
// retries: once a partial answer has reached the consumer, a retry
// replays the whole answer into the same stream. Truncation errors
// always arrive this way, so without the gate every capped response
// would both duplicate itself and pay for a second full generation.
if err == nil || al.Retry == nil || contentEmitted || !isRetryable(err) {
return finalContent, result, err
}
for attempt := 0; attempt < al.Retry.MaxRetries; attempt++ {
Expand Down
8 changes: 6 additions & 2 deletions pkg/agent/retry.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,10 +56,14 @@ func (r *RetryConfig) delay(attempt int) time.Duration {
return d
}

// isRetryable returns false for context-level errors that should not be retried.
// isRetryable returns false for errors that fail identically on every
// attempt: context-level cancellation, and a provider content-policy stop
// (deterministic for a given prompt — retrying only burns the budget).
func isRetryable(err error) bool {
if err == nil {
return false
}
return !errors.Is(err, context.Canceled) && !errors.Is(err, context.DeadlineExceeded)
return !errors.Is(err, context.Canceled) &&
!errors.Is(err, context.DeadlineExceeded) &&
!errors.Is(err, ErrLLMContentBlocked)
}
65 changes: 65 additions & 0 deletions pkg/agent/retry_hook_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package agent
import (
"context"
"errors"
"fmt"
"sync"
"testing"
"time"
Expand Down Expand Up @@ -114,3 +115,67 @@ func TestRetry_OnAttemptNilIsZeroCost(t *testing.T) {
// Compile-time check that history is wired through; the import is used by
// the flakyProvider receiver method.
var _ = history.Message{}

// truncatingProvider streams a partial answer and then reports it as
// truncated — the shape every provider returns when its output cap fires.
type truncatingProvider struct {
mu sync.Mutex
calls int
}

func (p *truncatingProvider) GenerateStream(_ context.Context, _ []history.Message, _ *tools.Registry, ch chan<- StreamEvent) (LLMResult, error) {
p.mu.Lock()
p.calls++
p.mu.Unlock()
ch <- Event(ContentEvent{Text: "partial"})
return LLMResult{Content: "partial"}, &IncompleteResponseError{
Provider: "test",
Reason: "max_tokens",
Kind: IncompleteTruncated,
}
}

func TestRetry_TruncationAfterStreamedContentIsNotRetried(t *testing.T) {
// Retrying would replay the whole answer into a stream the consumer has
// already read, and pay for a second cap-sized generation to do it.
prov := &truncatingProvider{}
loop, _ := setup(prov)
loop.Retry = &RetryConfig{MaxRetries: 3, BaseDelay: time.Millisecond, MaxDelay: time.Millisecond}

if _, err := loop.RunIteration(context.Background(), "s1", "go"); err == nil {
t.Fatal("a truncated response must surface as an error")
}

prov.mu.Lock()
defer prov.mu.Unlock()
if prov.calls != 1 {
t.Fatalf("provider called %d times, want 1 (no retry after streamed content)", prov.calls)
}
}

func TestIsRetryable_ProviderStopClassification(t *testing.T) {
tests := []struct {
name string
err error
want bool
}{
{"nil", nil, false},
{"cancelled", context.Canceled, false},
{"deadline", context.DeadlineExceeded, false},
{"generic provider error", errors.New("connection reset"), true},
// A length cut depends on this response, not the request, so the
// next attempt may well fit.
{"truncated", ErrLLMTruncated, true},
{"wrapped truncated", fmt.Errorf("gemini: %w", ErrLLMTruncated), true},
// A policy stop is deterministic: every retry reproduces it.
{"content blocked", ErrLLMContentBlocked, false},
{"wrapped content blocked", fmt.Errorf("gemini: %w", ErrLLMContentBlocked), false},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := isRetryable(tt.err); got != tt.want {
t.Fatalf("isRetryable(%v) = %v, want %v", tt.err, got, tt.want)
}
})
}
}
21 changes: 18 additions & 3 deletions pkg/llm/anthropic/anthropic.go
Original file line number Diff line number Diff line change
Expand Up @@ -275,8 +275,9 @@ func (p *Provider) GenerateStream(ctx context.Context, memory []history.Message,
// Surface the per-call MaxTokens truncation as a typed cap event so
// adopters can render "model truncated; raise MaxTokens" instead of
// silently shipping half-rendered code blocks. Skipped when soft
// truncation already fired the same event above.
if !truncated && string(accumulated.StopReason) == "max_tokens" {
// truncation already fired the same event above. The typed error is
// returned after the content is extracted, below.
if !truncated && accumulated.StopReason == anthropic.StopReasonMaxTokens {
streamChan <- agent.LimitExhaustedStreamEvent(agent.LimitKindProviderMaxTokens, int(p.MaxTokens), 0)
}

Expand Down Expand Up @@ -319,7 +320,21 @@ func (p *Provider) GenerateStream(ctx context.Context, memory []history.Message,
}
usage.TotalTokens = usage.PromptTokens + usage.CompletionTokens

return agent.LLMResult{Content: finalContent, ToolCalls: pendingCalls, Usage: usage}, nil
result := agent.LLMResult{Content: finalContent, ToolCalls: pendingCalls, Usage: usage}
// A capped or refused response is a prefix, not an answer. Returning it
// as a success is what turns a truncation into a decode error several
// layers up, naming the caller's schema instead of the real cause. The
// partial rides on the result so a host can still show what arrived.
// Soft truncation reports max_tokens too: the SDK could not finalize a
// tool_use block precisely because the cap fired mid-JSON.
stopReason := accumulated.StopReason
if truncated {
stopReason = anthropic.StopReasonMaxTokens
}
if err := stopReasonErr(stopReason); err != nil {
return result, err
}
return result, nil
}

// synthesizeStructuredTool builds a fake tool whose InputSchema
Expand Down
25 changes: 25 additions & 0 deletions pkg/llm/anthropic/errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,3 +38,28 @@ func classifyErr(err error) error {
}
return err
}

// stopReasonErr returns an *agent.IncompleteResponseError unless the
// generation ended cleanly, mapping Anthropic's stop reasons onto the
// provider-neutral classification consumers route on.
//
// "end_turn", "stop_sequence", and "tool_use" are complete responses.
// "pause_turn" is also complete — a long-running server tool asked the
// caller to continue the turn, and the content so far is intact. An empty
// reason means the API never reported one, treated as a clean stop so the
// default path is unchanged.
func stopReasonErr(r anthropic.StopReason) error {
var kind agent.IncompleteKind
switch r {
case "", anthropic.StopReasonEndTurn, anthropic.StopReasonStopSequence,
anthropic.StopReasonToolUse, anthropic.StopReasonPauseTurn:
return nil
case anthropic.StopReasonMaxTokens:
kind = agent.IncompleteTruncated
case anthropic.StopReasonRefusal:
kind = agent.IncompleteBlocked
default:
kind = agent.IncompleteOther
}
return &agent.IncompleteResponseError{Provider: "anthropic", Reason: string(r), Kind: kind}
}
Loading
Loading