diff --git a/trios/agent-server/apps/server/src/api/routes/queen-public-status.ts b/trios/agent-server/apps/server/src/api/routes/queen-public-status.ts index 3735667abe..0f7bce023c 100644 --- a/trios/agent-server/apps/server/src/api/routes/queen-public-status.ts +++ b/trios/agent-server/apps/server/src/api/routes/queen-public-status.ts @@ -20,11 +20,18 @@ * A count nobody can act on was the defect this projection shipped with: the * count said two issues lacked a boundary and could not say which two. The * number alone is inert - paths, holders, titles and prose never leave. + * + * `queue` is the same answer for the paid slots themselves: `workers`-style + * utilization says HOW MANY slots are busy, `queue.state` says why the rest + * are not - no eligible work, a full hive, work in flight, or evidence this + * page refuses to interpret. One closed word and one timestamp, nothing else, + * because the second anything else rides along it becomes a channel. */ import { Hono } from 'hono' import { Pool } from 'pg' import { logger } from '../../lib/logger' +import { configuredWorkerCapacity } from '../services/queen-dispatch' interface QueryResult { rowCount: number | null @@ -41,6 +48,10 @@ interface QueenPublicStatusDeps { createPool?: (url: string) => StatusPool tickIntervalSeconds?: () => number billingMode?: () => BillingMode + /** Paid worker slots, from the same authority /queen/public-research reads. */ + workerCapacity?: () => number + /** Clock, injectable so staleness is testable against fixed tick dates. */ + now?: () => number } type BillingMode = 'api_metered' | 'coding_plan' @@ -289,6 +300,147 @@ function classifySwarmState(facts: { return 'healthy_idle' } +/** + * Every value `queue.state` can carry. + * + * Closed like `swarmState` and the skip categories: a public consumer must be + * able to rely on every value it will ever read. These are the four honest + * answers to "why is a paid slot idle": + * + * capacity-full every configured slot is carrying a worker; there is + * no idle slot left to explain + * work-dispatched at least one slot is busy and at least one is free - + * the queue supplied work this round + * no-eligible-work every slot idle AND the latest decision explicitly + * found no eligible candidate + * unknown idle slots this page cannot account for: no readable + * or recent-enough decision, no configured capacity, a + * disabled scheduler, or a refusal this projection does + * not interpret + */ +type QueueState = + | 'capacity-full' + | 'work-dispatched' + | 'no-eligible-work' + | 'unknown' + +/** + * How many tick intervals a decision stays trustworthy, measured from + * `decided_at`. + * + * The tick timer guarantees a decision at least every interval while the loop + * lives, so a decision older than two intervals means nobody has decided for a + * whole extra period - the scheduler stopped, and the last thing it said must + * not go on explaining a present it never saw. One missed interval is + * tolerated because a round that runs long can straddle the timer. + */ +const TICK_STALENESS_INTERVALS = 2 + +/** + * `decided_at` as milliseconds, or null when the row never carried a date + * this process could read. + * + * Postgres returns a timestamptz as a Date and the test fixtures carry ISO + * strings; both are read here so the projection judges the same age whichever + * shape the row arrives in. Anything else - a corrupt value, a non-date - + * vouches for nothing. + */ +function decidedAtMs(value: unknown): number | null { + if (value instanceof Date) { + const ms = value.getTime() + return Number.isFinite(ms) ? ms : null + } + if (typeof value === 'string') { + const ms = Date.parse(value) + return Number.isFinite(ms) ? ms : null + } + return null +} + +/** + * The one closed word for what the paid slots are doing, from facts this + * endpoint already holds. + * + * No second scheduler and no second eligibility rule feeds this - only the + * dispatch counts the page already reads, the capacity the deployment already + * declares through `configuredWorkerCapacity` (the same single authority + * `/queen/public-research` reads), and the last tick row the page already + * publishes. The order is the contract, and it is the order of how directly + * each fact speaks about the present: + * + * 1. `capacity-full` - the dispatch table says every configured slot is + * busy, now. A full hive outranks everything the older tick might say + * about skipped candidates, because the skip story explains a past round + * while the table describes this instant; a hive that filled up after a + * no-choice decision must not keep explaining its busy slots as an empty + * queue. + * 2. `work-dispatched` - the table says work is in flight with room to + * spare. The queue is demonstrably supplying work; how many slots remain + * idle is a question about the round that has not happened yet. + * 3. `no-eligible-work` - every slot idle, and only now is the tick + * consulted: it must be readable, recent (a decision older than two + * intervals explains a scheduler that stopped, not a queue that is + * empty), produced under an enabled scheduler, with capacity configured, + * and it must be THE explicit no-candidate decision - `allowed: false` + * with queend's closed `nothing to choose` refusal. Any other refusal - + * a spent budget, a capacity limit - is work the queue HAD and money or + * slots refused, which is not this state's to claim. + * 4. `unknown` - everything else. Silence, staleness and unreadable + * evidence are never dressed up as an empty queue: an idle slot with no + * explanation is an operator's restart trigger, and this word exists so + * the trigger is pulled for a reason. + */ +function classifyQueueState(facts: { + activeWorkers: number + capacity: number + schedulerEnabled: boolean + tickDecidedAtMs: number | null + tickRefusedNothingToChoose: boolean + nowMs: number + intervalSeconds: number +}): QueueState { + if (facts.capacity > 0 && facts.activeWorkers >= facts.capacity) { + return 'capacity-full' + } + if (facts.activeWorkers > 0) return 'work-dispatched' + const freshWithinMs = facts.intervalSeconds * 1000 * TICK_STALENESS_INTERVALS + const tickVouchesForEmptyQueue = + facts.schedulerEnabled && + facts.capacity > 0 && + facts.tickDecidedAtMs != null && + facts.tickDecidedAtMs >= facts.nowMs - freshWithinMs && + facts.tickRefusedNothingToChoose + return tickVouchesForEmptyQueue ? 'no-eligible-work' : 'unknown' +} + +/** + * The public queue projection: a closed state and when the evidence behind it + * was observed, and not one byte more. + * + * `observedAt` dates the evidence, not the HTTP response: the tick's decision + * time whenever the state rests on the tick row (`no-eligible-work`, and + * `unknown` whose reason is a row that has gone quiet - its age is the + * finding), and the read time whenever the dispatch table is the evidence + * (`capacity-full`, `work-dispatched`) or nothing was ever recorded. A + * dashboard can therefore trust that a non-`unknown` state is either live or + * freshly decided, and measure exactly how stale an `unknown` is. + */ +function queueProjection( + state: QueueState, + tickDecidedAtMs: number | null, + readAtMs: number, +): { state: QueueState; observedAt: string } { + const restsOnTick = + state === 'no-eligible-work' || + (state === 'unknown' && tickDecidedAtMs != null) + return { + state, + observedAt: new Date( + restsOnTick && tickDecidedAtMs != null ? tickDecidedAtMs : readAtMs, + ).toISOString(), + } +} + export function createQueenPublicStatusRoute(deps: QueenPublicStatusDeps = {}) { const databaseUrl = deps.databaseUrl ?? configuredDatabaseUrl const createPool = @@ -297,6 +449,8 @@ export function createQueenPublicStatusRoute(deps: QueenPublicStatusDeps = {}) { const tickIntervalSeconds = deps.tickIntervalSeconds ?? configuredTickIntervalSeconds const billingMode = deps.billingMode ?? configuredBillingMode + const workerCapacity = deps.workerCapacity ?? configuredWorkerCapacity + const now = deps.now ?? Date.now return new Hono().get('/', async (c) => { c.header('Cache-Control', 'no-store') @@ -344,6 +498,10 @@ export function createQueenPublicStatusRoute(deps: QueenPublicStatusDeps = {}) { const countRow = counts.rows[0] ?? {} const latestRow = latest.rowCount ? latest.rows[0] : null const schedulerEnabled = intervalSeconds > 0 + const readAtMs = now() + const activeWorkers = asCount(countRow.running) + const capacity = workerCapacity() + const tickDecidedAtMs = decidedAtMs(tickRow?.decided_at) // Read once, quoted twice: `lastTick.refusal` and the swarmState // classification must be two readings of the same tick decision, or // the state could name a cause the tick never measured. @@ -353,7 +511,7 @@ export function createQueenPublicStatusRoute(deps: QueenPublicStatusDeps = {}) { return c.json({ status: 'ok', swarmState: classifySwarmState({ - running: asCount(countRow.running), + running: activeWorkers, unreviewed: asCount(countRow.unreviewed), schedulerEnabled, // A decision that is not a readable object cannot explain the @@ -368,6 +526,27 @@ export function createQueenPublicStatusRoute(deps: QueenPublicStatusDeps = {}) { // after accepted work finishes. decisionFoundNoEligibleCandidate: decision?.allowed === false, }), + // The queue reading of the same three rows. `running` is the active + // worker count: a dispatch that never started is written already + // finished, so an unfinished dispatch is a bee holding a paid slot. + queue: queueProjection( + classifyQueueState({ + activeWorkers, + capacity, + schedulerEnabled, + tickDecidedAtMs, + // queend writes this refusal from one closed template when its + // chooser found no candidate (main.swift); any other refusal - + // budget, capacity - is not this page's to reinterpret. + tickRefusedNothingToChoose: + decision?.allowed === false && + decision?.refusal === 'nothing to choose', + nowMs: readAtMs, + intervalSeconds, + }), + tickDecidedAtMs, + readAtMs, + ), scheduler: { enabled: schedulerEnabled, intervalSeconds, diff --git a/trios/agent-server/apps/server/tests/api/queen-public-status.test.ts b/trios/agent-server/apps/server/tests/api/queen-public-status.test.ts index 8d777ce7fc..880a5b9a4e 100644 --- a/trios/agent-server/apps/server/tests/api/queen-public-status.test.ts +++ b/trios/agent-server/apps/server/tests/api/queen-public-status.test.ts @@ -10,12 +10,19 @@ type QueryResult = { rowCount: number; rows: Array> } function fakePool(results: QueryResult[]) { let at = 0 let ended = false + const statements: string[] = [] return { - query: async () => results[at++] ?? { rowCount: 0, rows: [] }, + query: async (sql?: string) => { + // Counted, not parsed: the contract under test is HOW MANY queries a + // read issues, never their text. + if (typeof sql === 'string') statements.push(sql) + return results[at++] ?? { rowCount: 0, rows: [] } + }, end: async () => { ended = true }, wasEnded: () => ended, + statementCount: () => statements.length, } } @@ -30,6 +37,20 @@ const emptyDispatchCounts: QueryResult = { } const noLatestDispatch: QueryResult = { rowCount: 0, rows: [] } +/** + * A clock fixed inside the fixtures' own time, so whether a tick is fresh is + * a property of the fixture rather than of when the suite happens to run. + * READ_AT sits half an interval after FRESH_TICK_DECIDED_AT (interval 1800s, + * freshness window two intervals) and hours after every stale date below. + * + * Tick-backed queue states echo the decision time as `observedAt`; table- + * backed ones echo READ_AT - both pinned exactly by these constants. + */ +const FRESH_TICK_DECIDED_AT = '2026-09-01T04:17:42.983Z' +const READ_AT = '2026-09-01T04:47:42.983Z' +const STALE_TICK_DECIDED_AT = '2026-09-01T00:47:42.983Z' +const fixedNow = () => Date.parse(READ_AT) + /** * One skip reason per sentence `queend` writes, verbatim in shape * (queen-core/Sources/queend/main.swift): issue number, payload, and for the @@ -153,7 +174,7 @@ describe('GET /queen/status', () => { rowCount: 1, rows: [ { - decided_at: '2026-09-01T04:17:42.983Z', + decided_at: FRESH_TICK_DECIDED_AT, decision: { allowed: false, refusal: 'nothing to choose', @@ -188,6 +209,8 @@ describe('GET /queen/status', () => { createPool: () => pool, tickIntervalSeconds: () => 1800, billingMode: () => 'coding_plan', + workerCapacity: () => 4, + now: fixedNow, }).request('/') expect(response.status).toBe(200) @@ -200,6 +223,13 @@ describe('GET /queen/status', () => { // stay counted under dispatches.unreviewed; they just cannot name the // quiet while the tick that measured it says otherwise. swarmState: 'healthy_idle', + // Four paid slots, all idle, and a fresh explicit no-choice decision + // behind them: the queue is empty, which is the answer an operator + // needs before reaching for a restart. + queue: { + state: 'no-eligible-work', + observedAt: FRESH_TICK_DECIDED_AT, + }, scheduler: { enabled: true, intervalSeconds: 1800, @@ -207,7 +237,7 @@ describe('GET /queen/status', () => { estimatedUSDGateEnabled: false, }, lastTick: { - decidedAt: '2026-09-01T04:17:42.983Z', + decidedAt: FRESH_TICK_DECIDED_AT, allowed: false, refusal: 'nothing to choose', skippedCount: representativeSkipped.length, @@ -306,13 +336,20 @@ describe('GET /queen/status', () => { const leakBranch = 'queen-1291-leak-probe' const leakTitle = 'Teach the swarm to dream' const leakSecret = 'sk-trios-9f8e7d6c5b4a3210fedcba9876543210' + // The two values the queue criteria add: an issue body and a transcript + // excerpt, planted where a scheduler would actually hold them, to prove + // the newest field opens no newest channel. + const leakBody = + 'Boundary: the sacred scroll of queen-secret.ts and nothing else' + const leakTranscript = + 'worker murmurs the sacred scroll of queen-secret.ts' const pool = fakePool([ { rowCount: 1, rows: [ { - decided_at: '2026-09-01T04:17:42.983Z', + decided_at: FRESH_TICK_DECIDED_AT, decision: { allowed: false, refusal: 'nothing to choose', @@ -326,6 +363,8 @@ describe('GET /queen/status', () => { credentials: { token: leakSecret }, chosenPaths: [leakPath], strays: [{ issue: 1291, paths: [leakPath] }], + candidateBodies: { 1291: leakBody }, + transcript: leakTranscript, }, }, ], @@ -341,7 +380,7 @@ describe('GET /queen/status', () => { outcome: 'running', branch: leakBranch, title: leakTitle, - detail: leakSecret, + detail: leakTranscript, review_note: leakSecret, conversation_id: leakSecret, provider: leakSecret, @@ -355,6 +394,8 @@ describe('GET /queen/status', () => { databaseUrl: () => 'postgres://configured', createPool: () => pool, tickIntervalSeconds: () => 1800, + workerCapacity: () => 4, + now: fixedNow, }).request('/') expect(response.status).toBe(200) @@ -362,10 +403,20 @@ describe('GET /queen/status', () => { // The classification ran on planted rows without echoing any of them, // and an empty swarm under a live scheduler reads as health. expect(body.swarmState).toBe('healthy_idle') + // The queue projection is closed: one state word, one timestamp, and no + // third key a payload could ever ride. + expect(body.queue).toEqual({ + state: 'no-eligible-work', + observedAt: FRESH_TICK_DECIDED_AT, + }) + expect(Object.keys(body.queue as object).sort()).toEqual([ + 'observedAt', + 'state', + ]) // The categorisation still worked, and what it published is the issue // numbers - bare integers - and nothing else the reasons carried. expect(body.lastTick).toEqual({ - decidedAt: '2026-09-01T04:17:42.983Z', + decidedAt: FRESH_TICK_DECIDED_AT, allowed: false, refusal: 'nothing to choose', skippedCount: 3, @@ -386,6 +437,8 @@ describe('GET /queen/status', () => { expect(serialized).not.toContain(leakBranch) expect(serialized).not.toContain(leakTitle) expect(serialized).not.toContain(leakSecret) + expect(serialized).not.toContain(leakBody) + expect(serialized).not.toContain(leakTranscript) }) it('summarises a legacy decision with no skip array as empty', async () => { @@ -410,6 +463,8 @@ describe('GET /queen/status', () => { databaseUrl: () => 'postgres://configured', createPool: () => pool, tickIntervalSeconds: () => 1800, + workerCapacity: () => 4, + now: fixedNow, }).request('/') expect(response.status).toBe(200) @@ -419,6 +474,13 @@ describe('GET /queen/status', () => { // be the real recordTick -> recordDispatch window, so the snapshot is // unavailable until the row appears or a no-choice tick supersedes it. swarmState: 'unavailable', + // And the queue cannot claim dispatched work the table does not show, + // nor an empty queue a decision this old never said. The state is + // unknown, dated by the row that went quiet. + queue: { + state: 'unknown', + observedAt: '2026-08-30T09:00:00.000Z', + }, scheduler: { enabled: true, intervalSeconds: 1800, @@ -857,4 +919,250 @@ describe('skipSummary issue numbers', () => { }) expect(JSON.stringify(body)).not.toContain('1188') }) + + /** + * The queue fixtures below share one shape so each states only what it + * varies: four paid slots, a 1800s tick interval, and a read at READ_AT. + */ + const queueRoute = ( + pool: ReturnType, + opts?: { + workerCapacity?: () => number + tickIntervalSeconds?: () => number + }, + ) => + createQueenPublicStatusRoute({ + databaseUrl: () => 'postgres://configured', + createPool: () => pool, + tickIntervalSeconds: opts?.tickIntervalSeconds ?? (() => 1800), + workerCapacity: opts?.workerCapacity ?? (() => 4), + now: fixedNow, + }) + + const noChoiceTick = (decidedAt: string): QueryResult => ({ + rowCount: 1, + rows: [ + { + decided_at: decidedAt, + decision: { + allowed: false, + refusal: 'nothing to choose', + skipped: ['#1298: not first'], + }, + }, + ], + }) + + const runningCounts = (running: number): QueryResult => ({ + rowCount: 1, + rows: [ + { + total: String(running), + finished: '0', + running: String(running), + unreviewed: '0', + }, + ], + }) + + const queueOf = async ( + pool: ReturnType, + opts?: { + workerCapacity?: () => number + tickIntervalSeconds?: () => number + }, + ): Promise<{ state: string; observedAt: string }> => { + const body = (await (await queueRoute(pool, opts).request('/')).json()) as { + queue: { state: string; observedAt: string } + } + return body.queue + } + + it('reads idle paid slots as no-eligible-work from the explicit no-choice decision', async () => { + // Acceptance scenario 1: four idle slots, nothing running, and a last + // tick refused as `nothing to choose`. The skip summary explains WHY each + // candidate was passed over; the queue state is the one word a dashboard + // can print next to 0% utilization without guessing. + const pool = fakePool([ + noChoiceTick(FRESH_TICK_DECIDED_AT), + runningCounts(0), + noLatestDispatch, + ]) + expect(await queueOf(pool)).toEqual({ + state: 'no-eligible-work', + observedAt: FRESH_TICK_DECIDED_AT, + }) + }) + + it('reads a full worker pool as capacity-full regardless of skipped candidates', async () => { + // Acceptance scenario 2: active workers equal capacity, so there is no + // idle slot left to explain - and the tick's skip story must not re-explain + // a full hive as an empty queue. The tick here is fresh AND explicitly + // no-choice with a skipped candidate, exactly the rows that would read + // no-eligible-work if precedence were the other way round. + const fullWithSkips = fakePool([ + noChoiceTick(FRESH_TICK_DECIDED_AT), + runningCounts(4), + noLatestDispatch, + ]) + expect(await queueOf(fullWithSkips)).toEqual({ + state: 'capacity-full', + observedAt: READ_AT, + }) + + // And with no tick row at all: the table alone vouches for the full pool, + // dated by the read rather than by telemetry that does not exist. + const fullWithoutTick = fakePool([ + { rowCount: 0, rows: [] }, + runningCounts(4), + noLatestDispatch, + ]) + expect(await queueOf(fullWithoutTick)).toEqual({ + state: 'capacity-full', + observedAt: READ_AT, + }) + }) + + it('reads partial utilization as work-dispatched on the dispatch table alone', async () => { + // The fourth closed state: two of four slots busy. The tick is stale and + // no-choice besides, and none of that matters - unfinished dispatches are + // a fact about the present, read straight off the table, so the queue is + // demonstrably supplying work and the free slots belong to a round that + // has not happened yet. + const pool = fakePool([ + noChoiceTick(STALE_TICK_DECIDED_AT), + runningCounts(2), + noLatestDispatch, + ]) + expect(await queueOf(pool)).toEqual({ + state: 'work-dispatched', + observedAt: READ_AT, + }) + }) + + it('never calls a stale or missing tick an empty queue', async () => { + // Acceptance scenario 3: scheduler telemetry absent or stale reads as + // unknown, never as no-eligible-work. A decision four hours old explains + // a scheduler that stopped, not a queue that is empty. + const stale = fakePool([ + noChoiceTick(STALE_TICK_DECIDED_AT), + runningCounts(0), + noLatestDispatch, + ]) + expect(await queueOf(stale)).toEqual({ + state: 'unknown', + observedAt: STALE_TICK_DECIDED_AT, + }) + + // No tick row at all: nothing was ever observed, so the observation is + // the read itself. + const missing = fakePool([ + { rowCount: 0, rows: [] }, + runningCounts(0), + noLatestDispatch, + ]) + expect(await queueOf(missing)).toEqual({ + state: 'unknown', + observedAt: READ_AT, + }) + + // A row whose decision cannot be read says nothing, even at a fresh date; + // `lastTick` still publishes the row exactly as it always has. + const unreadable = fakePool([ + { + rowCount: 1, + rows: [{ decided_at: FRESH_TICK_DECIDED_AT, decision: 'corrupt' }], + }, + runningCounts(0), + noLatestDispatch, + ]) + expect(await queueOf(unreadable)).toEqual({ + state: 'unknown', + observedAt: FRESH_TICK_DECIDED_AT, + }) + }) + + it('reserves no-eligible-work for the closed nothing-to-choose refusal', async () => { + // Other refusals are work the queue HAD and something else refused. A + // spent budget (queend's other closed template, main.swift) with every + // slot idle is not an empty queue, and must not read as one. + const budgetRefusal = fakePool([ + { + rowCount: 1, + rows: [ + { + decided_at: FRESH_TICK_DECIDED_AT, + decision: { + allowed: false, + refusal: + 'the swarm has spent about $10.85 today, $0.85 past its ' + + '$10.00 daily limit (raise it with TRIOS_SWARM_DAILY_CAP_USD)', + }, + }, + ], + }, + runningCounts(0), + noLatestDispatch, + ]) + expect((await queueOf(budgetRefusal)).state).toBe('unknown') + + // A decision that chose work vouches for neither an empty queue nor a + // dispatch the table cannot show - the same window swarmState refuses to + // call healthy. + const choseWork = fakePool([ + { + rowCount: 1, + rows: [ + { + decided_at: FRESH_TICK_DECIDED_AT, + decision: { allowed: true, refusal: null, chosen: 1296 }, + }, + ], + }, + runningCounts(0), + noLatestDispatch, + ]) + expect((await queueOf(choseWork)).state).toBe('unknown') + + // A disabled scheduler: a perfectly fresh, perfectly explicit no-choice + // decision from a loop nobody restarted explains yesterday, not today. + const schedulerOff = fakePool([ + noChoiceTick(FRESH_TICK_DECIDED_AT), + runningCounts(0), + noLatestDispatch, + ]) + expect( + (await queueOf(schedulerOff, { tickIntervalSeconds: () => 0 })).state, + ).toBe('unknown') + + // No paid slot configured: whether a slot is idle is unknowable, and an + // unknowable idle must not borrow the empty-queue word. + const noCapacity = fakePool([ + noChoiceTick(FRESH_TICK_DECIDED_AT), + runningCounts(0), + noLatestDispatch, + ]) + expect((await queueOf(noCapacity, { workerCapacity: () => 0 })).state).toBe( + 'unknown', + ) + }) + + it('derives the queue projection without a fourth database query', async () => { + // The queue state is computed from the rows this route has always read - + // the tick row, the counts, the latest dispatch - plus the configured + // capacity, which comes from the environment and not the database. The + // count proves no query was added and none was skipped to compensate: the + // projection still derived a state. + const pool = fakePool([ + noChoiceTick(FRESH_TICK_DECIDED_AT), + runningCounts(4), + noLatestDispatch, + ]) + const response = await queueRoute(pool).request('/') + expect(response.status).toBe(200) + const body = (await response.json()) as Record + expect((body.queue as { state: string }).state).toBe('capacity-full') + expect(pool.statementCount()).toBe(3) + expect(pool.wasEnded()).toBe(true) + }) })