diff --git a/docs/tools/mcp/index.md b/docs/tools/mcp/index.md index 96fe2c261..27c02dab6 100644 --- a/docs/tools/mcp/index.md +++ b/docs/tools/mcp/index.md @@ -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 diff --git a/docs/tools/rag/index.md b/docs/tools/rag/index.md index ba90830e3..a50006288 100644 --- a/docs/tools/rag/index.md +++ b/docs/tools/rag/index.md @@ -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 @@ -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. diff --git a/pkg/rag/strategy/indexing_errors.go b/pkg/rag/strategy/indexing_errors.go index 50fd7ae29..5619a3856 100644 --- a/pkg/rag/strategy/indexing_errors.go +++ b/pkg/rag/strategy/indexing_errors.go @@ -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 @@ -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) +} diff --git a/pkg/rag/strategy/vector_store.go b/pkg/rag/strategy/vector_store.go index 82629e678..150afdad5 100644 --- a/pkg/rag/strategy/vector_store.go +++ b/pkg/rag/strategy/vector_store.go @@ -8,6 +8,7 @@ import ( "path/filepath" "strings" "sync" + "sync/atomic" "time" "github.com/fsnotify/fsnotify" @@ -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) @@ -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 } @@ -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) } diff --git a/pkg/rag/strategy/vector_store_test.go b/pkg/rag/strategy/vector_store_test.go index e86e2b5c5..0ca730df8 100644 --- a/pkg/rag/strategy/vector_store_test.go +++ b/pkg/rag/strategy/vector_store_test.go @@ -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") } @@ -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 @@ -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 @@ -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}, }) @@ -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{ diff --git a/pkg/tools/builtin/rag/rag_backoff_test.go b/pkg/tools/builtin/rag/rag_backoff_test.go index 4010e1eef..4e6d8d80b 100644 --- a/pkg/tools/builtin/rag/rag_backoff_test.go +++ b/pkg/tools/builtin/rag/rag_backoff_test.go @@ -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)