From 216116912a035e634078d15daec9b2b4b63bba31 Mon Sep 17 00:00:00 2001 From: Garry Tan Date: Tue, 21 Jul 2026 14:57:44 -0700 Subject: [PATCH] fix(cycle): drain PGLite synth subagents inline (takeover of #2699) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit PGLite holds an exclusive file lock on its embedded data-dir, so no separate Minions worker can serve the subagent children the synthesize phase enqueues — they sat in 'waiting' until waitForCompletion timed out. Drain a private per-run child queue inline (claim → run → complete/fail, plus the promote/stall/timeout housekeeping a worker would perform). No-op on Postgres, where children stay on the shared 'default' queue. Rebased onto the reworked synthesize (config.subagentTimeoutMs, #1586 source scoping): the inline job context now carries deadlineAtMs from the claim-time timeout_at stamp, and opts.yieldDuringPhase is ticked on a 60s keepalive while each child runs so the 5-min cycle lock TTL refreshes across long (up to 30-min) children. Co-authored-by: TheRealMrSystem Co-Authored-By: Claude Fable 5 --- src/core/cycle/synthesize.ts | 114 ++++++++++++++++++++++- test/e2e/dream-synthesize-pglite.test.ts | 77 ++++++++++++++- 2 files changed, 189 insertions(+), 2 deletions(-) diff --git a/src/core/cycle/synthesize.ts b/src/core/cycle/synthesize.ts index 2dc1a5f7a..13b55a119 100644 --- a/src/core/cycle/synthesize.ts +++ b/src/core/cycle/synthesize.ts @@ -28,6 +28,7 @@ import type Anthropic from '@anthropic-ai/sdk'; import { readFileSync, existsSync, writeFileSync, mkdirSync } from 'node:fs'; +import { randomUUID } from 'node:crypto'; import { chat as gatewayChat, validateModelId, type ChatResult } from '../ai/gateway.ts'; import { AIConfigError } from '../ai/errors.ts'; import { normalizeModelId } from '../model-id.ts'; @@ -37,7 +38,8 @@ import type { BrainEngine } from '../engine.ts'; import type { PhaseResult, PhaseError } from '../cycle.ts'; import { MinionQueue } from '../minions/queue.ts'; import { waitForCompletion, TimeoutError } from '../minions/wait-for-completion.ts'; -import type { MinionJobInput, SubagentHandlerData } from '../minions/types.ts'; +import { makeSubagentHandler } from '../minions/handlers/subagent.ts'; +import type { MinionJobInput, MinionJobContext, MinionHandler, SubagentHandlerData } from '../minions/types.ts'; import { discoverTranscripts, type DiscoveredTranscript } from './transcript-discovery.ts'; import { serializeMarkdown, serializePageToMarkdown } from '../markdown.ts'; import type { Page, PageType } from '../types.ts'; @@ -261,6 +263,102 @@ export interface SynthesizePhaseOpts { once?: boolean; } +const INLINE_PGLITE_LOCK_MS = 30_000; + +/** + * PGLite cannot be served by a separate Minions worker process: the embedded + * data-dir holds an exclusive file lock, so subagent children enqueued by the + * synth parent would sit in 'waiting' until waitForCompletion times out. + * Drive the same claim → run → complete/fail loop a worker would perform, + * inline, against this phase's private child queue. + * + * `yieldDuringPhase` is ticked on a 60s interval while a child runs so the + * 5-min cycle lock TTL keeps refreshing during long (up to 30-min) children. + */ +async function runPgliteSubagentsInline( + engine: BrainEngine, + queue: MinionQueue, + queueName: string, + yieldDuringPhase?: () => Promise, + handler: MinionHandler = makeSubagentHandler({ engine }), +): Promise { + if (engine.kind !== 'pglite') return; + + while (true) { + // Housekeeping a worker would normally perform, so child rows can reach + // terminal states (delayed retries promoted, timeouts dead-lettered) + // before the synth parent enters waitForCompletion polling. + await queue.promoteDelayed(); + await queue.handleStalled(); + await queue.handleTimeouts(); + await queue.handleWallClockTimeouts(INLINE_PGLITE_LOCK_MS); + + const lockToken = randomUUID(); + const job = await queue.claim(lockToken, INLINE_PGLITE_LOCK_MS, queueName, ['subagent']); + if (!job) return; + + const abort = new AbortController(); + const shutdown = new AbortController(); + const context: MinionJobContext = { + id: job.id, + name: job.name, + data: job.data, + attempts_made: job.attempts_made, + signal: abort.signal, + deadlineAtMs: job.timeout_at != null ? job.timeout_at.getTime() : null, + shutdownSignal: shutdown.signal, + updateProgress: async (progress: unknown) => { + await queue.updateProgress(job.id, lockToken, progress); + }, + updateTokens: async (tokens) => { + await queue.updateTokens(job.id, lockToken, tokens); + }, + log: async (message) => { + const value = typeof message === 'string' ? message : JSON.stringify(message); + await engine.executeRaw( + `UPDATE minion_jobs SET stacktrace = COALESCE(stacktrace, '[]'::jsonb) || to_jsonb($1::text), + updated_at = now() + WHERE id = $2 AND status = 'active' AND lock_token = $3`, + [value, job.id, lockToken], + ); + }, + isActive: async () => { + const rows = await engine.executeRaw<{ id: number }>( + `SELECT id FROM minion_jobs WHERE id = $1 AND status = 'active' AND lock_token = $2`, + [job.id, lockToken], + ); + return rows.length > 0; + }, + readInbox: async () => queue.readInbox(job.id, lockToken), + }; + + // Cycle-lock keepalive while the child runs (best-effort, never throws). + const keepalive = yieldDuringPhase + ? setInterval(() => { yieldDuringPhase().catch(() => { /* best-effort */ }); }, 60_000) + : null; + try { + const result = await handler(context); + await queue.completeJob( + job.id, + lockToken, + result != null ? (typeof result === 'object' ? result as Record : { value: result }) : undefined, + ); + } catch (e) { + const errorText = e instanceof Error ? e.message : String(e); + const attemptsExhausted = job.attempts_made + 1 >= job.max_attempts; + await queue.failJob( + job.id, + lockToken, + errorText, + attemptsExhausted ? 'dead' : 'delayed', + 0, + ); + } finally { + if (keepalive) clearInterval(keepalive); + } + } +} + export async function runPhaseSynthesize( engine: BrainEngine, opts: SynthesizePhaseOpts, @@ -427,6 +525,12 @@ export async function runPhaseSynthesize( } const queue = new MinionQueue(engine); + // PGLite children drain inline (no separate worker can open the embedded + // data-dir), so give them a private per-run queue: the inline drain must + // never claim unrelated 'default'-queue jobs a Postgres worker owns. + const childQueueName = engine.kind === 'pglite' + ? `dream-inline-${Date.now()}-${randomUUID().slice(0, 8)}` + : 'default'; const childIds: number[] = []; /** Map child job_id → chunk metadata for D6 orchestrator-side slug rewrite. */ const chunkInfo = new Map(); @@ -505,6 +609,7 @@ export async function runPhaseSynthesize( on_child_fail: 'continue', idempotency_key, timeout_ms: config.subagentTimeoutMs, + queue: childQueueName, }; const child = await queue.add( 'subagent', @@ -519,6 +624,12 @@ export async function runPhaseSynthesize( } } + // PGLite cannot run a separate Minions worker because the embedded DB + // holds an exclusive file lock. Drain this phase's private child queue + // inline so the parent observes terminal child states instead of polling + // waiters until subagentWaitTimeoutMs expires. No-op on Postgres. + await runPgliteSubagentsInline(engine, queue, childQueueName, opts.yieldDuringPhase); + // Wait for every child to reach a terminal state. Tick yieldDuringPhase // every 5 min so the cycle lock TTL refreshes. const childOutcomes: Array<{ jobId: number; status: string }> = []; @@ -1381,4 +1492,5 @@ export const __testing = { buildSynthesisPrompt, stampDreamProvenance, reverseWriteRefs, + runPgliteSubagentsInline, }; diff --git a/test/e2e/dream-synthesize-pglite.test.ts b/test/e2e/dream-synthesize-pglite.test.ts index 7cb2082bf..d58b03ddd 100644 --- a/test/e2e/dream-synthesize-pglite.test.ts +++ b/test/e2e/dream-synthesize-pglite.test.ts @@ -17,7 +17,7 @@ import { mkdtempSync, rmSync, writeFileSync } from 'fs'; import { tmpdir } from 'os'; import { join } from 'path'; import { PGLiteEngine } from '../../src/core/pglite-engine.ts'; -import { runPhaseSynthesize, renderPageToMarkdown } from '../../src/core/cycle/synthesize.ts'; +import { runPhaseSynthesize, renderPageToMarkdown, __testing as synthTesting } from '../../src/core/cycle/synthesize.ts'; interface TestRig { engine: PGLiteEngine; @@ -516,3 +516,78 @@ describe('E2E synthesize — verdict cache (Q-2)', () => { } }, 30_000); }); + +describe('E2E synthesize — PGLite inline subagent drain (takeover of #2699)', () => { + test('drains private subagent queue inline so the parent can observe completion', async () => { + const rig = await setupRig(); + try { + const { MinionQueue } = await import('../../src/core/minions/queue.ts'); + const queue = new MinionQueue(rig.engine); + const queueName = `dream-inline-test-${Date.now()}`; + const child = await queue.add( + 'subagent', + { prompt: 'test', model: 'anthropic:claude-sonnet-4-6', max_turns: 1 }, + { queue: queueName, max_attempts: 1 }, + { allowProtectedSubmit: true }, + ); + + let ticks = 0; + await synthTesting.runPgliteSubagentsInline( + rig.engine, + queue, + queueName, + async () => { ticks++; }, + async (ctx) => { + await ctx.log('inline child ran'); + await ctx.updateProgress({ step: 'done' }); + return { ok: true }; + }, + ); + expect(ticks).toBe(0); // 60s keepalive never fires for a fast child + + const final = await queue.getJob(child.id); + expect(final?.status).toBe('completed'); + expect(final?.result).toEqual({ ok: true }); + expect(final?.progress).toEqual({ step: 'done' }); + + const waiting = await rig.engine.executeRaw<{ count: string }>( + `SELECT COUNT(*)::text AS count FROM minion_jobs WHERE queue = $1 AND status = 'waiting'`, + [queueName], + ); + expect(waiting[0]?.count).toBe('0'); + } finally { + await rig.cleanup(); + } + }, 30_000); + + test('terminally marks failed inline children so synth parent will not hang', async () => { + const rig = await setupRig(); + try { + const { MinionQueue } = await import('../../src/core/minions/queue.ts'); + const queue = new MinionQueue(rig.engine); + const queueName = `dream-inline-test-fail-${Date.now()}`; + const child = await queue.add( + 'subagent', + { prompt: 'test', model: 'anthropic:claude-sonnet-4-6', max_turns: 1 }, + { queue: queueName, max_attempts: 1 }, + { allowProtectedSubmit: true }, + ); + + await synthTesting.runPgliteSubagentsInline( + rig.engine, + queue, + queueName, + undefined, + async () => { + throw new Error('synthetic child failure'); + }, + ); + + const final = await queue.getJob(child.id); + expect(final?.status).toBe('dead'); + expect(final?.error_text).toContain('synthetic child failure'); + } finally { + await rig.cleanup(); + } + }, 30_000); +});