mirror of
https://github.com/garrytan/gbrain.git
synced 2026-08-14 00:48:18 +00:00
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]>
This commit is contained in:
co-authored by
maxpetrusenkoagent <[REDACTED EMAIL]>
parent
447e57ec41
commit
f529eaa231
+57
-14
@@ -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<void> {
|
||||
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 <id>`). 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 <id>` 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;
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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);
|
||||
|
||||
|
||||
@@ -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');
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user