From b82f520314a75a18277c4cbf6df44b0c58f95dec Mon Sep 17 00:00:00 2001 From: Masa <98894508+Masashi-Ono0611@users.noreply.github.com> Date: Fri, 24 Jul 2026 06:21:34 +0900 Subject: [PATCH] fix(cycle): propagate all-provider-failed atom drains so durable jobs retry (#3218) (#3248) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit extract-atoms-drain's runBatch discarded runPhaseExtractAtoms's per-item failures/status, so a batch where EVERY provider call errored collapsed to {extracted: 0, skipped: 0} — indistinguishable from a legitimate no-op. The drain loop reported status: 'ok' regardless, the Minion handler returned normally, and the worker marked the durable job complete while the backlog sat untouched with no retry ever applied. - runBatch now derives providerFailure from the same counts the phase already returns (failures.length > 0 && transcripts_processed + pages_processed === 0 — every attempted item errored, zero succeeded). Partial success (>=1 item processed) is unaffected. - The pure loop surfaces this as status/stopped = 'provider_failure', breaking immediately (same hot-loop guard as no_progress) instead of letting a final remaining===0 recount silently overwrite it to 'drained'. - The extract-atoms-drain Minion handler throws when it sees status === 'provider_failure', so the worker's ordinary failJob path (attempt+backoff, dead-letter on exhaustion) takes over. The LockUnavailableError -> deferred path is unchanged. - autopilot's auto-drain submission bumps max_attempts from 1 to 3 (queue default) — with the handler now actually throwing, max_attempts:1 meant the first provider blip dead-lettered instantly with no backoff attempt. Tests: pure-loop provider_failure propagation (incl. the remaining===0 precedence case), runPhaseExtractAtoms's all-items-fail counts contract, and source-shape guards on the handler throw + autopilot max_attempts. Full suite deferred to CI per repo convention (targeted run: 104 pass / 0 fail across the touched + adjacent extract-atoms/drain/autopilot files; `bun run typecheck` clean). Two rounds of codex review (gpt-5.6-sol, high effort): round 1 flagged autopilot's max_attempts:1 and the stopped-precedence bug (both fixed above); round 2 confirmed no new issues. Thanks to @aaronkhawkins for the detailed report. Addresses the report in #3218. Co-authored-by: Claude Fable 5 --- src/commands/autopilot.ts | 5 +- src/commands/jobs.ts | 17 +++- src/core/cycle/extract-atoms-drain.ts | 69 ++++++++++++-- test/autopilot-auto-drain-wiring.test.ts | 13 +++ .../extract-atoms-synthesize-concepts.test.ts | 25 +++++ test/extract-atoms-drain.test.ts | 94 +++++++++++++++++++ 6 files changed, 214 insertions(+), 9 deletions(-) diff --git a/src/commands/autopilot.ts b/src/commands/autopilot.ts index 009c4ecbe..d5d68be34 100644 --- a/src/commands/autopilot.ts +++ b/src/commands/autopilot.ts @@ -828,7 +828,10 @@ export async function runAutopilot(engine: BrainEngine, args: string[]) { { queue: 'default', idempotency_key: idemKey, - max_attempts: 1, + // issue #3218: the handler now throws on an + // all-provider-failed batch, so give the queue's + // backoff a chance (was 1 — dead-lettered instantly). + max_attempts: 3, timeout_ms: timeoutMs, }, { allowProtectedSubmit: true }, diff --git a/src/commands/jobs.ts b/src/commands/jobs.ts index c5bbe55ad..599a934ff 100644 --- a/src/commands/jobs.ts +++ b/src/commands/jobs.ts @@ -2061,11 +2061,26 @@ export async function registerBuiltinHandlers( ? job.data.repoPath : ((await engine.getConfig('sync.repo_path')) ?? undefined); try { - return await runExtractAtomsDrainForSource(engine, { + const result = await runExtractAtomsDrainForSource(engine, { sourceId, windowSeconds, brainDir: repoPath, }); + // issue #3218: every item the drain attempted failed (0 succeeded, >=1 + // provider error) — completing this job normally would mark the + // durable job done while the backlog sits untouched, and no retry + // policy would ever fire on it again. Throw so the worker's ordinary + // failJob path (attempt+backoff, or dead-letter once exhausted) takes + // over instead — matching the existing behavior for every other + // handler failure. Partial success (>=1 item extracted) keeps + // completing normally, unchanged. + if (result.status === 'provider_failure') { + throw new Error( + `extract-atoms-drain: all provider calls failed this batch ` + + `(batches=${result.batches}, remaining=${result.remaining ?? '?'}) — retrying`, + ); + } + return result; } catch (e) { if (e instanceof LockUnavailableError) { return { phase: 'extract_atoms', status: 'skipped', deferred: true, reason: 'cycle_already_running' }; diff --git a/src/core/cycle/extract-atoms-drain.ts b/src/core/cycle/extract-atoms-drain.ts index 98a4bfa69..c6f474d2a 100644 --- a/src/core/cycle/extract-atoms-drain.ts +++ b/src/core/cycle/extract-atoms-drain.ts @@ -33,8 +33,14 @@ export interface ExtractAtomsDrainDeps { * routine cycle's skip contract. */ withLock: (work: () => Promise) => Promise; - /** Process one bounded batch (rediscovers eligibility). Returns counts. */ - runBatch: () => Promise<{ extracted: number; skipped: number }>; + /** + * Process one bounded batch (rediscovers eligibility). Returns counts, plus + * `providerFailure` (issue #3218) when EVERY item the batch attempted threw + * (zero items succeeded, at least one failure) — i.e. the batch's warning + * result was actually a total provider outage, not a partial/no-op batch. + * Omit/false for the ordinary partial-success or nothing-to-do cases. + */ + runBatch: () => Promise<{ extracted: number; skipped: number; providerFailure?: boolean }>; /** Count remaining eligible-but-unextracted pages, or null on query error. */ countRemaining: () => Promise; /** Injectable clock. Production: Date.now. */ @@ -52,15 +58,22 @@ export interface ExtractAtomsDrainOpts { export interface ExtractAtomsDrainResult { phase: 'extract_atoms'; - status: 'ok'; + /** + * issue #3218: 'provider_failure' when any batch reported `providerFailure` + * (every item it attempted errored). The Minion handler throws on this + * status so the durable job retries instead of completing over a backlog + * that made zero forward progress. Partial-success batches (>=1 item + * succeeded) always report 'ok', unchanged from before. + */ + status: 'ok' | 'provider_failure'; extracted: number; skipped: number; /** Eligible pages still pending after the window. null if the count errored. */ remaining: number | null; /** Batches actually processed. */ batches: number; - /** Why the loop stopped: drained | window | no_progress | max_batches. */ - stopped: 'drained' | 'window' | 'no_progress' | 'max_batches'; + /** Why the loop stopped: drained | window | no_progress | max_batches | provider_failure. */ + stopped: 'drained' | 'window' | 'no_progress' | 'max_batches' | 'provider_failure'; } export async function runExtractAtomsDrain( @@ -74,6 +87,10 @@ export async function runExtractAtomsDrain( let skipped = 0; let batches = 0; let stopped: ExtractAtomsDrainResult['stopped'] = 'window'; + // issue #3218: latched once any batch reports providerFailure — drives + // the returned `status`, independent of how `stopped` reads after the + // final (possibly overriding) remaining-count check below. + let providerFailure = false; while (deps.now() < deadline) { if (batches >= maxBatches) { stopped = 'max_batches'; break; } @@ -87,6 +104,17 @@ export async function runExtractAtomsDrain( batches++; deps.onBatch?.({ batch: batches, extracted: r.extracted, remaining: before }); + // issue #3218: every item this batch attempted failed (0 succeeded, >=1 + // error) — a total provider outage, not ordinary no-op/partial progress. + // Stop immediately (same hot-loop guard as no_progress below) and flag + // it so the caller can retry via its own policy instead of treating the + // drain as a clean completion. + if (r.providerFailure) { + providerFailure = true; + stopped = 'provider_failure'; + break; + } + // Stop if a batch made zero forward progress — extraction is failing or // everything left is ineligible (e.g. all skipped). Prevents a hot loop // that spends budget without draining. @@ -94,8 +122,22 @@ export async function runExtractAtomsDrain( } const remaining = await deps.countRemaining(); - if (remaining === 0) stopped = 'drained'; - return { phase: 'extract_atoms', status: 'ok', extracted, skipped, remaining, batches, stopped }; + // issue #3218 (codex P2): don't let a final remaining===0 recount + // overwrite 'provider_failure' back to 'drained' — that would report the + // contradictory {status: 'provider_failure', stopped: 'drained'} and + // mislead the CLI/JSON consumer (dream.ts prints both fields verbatim). + // status already takes precedence for the Minion handler's retry + // decision; keep `stopped` consistent with it once a failure latched. + if (!providerFailure && remaining === 0) stopped = 'drained'; + return { + phase: 'extract_atoms', + status: providerFailure ? 'provider_failure' : 'ok', + extracted, + skipped, + remaining, + batches, + stopped, + }; }); } @@ -157,9 +199,22 @@ export async function runExtractAtomsDrainForSource( brainDir: opts.brainDir, }); const d = (r.details ?? {}) as Record; + // issue #3218: `r.status` collapses to 'warn' whether ONE item failed + // (partial success — leave the drain's existing ok/no_progress path + // alone) or EVERY item failed (a total provider outage the drain + // adapter was silently swallowing). Re-derive the total-failure case + // from the per-item counts `runPhaseExtractAtoms` already returns: + // >=1 failure AND zero items successfully processed (transcripts_processed + // + pages_processed both 0 means every attempted `chat()` call threw — + // items that succeed with 0 atoms still count as processed, so this + // does not fire on "provider fine, nothing extractable"). + const failures = Array.isArray(d.failures) ? d.failures : []; + const itemsSucceeded = + Number(d.transcripts_processed ?? 0) + Number(d.pages_processed ?? 0); return { extracted: Number(d.atoms_extracted ?? 0), skipped: Number(d.duplicates_skipped ?? 0), + providerFailure: failures.length > 0 && itemsSucceeded === 0, }; }, countRemaining: () => countExtractAtomsBacklog(engine, extractionSourceId), diff --git a/test/autopilot-auto-drain-wiring.test.ts b/test/autopilot-auto-drain-wiring.test.ts index c8dd28c21..c800d829a 100644 --- a/test/autopilot-auto-drain-wiring.test.ts +++ b/test/autopilot-auto-drain-wiring.test.ts @@ -43,6 +43,19 @@ describe('autopilot auto-drain wiring', () => { expect(SRC).toMatch(/engine\.kind === 'postgres'[\s\S]{0,400}auto_drain/); }); + // issue #3218 (codex P1): with the handler now throwing on an + // all-provider-failed batch, max_attempts:1 made the queue's retry policy + // "dead-letter on the first failure, no backoff attempt" — regression-guard + // against silently reverting to 1. + test('issue #3218: submits with max_attempts 3 (not 1) so a retry can backoff before dead-lettering', () => { + // lastIndexOf: the queue.add(...) call site itself (the earlier occurrence + // is the unrelated created_at count query above it in the same function). + const callSite = SRC.lastIndexOf("'extract-atoms-drain'"); + const drainBlock = SRC.slice(callSite, callSite + 900); + expect(drainBlock).toContain('max_attempts: 3'); + expect(drainBlock).not.toContain('max_attempts: 1'); + }); + test('CODEX impl #4: no maxWaiting (it coalesces by name+queue, not source)', () => { // maxWaiting would return source A's waiting job for source B's submit, // never queuing B and over-counting the cap. The per-source idempotency key diff --git a/test/cycle/extract-atoms-synthesize-concepts.test.ts b/test/cycle/extract-atoms-synthesize-concepts.test.ts index d14102495..4b410fa9a 100644 --- a/test/cycle/extract-atoms-synthesize-concepts.test.ts +++ b/test/cycle/extract-atoms-synthesize-concepts.test.ts @@ -177,6 +177,31 @@ describe('v0.41 T5: runPhaseExtractAtoms via stubbed chat', () => { expect((result.details?.failures as unknown[]).length).toBe(1); }); + // issue #3218 — when EVERY item's chat() call throws (all-provider-failed), + // `transcripts_processed`/`pages_processed` must stay 0 while `failures` + // records one entry per item. This is the exact shape the + // extract-atoms-drain wiring (`runExtractAtomsDrainForSource`) uses to + // derive `providerFailure` (failures.length > 0 && itemsSucceeded === 0), + // distinguishing a total outage from the partial-success case above. + test('all items fail: transcripts_processed/pages_processed stay 0, every item recorded in failures', async () => { + const chat = async (_o: ChatOpts): Promise => { + throw new Error('provider unavailable'); + }; + const result = await runPhaseExtractAtoms(engine, { + _transcripts: [ + { filePath: '/a.txt', content: 'a', contentHash: 'ha' }, + { filePath: '/b.txt', content: 'b', contentHash: 'hb' }, + ], + _pages: [], + _chat: chat as typeof import('../../src/core/ai/gateway.ts').chat, + }); + expect(result.status).toBe('warn'); + expect(result.details?.atoms_extracted).toBe(0); + expect(result.details?.transcripts_processed).toBe(0); + expect(result.details?.pages_processed).toBe(0); + expect((result.details?.failures as unknown[]).length).toBe(2); + }); + // v0.41.2.1 regression case (D9 #14 wording): with _pages:[] and same // _transcripts, all PRE-EXISTING PhaseResult.details fields match // pre-fix values byte-for-byte. The new fields (pages_processed, diff --git a/test/extract-atoms-drain.test.ts b/test/extract-atoms-drain.test.ts index cb7dae821..fa8aaa19d 100644 --- a/test/extract-atoms-drain.test.ts +++ b/test/extract-atoms-drain.test.ts @@ -77,6 +77,62 @@ describe('runExtractAtomsDrain (issue #1678)', () => { expect(result.stopped).toBe('no_progress'); expect(batches).toBe(1); expect(result.remaining).toBe(5); + expect(result.status).toBe('ok'); + }); + + // issue #3218 — a batch where every attempted item errored (providerFailure) + // must surface distinctly from an ordinary no_progress/drained/window stop, + // so the Minion handler can retry instead of completing the durable job. + it('stops with status=provider_failure when a batch reports providerFailure', async () => { + let batches = 0; + const result = await runExtractAtomsDrain( + { + withLock: passThroughLock, + countRemaining: async () => 5, + runBatch: async () => { + batches++; + return { extracted: 0, skipped: 0, providerFailure: true }; + }, + now: () => 0, + }, + { windowMs: 1_000_000 }, + ); + expect(result.status).toBe('provider_failure'); + expect(result.stopped).toBe('provider_failure'); + expect(batches).toBe(1); + expect(result.remaining).toBe(5); + }); + + // issue #3218 (codex P2) — a final recount of 0 must NOT overwrite + // stopped='provider_failure' back to 'drained'. Otherwise the caller sees + // the contradictory {status: 'provider_failure', stopped: 'drained'}. + it('preserves stopped=provider_failure even when the final recount is 0', async () => { + const result = await runExtractAtomsDrain( + { + withLock: passThroughLock, + countRemaining: seq([3, 0]), // before-check: 3; final post-loop recount: 0 + runBatch: async () => ({ extracted: 0, skipped: 0, providerFailure: true }), + now: () => 0, + }, + { windowMs: 1_000_000 }, + ); + expect(result.status).toBe('provider_failure'); + expect(result.stopped).toBe('provider_failure'); + expect(result.remaining).toBe(0); + }); + + it('does not flag provider_failure for an ordinary partial-success batch', async () => { + const result = await runExtractAtomsDrain( + { + withLock: passThroughLock, + countRemaining: seq([3, 0, 0]), + runBatch: async () => ({ extracted: 1, skipped: 0, providerFailure: false }), + now: () => 0, + }, + { windowMs: 1_000_000 }, + ); + expect(result.status).toBe('ok'); + expect(result.stopped).toBe('drained'); }); it('propagates a busy-lock error (caller reports cycle_already_running)', async () => { @@ -133,4 +189,42 @@ describe('shared wiring helper holds the cycle lock (5A)', () => { expect(src).toContain('cycleLockIdFor(opts.sourceId)'); expect(src).toContain('withRefreshingLock(engine, lockId'); }); + + // issue #3218 — the wiring's `runBatch` must derive `providerFailure` from + // the SAME per-item counts pinned by + // `extract-atoms-synthesize-concepts.test.ts`'s "all items fail" case + // (failures.length > 0 && transcripts_processed + pages_processed === 0), + // not from `r.status` (which collapses partial and total failure into the + // same 'warn' value — the exact discard the issue reports). + it('runBatch derives providerFailure from failures.length + zero processed items, not r.status', () => { + const runBatchBlock = src.slice(src.indexOf('runBatch: async () => {')); + expect(runBatchBlock).toContain('d.failures'); + expect(runBatchBlock).toContain('transcripts_processed'); + expect(runBatchBlock).toContain('pages_processed'); + expect(runBatchBlock).toContain('providerFailure: failures.length > 0 && itemsSucceeded === 0'); + }); +}); + +// issue #3218 — the Minion handler must throw (not complete) when the drain +// reports status='provider_failure', so the worker's ordinary failJob path +// (attempt+backoff / dead-letter) retries the durable job instead of the +// backlog silently completing untouched. +describe('extract-atoms-drain Minion handler retries on provider_failure (issue #3218)', () => { + const jobsSrc = readFileSync(join(import.meta.dir, '../src/commands/jobs.ts'), 'utf8'); + const handlerBlock = jobsSrc.slice( + jobsSrc.indexOf("registerBuiltinJob(worker, engine, 'extract-atoms-drain'"), + jobsSrc.indexOf("registerBuiltinJob(worker, engine, 'extract-atoms-drain'") + 2200, + ); + + it("throws when result.status === 'provider_failure' instead of returning it", () => { + expect(handlerBlock).toMatch(/result\.status === 'provider_failure'/); + expect(handlerBlock).toMatch(/if \(result\.status === 'provider_failure'\) \{\s*throw new Error/); + }); + + it('still returns the deferred/skipped shape on LockUnavailableError (unchanged)', () => { + expect(handlerBlock).toContain('e instanceof LockUnavailableError'); + expect(handlerBlock).toContain( + "{ phase: 'extract_atoms', status: 'skipped', deferred: true, reason: 'cycle_already_running' }", + ); + }); });