From f63b339e894177723b22f891295b6af549e12433 Mon Sep 17 00:00:00 2001 From: Dmitrii Fedorov Date: Mon, 21 Sep 2026 21:14:30 -0300 Subject: [PATCH] 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 8658b2c10d..f43e580d16 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 @@ -883,6 +883,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. * @@ -4051,11 +4134,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', { @@ -4075,6 +4165,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 @@ -4097,8 +4188,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, ) @@ -4119,7 +4214,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], @@ -4127,6 +4222,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, @@ -4134,11 +4241,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 [] @@ -4156,10 +4265,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) } @@ -4245,6 +4356,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 @@ -4272,28 +4431,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 @@ -4306,7 +4532,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, @@ -4329,41 +4554,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, } } @@ -4589,6 +4876,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 @@ -4644,11 +4937,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, @@ -4668,6 +4962,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 @@ -4740,6 +5041,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 c74b3d7425..21ffd0bc4b 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']) + }) +})