From 3604c576dc3c06a7a948f6872fa81df7477e8229 Mon Sep 17 00:00:00 2001 From: Dmitrii Fedorov Date: Mon, 21 Sep 2026 21:14:30 -0300 Subject: [PATCH 1/2] feat(queen): a bee can run in its own container, so the swarm is no longer one machine wide A bee was a thread of the Queen. `dispatchBee` cut a worktree on her volume and sent the turn to her own `/chat`, so every bee shared her memory and her disk, and the widest the swarm could ever be was one container. Measured on the deployed one, 2026-09-18: about a gigabyte a bee against a 24 GB limit, and a 50 GB volume holding every worktree. Credentials stopped being the limit the day the key list was read without a counter; the container was the limit after that. So decision and execution are separated. The Queen still chooses everything she chose before - the issue, its boundary, the credential - and writes it down; a runner, in its own container, takes the order and does the work. Add replicas and the swarm is wider. THE ROW IS THE PROTOCOL. `queen_dispatch` already says which issue is in flight, under which boundary, on which credential, and a review sweep and two reapers read it. A second queue beside it would be a second answer to one question. With TRIOS_QUEEN_BEES_RUN_ELSEWHERE=on the Queen writes the row with `queued_at` and the brief and runs nothing; the row is in flight from that moment, so the boundary is held and the key is taken exactly as before. A runner (TRIOS_BEE_RUNNER_SECONDS set) claims one with an UPDATE over a row picked FOR UPDATE SKIP LOCKED, so two runners cannot take one order - proven against a real PostgreSQL 16 with eight runners reaching for one row at once. NO SECRET IN THE TABLE. The row carries `key_index` only. The runner resolves it with `workerProviderForKeyIndex` against the same worker variables the Queen reads; a runner whose variables differ refuses rather than reaching for the next key along, which is a different account. THE WORK STILL LEAVES WITHOUT A PUSH CREDENTIAL. A runner has no volume anyone can fetch from and may be gone minutes later, so it writes the bundle of base..queen-N into `queen_bundle`, and the export route serves it from there when the branch is not in its own checkout. `bundleOfBranch` is the existing bundle code, extracted so both can use it. (The export route file was already out of biome's format on the base; the pre-commit hook reformatted it, which is most of that file's diff.) THE REAPERS HAD TO LEARN WHERE A BEE RUNS, or the first restart would have been a disaster: - The boot reaper buried every unfinished row, because a restart of THIS container killed every bee in it. Runner bees do not die with the Queen. It now takes only rows the Queen ran herself (`queued_at IS NULL`). - A runner vouches for its bee by renewing `claimed_at` every fifteen seconds while it waits for the ending. The stall reaper releases a claimed row whose runner has been silent for ten minutes - a vanished container - and never salvages it, because nothing of that bee is on this disk. - The runner resets `dispatched_at` when the turn really starts, so the two-hour rule measures a turn and not the time the order waited. The shared half of starting a turn - cut the worktree, start the turn, hand back `begin` so the stream is read only after the caller's row exists - is one function, `cutAndStart`, used by both paths; the ledger around it differs (an insert for the Queen, an update for a runner) and stays with each caller. `markBeeRunningHere` from #505 moved into it with the rest of the start. Off by default on both sides. A deployment may be a Queen, a runner, or both. Co-Authored-By: Claude Opus 5 --- .../server/src/api/routes/queen-export.ts | 386 +++++++++------- .../server/src/api/services/queen-dispatch.ts | 422 +++++++++++++++--- .../server/src/api/services/queen-runner.ts | 293 ++++++++++++ .../apps/server/src/lib/db/pg-migrate.ts | 48 ++ trios/agent-server/apps/server/src/main.ts | 6 +- .../server/tests/api/queen-dispatch.test.ts | 200 +++++++++ .../server/tests/api/queen-runner.test.ts | 169 +++++++ .../tests/pglive/queen-runner-live.test.ts | 271 +++++++++++ 8 files changed, 1591 insertions(+), 204 deletions(-) create mode 100644 trios/agent-server/apps/server/src/api/services/queen-runner.ts create mode 100644 trios/agent-server/apps/server/tests/api/queen-runner.test.ts create mode 100644 trios/agent-server/apps/server/tests/pglive/queen-runner-live.test.ts diff --git a/trios/agent-server/apps/server/src/api/routes/queen-export.ts b/trios/agent-server/apps/server/src/api/routes/queen-export.ts index 53ed5d81bd..cc32bf0d0d 100644 --- a/trios/agent-server/apps/server/src/api/routes/queen-export.ts +++ b/trios/agent-server/apps/server/src/api/routes/queen-export.ts @@ -35,9 +35,10 @@ import { readFile, rm } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' import { Hono } from 'hono' +import { createQueenPool } from '../../lib/db/queen-pool' import { logger } from '../../lib/logger' import { baseRef, workspaceRoot } from '../services/queen-dispatch' - +import { queenLeaseDatabaseUrl } from '../services/queen-lease' /** * Git, run so that nothing in the checkout can steer it. @@ -128,166 +129,251 @@ function git( }) } -/** `queen-1234` -> 1234, and nothing else. */ -function issueOf(branch: string): number | null { - const m = /^queen-(\d+)$/.exec(branch.trim()) - return m ? Number(m[1]) : null +export interface BundledBranch { + ok: true + branch: string + base: string + commits: Array<{ sha: string; author: string; date: string; subject: string }> + files: string[] + bytes: Buffer } -export function createQueenExportRoute() { - return new Hono() - /** - * What is waiting to be published. - * - * Only branches that are AHEAD of the base: a branch level with it carries - * no work, and listing it would send a publisher to fetch an empty bundle. - */ - .get('/', async (c) => { - const root = workspaceRoot() - const base = baseRef() - const listed = await git(['-C', root, 'branch', '--list', 'queen-*']) - if (listed.code !== 0) { - return c.json( - { error: `could not list branches: ${listed.err.slice(0, 300)}` }, - 500, - ) - } - const branches = listed.out - .split('\n') - .map((l) => l.replace(/^[*+]?\s*/, '').trim()) - .filter((l) => l.length > 0) +/** + * One branch as a git bundle, made from the checkout this process can see. + * + * A bundle rather than a patch because it carries the commits themselves - + * their shas, authors and dates survive, so what lands upstream is what the bee + * actually wrote rather than a replay of it. It is emitted relative to the + * base, so it holds only this branch's work. + * + * Exported because the container that MAKES a bundle is no longer always the + * one that serves it: a runner bundles its own work before its worktree and its + * container go away. + */ +export async function bundleOfBranch( + issue: number, +): Promise { + const root = workspaceRoot() + const base = baseRef() + const branch = `queen-${issue}` - const waiting: Array<{ - issue: number - branch: string - commits: number - files: number - head: string - }> = [] - for (const branch of branches) { - const issue = issueOf(branch) - if (issue === null) continue - const count = await git([ - '-C', - root, - 'rev-list', - '--count', - `${base}..${branch}`, - ]) - const ahead = Number(count.out.trim()) - if (!Number.isInteger(ahead) || ahead <= 0) continue - const names = await git([ - '-C', - root, - 'diff', - '--name-only', - `${base}...${branch}`, - ]) - const head = await git(['-C', root, 'rev-parse', branch]) - waiting.push({ - issue, - branch, - commits: ahead, - files: names.out.split('\n').filter((l) => l.trim()).length, - head: head.out.trim(), - }) + const exists = await git(['-C', root, 'rev-parse', '--verify', branch]) + if (exists.code !== 0) { + return { + ok: false, + status: 404, + error: `no branch ${branch} in this checkout`, + } + } + const ahead = await git([ + '-C', + root, + 'rev-list', + '--count', + `${base}..${branch}`, + ]) + if (Number(ahead.out.trim()) <= 0) { + return { + ok: false, + status: 409, + error: `${branch} has no commits beyond ${base}`, + } + } + // Written to a path OUTSIDE the checkout: a bundle created inside a tree the + // agents can write is a file they can replace between creation and read. + // Written by the BEE (git drops to it), then read by this process as root. + const path = join(tmpdir(), `queen-export-${issue}-${randomUUID()}.bundle`) + try { + const made = await git( + ['-C', root, 'bundle', 'create', path, `${base}..${branch}`], + 120_000, + ) + if (made.code !== 0) { + return { + ok: false, + status: 500, + error: `bundle failed: ${made.err.slice(0, 400)}`, } - waiting.sort((a, b) => a.issue - b.issue) - return c.json({ base, count: waiting.length, branches: waiting }) - }) + } + const log = await git([ + '-C', + root, + 'log', + '--format=%H%x1f%an%x1f%aI%x1f%s', + `${base}..${branch}`, + ]) + const commits = log.out + .split('\n') + .filter((l) => l.trim()) + .map((l) => { + const [sha, author, date, subject] = l.split('\x1f') + return { sha, author, date, subject } + }) + const files = await git([ + '-C', + root, + 'diff', + '--name-only', + `${base}...${branch}`, + ]) + return { + ok: true, + branch, + base, + commits, + files: files.out.split('\n').filter((l) => l.trim()), + bytes: await readFile(path), + } + } finally { + await rm(path, { force: true }).catch(() => {}) + } +} - /** - * One branch, as a git bundle. - * - * A bundle rather than a patch because it carries the commits themselves - - * their shas, authors and dates survive, so what lands upstream is what the - * bee actually wrote rather than a replay of it. It is emitted relative to - * the base, so it holds only this branch's work and stays small enough to - * pass through JSON. - */ - .get('/:issue', async (c) => { - const issue = Number(c.req.param('issue')) - if (!Number.isInteger(issue) || issue <= 0) { - return c.json({ error: 'issue must be a positive integer' }, 400) - } - const root = workspaceRoot() - const base = baseRef() - const branch = `queen-${issue}` +/** The bundle a runner left behind, or null. Never a reason to fail a request. */ +async function storedBundle( + issue: number, +): Promise | null> { + const url = queenLeaseDatabaseUrl() + if (!url) return null + const pool = createQueenPool(url) + try { + const rows = await pool.query( + 'SELECT branch, base, runner, created_at, bytes FROM queen_bundle WHERE issue = $1', + [issue], + ) + const row = rows.rows?.[0] + if (!row) return null + return { + issue, + branch: String(row.branch), + base: String(row.base), + runner: String(row.runner), + at: row.created_at, + // The commit list is not stored: it is in queen_transcript and on the + // board, and a second copy is a second thing that can disagree. + bundleBase64: Buffer.from(row.bytes as Buffer).toString('base64'), + } + } catch (error) { + logger.warn('Queen export could not read a stored bundle', { + issue, + error: error instanceof Error ? error.message : String(error), + }) + return null + } finally { + await pool.end().catch(() => undefined) + } +} - const exists = await git(['-C', root, 'rev-parse', '--verify', branch]) - if (exists.code !== 0) { - return c.json({ error: `no branch ${branch} in this checkout` }, 404) - } - const ahead = await git([ - '-C', - root, - 'rev-list', - '--count', - `${base}..${branch}`, - ]) - if (Number(ahead.out.trim()) <= 0) { - return c.json( - { error: `${branch} has no commits beyond ${base}` }, - 409, - ) - } +/** `queen-1234` -> 1234, and nothing else. */ +function issueOf(branch: string): number | null { + const m = /^queen-(\d+)$/.exec(branch.trim()) + return m ? Number(m[1]) : null +} - // Written to a path OUTSIDE the checkout: a bundle created inside a tree - // the agents can write is a file they can replace between creation and - // read. - // Written by the BEE (git drops to it), then read by this process as root: - // a path both can reach, and outside the checkout so the agents cannot - // swap the file between its creation and its read. - const path = join(tmpdir(), `queen-export-${issue}-${randomUUID()}.bundle`) - try { - const made = await git( - ['-C', root, 'bundle', 'create', path, `${base}..${branch}`], - 120_000, - ) - if (made.code !== 0) { +export function createQueenExportRoute() { + return ( + new Hono() + /** + * What is waiting to be published. + * + * Only branches that are AHEAD of the base: a branch level with it carries + * no work, and listing it would send a publisher to fetch an empty bundle. + */ + .get('/', async (c) => { + const root = workspaceRoot() + const base = baseRef() + const listed = await git(['-C', root, 'branch', '--list', 'queen-*']) + if (listed.code !== 0) { return c.json( - { error: `bundle failed: ${made.err.slice(0, 400)}` }, + { error: `could not list branches: ${listed.err.slice(0, 300)}` }, 500, ) } - const log = await git([ - '-C', - root, - 'log', - '--format=%H%x1f%an%x1f%aI%x1f%s', - `${base}..${branch}`, - ]) - const commits = log.out + const branches = listed.out .split('\n') - .filter((l) => l.trim()) - .map((l) => { - const [sha, author, date, subject] = l.split('\x1f') - return { sha, author, date, subject } + .map((l) => l.replace(/^[*+]?\s*/, '').trim()) + .filter((l) => l.length > 0) + + const waiting: Array<{ + issue: number + branch: string + commits: number + files: number + head: string + }> = [] + for (const branch of branches) { + const issue = issueOf(branch) + if (issue === null) continue + const count = await git([ + '-C', + root, + 'rev-list', + '--count', + `${base}..${branch}`, + ]) + const ahead = Number(count.out.trim()) + if (!Number.isInteger(ahead) || ahead <= 0) continue + const names = await git([ + '-C', + root, + 'diff', + '--name-only', + `${base}...${branch}`, + ]) + const head = await git(['-C', root, 'rev-parse', branch]) + waiting.push({ + issue, + branch, + commits: ahead, + files: names.out.split('\n').filter((l) => l.trim()).length, + head: head.out.trim(), }) - const files = await git([ - '-C', - root, - 'diff', - '--name-only', - `${base}...${branch}`, - ]) - const bytes = await readFile(path) - logger.info('Queen export served', { - issue, - branch, - commits: commits.length, - bundleBytes: bytes.length, - }) - return c.json({ - issue, - branch, - base, - commits, - files: files.out.split('\n').filter((l) => l.trim()), - bundleBase64: bytes.toString('base64'), - }) - } finally { - await rm(path, { force: true }).catch(() => {}) - } - }) + } + waiting.sort((a, b) => a.issue - b.issue) + return c.json({ base, count: waiting.length, branches: waiting }) + }) + + /** + * One branch, as a git bundle. + * + * A bundle rather than a patch because it carries the commits themselves - + * their shas, authors and dates survive, so what lands upstream is what the + * bee actually wrote rather than a replay of it. It is emitted relative to + * the base, so it holds only this branch's work and stays small enough to + * pass through JSON. + */ + .get('/:issue', async (c) => { + const issue = Number(c.req.param('issue')) + if (!Number.isInteger(issue) || issue <= 0) { + return c.json({ error: 'issue must be a positive integer' }, 400) + } + const made = await bundleOfBranch(issue) + if (made.ok) { + logger.info('Queen export served', { + issue, + branch: made.branch, + commits: made.commits.length, + bundleBytes: made.bytes.length, + }) + return c.json({ + issue, + branch: made.branch, + base: made.base, + commits: made.commits, + files: made.files, + bundleBase64: made.bytes.toString('base64'), + }) + } + // Not in THIS checkout - which is the normal case once bees run in their + // own containers. A runner has no volume anyone can fetch from, so it + // leaves the bundle in the database on its way out; here is where it is + // collected. Looked at only after the local branch was not found, so a + // container that still runs its own bees behaves exactly as it did. + if (made.status === 404) { + const stored = await storedBundle(issue) + if (stored) return c.json(stored) + } + return c.json({ error: made.error }, made.status as 404 | 409 | 500) + }) + ) } diff --git a/trios/agent-server/apps/server/src/api/services/queen-dispatch.ts b/trios/agent-server/apps/server/src/api/services/queen-dispatch.ts index 722e6ea12b..f440952c87 100644 --- a/trios/agent-server/apps/server/src/api/services/queen-dispatch.ts +++ b/trios/agent-server/apps/server/src/api/services/queen-dispatch.ts @@ -917,6 +917,89 @@ function configuredEndpointProvider( } } +/** + * The credential a durable `key_index` names, for whoever has to USE it later. + * + * `resolveWorkerProvider` chooses; this one remembers. A bee that runs outside + * the Queen is handed its issue through `queen_dispatch`, and the row carries + * the index and never the key: a secret in a table is a secret in every backup, + * every logical replica and every `SELECT *` a dashboard ever runs. The runner + * reads the same variables the Queen reads and resolves the index against them, + * so the credential exists in exactly the two places it already did - the + * environment and the request - and nowhere in between. + * + * Null when the index names nothing: keys were removed, or the runner's + * environment is not the Queen's. Refusing is the only safe answer, because the + * next index along is a DIFFERENT account with a different rate limit, and + * quietly using it would put two bees on one credential while the ledger says + * otherwise. + */ +export function workerProviderForKeyIndex( + keyIndex: number, +): WorkerProvider | null { + const override = process.env.TRIOS_QUEEN_WORKER_MODEL + if (configuredWorkerBaseUrl()) { + const pools = configuredEndpointPools() + if (pools.length === 0) return null + if (pools.length > 1) { + for (const pool of pools) { + const position = keyIndex - (pool.number - 1) * POOL_KEY_STRIDE + if (position < 0 || position >= pool.keys.length) continue + return { + provider: pool.provider, + model: pool.model, + baseUrl: pool.baseUrl, + apiKey: pool.keys[position], + keyIndex, + keyCount: pool.keys.length, + poolNumber: pool.number, + poolCount: pools.length, + contextWindow: pool.contextWindow, + } + } + return null + } + const only = pools[0] + if (keyIndex < 0 || keyIndex >= only.keys.length) { + // Ollama is one inference slot with one token name, and its index is 0. + return only.provider === 'ollama' && keyIndex === 0 + ? { + provider: 'ollama', + model: only.model, + baseUrl: only.baseUrl, + apiKey: only.keys[0] || 'local', + keyIndex: 0, + keyCount: 1, + contextWindow: only.contextWindow, + } + : null + } + return { + provider: only.provider, + model: only.model, + baseUrl: only.baseUrl, + apiKey: only.keys[keyIndex], + keyIndex, + keyCount: only.keys.length, + contextWindow: only.contextWindow, + } + } + for (const candidate of WORKER_PROVIDERS) { + const keys = keysFor(candidate.envVar) + if (keys.length === 0) continue + if (keyIndex < 0 || keyIndex >= keys.length) return null + return { + provider: candidate.provider, + model: override || candidate.model, + apiKey: keys[keyIndex], + keyIndex, + keyCount: keys.length, + laneCount: workerLanesFor(candidate.provider), + } + } + return null +} + /** * A bounded number of concurrent lanes per distinct credential. * @@ -4258,11 +4341,18 @@ async function salvageBeforeRelease( const salvage = deps.salvage ?? salvageDispatch const budgetMs = deps.budgetMs ?? SALVAGE_SWEEP_BUDGET_MS const due: number[] = [] + // Released, never salvaged: a bee that ran - or was to run - in a runner's + // container left nothing on THIS disk, and a salvage here would be git on a + // worktree that does not exist. + const elsewhere = new Set() try { const rows = await pool.query(sql, params) for (const row of rows.rows ?? []) { const issue = Number((row as { issue?: unknown }).issue) - if (Number.isFinite(issue)) due.push(issue) + if (!Number.isFinite(issue)) continue + due.push(issue) + if ((row as { queued_at?: unknown }).queued_at != null) + elsewhere.add(issue) } } catch (error) { logger.warn('Queen could not read the rows a reaper is about to release', { @@ -4282,6 +4372,7 @@ async function salvageBeforeRelease( }) continue } + if (elsewhere.has(issue)) continue // AND NEVER FOR LONGER THAN THE BUDGET. The release is what the reaper is // for; the repair is what it would like to do on the way. A sweep inside // the round gate that spends unbounded git time is a swarm that dispatches @@ -4304,8 +4395,12 @@ export async function reapDispatchesFromPreviousBoot( ): Promise { const due = await salvageBeforeRelease( pool, - `SELECT issue FROM queen_dispatch - WHERE started = true AND finished_at IS NULL`, + // ONLY BEES THIS CONTAINER RAN. A row the Queen QUEUED belongs to a runner, + // and a runner does not die when the Queen restarts: releasing its row + // would free the boundary and the key of a bee that is still writing, and + // the next round would dispatch the same issue into a second container. + `SELECT issue, queued_at FROM queen_dispatch + WHERE started = true AND finished_at IS NULL AND queued_at IS NULL`, [], deps, ) @@ -4326,7 +4421,7 @@ export async function reapDispatchesFromPreviousBoot( `UPDATE queen_dispatch SET finished_at = now(), outcome = '${DISPATCH_OUTCOME_LABELS.reapedAtBoot}: the container running this turn was replaced' - WHERE started = true AND finished_at IS NULL + WHERE started = true AND finished_at IS NULL AND queued_at IS NULL AND issue = ANY($1::int[]) RETURNING issue`, [due], @@ -4334,6 +4429,18 @@ export async function reapDispatchesFromPreviousBoot( return reaped.rows.map((r) => r.issue as number) } +/** + * How long a runner may go without vouching for its bee before the bee is + * presumed dead with it. + * + * A runner renews `claimed_at` every few seconds while its turn streams. Its + * container can vanish without a word - a redeploy, a crash, a scale-down - and + * nothing else would ever end the row: the Queen does not see that process and + * the two-hour rule would hold the boundary and the key for two hours. Ten + * minutes is forty missed renewals, which is not a slow network. + */ +export const RUNNER_SILENT_MINUTES = 10 + export async function reapStalledDispatches( pool: Pool, stallMinutes = 120, @@ -4341,11 +4448,13 @@ export async function reapStalledDispatches( ): Promise { const due = await salvageBeforeRelease( pool, - `SELECT issue FROM queen_dispatch + `SELECT issue, queued_at FROM queen_dispatch WHERE started = true AND finished_at IS NULL - AND dispatched_at < now() - make_interval(mins => $1)`, - [stallMinutes], + AND (dispatched_at < now() - make_interval(mins => $1) + OR (claimed_by IS NOT NULL + AND claimed_at < now() - make_interval(mins => $2)))`, + [stallMinutes, RUNNER_SILENT_MINUTES], deps, ) if (due.length === 0) return [] @@ -4363,10 +4472,12 @@ export async function reapStalledDispatches( outcome = '${DISPATCH_OUTCOME_LABELS.reapedStalled}: no completion within ' || $1 || ' minutes' WHERE started = true AND finished_at IS NULL - AND dispatched_at < now() - make_interval(mins => $1) + AND (dispatched_at < now() - make_interval(mins => $1) + OR (claimed_by IS NOT NULL + AND claimed_at < now() - make_interval(mins => $3))) AND issue = ANY($2::int[]) RETURNING issue`, - [stallMinutes, due], + [stallMinutes, due, RUNNER_SILENT_MINUTES], ) return reaped.rows.map((r) => r.issue as number) } @@ -4453,6 +4564,54 @@ export async function dispatchBee( return { started: false, issue, branch, detail } } + // THE HANDOVER, when bees run in their own containers. + // + // The Queen's work ends here: she has chosen the issue, its boundary and its + // credential, and the row IS the order. She asks nothing of her own container + // because the bee will not run in it, and cuts no worktree because the runner + // cuts its own. The row is in flight from this moment - started, unfinished - + // so the boundary is held and the key is taken exactly as before, and every + // reader of the table keeps working without knowing where the bee runs. + if (beesRunElsewhere()) { + const conversationId = randomUUID() + const detail = + `queued for a runner; ${chosen.provider}/${chosen.model}` + + (chosen.poolCount && chosen.poolCount > 1 + ? ` pool ${chosen.poolNumber}` + : '') + + (chosen.keyCount && chosen.keyCount > 1 + ? ` key ${((chosen.keyIndex ?? 0) % POOL_KEY_STRIDE) + 1}/${chosen.keyCount}` + : '') + await recordDispatch( + pool, + issue, + branch, + true, + detail, + ownedPaths, + conversationId, + chosen.keyIndex, + criteria, + criteriaSource, + chosen.provider, + chosen.model, + { brief }, + ) + logger.info('Queen queued a bee for a runner', { + issue, + branch, + keyIndex: chosen.keyIndex, + }) + return { + started: true, + issue, + branch, + detail, + conversationId, + keyIndex: chosen.keyIndex, + } + } + // A key is free and a provider answers. The remaining question is the one // nothing used to ask: can the CONTAINER carry another bee? Asked here, before // a worktree is cut, so a refusal costs a measurement and nothing else. The @@ -4480,28 +4639,95 @@ export async function dispatchBee( } } + const cut = await cutAndStart( + pool, + { issue, branch, brief, ownedPaths, conversationId: randomUUID(), chosen }, + deps, + ) + if (!cut.ok) { + await recordDispatch(pool, issue, branch, false, cut.detail, ownedPaths) + return { started: false, issue, branch, detail: cut.detail } + } + + await recordDispatch( + pool, + issue, + branch, + cut.started, + cut.detail, + ownedPaths, + cut.conversationId, + chosen.keyIndex, + criteria, + criteriaSource, + chosen.provider, + chosen.model, + ) + cut.begin() + logger.info('Queen dispatch', { + issue, + branch, + started: cut.started, + detail: cut.detail, + }) + return { + started: cut.started, + issue, + branch, + detail: cut.detail, + conversationId: cut.conversationId, + keyIndex: chosen.keyIndex, + } +} + +interface BeePlan { + issue: number + branch: string + brief: string + ownedPaths: string[] + conversationId: string + chosen: WorkerProvider +} + +/** + * Cut the worktree and hand the turn to the agent, wherever this is running. + * + * The half of a dispatch that is the same whether the Queen is starting the bee + * herself or a runner is starting one she ordered. What differs is the ledger + * around it - an insert here, an update there - and that stays with the caller, + * because a row written by the wrong rule is a bee nobody can account for. + * + * `begin` is handed back rather than called: the stream must not be read until + * the caller's row exists, or the reader that closes the row can outrun its + * creation. That race was measured, and the phantom it leaves is a bee that has + * stopped and looks like it is running until the reaper comes. + */ +async function cutAndStart( + pool: Pool, + plan: BeePlan, + deps: BeeRoomDeps, +): Promise< + | { ok: false; detail: string } + | { + ok: true + started: boolean + detail: string + conversationId: string + begin: () => void + } +> { + const { issue, brief, ownedPaths, conversationId, chosen } = plan const worktree = await prepareWorktree(issue, { running: () => runningBeeBranches(pool), volumeUsed: deps.volumeUsed, }) - if (!worktree.ok) { - await recordDispatch( - pool, - issue, - branch, - false, - worktree.detail, - ownedPaths, - ) - return { started: false, issue, branch, detail: worktree.detail } - } + if (!worktree.ok) return { ok: false, detail: worktree.detail } // The PROJECT inside the checkout, not the checkout root. A worktree is a // clone of the repository and this project is a directory inside it, so // standing at the root makes every project-relative boundary resolve one // level too high - the bee writes `/docs/x.md` where the committer // looks for `trios/docs/x.md`, and its work reads as no work at all. - // The PROJECT inside the checkout — which is not always a subdirectory. // // This was hardcoded to `/trios`, which is right for BrowserOS and wrong for // every other repository: aimed at a repo whose code sits at its root, the @@ -4514,7 +4740,6 @@ export async function dispatchBee( const subdir = repoSubdir() const workingDirectory = subdir ? `${worktree.path}/${subdir}` : worktree.path - const conversationId = randomUUID() const turn = await startTurn( pool, issue, @@ -4539,41 +4764,103 @@ export async function dispatchBee( : '') + (chosen.rehearsal ? ' (REHEARSAL - a recorded stream, not a model)' : '') : turn.detail - - await recordDispatch( - pool, - issue, - branch, - turn.ok, + return { + ok: true, + started: turn.ok, detail, - ownedPaths, conversationId, - chosen.keyIndex, - criteria, - criteriaSource, - chosen.provider, - chosen.model, + begin: () => { + // A bee that just started has not allocated what it will hold. The next + // dispatch of this round must not read the container as empty because of + // it. + if (turn.ok) noteBeeStarted(conversationId, (deps.now ?? Date.now)()) + // THIS PROCESS NOW KNOWS THIS BEE IS ALIVE, which is a question the + // database cannot answer: a row that has not been closed for two hours is + // a row, not a corpse. The stall sweep reads this before it salvages, so + // a long turn's worktree is never committed from under it. + // `closeDispatch` clears it. + if (turn.ok) markBeeRunningHere(issue, conversationId) + // ONLY NOW may the stream be read. Everything that reads the bee's output + // eventually writes to the caller's row, and a writer that can outrun the + // row's creation is a writer that silently updates nothing. + turn.beginDrain?.() + }, + } +} + +/** + * Run a bee the Queen ordered and a runner claimed. + * + * Same turn, same worktree, same transcript; what is different is only where it + * happens and what the ledger already holds. The row exists and is in flight, + * so a success UPDATES its detail rather than writing it again - a second + * insert would archive a history row and clear the criteria the bee is to be + * judged against. A failure is recorded as a refusal, which ends the row and + * hands the issue and its boundary back to the next round; nothing else here + * could release them, and an order nobody can retry is worse than one nobody + * took. + */ +export async function runClaimedBee( + pool: Pool, + order: { + issue: number + branch: string + brief: string + ownedPaths: string[] + conversationId: string + keyIndex: number + }, + deps: BeeRoomDeps = {}, +): Promise { + const { issue, branch, ownedPaths } = order + const refuse = async (detail: string): Promise => { + await recordDispatch(pool, issue, branch, false, detail, ownedPaths) + return { started: false, issue, branch, detail } + } + + const chosen = workerProviderForKeyIndex(order.keyIndex) + if (!chosen?.apiKey) { + // The runner's environment is not the Queen's, or the key list changed + // under the order. Taking the next credential along would put two bees on + // one account while the ledger says otherwise. + return refuse( + `this runner cannot resolve key_index ${order.keyIndex}: its worker ` + + 'variables differ from the ones the order was written against', + ) + } + const noRoom = await beeRoomRefusal(pool, branch, deps) + if (noRoom) { + logger.warn('Runner claimed a bee but its container has no room', { + issue, + resource: noRoom.resource, + detail: noRoom.detail, + }) + return refuse(noRoom.detail) + } + + const cut = await cutAndStart(pool, { ...order, chosen }, deps) + if (!cut.ok || !cut.started) { + return refuse(cut.ok ? cut.detail : cut.detail) + } + // `dispatched_at` too: the two-hour rule measures a TURN, and until this + // moment the row was only an order waiting in a queue. Left at the queue time + // a bee that waited an hour for a runner would be reaped an hour into its + // work. + await pool.query( + `UPDATE queen_dispatch + SET detail = $2, claimed_at = now(), dispatched_at = now() + WHERE issue = $1`, + [issue, cut.detail], ) - // A bee that just started has not allocated what it will hold. The next - // dispatch of this round must not read the container as empty because of it. - if (turn.ok) noteBeeStarted(conversationId, (deps.now ?? Date.now)()) - // THIS PROCESS NOW KNOWS THIS BEE IS ALIVE, which is a question the database - // cannot answer: a row that has not been closed for two hours is a row, not - // a corpse. The stall sweep reads this before it salvages, so a long turn's - // worktree is never committed from under it. `closeDispatch` clears it. - if (turn.ok) markBeeRunningHere(issue, conversationId) - // ONLY NOW may the stream be read. Everything that reads the bee's output - // eventually writes to the row above, and a writer that can outrun the row's - // creation is a writer that silently updates nothing. - turn.beginDrain?.() - logger.info('Queen dispatch', { issue, branch, started: turn.ok, detail }) + cut.begin() + logger.info('Runner started a bee', { issue, branch, detail: cut.detail }) return { - started: turn.ok, + started: true, issue, branch, - detail, - conversationId, - keyIndex: chosen.keyIndex, + detail: cut.detail, + conversationId: cut.conversationId, + keyIndex: order.keyIndex, } } @@ -4799,6 +5086,12 @@ export async function recordDispatch( */ provider?: string, model?: string, + /** + * Set only when this row is an ORDER for a runner rather than a bee already + * running here. The brief travels with it because a runner has no issue body + * to build one from, and `queued_at` is what a runner looks for. + */ + queue?: { brief: string }, ): Promise { // #1360. A dispatch that never started is recorded with its ending, and the // ending is ONE WORD. It used to be the refusal detail verbatim, which is @@ -4854,11 +5147,12 @@ export async function recordDispatch( `INSERT INTO queen_dispatch (issue, branch, started, detail, owned_paths, conversation_id, dispatched_at, finished_at, outcome, key_index, - criteria, criteria_source, provider, model) + criteria, criteria_source, provider, model, + queued_at, brief, claimed_by, claimed_at) VALUES ($1, $2, $3, $4, $5::jsonb, $6, now(), CASE WHEN $3 THEN NULL ELSE now() END, CASE WHEN $3 THEN NULL ELSE $12 END, - $7, $8::jsonb, $9, $10, $11) + $7, $8::jsonb, $9, $10, $11, $13, $14, NULL, NULL) ON CONFLICT (issue) DO UPDATE SET branch = EXCLUDED.branch, started = EXCLUDED.started, @@ -4878,6 +5172,13 @@ export async function recordDispatch( output_tokens = NULL, criteria = EXCLUDED.criteria, criteria_source = EXCLUDED.criteria_source, + -- A re-dispatch is a NEW order: whoever claimed the last one is not + -- working on this one, and a stale claim would keep every runner off + -- the row for ever. + queued_at = EXCLUDED.queued_at, + brief = EXCLUDED.brief, + claimed_by = NULL, + claimed_at = NULL, provider = EXCLUDED.provider, model = EXCLUDED.model, -- A dispatch that starts again is new work, so last turn's verdict @@ -4950,6 +5251,21 @@ export async function recordDispatch( provider ?? null, model ?? null, outcome, + queue ? new Date() : null, + queue?.brief ?? null, ], ) } + +/** + * Whether this deployment runs its bees somewhere else. + * + * Off by default, and deliberately: the swarm this was written for is one + * container that has always run its own bees, and a supervisor that silently + * stops doing the work and starts writing orders nobody collects is a swarm + * that looks busy and moves nothing. Turning it on is a deployment decision, + * made once, alongside the runners that answer it. + */ +export function beesRunElsewhere(): boolean { + return process.env.TRIOS_QUEEN_BEES_RUN_ELSEWHERE === 'on' +} diff --git a/trios/agent-server/apps/server/src/api/services/queen-runner.ts b/trios/agent-server/apps/server/src/api/services/queen-runner.ts new file mode 100644 index 0000000000..e013edf404 --- /dev/null +++ b/trios/agent-server/apps/server/src/api/services/queen-runner.ts @@ -0,0 +1,293 @@ +/** + * @license + * Copyright 2025 BrowserOS + * SPDX-License-Identifier: AGPL-3.0-or-later + * + * A BEE THAT DOES NOT LIVE INSIDE THE QUEEN. + * + * Until now a bee was a thread of the supervisor: `dispatchBee` cut a worktree + * on her volume and sent the turn to her own `/chat`. Everything the swarm + * could ever be was therefore one container. Measured on the deployed one, + * 2026-09-18: 24 GB of memory, about a gigabyte a bee, a 50 GB volume holding + * every worktree, and a ceiling the operator kept moving up and down against + * numbers that were really the container's and not the swarm's. + * + * So the work is handed over instead. The Queen still decides everything she + * decided before - which issue, which boundary, which credential - and writes + * it as a row. This file is the other half: a process that takes one such row + * and does the work in ITS container, with its own memory, its own disk and its + * own checkout. Add a replica and the swarm is wider; take one away and the + * rows it had not claimed are still there. + * + * WHY THE ROW IS THE PROTOCOL. There is already a table that says which issue + * is in flight, under which boundary, on which credential, and a review sweep + * and two reapers that read it. A queue beside it would be a second answer to + * the same question, and the two would disagree the first time a process died + * between them. `UPDATE ... WHERE claimed_by IS NULL` over a row picked `FOR + * UPDATE SKIP LOCKED` is the whole of the mutual exclusion: two runners cannot + * take one issue, and a runner that dies without finishing leaves exactly what + * a Queen that died mid-bee always left, which the boot reaper already clears. + * + * WHAT IS NOT SENT. The credential. The row carries `key_index`, and the runner + * resolves it against the same worker variables the Queen reads, so a secret + * never enters the table, the backups or anyone's `SELECT *`. A runner whose + * variables differ cannot resolve the index and says so rather than reaching + * for the next key along, which would be a different account. + */ +import { hostname } from 'node:os' +import type { Pool } from 'pg' +import { createQueenPool } from '../../lib/db/queen-pool' +import { logger } from '../../lib/logger' +import { bundleOfBranch } from '../routes/queen-export' +import { runClaimedBee } from './queen-dispatch' +import { queenLeaseDatabaseUrl } from './queen-lease' + +export interface BeeOrder { + issue: number + branch: string + brief: string + ownedPaths: string[] + conversationId: string + keyIndex: number +} + +/** Who this runner is, in logs and in `claimed_by`. */ +export function runnerName(): string { + const stated = process.env.TRIOS_BEE_RUNNER_NAME?.trim() + if (stated) return stated.slice(0, 80) + // Railway gives every replica the same service name, so the pid is what + // separates two runners in one image. A name that collides is not a + // correctness problem - the claim is - but it is an unreadable log. + return `${hostname()}:${process.pid}` +} + +/** How many bees this replica carries at once. One by default, and small on purpose. */ +export function runnerSlots(): number { + const parsed = Number(process.env.TRIOS_BEE_RUNNER_SLOTS) + if (!Number.isInteger(parsed) || parsed < 1) return 1 + return Math.min(parsed, 64) +} + +function pollSeconds(): number { + const parsed = Number(process.env.TRIOS_BEE_RUNNER_SECONDS ?? '0') + if (!Number.isInteger(parsed) || parsed < 5) return 0 + return Math.min(parsed, 3600) +} + +/** + * Take one order, or nothing. + * + * `FOR UPDATE SKIP LOCKED` rather than a transaction two runners queue behind: + * a runner that waits for another runner's lock is a runner doing nothing while + * an unclaimed row sits next to it. Oldest first, because an order that has + * waited longest is the one whose issue has been held longest. + */ +export async function claimQueuedBee( + pool: Pool, + runner: string, +): Promise { + const rows = await pool.query( + `UPDATE queen_dispatch + SET claimed_by = $1, claimed_at = now() + WHERE issue = ( + SELECT issue FROM queen_dispatch + WHERE queued_at IS NOT NULL + AND claimed_by IS NULL + AND finished_at IS NULL + AND started = true + ORDER BY queued_at + FOR UPDATE SKIP LOCKED + LIMIT 1) + RETURNING issue, branch, brief, owned_paths, conversation_id, key_index`, + [runner], + ) + const row = rows.rows?.[0] + if (!row) return null + const keyIndex = Number(row.key_index) + const conversationId = String(row.conversation_id ?? '') + if (!Number.isInteger(keyIndex) || !conversationId) { + // An order the Queen could not complete. Left claimed rather than retried: + // a row without a credential or a conversation is not work this runner can + // do, and releasing it would hand every runner the same unusable order. + logger.warn('Runner claimed an order it cannot run', { + issue: row.issue, + hasConversation: Boolean(conversationId), + }) + return null + } + return { + issue: Number(row.issue), + branch: String(row.branch), + brief: String(row.brief ?? ''), + ownedPaths: Array.isArray(row.owned_paths) + ? (row.owned_paths as string[]) + : [], + conversationId, + keyIndex, + } +} + +/** + * Wait for the ending the drain writes, vouching for the bee while it runs. + * + * Each poll renews `claimed_at`, which is how the Queen tells a long turn from a + * dead runner: she cannot see this process, and without the renewal a runner + * that vanished mid-turn would hold its issue's boundary and credential until + * the two-hour rule. One statement does both jobs - a row that no longer + * matches (finished, reaped, or re-dispatched to somebody else) is the answer + * that this turn is no longer the table's business. + */ +export async function waitForEnding( + pool: Pool, + order: BeeOrder, + runner: string, + everyMs = 15_000, +): Promise { + for (;;) { + await new Promise((resolve) => setTimeout(resolve, everyMs)) + try { + const rows = await pool.query( + `UPDATE queen_dispatch SET claimed_at = now() + WHERE issue = $1 AND conversation_id = $2 AND claimed_by = $3 + AND finished_at IS NULL + RETURNING issue`, + [order.issue, order.conversationId, runner], + ) + if (!rows.rows?.length) return + } catch (error) { + // One failed renewal is not a dead runner; the Queen gives it forty. + logger.warn('Runner could not vouch for its bee', { + issue: order.issue, + error: error instanceof Error ? error.message : String(error), + }) + } + } +} + +/** + * Leave the work where somebody can reach it. + * + * The container publishes nothing - that rule is older than this file and the + * reason for it is in queen-export.ts - but a runner cannot even be VISITED: + * replicas share one address, its volume is its own, and it may be gone before + * anyone asks. So the bundle goes into the table the export route now reads. + * A turn that committed nothing has nothing to store, and that is not a + * failure: plenty of bees end with a verdict and no commit. + */ +export async function storeBundle( + pool: Pool, + issue: number, + runner: string, +): Promise { + const made = await bundleOfBranch(issue) + if (!made.ok) { + if (made.status !== 409) { + logger.warn('Runner could not bundle its work', { + issue, + error: made.error, + }) + } + return + } + await pool.query( + `INSERT INTO queen_bundle (issue, branch, base, runner, bytes) + VALUES ($1, $2, $3, $4, $5) + ON CONFLICT (issue) DO UPDATE + SET branch = EXCLUDED.branch, base = EXCLUDED.base, + runner = EXCLUDED.runner, bytes = EXCLUDED.bytes, + created_at = now()`, + [issue, made.branch, made.base, runner, made.bytes], + ) + logger.info('Runner stored the bundle of a finished bee', { + issue, + branch: made.branch, + commits: made.commits.length, + bundleBytes: made.bytes.length, + }) +} + +/** + * One order, start to finish. Resolves when the bee has ended and its work is + * where the publisher can fetch it, so the caller may then take another. + */ +export async function runOneOrder( + pool: Pool, + runner: string, +): Promise { + const order = await claimQueuedBee(pool, runner) + if (!order) return null + logger.info('Runner claimed a bee', { + runner, + issue: order.issue, + branch: order.branch, + }) + const outcome = await runClaimedBee(pool, order) + if (!outcome.started) return order + await waitForEnding(pool, order, runner) + await storeBundle(pool, order.issue, runner).catch((error) => { + logger.warn('Runner could not store the bundle of a finished bee', { + issue: order.issue, + error: error instanceof Error ? error.message : String(error), + }) + }) + return order +} + +/** + * Start taking orders, or explain why not. + * + * Off unless `TRIOS_BEE_RUNNER_SECONDS` is set, for the same reason the Queen's + * own loop is off by default: a server started on a laptop, in a test or beside + * the app must not quietly join the hive and take work nobody can then find. + */ +export function startBeeRunner(): void { + const every = pollSeconds() + if (!every) return + const url = queenLeaseDatabaseUrl() + if (!url) { + logger.warn('Bee runner requested but no database is configured') + return + } + const runner = runnerName() + const slots = runnerSlots() + const pool = createQueenPool(url) + let busy = 0 + let stopped = false + logger.info('Bee runner starting', { runner, slots, everySeconds: every }) + + const take = (): void => { + if (stopped || busy >= slots) return + busy += 1 + runOneOrder(pool, runner) + .then((order) => { + // A claim that found work asks again at once: a replica that waits out + // its interval between bees is a replica idle for no reason, and the + // queue is where the pacing belongs. + if (order && !stopped) setImmediate(take) + }) + .catch((error) => { + logger.warn('Bee runner round failed', { + runner, + error: error instanceof Error ? error.message : String(error), + }) + }) + .finally(() => { + busy -= 1 + }) + } + + const timer = setInterval(() => { + // Every free slot, not one: a replica with four slots and four orders + // waiting should be carrying four bees, not one a minute. + for (let slot = busy; slot < slots; slot++) take() + }, every * 1000) + take() + + process.once('SIGTERM', () => { + stopped = true + clearInterval(timer) + // The bees in flight are NOT abandoned here. Their rows stay claimed and + // unfinished, which is exactly what a Queen that died mid-bee always left, + // and the boot reaper clears it the same way. + }) +} diff --git a/trios/agent-server/apps/server/src/lib/db/pg-migrate.ts b/trios/agent-server/apps/server/src/lib/db/pg-migrate.ts index 29ac242525..f896284278 100644 --- a/trios/agent-server/apps/server/src/lib/db/pg-migrate.ts +++ b/trios/agent-server/apps/server/src/lib/db/pg-migrate.ts @@ -202,6 +202,54 @@ ALTER TABLE queen_dispatch ALTER TABLE queen_dispatch ADD COLUMN IF NOT EXISTS output_tokens bigint; +-- A dispatch that is WAITING for a runner, and the runner that took it. +-- +-- A bee used to be a thread of the Queen: dispatch cut a worktree on her volume +-- and sent the turn to her own /chat, so every bee shared her memory and her +-- disk and the swarm could never be wider than one container. Measured +-- 2026-09-18: about a gigabyte of memory a bee against a 24 GB limit. +-- +-- So the row becomes the handover. The Queen writes it with queued_at set and +-- runs nothing; a runner claims it with an UPDATE that skips locked rows, so +-- two runners cannot take one issue, and does the work in its own container. +-- Everything else about the row is unchanged, which is the point: it is in +-- flight from the moment it is queued, so its boundary is held and its key is +-- taken exactly as before, and every reader of this table - the board, the +-- reapers, the review sweep - keeps working without knowing where the bee runs. +-- +-- The brief travels with it because the runner has no issue body to build one +-- from. It is prompt text, and no secret is stored here: the runner resolves +-- its credential from key_index against the same environment the Queen reads. +ALTER TABLE queen_dispatch + ADD COLUMN IF NOT EXISTS queued_at timestamptz; +ALTER TABLE queen_dispatch + ADD COLUMN IF NOT EXISTS claimed_by text; +ALTER TABLE queen_dispatch + ADD COLUMN IF NOT EXISTS claimed_at timestamptz; +ALTER TABLE queen_dispatch + ADD COLUMN IF NOT EXISTS brief text; + +CREATE INDEX IF NOT EXISTS idx_queen_dispatch_queued + ON queen_dispatch (queued_at) + WHERE claimed_by IS NULL AND finished_at IS NULL; + +-- Work that finished in a container nobody can reach afterwards. +-- +-- Publishing is done from outside by whoever holds the token, against a bundle +-- this container hands out; that is the rule queen-export.ts exists to keep. A +-- runner replica has no volume anyone can fetch from and may be gone minutes +-- later, so it writes the bundle here, where the export route looks for it. The +-- bytes are a git bundle of base..queen-N and nothing else - no credential ever +-- travels this way, in either direction. +CREATE TABLE IF NOT EXISTS queen_bundle ( + issue int PRIMARY KEY, + branch text NOT NULL, + base text NOT NULL, + runner text NOT NULL, + created_at timestamptz NOT NULL DEFAULT now(), + bytes bytea NOT NULL +); + -- The attempts that queen_dispatch overwrote. -- -- queen_dispatch is keyed by issue alone, so dispatching an issue a second diff --git a/trios/agent-server/apps/server/src/main.ts b/trios/agent-server/apps/server/src/main.ts index f1349a94ed..c1ac716b3f 100644 --- a/trios/agent-server/apps/server/src/main.ts +++ b/trios/agent-server/apps/server/src/main.ts @@ -8,15 +8,16 @@ * Manages server lifecycle: initialization, startup, and shutdown. */ -import { markBrowserlessByConfiguration } from './agent/browserless' import fs from 'node:fs' import path from 'node:path' import { EXIT_CODES } from '@browseros/shared/constants/exit-codes' +import { markBrowserlessByConfiguration } from './agent/browserless' import { createHttpServer } from './api/server' import { configureVmRuntime, peekOpenClawService, } from './api/services/openclaw/openclaw-service' +import { startBeeRunner } from './api/services/queen-runner' import { startQueenTick } from './api/services/queen-tick' import { CdpBackend } from './browser/backends/cdp' import { Browser } from './browser/browser' @@ -86,6 +87,9 @@ export class Application { // After migrations: the tick's first act is a lease query against a table // the migration above creates. startQueenTick() + // A container that takes bees the Queen ordered. Off unless asked; a + // deployment may be a Queen, a runner, or both. + startBeeRunner() // The same asymmetry the comment below describes, one level up: this used // to exit when no CDP port was *configured*, while tolerating a configured diff --git a/trios/agent-server/apps/server/tests/api/queen-dispatch.test.ts b/trios/agent-server/apps/server/tests/api/queen-dispatch.test.ts index 2c06d7eb5a..b8dc8f8049 100644 --- a/trios/agent-server/apps/server/tests/api/queen-dispatch.test.ts +++ b/trios/agent-server/apps/server/tests/api/queen-dispatch.test.ts @@ -17,6 +17,7 @@ import { tmpdir } from 'node:os' import { join } from 'node:path' import type { Pool } from 'pg' import { + beesRunElsewhere, classifyQuotaExhaustion, closeDispatch, committedFileCount, @@ -36,8 +37,10 @@ import { reapWorktrees, recordDispatch, resolveWorkerProvider, + runClaimedBee, setDurableCloseListener, workerCapacityBreakdown, + workerProviderForKeyIndex, workspaceRoot, } from '../../src/api/services/queen-dispatch' import { resetYoungBees } from '../../src/api/services/queen-resources' @@ -94,6 +97,7 @@ const KEYS = [ 'TRIOS_QUEEN_BEE_DISK_MB', 'TRIOS_QUEEN_DISK_HEADROOM_PERCENT', 'TRIOS_QUEEN_MEMORY_LIMIT_MB', + 'TRIOS_QUEEN_BEES_RUN_ELSEWHERE', // The far end of the key list (MAX_KEYS_PER_POOL) and the first name past it. 'TRIOS_QUEEN_WORKER_API_KEY_17', 'TRIOS_QUEEN_WORKER_API_KEY_1024', @@ -1213,6 +1217,44 @@ describe('the container is asked before a worktree is cut', () => { } }) + it('lets a runner start an order it claimed, and books it as an update, not a new row', async () => { + // The order already holds the issue, the boundary, the criteria and the + // key. A success must say what happened on it - and reset the clock the + // two-hour rule reads, or a bee that waited an hour for a runner would be + // reaped an hour into its work - without inserting it again, which would + // archive a history row and clear what it is to be judged against. + oneFreeKey() + hive() + const turns = realStarts() + const { pool, queries } = recordingPool() + try { + const outcome = await runClaimedBee( + pool, + { + issue: 1530, + branch: 'queen-1530', + brief: 'the order the Queen wrote', + ownedPaths: ['docs/a.md'], + conversationId: 'conv-ordered-by-the-queen', + keyIndex: 0, + }, + { memory: calm, volume: roomy, volumeUsed: () => 40 }, + ) + expect(outcome.started).toBe(true) + expect(outcome.conversationId).toBe('conv-ordered-by-the-queen') + expect(outcome.detail).toContain('zai/') + const touched = queries.filter((q) => /queen_dispatch\b/.test(q.text)) + expect( + touched.some((q) => /INSERT INTO queen_dispatch\b/.test(q.text)), + ).toBe(false) + const update = touched.find((q) => q.text.includes('SET detail = $2')) + expect(update?.text).toContain('dispatched_at = now()') + expect(update?.values).toEqual([1530, outcome.detail]) + } finally { + await turns.restore() + } + }) + it('never reaps the tree of the issue it is about to dispatch again', async () => { // A redeploy killed the swarm; #1520 was reaped at boot and its tree holds // one committed, unpushed attempt. The container carries no push @@ -1313,6 +1355,164 @@ describe('the container is asked before a worktree is cut', () => { }) }) +/** + * A bee that runs somewhere else. + * + * The Queen decides exactly what she decided before and writes it down; the + * turn happens in another container. Measured 2026-09-18, which is why: her own + * container held about a gigabyte a bee against a 24 GB limit, so the swarm + * could never be wider than one machine however many credentials it had. + */ +describe('handing a bee to a runner', () => { + const pool = () => { + const queries: Array<{ text: string; values?: unknown[] }> = [] + const fake = { + query: async (text: string, values?: unknown[]) => { + queries.push({ text, values }) + return { rowCount: 0, rows: [] } + }, + } as unknown as Pool + return { fake, queries } + } + const oneKey = () => { + process.env.TRIOS_QUEEN_WORKER_PROVIDER = 'zai' + process.env.TRIOS_QUEEN_WORKER_BASE_URL = 'https://api.z.ai/api/paas/v4' + process.env.TRIOS_QUEEN_WORKER_API_KEY = 'first' + process.env.TRIOS_QUEEN_WORKER_API_KEY_2 = 'second' + } + + it('writes the order and starts nothing here', async () => { + oneKey() + process.env.TRIOS_QUEEN_BEES_RUN_ELSEWHERE = 'on' + const { fake, queries } = pool() + const realFetch = globalThis.fetch + let calls = 0 + globalThis.fetch = (async () => { + calls += 1 + return new Response('{}') + }) as unknown as typeof fetch + let measured = 0 + try { + const outcome = await dispatchBee( + fake, + 1600, + 'the brief for this bee', + ['docs/a.md'], + [], + undefined, + ['the tab opens'], + 'stated', + { + memory: () => { + measured += 1 + return { kind: 'unsupported', platform: 'darwin' } + }, + }, + ) + // In flight from this moment: the boundary is held and the key is taken + // exactly as they were when she ran the bee herself. + expect(outcome.started).toBe(true) + expect(outcome.keyIndex).toBe(0) + expect(outcome.conversationId).toMatch(/[0-9a-f-]{36}/) + expect(outcome.detail).toContain('queued for a runner; zai/') + // Nothing was run here: no turn asked for, and the container not even + // measured, because the bee will not live in it. + expect(calls).toBe(0) + expect(measured).toBe(0) + } finally { + globalThis.fetch = realFetch + } + const insert = queries.find((q) => + /INSERT INTO queen_dispatch\b/.test(q.text), + ) + expect(insert?.text).toContain('queued_at') + // The brief travels with the order: a runner has no issue body to build one + // from. The credential does not - only its index. + expect(insert?.values?.[13]).toBe('the brief for this bee') + expect(insert?.values?.[12]).toBeInstanceOf(Date) + expect(JSON.stringify(insert?.values)).not.toContain('first') + }) + + it('runs the bee here when nobody was told otherwise', async () => { + oneKey() + const { fake } = pool() + let measured = 0 + const outcome = await dispatchBee( + fake, + 1601, + 'brief', + [], + [], + undefined, + [], + 'none', + { + memory: () => { + measured += 1 + return { + kind: 'measured', + usedBytes: 23 * 1_000_000_000, + limitBytes: 24 * 1_000_000_000, + source: 'cgroup v2', + limitSource: 'cgroup', + } + }, + volume: () => ({ totalBytes: 50e9, freeBytes: 19e9 }), + }, + ) + expect(measured).toBeGreaterThan(0) + expect(outcome.detail).not.toContain('queued for a runner') + expect(beesRunElsewhere()).toBe(false) + }) + + it('remembers which credential an index names, without storing one', () => { + oneKey() + process.env.TRIOS_QUEEN_WORKER_POOL_2_BASE_URL = + 'https://integrate.api.nvidia.com/v1' + process.env.TRIOS_QUEEN_WORKER_POOL_2_MODEL = 'nvidia/nemotron' + process.env.TRIOS_QUEEN_WORKER_POOL_2_API_KEY = 'nv-one' + process.env.TRIOS_QUEEN_WORKER_POOL_2_API_KEY_2 = 'nv-two' + // Whatever the allocator hands out, the index it stamps on the row must + // name the same credential when a different process reads it back. + for (const taken of [[], [0], [0, 1], [0, 1, 10_000]]) { + const chosen = resolveWorkerProvider(taken) + const again = workerProviderForKeyIndex(chosen?.keyIndex ?? -1) + expect(again?.apiKey).toBe(chosen?.apiKey) + expect(again?.model).toBe(chosen?.model) + expect(again?.baseUrl).toBe(chosen?.baseUrl) + } + expect(workerProviderForKeyIndex(0)?.apiKey).toBe('first') + expect(workerProviderForKeyIndex(10_001)?.apiKey).toBe('nv-two') + // An index this environment cannot explain. The next key along is a + // different account, so there is no answer but none. + expect(workerProviderForKeyIndex(7)).toBeNull() + expect(workerProviderForKeyIndex(20_000)).toBeNull() + }) + + it('hands the issue back when the runner cannot resolve the credential', async () => { + // The runner's variables are not the Queen's. Reaching for the next key + // would put two bees on one account while the ledger says otherwise. + oneKey() + const { fake, queries } = pool() + const outcome = await runClaimedBee(fake, { + issue: 1602, + branch: 'queen-1602', + brief: 'brief', + ownedPaths: [], + conversationId: 'conv-1602', + keyIndex: 4242, + }) + expect(outcome.started).toBe(false) + expect(outcome.detail).toContain('cannot resolve key_index 4242') + // Recorded as a refusal, which ends the row and gives the issue and its + // boundary back to the next round. + const insert = queries.find((q) => + /INSERT INTO queen_dispatch\b/.test(q.text), + ) + expect(insert?.values?.[2]).toBe(false) + }) +}) + /** * #1308. `workers.capacity` answers a number; this breakdown answers what the * number is MADE of. An operator seeing capacity 4 cannot act on it without diff --git a/trios/agent-server/apps/server/tests/api/queen-runner.test.ts b/trios/agent-server/apps/server/tests/api/queen-runner.test.ts new file mode 100644 index 0000000000..e233495149 --- /dev/null +++ b/trios/agent-server/apps/server/tests/api/queen-runner.test.ts @@ -0,0 +1,169 @@ +/** + * @license + * Copyright 2025 BrowserOS + * SPDX-License-Identifier: AGPL-3.0-or-later + * + * The parts of the handover that are this code's own logic rather than + * PostgreSQL's. Whether two runners can take one order is the database's answer + * and is asked of a real one in tests/pglive/queen-runner-live.test.ts. + */ +import { describe, expect, it } from 'bun:test' +import type { Pool } from 'pg' +import { + RUNNER_SILENT_MINUTES, + reapStalledDispatches, +} from '../../src/api/services/queen-dispatch' +import { + runnerName, + runnerSlots, + waitForEnding, +} from '../../src/api/services/queen-runner' + +/** A database that returns the given rows to the reaper's SELECT and records everything. */ +function reaperPool(rows: Array<{ issue: number; queued_at: Date | null }>) { + const asked: Array<{ sql: string; params: unknown[] }> = [] + const pool = { + query: async (sql: string, params: unknown[] = []) => { + asked.push({ sql: String(sql), params }) + if ( + String(sql).startsWith('SELECT issue, queued_at FROM queen_dispatch') + ) { + return { rowCount: rows.length, rows } + } + return { + rowCount: rows.length, + rows: rows.map(({ issue }) => ({ issue })), + } + }, + } as unknown as Pool + return { pool, asked } +} + +describe('the stall reaper and bees that run elsewhere', () => { + it('releases a runner bee without reaching for a worktree this disk does not have', async () => { + // #7 ran here; #8 was an order a runner took. Both are past the line. + const { pool, asked } = reaperPool([ + { issue: 7, queued_at: null }, + { issue: 8, queued_at: new Date() }, + ]) + const salvaged: number[] = [] + const released = await reapStalledDispatches(pool, 120, { + salvage: async (_pool, issue) => { + salvaged.push(issue) + return { + committed: false, + detail: 'nothing', + left: [], + sha: null, + files: [], + } as never + }, + }) + // Both rows are released - a dead runner's issue must go back to the swarm + // - but only the bee that ran HERE is salvaged: git on the other one's + // worktree would be git on a directory that does not exist. + expect(salvaged).toEqual([7]) + expect(released).toEqual([7, 8]) + const release = asked.find((q) => q.sql.includes('UPDATE queen_dispatch')) + expect(release?.params).toEqual([120, [7, 8], RUNNER_SILENT_MINUTES]) + }) + + it('counts a runner that stopped vouching as dead long before two hours', () => { + // Forty missed renewals at fifteen seconds each. Not a slow network. + expect(RUNNER_SILENT_MINUTES).toBeGreaterThanOrEqual(5) + expect(RUNNER_SILENT_MINUTES).toBeLessThan(120) + }) +}) + +describe('a runner vouching for its bee', () => { + it('renews its claim until the row stops matching, then stops', async () => { + const asked: Array<{ sql: string; params: unknown[] }> = [] + let alive = 3 + const pool = { + query: async (sql: string, params: unknown[] = []) => { + asked.push({ sql: String(sql), params }) + alive -= 1 + // Three renewals find the row; then the drain writes the ending and + // the renewal matches nothing. + return alive > 0 ? { rows: [{ issue: 9 }] } : { rows: [] } + }, + } as unknown as Pool + await waitForEnding( + pool, + { + issue: 9, + branch: 'queen-9', + brief: '', + ownedPaths: [], + conversationId: 'conv-9', + keyIndex: 0, + }, + 'runner-a', + 1, + ) + expect(asked).toHaveLength(3) + for (const q of asked) { + expect(q.sql).toContain('SET claimed_at = now()') + // Only its OWN claim, on its own turn: a re-dispatched issue belongs to + // whoever claimed the new order. + expect(q.params).toEqual([9, 'conv-9', 'runner-a']) + } + }) + + it('keeps vouching through a database that blips', async () => { + let calls = 0 + const pool = { + query: async () => { + calls += 1 + if (calls === 1) throw new Error('Connection terminated unexpectedly') + return { rows: [] } + }, + } as unknown as Pool + await waitForEnding( + pool, + { + issue: 10, + branch: 'queen-10', + brief: '', + ownedPaths: [], + conversationId: 'conv-10', + keyIndex: 0, + }, + 'runner-a', + 1, + ) + // One failed renewal is not a dead runner; it tried again. + expect(calls).toBe(2) + }) +}) + +describe('what a runner calls itself and how much it carries', () => { + it('says who it is, and carries one bee unless told otherwise', () => { + const previous = { + name: process.env.TRIOS_BEE_RUNNER_NAME, + slots: process.env.TRIOS_BEE_RUNNER_SLOTS, + } + try { + delete process.env.TRIOS_BEE_RUNNER_NAME + delete process.env.TRIOS_BEE_RUNNER_SLOTS + expect(runnerName()).toContain(String(process.pid)) + expect(runnerSlots()).toBe(1) + process.env.TRIOS_BEE_RUNNER_NAME = 'replica-3' + process.env.TRIOS_BEE_RUNNER_SLOTS = '4' + expect(runnerName()).toBe('replica-3') + expect(runnerSlots()).toBe(4) + for (const bad of ['0', '-2', 'many', '1.5']) { + process.env.TRIOS_BEE_RUNNER_SLOTS = bad + expect(runnerSlots()).toBe(1) + } + process.env.TRIOS_BEE_RUNNER_SLOTS = '1000' + expect(runnerSlots()).toBe(64) + } finally { + if (previous.name === undefined) delete process.env.TRIOS_BEE_RUNNER_NAME + else process.env.TRIOS_BEE_RUNNER_NAME = previous.name + if (previous.slots === undefined) + delete process.env.TRIOS_BEE_RUNNER_SLOTS + else process.env.TRIOS_BEE_RUNNER_SLOTS = previous.slots + } + }) +}) diff --git a/trios/agent-server/apps/server/tests/pglive/queen-runner-live.test.ts b/trios/agent-server/apps/server/tests/pglive/queen-runner-live.test.ts new file mode 100644 index 0000000000..f9b9b4a731 --- /dev/null +++ b/trios/agent-server/apps/server/tests/pglive/queen-runner-live.test.ts @@ -0,0 +1,271 @@ +/** + * @license + * Copyright 2025 BrowserOS + * SPDX-License-Identifier: AGPL-3.0-or-later + * + * TWO RUNNERS, ONE ORDER: only PostgreSQL can answer this. + * + * The whole of the handover's mutual exclusion is one statement - an UPDATE + * over a row picked `FOR UPDATE SKIP LOCKED` - and nothing outside a real + * server tells you whether it holds. A fake pool would answer whatever it was + * written to answer, which is the shape of test that lets a double-claim ship: + * two runners on one issue means two bees writing one branch and two rows' + * worth of ledger for one boundary. + * + * So this runs against a scratch database, in the pglive GROUP for the reason + * that group exists: `mock.module` is process-global in bun, and the api group + * binds `pg` to a FakePool at module scope. + * + * Like its neighbour, this FAILS when no server is reachable rather than + * skipping, unless TRIOS_PG_MIGRATE_GATE=offline says the absence is deliberate + * - a silent skip is how a gate comes to report a success it never earned. + */ + +import { afterEach, beforeEach, describe, expect, it } from 'bun:test' +import { randomBytes } from 'node:crypto' +import { userInfo } from 'node:os' +import { Pool } from 'pg' +import { + reapDispatchesFromPreviousBoot, + reapStalledDispatches, +} from '../../src/api/services/queen-dispatch' +import { claimQueuedBee } from '../../src/api/services/queen-runner' +import { runPgMigrations } from '../../src/lib/db/pg-migrate' +import { createQueenPool, queenSchema } from '../../src/lib/db/queen-pool' + +const OFFLINE_KEY = 'TRIOS_PG_MIGRATE_GATE' +const URL_KEY = 'TRIOS_PG_TEST_URL' + +function offlineRequested(): boolean { + return (process.env[OFFLINE_KEY] ?? '').toLowerCase() === 'offline' +} + +function adminUrl(): string { + return ( + process.env[URL_KEY] ?? + `postgres://${userInfo().username}@127.0.0.1:5432/postgres` + ) +} + +/** A fresh database per run, dropped afterwards. */ +async function scratchDatabase(): Promise<{ + url: string + drop: () => Promise +} | null> { + const name = `queen_runner_${randomBytes(6).toString('hex')}` + const admin = new Pool({ connectionString: adminUrl(), max: 1 }) + try { + await admin.query(`CREATE DATABASE ${name}`) + } catch (error) { + await admin.end().catch(() => undefined) + if (offlineRequested()) return null + throw error + } + const url = new URL(adminUrl()) + url.pathname = `/${name}` + // The schema every Queen connection selects. Production has it because a role + // setting put it there years ago; a scratch database has to be told, and a + // pool that selects a schema nobody created cannot create a table in it. + const fresh = new Pool({ connectionString: url.toString(), max: 1 }) + try { + await fresh.query(`CREATE SCHEMA IF NOT EXISTS ${queenSchema()}`) + } finally { + await fresh.end().catch(() => undefined) + } + return { + url: url.toString(), + drop: async () => { + await admin + .query(`DROP DATABASE IF EXISTS ${name} WITH (FORCE)`) + .catch(() => undefined) + await admin.end().catch(() => undefined) + }, + } +} + +describe('a queued bee is claimed by exactly one runner', () => { + let scratch: { url: string; drop: () => Promise } | null = null + let pool: Pool | null = null + const previousUrl = process.env.DATABASE_URL + + beforeEach(async () => { + scratch = await scratchDatabase() + if (!scratch) return + // `runPgMigrations` reads DATABASE_URL, as the server does at boot. + process.env.DATABASE_URL = scratch.url + await runPgMigrations() + // The pool the server itself builds: it selects the Queen's schema on every + // connection, and a plain `new Pool` here would read a different namespace + // from the one the migration just wrote into - which is precisely the decoy + // queen-pool.ts exists to end. + pool = createQueenPool(scratch.url) + }) + + afterEach(async () => { + await pool?.end().catch(() => undefined) + pool = null + await scratch?.drop() + scratch = null + if (previousUrl === undefined) delete process.env.DATABASE_URL + else process.env.DATABASE_URL = previousUrl + }) + + const order = async (issue: number, keyIndex = 0): Promise => { + await pool?.query( + `INSERT INTO queen_dispatch + (issue, branch, started, detail, owned_paths, conversation_id, + key_index, queued_at, brief) + VALUES ($1, $2, true, 'queued for a runner', '["docs/a.md"]'::jsonb, + $3, $4, now(), 'do the thing')`, + [issue, `queen-${issue}`, `conv-${issue}`, keyIndex], + ) + } + + it('hands one order to one runner, whatever else is asking', async () => { + if (!pool) return expect(offlineRequested()).toBe(true) + await order(4100) + + // Eight runners reaching for one row at the same moment. + const claims = await Promise.all( + Array.from({ length: 8 }, (_, i) => + claimQueuedBee(pool as Pool, `runner-${i}`), + ), + ) + const took = claims.filter(Boolean) + expect(took).toHaveLength(1) + expect(took[0]).toMatchObject({ + issue: 4100, + branch: 'queen-4100', + brief: 'do the thing', + ownedPaths: ['docs/a.md'], + conversationId: 'conv-4100', + keyIndex: 0, + }) + const row = await pool.query( + 'SELECT claimed_by, claimed_at FROM queen_dispatch WHERE issue = 4100', + ) + expect(row.rows[0].claimed_by).toMatch(/^runner-\d$/) + expect(row.rows[0].claimed_at).not.toBeNull() + // And nobody gets it twice. + expect(await claimQueuedBee(pool, 'runner-late')).toBeNull() + }) + + it('gives each runner a different order, oldest first, and skips what is taken', async () => { + if (!pool) return expect(offlineRequested()).toBe(true) + await order(4201, 1) + await order(4202, 2) + await order(4203, 3) + // Reaching at the same moment: SKIP LOCKED means nobody waits behind + // another runner's lock while an unclaimed row sits next to it. + const claims = await Promise.all([ + claimQueuedBee(pool, 'a'), + claimQueuedBee(pool, 'b'), + claimQueuedBee(pool, 'c'), + ]) + const issues = claims + .filter(Boolean) + .map((claim) => (claim as { issue: number }).issue) + .sort() + expect(issues).toEqual([4201, 4202, 4203]) + expect(await claimQueuedBee(pool, 'd')).toBeNull() + }) + + it('takes only what the Queen queued and left unclaimed', async () => { + if (!pool) return expect(offlineRequested()).toBe(true) + // A bee running inside the Queen herself: no queued_at, so no order. + await pool.query( + `INSERT INTO queen_dispatch (issue, branch, started, detail, conversation_id, key_index) + VALUES (4301, 'queen-4301', true, 'running here', 'c1', 0)`, + ) + // An order already finished, and one already claimed. + await order(4302) + await pool.query( + "UPDATE queen_dispatch SET finished_at = now(), outcome = 'finished' WHERE issue = 4302", + ) + await order(4303) + await pool.query( + "UPDATE queen_dispatch SET claimed_by = 'somebody', claimed_at = now() WHERE issue = 4303", + ) + expect(await claimQueuedBee(pool, 'fresh')).toBeNull() + }) + + it('re-dispatching an order releases the claim the last one left', async () => { + if (!pool) return expect(offlineRequested()).toBe(true) + await order(4400) + expect(await claimQueuedBee(pool, 'first')).not.toBeNull() + // The runner died. The reaper ends the row, and the next round queues the + // issue again - which is the same upsert `recordDispatch` runs. + await pool.query( + `UPDATE queen_dispatch SET finished_at = now(), outcome = 'reaped' WHERE issue = 4400`, + ) + await pool.query( + `UPDATE queen_dispatch + SET started = true, finished_at = NULL, outcome = NULL, + queued_at = now(), claimed_by = NULL, claimed_at = NULL + WHERE issue = 4400`, + ) + const again = await claimQueuedBee(pool, 'second') + expect(again?.issue).toBe(4400) + }) + + // A salvage runs git on this machine's volume; these cases are about which + // ROWS a reaper takes, so it is replaced by a no-op that says nothing moved. + const noSalvage = async () => + ({ + committed: false, + detail: 'test', + left: [], + sha: null, + files: [], + }) as never + + it('the Queen restarting does not bury the bees that run on runners', async () => { + if (!pool) return expect(offlineRequested()).toBe(true) + // One bee the Queen ran herself - it died with her container. + await pool.query( + `INSERT INTO queen_dispatch (issue, branch, started, detail, conversation_id, key_index) + VALUES (4501, 'queen-4501', true, 'running here', 'c-4501', 0)`, + ) + // An order nobody has taken yet, and one a runner is working on. + await order(4502) + await order(4503) + expect(await claimQueuedBee(pool, 'runner-x')).not.toBeNull() + + const reaped = await reapDispatchesFromPreviousBoot(pool, { + salvage: noSalvage, + }) + expect(reaped).toEqual([4501]) + const open = await pool.query( + 'SELECT issue FROM queen_dispatch WHERE finished_at IS NULL ORDER BY issue', + ) + expect(open.rows.map((row) => row.issue)).toEqual([4502, 4503]) + }) + + it('a runner that stops vouching loses its bee after ten minutes, not two hours', async () => { + if (!pool) return expect(offlineRequested()).toBe(true) + await order(4601) + await order(4602) + expect(await claimQueuedBee(pool, 'silent')).not.toBeNull() + expect(await claimQueuedBee(pool, 'alive')).not.toBeNull() + // Both started a minute ago; one runner went quiet eleven minutes ago - its + // container is gone - and the other renewed a second ago. + await pool.query( + `UPDATE queen_dispatch SET dispatched_at = now() - interval '1 minute'`, + ) + await pool.query( + `UPDATE queen_dispatch SET claimed_at = now() - interval '11 minutes' + WHERE claimed_by = 'silent'`, + ) + const reaped = await reapStalledDispatches(pool, 120, { + salvage: noSalvage, + }) + const silent = await pool.query( + "SELECT issue FROM queen_dispatch WHERE claimed_by = 'silent'", + ) + expect(reaped).toEqual([silent.rows[0].issue]) + const left = await pool.query( + 'SELECT claimed_by FROM queen_dispatch WHERE finished_at IS NULL', + ) + expect(left.rows.map((row) => row.claimed_by)).toEqual(['alive']) + }) +}) From c1bcbfb32d6ca3e6d0183ff21eb2e586f4585a9f Mon Sep 17 00:00:00 2001 From: Dmitrii Fedorov Date: Fri, 2 Oct 2026 21:03:07 -0300 Subject: [PATCH 2/2] feat(queen): runners take keys from the contributor registry, and the Queen judges a runner's branch Rebased #506 onto the deploy branch, where keys and models also come from the contributor registry since #522/#524. A runner resolved an order's key_index from its environment only, so a key added in the account (negative index) or an owner's model choice would have been refused or ignored; it now resolves through the same candidates the Queen allocated from. The review reads queen-N in the Queen's checkout - diff, head, the specs at that commit, criteria in a worktree cut from it - and a runner left that branch only in queen_bundle, so every runner bee would have waited forever unjudged. The sweep now brings a claimed row's bundle into the checkout first, fetching origin when the runner cut from a newer base, and replacing a stale worktree of an earlier local attempt that holds the branch name. With bees running elsewhere the Queen's own memory says nothing about the swarm's width, so the entrypoint takes the operator ceiling (else the lane count) instead of dividing this container's memory. --- .../server/src/api/routes/queen-export.ts | 94 ++++++++- .../server/src/api/services/queen-dispatch.ts | 17 +- .../server/src/api/services/queen-reviewer.ts | 7 + .../server/src/api/services/queen-tick.ts | 15 ++ .../server/tests/api/queen-runner.test.ts | 178 +++++++++++++++++- trios/agent-server/docker-entrypoint.sh | 23 +++ 6 files changed, 331 insertions(+), 3 deletions(-) diff --git a/trios/agent-server/apps/server/src/api/routes/queen-export.ts b/trios/agent-server/apps/server/src/api/routes/queen-export.ts index cc32bf0d0d..21cf441daf 100644 --- a/trios/agent-server/apps/server/src/api/routes/queen-export.ts +++ b/trios/agent-server/apps/server/src/api/routes/queen-export.ts @@ -31,10 +31,11 @@ import { spawn } from 'node:child_process' import { randomUUID } from 'node:crypto' -import { readFile, rm } from 'node:fs/promises' +import { readFile, rm, writeFile } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' import { Hono } from 'hono' +import type { Pool } from 'pg' import { createQueenPool } from '../../lib/db/queen-pool' import { logger } from '../../lib/logger' import { baseRef, workspaceRoot } from '../services/queen-dispatch' @@ -229,6 +230,97 @@ export async function bundleOfBranch( } } +/** + * A runner's branch, brought into this checkout before the Queen judges it. + * + * The review reads `queen-N` here - its diff, its head, the specs at that + * commit, the criteria in a worktree cut from it - and a bee that ran in a + * runner's container left that branch only in `queen_bundle`. Without this the + * diff of every runner bee would fail to read and the row would wait forever. + * The bundle is `base..queen-N`, so a base newer than this checkout has seen is + * fetched from origin first; a leftover worktree of an earlier attempt on this + * disk holds the branch name and is removed (the bee there has finished). + */ +export async function importRunnerBranch( + pool: Pick, + issue: number, +): Promise<{ ok: true; imported: boolean } | { ok: false; error: string }> { + const rows = await pool.query( + 'SELECT bytes FROM queen_bundle WHERE issue = $1', + [issue], + ) + const bytes = rows.rows?.[0]?.bytes as Buffer | undefined + if (!bytes?.length) return { ok: true, imported: false } + const root = workspaceRoot() + const branch = `queen-${issue}` + const path = join(tmpdir(), `queen-import-${issue}-${randomUUID()}.bundle`) + await writeFile(path, bytes, { mode: 0o644 }) + try { + const heads = await git([ + '-C', + root, + 'bundle', + 'list-heads', + path, + `refs/heads/${branch}`, + ]) + const head = heads.out.trim().split(/\s+/)[0] ?? '' + if (heads.code !== 0 || !/^[0-9a-f]{40}$/.test(head)) { + return { ok: false, error: `the stored bundle names no ${branch}` } + } + const local = await git([ + '-C', + root, + 'rev-parse', + '--verify', + '--quiet', + branch, + ]) + if (local.code === 0 && local.out.trim() === head) { + return { ok: true, imported: false } + } + const fetch = () => + git( + [ + '-C', + root, + 'fetch', + '--no-tags', + '--quiet', + path, + `+refs/heads/${branch}:refs/heads/${branch}`, + ], + 120_000, + ) + let fetched = await fetch() + if (fetched.code !== 0 && /checked out/i.test(fetched.err)) { + const listed = await git(['-C', root, 'worktree', 'list', '--porcelain']) + const holder = listed.out + .split('\n\n') + .find((entry) => entry.includes(`branch refs/heads/${branch}`)) + ?.match(/^worktree (.+)$/m)?.[1] + if (holder && holder !== root) { + await git(['-C', root, 'worktree', 'remove', '--force', holder]) + fetched = await fetch() + } + } + if (fetched.code !== 0) { + await git(['-C', root, 'fetch', '--quiet', 'origin'], 120_000) + fetched = await fetch() + } + if (fetched.code !== 0) { + return { ok: false, error: fetched.err.trim().slice(0, 300) } + } + logger.info('Queen brought a runner branch into her checkout', { + issue, + head: head.slice(0, 12), + }) + return { ok: true, imported: true } + } finally { + await rm(path, { force: true }).catch(() => {}) + } +} + /** The bundle a runner left behind, or null. Never a reason to fail a request. */ async function storedBundle( issue: number, diff --git a/trios/agent-server/apps/server/src/api/services/queen-dispatch.ts b/trios/agent-server/apps/server/src/api/services/queen-dispatch.ts index f440952c87..a9a39cf12f 100644 --- a/trios/agent-server/apps/server/src/api/services/queen-dispatch.ts +++ b/trios/agent-server/apps/server/src/api/services/queen-dispatch.ts @@ -936,7 +936,19 @@ function configuredEndpointProvider( */ export function workerProviderForKeyIndex( keyIndex: number, + runtime?: ContributorRuntime, ): WorkerProvider | null { + // With the contributor registry in play the Queen allocated from + // `contributorWorkerCandidates`: a managed key (negative index) lives only + // there, a disabled key is absent, and an owner's model replaces the pool's. + // A runner must resolve the order exactly as it was written. + if (contributorRuntimeActive(runtime)) { + return ( + contributorWorkerCandidates(runtime).find( + (lane) => lane.keyIndex === keyIndex, + ) ?? null + ) + } const override = process.env.TRIOS_QUEEN_WORKER_MODEL if (configuredWorkerBaseUrl()) { const pools = configuredEndpointPools() @@ -4818,7 +4830,10 @@ export async function runClaimedBee( return { started: false, issue, branch, detail } } - const chosen = workerProviderForKeyIndex(order.keyIndex) + const chosen = workerProviderForKeyIndex( + order.keyIndex, + await contributorRuntime(pool, environmentContributorKeys()), + ) if (!chosen?.apiKey) { // The runner's environment is not the Queen's, or the key list changed // under the order. Taking the next credential along would put two bees on diff --git a/trios/agent-server/apps/server/src/api/services/queen-reviewer.ts b/trios/agent-server/apps/server/src/api/services/queen-reviewer.ts index 5217fd79f1..c1daa6858a 100644 --- a/trios/agent-server/apps/server/src/api/services/queen-reviewer.ts +++ b/trios/agent-server/apps/server/src/api/services/queen-reviewer.ts @@ -789,6 +789,13 @@ export interface ReviewDeps { reviewsPerRound: () => number /** Criteria measurements one sweep may buy. Optional for injected fakes. */ measurementsPerRound?: () => number + /** + * Brings a branch a runner left in `queen_bundle` into this checkout, so + * everything above reads it as if the bee had run here. Optional for fakes. + */ + importRunnerBranch?: ( + issue: number, + ) => Promise<{ ok: true; imported: boolean } | { ok: false; error: string }> } export function defaultReviewDeps(): ReviewDeps { diff --git a/trios/agent-server/apps/server/src/api/services/queen-tick.ts b/trios/agent-server/apps/server/src/api/services/queen-tick.ts index 04202fae4c..a51d845aa9 100644 --- a/trios/agent-server/apps/server/src/api/services/queen-tick.ts +++ b/trios/agent-server/apps/server/src/api/services/queen-tick.ts @@ -41,6 +41,7 @@ import type { Pool } from 'pg' import { createQueenPool } from '../../lib/db/queen-pool' import { logger } from '../../lib/logger' import { startModelProbes, workerModelRanking } from '../../lib/model-ranking' +import { importRunnerBranch } from '../routes/queen-export' import { outstandingEscalations } from '../routes/queen-needs-you' import { githubCiDeps, takeBackRefusedAcceptances } from './queen-ci-verdict' import { contributorRuntime } from './queen-contributor-keys' @@ -2649,6 +2650,7 @@ export async function reviewFinishedDispatches( taken, await contributorRuntime(pool, environmentContributorKeys()), ), + importRunnerBranch: (issue) => importRunnerBranch(pool, issue), ...overrides, } // TWO THINGS ABOUT THIS QUERY, BOTH MEASURED ON 2026-09-03. @@ -2702,6 +2704,7 @@ export async function reviewFinishedDispatches( -- work was salvaged rather than written; it changes NOTHING about -- how the work is judged. d.salvaged_at, d.salvaged_sha, d.salvaged_files, d.salvage_left, + d.claimed_by, (SELECT string_agg(t.text, '' ORDER BY t.seq) FROM queen_transcript t WHERE t.conversation_id = d.conversation_id AND t.kind = 'say') @@ -2784,6 +2787,18 @@ export async function reviewFinishedDispatches( const conversation = row.conversation_id == null ? null : String(row.conversation_id) + // A bee a runner claimed ran in another container; its branch is in + // `queen_bundle` until it is brought here. Failing that is a wait, the + // same as a diff that could not be read: nothing is counted against it. + if (row.claimed_by && deps.importRunnerBranch) { + const imported = await deps.importRunnerBranch(issue) + if (!imported.ok) { + logger.warn('Queen could not bring a runner branch into her checkout', { + issue, + error: imported.error, + }) + } + } // ONE git diff, asked once and used twice. The count is what the review // policy weighs; the names are what the boundary rule compares. const diff = await deps.committedFilesResult(issue) diff --git a/trios/agent-server/apps/server/tests/api/queen-runner.test.ts b/trios/agent-server/apps/server/tests/api/queen-runner.test.ts index e233495149..99dbe904ec 100644 --- a/trios/agent-server/apps/server/tests/api/queen-runner.test.ts +++ b/trios/agent-server/apps/server/tests/api/queen-runner.test.ts @@ -7,11 +7,17 @@ * PostgreSQL's. Whether two runners can take one order is the database's answer * and is asked of a real one in tests/pglive/queen-runner-live.test.ts. */ -import { describe, expect, it } from 'bun:test' +import { afterEach, beforeEach, describe, expect, it } from 'bun:test' +import { execFileSync } from 'node:child_process' +import { mkdtempSync, readFileSync, rmSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' import type { Pool } from 'pg' +import { importRunnerBranch } from '../../src/api/routes/queen-export' import { RUNNER_SILENT_MINUTES, reapStalledDispatches, + workerProviderForKeyIndex, } from '../../src/api/services/queen-dispatch' import { runnerName, @@ -167,3 +173,173 @@ describe('what a runner calls itself and how much it carries', () => { } }) }) + +describe('a runner resolves an order the way the Queen allocated it', () => { + const names = [ + 'TRIOS_QUEEN_WORKER_BASE_URL', + 'TRIOS_QUEEN_WORKER_PROVIDER', + 'TRIOS_QUEEN_WORKER_MODEL', + 'TRIOS_QUEEN_WORKER_API_KEY', + 'TRIOS_QUEEN_WORKER_API_KEY_2', + ] + const saved = new Map() + beforeEach(() => { + for (const name of names) saved.set(name, process.env[name]) + process.env.TRIOS_QUEEN_WORKER_BASE_URL = + 'https://integrate.api.nvidia.com/v1' + process.env.TRIOS_QUEEN_WORKER_PROVIDER = 'openai-compatible' + process.env.TRIOS_QUEEN_WORKER_MODEL = 'pool-model' + process.env.TRIOS_QUEEN_WORKER_API_KEY = 'env-key-a' + process.env.TRIOS_QUEEN_WORKER_API_KEY_2 = 'env-key-b' + }) + afterEach(() => { + for (const name of names) { + const value = saved.get(name) + if (value === undefined) delete process.env[name] + else process.env[name] = value + } + }) + const runtime = { + managed: [ + { + id: -2, + provider: 'nvidia' as const, + model: 'owner-model', + apiKey: 'managed-key', + baseUrl: 'https://integrate.api.nvidia.com/v1', + }, + ], + disabled: [0], + models: { 1: 'owner-model' }, + } + it('finds a managed key, applies the owner model and refuses a disabled key', () => { + expect(workerProviderForKeyIndex(-2, runtime)).toMatchObject({ + apiKey: 'managed-key', + model: 'owner-model', + keyIndex: -2, + }) + expect(workerProviderForKeyIndex(1, runtime)).toMatchObject({ + apiKey: 'env-key-b', + model: 'owner-model', + keyIndex: 1, + }) + expect(workerProviderForKeyIndex(0, runtime)).toBeNull() + expect(workerProviderForKeyIndex(-9, runtime)).toBeNull() + }) + it('without the registry, reads the environment exactly as before', () => { + expect(workerProviderForKeyIndex(0)).toMatchObject({ + apiKey: 'env-key-a', + model: 'pool-model', + }) + expect( + workerProviderForKeyIndex(0, { managed: [], disabled: [], models: {} }), + ).toMatchObject({ apiKey: 'env-key-a', model: 'pool-model' }) + }) +}) + +describe('the Queen brings a runner branch into her checkout before judging it', () => { + const sh = (cwd: string, ...args: string[]) => + execFileSync('git', ['-C', cwd, ...args], { + encoding: 'utf8', + env: { + ...process.env, + GIT_AUTHOR_NAME: 't', + GIT_AUTHOR_EMAIL: 't@t', + GIT_COMMITTER_NAME: 't', + GIT_COMMITTER_EMAIL: 't@t', + }, + }).trim() + let scratch = '' + const saved: Record = {} + beforeEach(() => { + scratch = mkdtempSync(join(tmpdir(), 'queen-import-')) + for (const name of [ + 'WORKSPACE_DIR', + 'TRIOS_REPO_URL', + 'TRIOS_TOOL_SHELL_USER', + ]) + saved[name] = process.env[name] + delete process.env.TRIOS_TOOL_SHELL_USER + // origin, the Queen's checkout, and a runner's checkout of the same repo. + execFileSync('git', ['init', '-q', '-b', 'master', join(scratch, 'origin')]) + sh(join(scratch, 'origin'), 'commit', '-q', '--allow-empty', '-m', 'base') + execFileSync('git', [ + 'clone', + '-q', + join(scratch, 'origin'), + join(scratch, 'ws', 't27'), + ]) + execFileSync('git', [ + 'clone', + '-q', + join(scratch, 'origin'), + join(scratch, 'runner'), + ]) + process.env.WORKSPACE_DIR = join(scratch, 'ws') + process.env.TRIOS_REPO_URL = 'https://github.com/gHashTag/t27.git' + }) + afterEach(() => { + for (const [name, value] of Object.entries(saved)) { + if (value === undefined) delete process.env[name] + else process.env[name] = value + } + rmSync(scratch, { recursive: true, force: true }) + }) + /** The runner's work after the origin moved on, as `queen_bundle` would hold it. */ + const runnerBundle = () => { + const origin = join(scratch, 'origin') + sh(origin, 'commit', '-q', '--allow-empty', '-m', 'master moved') + const runner = join(scratch, 'runner') + sh(runner, 'fetch', '-q', 'origin') + sh(runner, 'checkout', '-q', '-b', 'queen-7', 'origin/master') + sh(runner, 'commit', '-q', '--allow-empty', '-m', 'bee work') + const file = join(scratch, 'b.bundle') + sh(runner, 'bundle', 'create', file, 'origin/master..queen-7') + return { + bytes: readFileSync(file), + head: sh(runner, 'rev-parse', 'queen-7'), + } + } + const poolWith = (bytes: Buffer | null) => + ({ + query: async () => ({ rows: bytes ? [{ bytes }] : [] }), + }) as unknown as Pick + it('imports the branch even when the runner cut from a base this checkout has not fetched', async () => { + const { bytes, head } = runnerBundle() + const queen = join(scratch, 'ws', 't27') + expect(await importRunnerBranch(poolWith(bytes), 7)).toEqual({ + ok: true, + imported: true, + }) + expect(sh(queen, 'rev-parse', 'queen-7')).toBe(head) + // Already here: nothing to do the second time. + expect(await importRunnerBranch(poolWith(bytes), 7)).toEqual({ + ok: true, + imported: false, + }) + // No bundle: a bee that ran here, untouched. + expect(await importRunnerBranch(poolWith(null), 8)).toEqual({ + ok: true, + imported: false, + }) + }) + it('replaces an earlier attempt whose worktree still holds the branch', async () => { + const queen = join(scratch, 'ws', 't27') + sh( + queen, + 'worktree', + 'add', + '-q', + '-b', + 'queen-7', + join(queen, '.worktrees', 'queen-7'), + 'master', + ) + const { bytes, head } = runnerBundle() + expect(await importRunnerBranch(poolWith(bytes), 7)).toEqual({ + ok: true, + imported: true, + }) + expect(sh(queen, 'rev-parse', 'queen-7')).toBe(head) + }) +}) diff --git a/trios/agent-server/docker-entrypoint.sh b/trios/agent-server/docker-entrypoint.sh index 4e87054c6c..f3ef2d40fc 100755 --- a/trios/agent-server/docker-entrypoint.sh +++ b/trios/agent-server/docker-entrypoint.sh @@ -267,6 +267,29 @@ derive_worker_cap() { [ "$lanes" -le 4 ] || lanes=4 from_keys=$((credentials * lanes)) + + # BEES RUN ELSEWHERE: the Queen only writes orders and runners in their own + # containers run them, so this container's memory says nothing about how + # wide the swarm is. The width is what the runners hold, which only the + # operator knows: the named ceiling, else the lanes counted here. (Only the + # first pool is counted above; the TypeScript side bounds by every lane.) + case "${TRIOS_QUEEN_BEES_RUN_ELSEWHERE:-}" in + on|true|1) + ceiling=${TRIOS_QUEEN_MAX_WORKERS_CEILING:-} + case "$ceiling" in + ''|*[!0-9]*) ceiling="" ;; + esac + if [ -n "$ceiling" ] && [ "$ceiling" -ge 1 ]; then + echo "[entrypoint] bees run elsewhere: width is the operator ceiling $ceiling" >&2 + echo "$ceiling" + else + echo "[entrypoint] bees run elsewhere: width is the lane count $from_keys" >&2 + echo "$from_keys" + fi + return + ;; + esac + bee_mb=${TRIOS_QUEEN_BEE_MEMORY_MB:-1024} case "$bee_mb" in ''|*[!0-9]*) bee_mb=1024 ;;