fix(cycle): bound propose_takes with per-call timeout + phase deadline

Takeover of #2262. The extractor's gateway.chat call had no abortSignal, so
one stalled provider socket could pin the phase for the 300s gateway default
per page; the nightly wrapper then SIGTERMed the whole phase mid-run. Each
extractor call is now bounded at 90s (per-page failure already logs a warning
and continues), and the page loop carries a 30-min wall-clock deadline that
breaks cleanly into a partial result with deadline_hit:true + warn status.

Unlike the original PR, the default pageLimit stays at 100 — shrinking it to
30 was an unrelated product-knob change that would permanently cut nightly
take coverage.

Co-authored-by: tschew72 <tschew72@users.noreply.github.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Garry Tan
2026-07-21 14:45:47 -07:00
co-authored by tschew72 Claude Fable 5
parent 5df66c3855
commit 17fa84cda5
2 changed files with 68 additions and 1 deletions
+34 -1
View File
@@ -145,6 +145,8 @@ export interface ProposeTakesOpts extends BasePhaseOpts {
model?: string;
/** Skip pages that already have a complete takes fence. Default: true. */
skipPagesWithFence?: boolean;
/** Override the phase wall-clock deadline (tests). Default: 30 min. */
deadlineMs?: number;
}
export interface ProposeTakesResult {
@@ -153,6 +155,8 @@ export interface ProposeTakesResult {
cache_misses: number;
proposals_inserted: number;
budget_exhausted: boolean;
/** True when the phase deadline fired before the page loop completed (partial result). */
deadline_hit?: boolean;
warnings: string[];
}
@@ -210,6 +214,9 @@ export function extractExistingTakesForDedup(pageBody: string): Array<{
return rows;
}
/** Per-call wall-clock timeout for the extractor LLM call. */
const EXTRACTOR_CALL_TIMEOUT_MS = 90_000;
/**
* Production extractor — calls gateway.chat with the EXTRACT_TAKES_PROMPT
* and parses the JSON array output. Returns [] on parse failure (logged as
@@ -227,10 +234,14 @@ export async function defaultExtractor(
.replace('{EXISTING_TAKES_JSON}', JSON.stringify(input.existingTakes, null, 2))
.replace('{PAGE_BODY}', input.pageBody);
// Bound each call so one stalled provider socket can't pin the phase for the
// full gateway default (GBRAIN_AI_CHAT_TIMEOUT_MS, 300s) x pageLimit. The
// caller already catches per-page errors, logs a warning, and continues.
const result = await gatewayChat({
messages: [{ role: 'user', content: prompt }],
...(input.modelHint ? { model: input.modelHint } : {}),
maxTokens: 2048,
abortSignal: AbortSignal.timeout(EXTRACTOR_CALL_TIMEOUT_MS),
});
// ChatResult.text is already the concatenated text content.
@@ -287,6 +298,14 @@ class ProposeTakesPhase extends BaseCyclePhase {
readonly name = 'propose_takes' as CyclePhase;
protected readonly budgetUsdKey = 'cycle.propose_takes.budget_usd';
protected readonly budgetUsdDefault = 5.0;
/**
* Hard wall-clock deadline for the phase. Even with the per-call timeout in
* defaultExtractor, a long tail of slow-but-completing calls can accumulate.
* The phase breaks cleanly and returns a partial result with
* `deadline_hit: true` instead of being killed mid-write by an outer
* `timeout` wrapper (the recurring SIGTERM in nightly dream runs).
*/
private static readonly PHASE_DEADLINE_MS = 30 * 60 * 1000;
protected override mapErrorCode(err: unknown): string {
if (err instanceof GBrainError) return err.problem;
@@ -307,6 +326,8 @@ class ProposeTakesPhase extends BaseCyclePhase {
const promptVersion = opts.promptVersion ?? PROPOSE_TAKES_PROMPT_VERSION;
const pageLimit = opts.pageLimit ?? 100;
const skipPagesWithFence = opts.skipPagesWithFence ?? false;
const deadlineMs = opts.deadlineMs ?? ProposeTakesPhase.PHASE_DEADLINE_MS;
const phaseStartMs = Date.now();
const proposalRunId = `propose-${new Date().toISOString().slice(0, 19).replace(/[-:T]/g, '')}-${randomUUID().slice(0, 8)}`;
const result: ProposeTakesResult = {
@@ -333,6 +354,18 @@ class ProposeTakesPhase extends BaseCyclePhase {
const modelId = opts.model ?? getChatModel();
for (const page of pages) {
// Phase deadline check. Break (not throw) so the phase returns a
// partial result with deadline_hit:true; work already banked stays.
const elapsedMs = Date.now() - phaseStartMs;
if (elapsedMs > deadlineMs) {
result.warnings.push(
`phase deadline hit at page ${result.pages_scanned}/${pages.length} ` +
`after ${(elapsedMs / 1000).toFixed(0)}s (cap ${(deadlineMs / 1000).toFixed(0)}s); partial completion`,
);
result.deadline_hit = true;
break;
}
result.pages_scanned += 1;
this.tick(opts);
@@ -450,7 +483,7 @@ class ProposeTakesPhase extends BaseCyclePhase {
return {
summary: `propose_takes: scanned ${result.pages_scanned} pages, ${result.cache_hits} cached, ${result.proposals_inserted} new proposals (run ${proposalRunId})`,
details: { ...result, proposal_run_id: proposalRunId, prompt_version: promptVersion },
status: result.budget_exhausted ? 'warn' : 'ok',
status: result.budget_exhausted || result.deadline_hit ? 'warn' : 'ok',
};
}
}
+34
View File
@@ -368,6 +368,40 @@ New prose appended here.`;
expect(extractorCalls).toBe(1);
});
test('phase deadline breaks the page loop with a partial result (deadline_hit)', async () => {
const pages = [
buildPage({ slug: 'wiki/slow-a', body: 'page a' }),
buildPage({ slug: 'wiki/slow-b', body: 'page b' }),
];
const { engine } = buildMockEngine({ pages });
let extractorCalls = 0;
const extractor: ProposeTakesExtractor = async () => {
extractorCalls++;
await new Promise((r) => setTimeout(r, 10));
return [{ claim_text: 'x', kind: 'take', holder: 'brain', weight: 0.5 }];
};
// 5ms deadline: page 1 processes (elapsed 0 at check), the 10ms extractor
// call pushes elapsed past the cap, page 2 is never scanned.
const result = await runPhaseProposeTakes(buildCtx(engine), { extractor, deadlineMs: 5 });
expect(result.status).toBe('warn');
const details = result.details as Record<string, unknown>;
expect(details.deadline_hit).toBe(true);
expect(details.pages_scanned).toBe(1);
expect(extractorCalls).toBe(1);
expect((details.warnings as string[]).some(w => w.includes('phase deadline hit'))).toBe(true);
});
test('default deadline does not fire on a fast run', async () => {
const pages = [buildPage({ slug: 'wiki/fast', body: 'quick page' })];
const { engine } = buildMockEngine({ pages });
const extractor: ProposeTakesExtractor = async () => [];
const result = await runPhaseProposeTakes(buildCtx(engine), { extractor });
const details = result.details as Record<string, unknown>;
expect(details.deadline_hit).toBeUndefined();
expect(details.pages_scanned).toBe(1);
});
test('proposal_run_id is stable across all proposals from one phase invocation', async () => {
const pages = [
buildPage({ slug: 'wiki/a', body: 'page a' }),