From f529eaa231c96a22708a69c56a2df0e6d678f2cd Mon Sep 17 00:00:00 2001 From: maxpetrusenkoagent Date: Tue, 21 Jul 2026 16:21:51 -0400 Subject: [PATCH] fix(jobs): refresh gateway config for queued AI work (#2125) Long-lived minion workers can outlive DB-backed model config changes. Refresh the AI gateway before gateway-backed handlers run so queued cycle/propose_takes work does not fall back to a stale Anthropic default when the operator configured another provider. Also record the active gateway chat model in propose_takes budget/proposal metadata instead of hardcoding claude-sonnet-4-6, and keep provider:model IDs intact for budget pricing. Regression coverage verifies queued worker refresh, propose_takes model metadata, nested provider IDs, skipFence threading, and the updated autopilot signal source guard. Co-authored-by: maxpetrusenkoagent <[REDACTED EMAIL]> --- src/commands/jobs.ts | 71 ++++++++++++++++++++++++++------- src/core/cycle/propose-takes.ts | 8 ++-- test/cycle-abort.test.ts | 2 +- test/handlers.test.ts | 60 ++++++++++++++++++++++++++++ test/propose-takes.test.ts | 51 +++++++++++++++++++++++ 5 files changed, 174 insertions(+), 18 deletions(-) diff --git a/src/commands/jobs.ts b/src/commands/jobs.ts index cb09f9aa9..31aa24e14 100644 --- a/src/commands/jobs.ts +++ b/src/commands/jobs.ts @@ -7,7 +7,7 @@ import type { BrainEngine } from '../core/engine.ts'; import { MinionQueue } from '../core/minions/queue.ts'; import { MinionWorker } from '../core/minions/worker.ts'; import { WORKER_EXIT_RSS_WATCHDOG } from '../core/minions/worker-exit-codes.ts'; -import type { MinionJob, MinionJobStatus } from '../core/minions/types.ts'; +import type { MinionHandler, MinionJob, MinionJobStatus } from '../core/minions/types.ts'; import type { PaceKeyOverrides } from '../core/pace-mode.ts'; import { loadConfig, isThinClient } from '../core/config.ts'; import { callRemoteTool, unpackToolResult } from '../core/mcp-client.ts'; @@ -22,6 +22,49 @@ function hasFlag(args: string[], flag: string): boolean { return args.includes(flag); } +/** + * Long-lived workers outlive operator config changes. Re-stamp the AI gateway + * from DB-backed model config immediately before queued jobs enter gateway-backed + * paths, so a stale process-level default cannot route new work to the wrong + * provider. + */ +async function refreshGatewayForJob(engine: BrainEngine): Promise { + const { reconfigureGatewayWithEngine } = await import('../core/ai/gateway.ts'); + await reconfigureGatewayWithEngine(engine); +} + +const GATEWAY_REFRESH_JOB_NAMES = new Set([ + 'embed', + 'extract-conversation-facts', + 'enrich', + 'contextual_reindex_per_chunk', + 'autopilot-cycle', + 'synthesize', + 'patterns', + 'consolidate', + 'extract_facts', + 'extract-atoms-drain', + 'embed-backfill', + 'extract-takes-from-pages', + 'embed-catch-up', +]); + +function registerBuiltinJob( + worker: MinionWorker, + engine: BrainEngine, + name: string, + handler: MinionHandler, +): void { + if (!GATEWAY_REFRESH_JOB_NAMES.has(name)) { + worker.register(name, handler); + return; + } + worker.register(name, async (job) => { + await refreshGatewayForJob(engine); + return await handler(job); + }); +} + /** Parse `--max-waiting N` from CLI args. Returns undefined if absent. * Throws on malformed input (caller should surface the error and exit). * Clamps to [1, 100] to match the queue-layer clamp in MinionQueue.add. @@ -1439,7 +1482,7 @@ export async function registerBuiltinHandlers( return { ...result, embed_job_id: embedJobId, embed_skip_reason: embedSkipReason }; }); - worker.register('embed', async (job) => { + registerBuiltinJob(worker, engine, 'embed', async (job) => { const { runEmbedCore } = await import('./embed.ts'); // Primary Minion progress channel is job.updateProgress (DB-backed, // readable via `gbrain jobs get `). Stderr from the worker daemon @@ -1486,7 +1529,7 @@ export async function registerBuiltinHandlers( // BudgetTracker inside its own process. BudgetExhausted is caught at // the core level and returned as `result.budget_exhausted: true` (NOT // a job failure) so the user can resume with a higher cap. - worker.register('extract-conversation-facts', async (job) => { + registerBuiltinJob(worker, engine, 'extract-conversation-facts', async (job) => { const { runExtractConversationFactsCore } = await import('./extract-conversation-facts.ts'); const sourceId = typeof job.data.sourceId === 'string' ? job.data.sourceId : undefined; if (!sourceId) { @@ -1545,7 +1588,7 @@ export async function registerBuiltinHandlers( // at the core level and returned as result.budget_exhausted (NOT a failure). // Strict per-source: the CLI fans out one job per source when --source is // omitted, so a job ALWAYS carries data.sourceId. - worker.register('enrich', async (job) => { + registerBuiltinJob(worker, engine, 'enrich', async (job) => { const { runEnrichCore } = await import('./enrich.ts'); const sourceId = typeof job.data.sourceId === 'string' ? job.data.sourceId : undefined; if (!sourceId) { @@ -1685,13 +1728,13 @@ export async function registerBuiltinHandlers( const { makeContextualReindexHandler } = await import( '../core/minions/handlers/contextual-reindex-per-chunk.ts' ); - worker.register('contextual_reindex_per_chunk', makeContextualReindexHandler({ engine })); + registerBuiltinJob(worker, engine, 'contextual_reindex_per_chunk', makeContextualReindexHandler({ engine })); } // derivation); the handler returns { partial, status, report } so // `gbrain jobs get ` shows the full structured report. Does NOT // throw on partial: a flaky phase must not block every future cycle. - worker.register('autopilot-cycle', async (job) => { + registerBuiltinJob(worker, engine, 'autopilot-cycle', async (job) => { const { runCycle } = await import('../core/cycle.ts'); // v0.41.30 (T2): fall back to null (NOT cwd '.') when no repo is configured. // The queued cycle is the same primitive `gbrain dream` uses; a checkout-less @@ -1986,12 +2029,12 @@ export async function registerBuiltinHandlers( }; // PROTECTED — internally spawn subagent children - worker.register('synthesize', makePhaseHandler('synthesize')); - worker.register('patterns', makePhaseHandler('patterns')); - worker.register('consolidate', makePhaseHandler('consolidate')); + registerBuiltinJob(worker, engine, 'synthesize', makePhaseHandler('synthesize')); + registerBuiltinJob(worker, engine, 'patterns', makePhaseHandler('patterns')); + registerBuiltinJob(worker, engine, 'consolidate', makePhaseHandler('consolidate')); // Open — DB writes only, no LLM spend - worker.register('extract_facts', makePhaseHandler('extract_facts')); + registerBuiltinJob(worker, engine, 'extract_facts', makePhaseHandler('extract_facts')); worker.register('resolve_symbol_edges', makePhaseHandler('resolve_symbol_edges')); worker.register('recompute_emotional_weight', makePhaseHandler('recompute_emotional_weight')); @@ -2001,7 +2044,7 @@ export async function registerBuiltinHandlers( // window / defer behavior. On LockUnavailableError (the routine cycle holds // the per-source lock) the job completes `{ deferred: true }` and retries // next tick instead of failing — cooperative interleave (CODEX accepted). - worker.register('extract-atoms-drain', async (job) => { + registerBuiltinJob(worker, engine, 'extract-atoms-drain', async (job) => { const { runExtractAtomsDrainForSource } = await import('../core/cycle/extract-atoms-drain.ts'); const { LockUnavailableError } = await import('../core/db-lock.ts'); const sourceId = typeof job.data.sourceId === 'string' ? job.data.sourceId : undefined; @@ -2029,7 +2072,7 @@ export async function registerBuiltinHandlers( // Cost-bounded via D6 ($10/job BudgetTracker) + D19 (source-level cooldown // + 24h rolling cap, gated at submit time). NOT in PROTECTED_JOB_NAMES — // embedding-only spend, no API-by-the-minute risk like subagent. - worker.register('embed-backfill', async (job) => { + registerBuiltinJob(worker, engine, 'embed-backfill', async (job) => { const { makeEmbedBackfillHandler } = await import('../core/minions/handlers/embed-backfill.ts'); return await makeEmbedBackfillHandler(engine)(job); }); @@ -2050,7 +2093,7 @@ export async function registerBuiltinHandlers( // (LLM-bearing). Two-gate consent enforced at the handler boundary: // refuses to run unless takes.bootstrap_enabled config is true, even // when allowProtectedSubmit was set at queue.add time. - worker.register('extract-takes-from-pages', async (job) => { + registerBuiltinJob(worker, engine, 'extract-takes-from-pages', async (job) => { const { extractTakesFromPages } = await import('../core/extract-takes-from-pages.ts'); const data = (job.data ?? {}) as { sourceId?: string; maxPages?: number }; const bootstrapCfg = await engine.getConfig('takes.bootstrap_enabled'); @@ -2077,7 +2120,7 @@ export async function registerBuiltinHandlers( // remediation pipeline. Wraps runEmbedCore with stale + catchUp + the // priority/batchSize the recommendation supplies. NOT in // PROTECTED_JOB_NAMES (embedding spend only). - worker.register('embed-catch-up', async (job) => { + registerBuiltinJob(worker, engine, 'embed-catch-up', async (job) => { const { runEmbedCore } = await import('./embed.ts'); const data = (job.data ?? {}) as { sourceId?: string; diff --git a/src/core/cycle/propose-takes.ts b/src/core/cycle/propose-takes.ts index 33fbe054c..63ada141e 100644 --- a/src/core/cycle/propose-takes.ts +++ b/src/core/cycle/propose-takes.ts @@ -39,7 +39,7 @@ import { randomUUID, createHash } from 'node:crypto'; import { BaseCyclePhase, type ScopedReadOpts, type BasePhaseOpts } from './base-phase.ts'; -import { chat as gatewayChat } from '../ai/gateway.ts'; +import { chat as gatewayChat, getChatModel } from '../ai/gateway.ts'; import { writeReceipt } from '../extract/receipt-writer.ts'; import { upsertExtractRollup } from '../extract/rollup-writer.ts'; import { GBrainError } from '../types.ts'; @@ -330,6 +330,8 @@ class ProposeTakesPhase extends BaseCyclePhase { opts.reporter.start('propose_takes.pages' as never, pages.length); } + const modelId = opts.model ?? getChatModel(); + for (const page of pages) { result.pages_scanned += 1; this.tick(opts); @@ -359,7 +361,7 @@ class ProposeTakesPhase extends BaseCyclePhase { // Budget pre-check before the LLM call. Estimate: ~1500 input tokens + 500 output. const budget = this.checkBudget({ - modelId: opts.model ?? 'claude-sonnet-4-6', + modelId, estimatedInputTokens: 1500, maxOutputTokens: 500, }); @@ -408,7 +410,7 @@ class ProposeTakesPhase extends BaseCyclePhase { p.weight, p.domain ?? null, JSON.stringify(existingTakes), - opts.model ?? 'claude-sonnet-4-6', + modelId, ], ); result.proposals_inserted += 1; diff --git a/test/cycle-abort.test.ts b/test/cycle-abort.test.ts index 8b1dec4b1..9c49fe52f 100644 --- a/test/cycle-abort.test.ts +++ b/test/cycle-abort.test.ts @@ -113,7 +113,7 @@ describe('autopilot-cycle handler contract (v0.20.5)', () => { // the original 2000-char ceiling. The intent of the guard is unchanged: // "the autopilot-cycle handler passes job.signal to runCycle." The // window just needs to be wide enough to span any reasonable handler. - const handlerStart = jobsSource.indexOf("worker.register('autopilot-cycle'"); + const handlerStart = jobsSource.indexOf("registerBuiltinJob(worker, engine, 'autopilot-cycle'"); expect(handlerStart).toBeGreaterThan(-1); const handlerBlock = jobsSource.slice(handlerStart, handlerStart + 6000); diff --git a/test/handlers.test.ts b/test/handlers.test.ts index 487e2a3eb..7a4ec3d4f 100644 --- a/test/handlers.test.ts +++ b/test/handlers.test.ts @@ -13,6 +13,7 @@ import { describe, test, expect, beforeAll, afterAll, mock } from 'bun:test'; import { PGLiteEngine } from '../src/core/pglite-engine.ts'; import { MinionWorker } from '../src/core/minions/worker.ts'; import { registerBuiltinHandlers } from '../src/commands/jobs.ts'; +import { configureGateway, getChatModel, resetGateway } from '../src/core/ai/gateway.ts'; let engine: PGLiteEngine; let worker: MinionWorker; @@ -122,6 +123,65 @@ describe('autopilot-cycle handler — partial failure does NOT throw', () => { }); describe('autopilot-cycle handler — phase passthrough', () => { + test('refreshes DB-backed chat model config before a queued cycle runs', async () => { + const handler = (worker as any).handlers.get('autopilot-cycle'); + expect(handler).toBeDefined(); + + const oldModel = await engine.getConfig('models.chat'); + configureGateway({ + chat_model: 'anthropic:claude-sonnet-4-6', + env: { ANTHROPIC_API_KEY: 'stale-key', OPENAI_API_KEY: 'fresh-key' }, + }); + await engine.setConfig('models.chat', 'openai:gpt-5'); + + try { + const result = await handler({ + data: { phases: ['orphans'], pull: false }, + signal: { aborted: false } as any, + job: { id: 9, name: 'autopilot-cycle' } as any, + }); + + expect(result).toBeDefined(); + expect(getChatModel()).toBe('openai:gpt-5'); + } finally { + resetGateway(); + if (oldModel === null) { + await engine.unsetConfig('models.chat'); + } else { + await engine.setConfig('models.chat', oldModel); + } + } + }); + + test('refreshes DB-backed chat model config before gateway-backed handlers validate job data', async () => { + const handler = (worker as any).handlers.get('enrich'); + expect(handler).toBeDefined(); + + const oldModel = await engine.getConfig('models.chat'); + configureGateway({ + chat_model: 'anthropic:claude-sonnet-4-6', + env: { ANTHROPIC_API_KEY: 'stale-key', OPENAI_API_KEY: 'fresh-key' }, + }); + await engine.setConfig('models.chat', 'openai:gpt-5'); + + try { + await expect(handler({ + data: {}, + signal: { aborted: false } as any, + job: { id: 10, name: 'enrich' } as any, + })).rejects.toThrow('enrich Minion job requires data.sourceId'); + + expect(getChatModel()).toBe('openai:gpt-5'); + } finally { + resetGateway(); + if (oldModel === null) { + await engine.unsetConfig('models.chat'); + } else { + await engine.setConfig('models.chat', oldModel); + } + } + }); + test('job.data.phases restricts which phases run', async () => { const fs = await import('fs'); const { execSync } = await import('child_process'); diff --git a/test/propose-takes.test.ts b/test/propose-takes.test.ts index a49c317eb..3c0ccb68d 100644 --- a/test/propose-takes.test.ts +++ b/test/propose-takes.test.ts @@ -25,6 +25,8 @@ import { type ProposeTakesExtractor, type ProposedTake, } from '../src/core/cycle/propose-takes.ts'; +import { configureGateway, resetGateway } from '../src/core/ai/gateway.ts'; +import { BudgetMeter } from '../src/core/cycle/budget-meter.ts'; import type { OperationContext } from '../src/core/operations.ts'; import type { BrainEngine } from '../src/core/engine.ts'; import type { Page } from '../src/core/types.ts'; @@ -384,4 +386,53 @@ New prose appended here.`; expect(typeof runIdA).toBe('string'); expect((runIdA as string).startsWith('propose-')).toBe(true); }); + + test('records the configured gateway chat model when no phase model override is passed', async () => { + configureGateway({ + chat_model: 'openai:gpt-5', + env: { OPENAI_API_KEY: 'test-key' }, + }); + try { + const pages = [buildPage({ slug: 'wiki/model-default', body: 'configured model should be recorded' })]; + const { engine, captured } = buildMockEngine({ pages }); + const extractor: ProposeTakesExtractor = async () => [ + { claim_text: 'configured model should be recorded', kind: 'take', holder: 'brain', weight: 0.5 }, + ]; + + await runPhaseProposeTakes(buildCtx(engine), { extractor }); + + const insert = captured.find(c => c.sql.includes('INSERT INTO take_proposals')); + expect(insert).toBeDefined(); + expect(insert!.params[11]).toBe('openai:gpt-5'); + } finally { + resetGateway(); + } + }); + + test('keeps nested provider model ids intact for budget checks and proposal records', async () => { + configureGateway({ + chat_model: 'openrouter:anthropic/claude-sonnet-4-6', + env: { OPENROUTER_API_KEY: 'test-key' }, + }); + try { + const pages = [buildPage({ slug: 'wiki/openrouter-model', body: 'nested provider model should stay intact' })]; + const { engine, captured } = buildMockEngine({ pages }); + const extractor: ProposeTakesExtractor = async () => [ + { claim_text: 'nested provider model should stay intact', kind: 'take', holder: 'brain', weight: 0.5 }, + ]; + + const result = await runPhaseProposeTakes(buildCtx(engine), { + extractor, + meter: new BudgetMeter({ budgetUsd: 0.000001, phase: 'propose_takes' }), + }); + + expect(result.status).toBe('ok'); + expect(result.details.budget_exhausted).toBe(false); + const insert = captured.find(c => c.sql.includes('INSERT INTO take_proposals')); + expect(insert).toBeDefined(); + expect(insert!.params[11]).toBe('openrouter:anthropic/claude-sonnet-4-6'); + } finally { + resetGateway(); + } + }); });