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
2 changes: 1 addition & 1 deletion docs/tools/mcp/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -258,7 +258,7 @@ toolsets:

See [Toolset Lifecycle](../../configuration/tools/index.md#toolset-lifecycle) for all profiles and tuning knobs, and [`/toolset-restart`](../../features/tui/index.md) to force a reconnect from the TUI.

**Startup failure behaviour:** local MCP failures (missing binary, connection refused, bad auth) fail fast — each turn retries immediately with no artificial delay. Remote MCP servers (Streamable HTTP / SSE) that respond with one of a fixed set of retryable HTTP statuses — 429 Too Many Requests, 408 Request Timeout, 500/502/503/504, or 529 (Anthropic-style "overloaded") — are paced by the same [bounded exponential backoff gate](../rag/index.md#indexing-failures-retries-and-backoff) that RAG embedding calls use, so a temporarily-overloaded remote MCP server does not trigger a new connect attempt on every agent turn. This is a fixed enumeration, not a full 5xx range: less-common codes such as 501, 505, or the Cloudflare 520–527 family do not arm the gate. Note MCP's trigger set is broader than RAG's current 429-only pacing (see the linked page and [#4097](https://github.com/docker/docker-agent/issues/4097)). This pacing also applies when toolsets are wrapped in [code mode](../../features/code-mode/index.md#limits--security-notes): a retryable failure in the degraded subset paces that subset's retry the same way, while the composite's healthy tools stay available.
**Startup failure behaviour:** local MCP failures (missing binary, connection refused, bad auth) fail fast — each turn retries immediately with no artificial delay. Remote MCP servers (Streamable HTTP / SSE) that respond with one of a fixed set of retryable HTTP statuses — 429 Too Many Requests, 408 Request Timeout, 500/502/503/504, or 529 (Anthropic-style "overloaded") — are paced by the same [bounded exponential backoff gate](../rag/index.md#indexing-failures-retries-and-backoff) that RAG embedding calls use, so a temporarily-overloaded remote MCP server does not trigger a new connect attempt on every agent turn. This is a fixed enumeration, not a full 5xx range: less-common codes such as 501, 505, or the Cloudflare 520–527 family do not arm the gate. MCP paces every connection attempt on this set; RAG's own trigger set differs slightly — 429 arms the gate on the very first failure, while 408/5xx only arm it once every file in an indexing run has hit one of these retryable statuses with none indexed successfully (see [RAG's retry-policy table](../rag/index.md#what-triggers-backoff)). This pacing also applies when toolsets are wrapped in [code mode](../../features/code-mode/index.md#limits--security-notes): a retryable failure in the degraded subset paces that subset's retry the same way, while the composite's healthy tools stay available.

## Combined Example

Expand Down
45 changes: 27 additions & 18 deletions docs/tools/rag/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -183,29 +183,38 @@ every agent turn.

### What triggers backoff

For **RAG indexing**, backoff applies only to **HTTP 429 rate-limit** responses
from the embedding or model provider — the one signal that reliably reaches the
toolset gate. Other errors (5xx, 408) are handled per-file within the indexing
run and do not arm the gate (tracked as a gap in
[#4097](https://github.com/docker/docker-agent/issues/4097)). These are the
current gate triggers for RAG indexing specifically:
For **RAG indexing**, the toolset gate arms on two different signals depending
on the failure:

- **HTTP 429 (rate limit)** aborts the whole indexing run at the *first*
failure — continuing would just keep hammering a provider that asked for
backoff — and that abort error reaches the gate immediately.
- **HTTP 408 (request timeout) and a fixed 5xx set** (`500, 502, 503, 504,
529`) are otherwise handled per-file: a single file's transient
failure is skipped so the run can keep indexing the rest. But if **no
file in the run is successfully indexed** — every attempted file hit one
of these retryable statuses — the run treats that as a sustained backend
failure rather than a one-off hiccup, and surfaces the error so the gate
arms on the next turn (fixed in
[#4097](https://github.com/docker/docker-agent/issues/4097); previously
only 429 reached the gate).

| Failure kind | Behaviour |
|---|---|
| HTTP 429 (rate limit) | Backoff: next attempt delayed |
| Other failures (5xx, 408, config errors, auth) | Fail fast: retried every turn with no added delay |
| HTTP 429 (rate limit) | Aborts the run on the first failure; backoff: next attempt delayed |
| HTTP 408 or 5xx, isolated to some files | Per-file skip; run succeeds, no backoff (indexed files persist, failures retried next run) |
| HTTP 408 or 5xx, affecting every file | Run fails; backoff: next attempt delayed |
| Other failures (config errors, auth, unrecognized 4xx) | Fail fast: retried every turn with no added delay |
| Context cancellation or agent shutdown | Immediate: no delay |

> [!NOTE]
> 5xx and 408 errors from the embedding provider are retried per-file and do not
> propagate to the toolset gate. Only 429 (rate-limit) terminates the indexing run
> early and surfaces the gate so Docker Agent can pace the next attempt.
>
> This 429-only trigger set is specific to the RAG/embedding path. Other toolset
> types have their own trigger sets against the same gate — for example, remote
> MCP toolsets also pace on 408 and a fixed set of 5xx-family statuses (see
> This trigger set (429 always, 408/5xx when sustained across every file) is
> specific to the RAG/embedding path. Other toolset types have their own
> trigger sets against the same gate — for example, remote MCP toolsets pace
> every connection attempt (not just a sustained run) on 408 and the same
> fixed 5xx set (see
> [MCP startup failure behaviour](../mcp/index.md#lifecycle-auto-restart-profiles)),
> and the A2A toolset paces its agent-card fetch on the same fixed set (see
> and the A2A toolset paces its agent-card fetch the same way (see
> [A2A startup failure behaviour](../a2a/index.md#startup-failure-behaviour)).

### Retry policy and parameters
Expand Down Expand Up @@ -246,9 +255,9 @@ are not affected.
- The knowledge-base tool does not appear in the agent's tool list until indexing
succeeds. A successful start is silent — the tool is listed and the agent uses it.

### Troubleshooting repeated 429 errors
### Troubleshooting repeated 429/5xx/408 errors

If you see persistent `429` errors in the logs:
If you see persistent `429`, `5xx`, or `408` errors in the logs:

1. **Check provider rate limits.** Your embedding API key may have a low requests-per-minute
quota. Upgrading the plan or using a different API key can help.
Expand Down
19 changes: 18 additions & 1 deletion pkg/rag/strategy/indexing_errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,11 @@ var errIndexingAborted = errors.New("indexing aborted due to non-retryable model
// call made during indexing. Permanent failures are wrapped with
// errIndexingAborted so callers can abort the run; transient failures (5xx,
// timeouts) and context cancellation are returned unchanged so callers can
// skip the current file and continue.
// skip the current file and continue. Per-file skipping is still correct for
// an isolated transient failure, but if every file in a run fails the same
// way, the caller uses isGateArmingTransientError to detect that and
// propagate instead of silently returning success (see vector_store.go's
// Initialize).
func classifyModelCallError(err error) error {
if err == nil {
return nil
Expand All @@ -40,3 +44,16 @@ func classifyModelCallError(err error) error {
func isIndexingAborted(err error) bool {
return errors.Is(err, errIndexingAborted)
}

// isGateArmingTransientError reports whether a transient (non-aborted) model
// error carries an HTTP status that StartableToolSet's backoff gate paces on
// (429, 408, or a fixed 5xx set — see startBackoffRetryable). Mirrors that
// function's own *modelerrors.StatusError pre-filter so plain-text errors
// (port numbers, chunk counters) can never arm the gate here either.
func isGateArmingTransientError(err error) bool {
var se *modelerrors.StatusError
if !errors.As(err, &se) {
return false
}
return modelerrors.RetryableHTTPStatus(se)
}
23 changes: 23 additions & 0 deletions pkg/rag/strategy/vector_store.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (
"path/filepath"
"strings"
"sync"
"sync/atomic"
"time"

"github.com/fsnotify/fsnotify"
Expand Down Expand Up @@ -302,6 +303,11 @@ func (s *VectorStore) Initialize(ctx context.Context, docPaths []string, chunkin
// Index files that need it in parallel
var indexed int
var indexedMu sync.Mutex
// firstGateArmingErr records the first transient failure whose HTTP
// status the StartableToolSet backoff gate paces on (see
// isGateArmingTransientError). CompareAndSwap makes "first" well-defined
// across concurrent goroutines without a separate mutex.
var firstGateArmingErr atomic.Pointer[error]

g, gctx := errgroup.WithContext(ctx)
g.SetLimit(s.fileIndexConcurrency)
Expand Down Expand Up @@ -331,6 +337,9 @@ func (s *VectorStore) Initialize(ctx context.Context, docPaths []string, chunkin
return err
}
slog.ErrorContext(ctx, "Failed to index file", "path", status.path, "error", err)
if isGateArmingTransientError(err) {
firstGateArmingErr.CompareAndSwap(nil, &err)
}
// Transient/local failure - continue indexing other files
return nil
}
Expand Down Expand Up @@ -360,6 +369,20 @@ func (s *VectorStore) Initialize(ctx context.Context, docPaths []string, chunkin
return err
}

// A gate-arming transient error (429/408/5xx) on every attempted file,
// with none indexed, means the provider is sustaining the same failure
// rather than hiccuping on one file - propagate it so the StartableToolSet
// backoff gate paces the next turn instead of retrying at full speed
// (see issue #4097). Partial progress still returns nil: indexed files
// persist via metadata, so the next run only retries the failures.
if indexed == 0 {
if errPtr := firstGateArmingErr.Load(); errPtr != nil {
err := fmt.Errorf("indexing failed for all %d file(s): %w", filesToIndex, *errPtr)
s.emitEvent(types.Event{Type: types.EventTypeError, Error: err})
return err
}
}

if err := s.cleanupOrphanedDocuments(ctx, seenFiles); err != nil {
slog.ErrorContext(ctx, "Failed to cleanup orphaned documents", "error", err)
}
Expand Down
96 changes: 89 additions & 7 deletions pkg/rag/strategy/vector_store_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,10 +53,14 @@ func TestClassifyModelCallError(t *testing.T) {
}
}

// fakeEmbeddingProvider counts embedding calls and always fails with a fixed error.
// fakeEmbeddingProvider counts embedding calls and fails with a fixed error.
// By default (failFirst == 0) it fails every call; when failFirst > 0 it
// fails only the first failFirst calls and succeeds afterward, simulating a
// backend that recovers partway through a run.
type fakeEmbeddingProvider struct {
calls atomic.Int64
err error
calls atomic.Int64
err error
failFirst int64
}

func (f *fakeEmbeddingProvider) ID() modelsdev.ID { return modelsdev.NewID("test", "fake-embed") }
Expand All @@ -67,8 +71,8 @@ func (f *fakeEmbeddingProvider) CreateChatCompletionStream(context.Context, []ch
}

func (f *fakeEmbeddingProvider) CreateEmbedding(context.Context, string) (*base.EmbeddingResult, error) {
f.calls.Add(1)
if f.err != nil {
n := f.calls.Add(1)
if f.err != nil && (f.failFirst <= 0 || n <= f.failFirst) {
return nil, f.err
}
return &base.EmbeddingResult{Embedding: []float64{0.1, 0.2}, TotalTokens: 1}, nil
Expand Down Expand Up @@ -131,6 +135,15 @@ func (db *fakeVectorDB) Close() error { return nil }

func newTestVectorStore(t *testing.T, embedErr error) (*VectorStore, *fakeEmbeddingProvider, []string) {
t.Helper()
return newTestVectorStoreWithConcurrency(t, embedErr, 1)
}

// newTestVectorStoreWithConcurrency is like newTestVectorStore but lets a
// test choose FileIndexConcurrency, so tests can exercise the concurrent
// errgroup path (e.g. racing writes to firstGateArmingErr) instead of the
// strictly sequential default.
func newTestVectorStoreWithConcurrency(t *testing.T, embedErr error, fileIndexConcurrency int) (*VectorStore, *fakeEmbeddingProvider, []string) {
t.Helper()

dir := t.TempDir()
const fileCount = 5
Expand All @@ -147,7 +160,7 @@ func newTestVectorStore(t *testing.T, embedErr error) (*VectorStore, *fakeEmbedd
Database: newFakeVectorDB(),
Embedder: embed.New(fake),
EmbeddingConcurrency: 1,
FileIndexConcurrency: 1,
FileIndexConcurrency: fileIndexConcurrency,
Chunking: ChunkingConfig{Size: 1024, Overlap: 0},
})

Expand Down Expand Up @@ -176,13 +189,82 @@ func TestInitializeContinuesOnTransientModelError(t *testing.T) {
Err: errors.New("internal server error"),
}
store, fake, docPaths := newTestVectorStore(t, embedErr)
fake.failFirst = 1 // only the first file fails; the rest succeed

err := store.Initialize(t.Context(), docPaths, ChunkingConfig{Size: 1024})
require.NoError(t, err, "transient errors skip the file and keep indexing")
require.NoError(t, err, "an isolated transient failure skips the file and keeps indexing")
assert.Equal(t, int64(len(docPaths)), fake.calls.Load(),
"every file should still be attempted on transient errors")
}

// TestInitializeSurfacesSustainedTransientModelError proves the fix for
// issue #4097: when every file in a run fails with the same gate-arming
// HTTP status (429, 408 or a retryable 5xx), Initialize propagates the
// error instead of silently returning nil, so StartableToolSet's backoff
// gate (startBackoffRetryable) arms on the next turn.
func TestInitializeSurfacesSustainedTransientModelError(t *testing.T) {
t.Parallel()
tests := []struct {
name string
statusCode int
}{
{name: "408 request timeout", statusCode: http.StatusRequestTimeout},
{name: "500 internal server error", statusCode: http.StatusInternalServerError},
{name: "502 bad gateway", statusCode: http.StatusBadGateway},
{name: "503 service unavailable", statusCode: http.StatusServiceUnavailable},
{name: "504 gateway timeout", statusCode: http.StatusGatewayTimeout},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
embedErr := &modelerrors.StatusError{
StatusCode: tt.statusCode,
Err: fmt.Errorf("HTTP %d error from provider", tt.statusCode),
}
store, fake, docPaths := newTestVectorStore(t, embedErr)

err := store.Initialize(t.Context(), docPaths, ChunkingConfig{Size: 1024})
require.Error(t, err, "a sustained transient failure across every file must be surfaced")
assert.Equal(t, int64(len(docPaths)), fake.calls.Load(),
"every file should still be attempted before the run gives up")
assert.False(t, isIndexingAborted(err),
"this is not the 400/401/404 abort path, just a propagated transient failure")

// startBackoffRetryable's own predicate: a *modelerrors.StatusError in
// the chain whose status RetryableHTTPStatus accepts. Asserted here
// (startBackoffRetryable itself is unexported in pkg/tools) to pin
// that the returned error actually arms the gate.
var statusErr *modelerrors.StatusError
require.ErrorAs(t, err, &statusErr)
assert.True(t, modelerrors.RetryableHTTPStatus(statusErr))
})
}
}

// TestInitializeSurfacesSustainedTransientModelError_ConcurrentFailures runs
// the same sustained-failure scenario with FileIndexConcurrency > 1, so
// multiple goroutines race to record into firstGateArmingErr via
// atomic.Pointer[error].CompareAndSwap concurrently — not just the strictly
// sequential path exercised above. Run with -race to catch any data race.
func TestInitializeSurfacesSustainedTransientModelError_ConcurrentFailures(t *testing.T) {
t.Parallel()
embedErr := &modelerrors.StatusError{
StatusCode: http.StatusServiceUnavailable,
Err: errors.New("service unavailable"),
}
store, fake, docPaths := newTestVectorStoreWithConcurrency(t, embedErr, 5) // fileCount in the helper is 5

err := store.Initialize(t.Context(), docPaths, ChunkingConfig{Size: 1024})
require.Error(t, err, "a sustained transient failure across every concurrently-indexed file must be surfaced")
assert.Equal(t, int64(len(docPaths)), fake.calls.Load())
assert.False(t, isIndexingAborted(err))

var statusErr *modelerrors.StatusError
require.ErrorAs(t, err, &statusErr)
assert.True(t, modelerrors.RetryableHTTPStatus(statusErr))
}

func TestCheckAndReindexAbortsOnNonRetryableModelError(t *testing.T) {
t.Parallel()
embedErr := &modelerrors.StatusError{
Expand Down
72 changes: 41 additions & 31 deletions pkg/tools/builtin/rag/rag_backoff_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -75,40 +75,50 @@ func buildRAGToolSet(t *testing.T, impl strategy.Strategy) *ToolSet {
// A frozen fake clock (WithStartRetryClock) and identity jitter give exact,
// deterministic windows without synctest complications from the RAG manager's
// internal goroutines.
//
// Subtests cover 429 (rate limit), 408 (request timeout) and 503 (service
// unavailable) - the family of statuses that, per issue #4097, must all
// reach the gate from RAG indexing, not just 429.
func TestRAGStartableBackoff_StatusErrorEngagesGate(t *testing.T) {
t.Parallel()

counting := &countingStatusErrStrategy{statusCode: 429}
toolset := buildRAGToolSet(t, counting)

// Frozen clock: window = base + 0% jitter = exactly base.
// We advance the clock manually to control expiry.
now := time.Unix(1_000_000, 0)
s := tools.NewStartable(toolset,
tools.WithStartRetryJitter(func(d time.Duration) time.Duration { return d }), // identity
tools.WithStartRetryClock(func() time.Time { return now }),
)

// Attempt 1: TryStart invokes Initialize and arms the gate.
_, err := s.TryStart(t.Context())
require.Error(t, err, "expected failure on attempt 1")
assert.Equal(t, int32(1), counting.calls.Load(), "Initialize must be called on attempt 1")

// Immediate TryStart: gate must suppress it (clock not advanced).
_, err = s.TryStart(t.Context())
require.Error(t, err)
assert.Equal(t, int32(1), counting.calls.Load(),
"gate must suppress TryStart within the window")

// Advance clock past base window (base × 1.0 with identity jitter).
// Use 6 minutes to exceed any possible jittered window for any attempt.
now = now.Add(6 * time.Minute)

// Gate expired: TryStart must invoke Initialize again.
_, err = s.TryStart(t.Context())
require.Error(t, err, "still failing — expected error")
assert.Equal(t, int32(2), counting.calls.Load(),
"Initialize must be called again once the backoff window expires")
for _, statusCode := range []int{429, 408, 503} {
t.Run(fmt.Sprintf("status_%d", statusCode), func(t *testing.T) {
t.Parallel()

counting := &countingStatusErrStrategy{statusCode: statusCode}
toolset := buildRAGToolSet(t, counting)

// Frozen clock: window = base + 0% jitter = exactly base.
// We advance the clock manually to control expiry.
now := time.Unix(1_000_000, 0)
s := tools.NewStartable(toolset,
tools.WithStartRetryJitter(func(d time.Duration) time.Duration { return d }), // identity
tools.WithStartRetryClock(func() time.Time { return now }),
)

// Attempt 1: TryStart invokes Initialize and arms the gate.
_, err := s.TryStart(t.Context())
require.Error(t, err, "expected failure on attempt 1")
assert.Equal(t, int32(1), counting.calls.Load(), "Initialize must be called on attempt 1")

// Immediate TryStart: gate must suppress it (clock not advanced).
_, err = s.TryStart(t.Context())
require.Error(t, err)
assert.Equal(t, int32(1), counting.calls.Load(),
"gate must suppress TryStart within the window")

// Advance clock past base window (base x 1.0 with identity jitter).
// Use 6 minutes to exceed any possible jittered window for any attempt.
now = now.Add(6 * time.Minute)

// Gate expired: TryStart must invoke Initialize again.
_, err = s.TryStart(t.Context())
require.Error(t, err, "still failing — expected error")
assert.Equal(t, int32(2), counting.calls.Load(),
"Initialize must be called again once the backoff window expires")
})
}
}

// TestRAGStartableBackoff_PlainErrorNoGate proves that a plain (non-StatusError)
Expand Down
Loading