diff --git a/trios/.trinity/events/akashic-log.jsonl b/trios/.trinity/events/akashic-log.jsonl index 953e4a07fd..5ab18065a8 100644 --- a/trios/.trinity/events/akashic-log.jsonl +++ b/trios/.trinity/events/akashic-log.jsonl @@ -9,3 +9,10 @@ {"ts":"2026-09-01T13:07:35Z","event":"loop.handoff","agent":"codex-root","handoff_id":"d676c4ce-b960-438c-8f0c-05aa1f6f7396","task_id":"GH-1300","resume_point":"Watch live dispatches #1301 and #1302 through Queen review; production is on coding_plan with two active unique credentials.","next_options":["Add two remaining unique credentials securely and prove four-way capacity.","Implement provider quota-state and reset-time visibility without exposing response bodies.","Reconcile the seven legacy review cards against the runtime ledger that reports zero unreviewed dispatches."]} {"ts":"2026-09-01T17:42:12Z","event":"task.intent","agent":"codex-root","task_id":"GH-1313","spec_path":".trinity/specs/queen-public-hardware-registry.md","graph_node":"queen_public_hardware_registry","priority":"P1"} {"ts":"2026-09-01T17:42:12Z","event":"claim.acquire","agent":"codex-root","claim_id":"5A13B719-47C0-4490-B68F-8A1D328A8401","resource":".trinity/specs/queen-public-hardware-registry.md","ttl_sec":3600} +{"ts":"2026-10-01T21:48:30.035Z","event":"task.intent","agent":"codex-contributor-keys","task_id":"contributor-keys","spec_path":"agent-server/specs/automation/queen-contributor-keys.t27","graph_node":"queen_contributor_keys"} +{"ts":"2026-10-01T22:06:32.008Z","event":"claim.heartbeat","agent":"codex-contributor-keys","claim_id":"2d5b4461-adaa-4135-9b39-24296a6456a5"} +{"ts":"2026-10-01T22:09:25.786Z","event":"claim.release","agent":"codex-contributor-keys","claim_id":"2d5b4461-adaa-4135-9b39-24296a6456a5","result":"clean","validation":"268 Bun tests and 9 native spec tests passed; deployment pending"} +{"ts":"2026-10-01T22:10:35.338Z","event":"task.intent","agent":"codex-contributor-keys","task_id":"contributor-keys-generation","spec_path":"agent-server/specs/automation/queen-contributor-keys.t27","detail":"Preserve byte-identical t27c output through the existing formatting hook"} +{"ts":"2026-10-01T22:11:04.287Z","event":"claim.release","agent":"codex-contributor-keys","claim_id":"b00e2190-4289-4c02-bb2b-04fa398f29e4","result":"clean","validation":"Generated output remains byte-identical after formatter; 12 host tests pass"} +{"ts":"2026-10-01T22:12:34.391Z","event":"task.intent","agent":"codex-contributor-keys","task_id":"contributor-keys-ci","spec_path":"agent-server/specs/automation/queen-contributor-keys.t27","detail":"Connect live contributor tests to the existing disposable CI PostgreSQL"} +{"ts":"2026-10-01T22:14:26.087Z","event":"claim.release","agent":"codex-contributor-keys","claim_id":"c15a0cee-a9c3-4d66-ac7b-bc4cde3a67fc","result":"clean","validation":"Existing CI pglive command: 12 tests pass, 75 assertions, no skips, with TRIOS_PG_TEST_URL"} diff --git a/trios/.trinity/experience/2026-10-01_contributor-keys.json b/trios/.trinity/experience/2026-10-01_contributor-keys.json new file mode 100644 index 0000000000..20cfae19b5 --- /dev/null +++ b/trios/.trinity/experience/2026-10-01_contributor-keys.json @@ -0,0 +1,27 @@ +{ + "timestamp": "2026-10-01T22:12:00Z", + "task_id": "CONTRIBUTOR-KEYS-001", + "tracking_issues": [ + "https://github.com/gHashTag/999-multibots-telegraf/issues/3251", + "https://github.com/gHashTag/t27/issues/5472" + ], + "contract": "Own-account Queen key management preserves immutable credential ownership, existing environment indices and dispatch-derived XP. Management authentication never overrides persisted consent.", + "red": [ + "Removing SQL owner filtering exposed a foreign credential to the actual route test.", + "Removing allocation filtering selected a disabled environment key.", + "Removing owner AAD made the wrong-owner decrypt test fail." + ], + "green": { + "bun_regression": "268 tests passed, 0 failed, 0 skipped, 1015 assertions across 11 files; includes 9 isolated PostgreSQL tests and compiled Swift core", + "typescript": "tsc --noEmit passed", + "spec_runtime": "9 native t27 tests passed with Zig 0.16.0", + "spec_negative_controls": "9 controls killed; verified by the canonical spec agent", + "generated_bindings": "Real t27c gen-ts output; every emitted constant compared with compiler AST in host test", + "review_fixes": "Disable does not decrypt; corrupt managed entries do not stop other credentials; removing the management capability retains consent and attribution; Ollama and legacy reviewers retain their paths" + }, + "external_gates": [ + "Deployment, live provider probes and XP visibility are not proven by these local checks.", + "Generated policy contains declarations; protocol, SQL and cryptographic operations remain explicit TypeScript host adapters." + ], + "learning": "A service capability authenticates management calls, while durable consent remains authoritative without that capability. Credential identity must be assigned before filtering, and disable must remain possible when decryption fails. Concurrent t27 test-report processes need separate TMPDIR directories because the compiler currently derives its work directory from the spec basename." +} diff --git a/trios/agent-server/apps/server/src/api/routes/queen-contributor-keys.ts b/trios/agent-server/apps/server/src/api/routes/queen-contributor-keys.ts new file mode 100644 index 0000000000..040dae4144 --- /dev/null +++ b/trios/agent-server/apps/server/src/api/routes/queen-contributor-keys.ts @@ -0,0 +1,209 @@ +import { Hono } from 'hono' +import type { Pool } from 'pg' +import { createQueenPool } from '../../lib/db/queen-pool' +import { + addContributorKey, + CONTRIBUTOR_PROVIDERS, + ContributorError, + changeContributorKey, + contributorGithub, + type EnvironmentKey, + listContributorKeys, + trustedContributor, +} from '../services/queen-contributor-keys' +import { CONTRIBUTOR_POLICY } from '../services/queen-contributor-policy' +import { environmentContributorKeys } from '../services/queen-dispatch' +import { + ACCEPTED_XP, + HOUR_XP, + keyWork, + parseOwners, + rank, + SPEC_XP, +} from '../services/queen-leaderboard' + +export interface ContributorRouteDeps { + pool?: () => Pool + environment?: () => EnvironmentKey[] + owners?: () => Record + fetcher?: typeof fetch +} +let productionPool: Pool | undefined +function poolForKeys(): Pool { + const url = process.env.DATABASE_URL + if (!url) throw new ContributorError('contributor_keys_unavailable', 503) + productionPool ??= createQueenPool(url) + return productionPool +} +const emptyContribution = () => ({ + xp: 0, + accepted: 0, + specs: 0, + finished: 0, + hours: 0, +}) +async function contributionOf(pool: Pool, id: number) { + const work = (await keyWork(pool)).filter((entry) => entry.keyIndex === id) + const total = rank(work, {})[0] + return total + ? { + xp: total.xp, + accepted: total.accepted, + specs: total.specs ?? 0, + finished: total.finished, + hours: total.hours, + } + : emptyContribution() +} + +/** Render attests a verified app subject. No browser can self-assert ownership. */ +export function createQueenContributorKeysRoute( + deps: ContributorRouteDeps = {}, +) { + const app = new Hono<{ Variables: { contributor: string; body: string } }>() + app.use('/*', async (c, next) => { + c.header('Cache-Control', 'no-store') + try { + c.set( + 'contributor', + trustedContributor( + c.req.header('authorization'), + c.req.header('x-queen-contributor-id'), + ), + ) + return await next() + } catch (error) { + const known = + error instanceof ContributorError + ? error + : new ContributorError('contributor_keys_unavailable', 503) + return c.json({ error: known.code }, known.status as 400) + } + }) + app.use('/*', async (c, next) => { + if (c.req.method !== 'POST') return next() + const reader = c.req.raw.body?.getReader() + const chunks: Uint8Array[] = [] + let bytes = 0 + if (reader) { + try { + for (;;) { + const { done, value } = await reader.read() + if (done) break + bytes += value.byteLength + if (bytes > CONTRIBUTOR_POLICY.MAX_BODY_BYTES) { + await reader.cancel() + throw new ContributorError('request_too_large', 413) + } + chunks.push(value) + } + } finally { + reader.releaseLock() + } + } + c.set('body', Buffer.concat(chunks).toString('utf8')) + return await next() + }) + app.get('/', async (c) => { + const pool = (deps.pool ?? poolForKeys)() + const subject = c.get('contributor') + const environment = (deps.environment ?? environmentContributorKeys)() + const owners = ( + deps.owners ?? (() => parseOwners(process.env.TRIOS_KEY_OWNERS)) + )() + const keys = await listContributorKeys(pool, subject, environment, owners) + const work = await keyWork(pool) + const byKey = new Map(work.map((item) => [item.keyIndex, item])) + const own = work.filter((item) => + keys.some((key) => key.id === item.keyIndex), + ) + const totals = rank( + own, + Object.fromEntries(keys.map((key) => [key.id, 'own'])), + )[0] + const contribution = (entry: typeof totals) => + entry + ? { + xp: entry.xp, + accepted: entry.accepted, + specs: entry.specs ?? 0, + finished: entry.finished, + hours: entry.hours, + } + : emptyContribution() + return c.json({ + keys: keys.map((key) => { + const entry = byKey.get(key.id) + return { + ...key, + contribution: contribution(rank(entry ? [entry] : [], {})[0]), + } + }), + providers: CONTRIBUTOR_PROVIDERS.map(({ id, label, model }) => ({ + id, + label, + model, + })), + contribution: contribution(totals), + attribution: { subject, github: contributorGithub(subject) }, + scoring: { acceptedXp: ACCEPTED_XP, specXp: SPEC_XP, hourXp: HOUR_XP }, + }) + }) + app.post('/', async (c) => { + let body: unknown + try { + body = JSON.parse(c.get('body')) + } catch { + throw new ContributorError('invalid_json') + } + const pool = (deps.pool ?? poolForKeys)() + const key = await addContributorKey( + pool, + c.get('contributor'), + body, + (deps.environment ?? environmentContributorKeys)(), + deps.fetcher, + ) + return c.json( + { key: { ...key, contribution: await contributionOf(pool, key.id) } }, + 201, + ) + }) + app.post('/:id/:action', async (c) => { + const raw = c.req.param('id') + const id = Number(raw) + if ( + raw !== String(id) || + !/^-?(?:0|[1-9][0-9]*)$/.test(raw) || + !Number.isInteger(id) || + id < -2147483647 || + id > 2147483647 + ) { + throw new ContributorError('invalid_key_id') + } + const action = c.req.param('action') + if (action !== 'probe' && action !== 'enable' && action !== 'disable') + throw new ContributorError('invalid_action') + const pool = (deps.pool ?? poolForKeys)() + const key = await changeContributorKey( + pool, + c.get('contributor'), + id, + action, + (deps.environment ?? environmentContributorKeys)(), + (deps.owners ?? (() => parseOwners(process.env.TRIOS_KEY_OWNERS)))(), + deps.fetcher, + ) + return c.json({ + key: { ...key, contribution: await contributionOf(pool, key.id) }, + }) + }) + app.onError((error, c) => { + const known = + error instanceof ContributorError + ? error + : new ContributorError('contributor_keys_unavailable', 503) + return c.json({ error: known.code }, known.status as 400) + }) + return app +} diff --git a/trios/agent-server/apps/server/src/api/routes/queen-kanban.ts b/trios/agent-server/apps/server/src/api/routes/queen-kanban.ts index 742dad1e8a..7593382d2d 100644 --- a/trios/agent-server/apps/server/src/api/routes/queen-kanban.ts +++ b/trios/agent-server/apps/server/src/api/routes/queen-kanban.ts @@ -46,9 +46,11 @@ */ import { Hono } from 'hono' +import type { Pool } from 'pg' import { createQueenPool } from '../../lib/db/queen-pool' import { logger } from '../../lib/logger' import { + liveWorkerCapacity, type WorkerCapacityBreakdown, workerCapacityBreakdown, } from '../services/queen-dispatch' @@ -560,7 +562,7 @@ async function build( (lastTick.rows[0]?.decision as { refusal?: string } | undefined) ?.refusal ?? null, roundSeconds: Number(process.env.TRIOS_QUEEN_TICK_SECONDS ?? '0') || null, - ...boardWorkerCapacity(), + ...boardWorkerCapacity(await liveWorkerCapacity(pool as Pool)), }, } } diff --git a/trios/agent-server/apps/server/src/api/routes/queen-public-research.ts b/trios/agent-server/apps/server/src/api/routes/queen-public-research.ts index 0a5d459aee..bb1655a0a3 100644 --- a/trios/agent-server/apps/server/src/api/routes/queen-public-research.ts +++ b/trios/agent-server/apps/server/src/api/routes/queen-public-research.ts @@ -15,9 +15,11 @@ */ import { Hono } from 'hono' +import type { Pool } from 'pg' import { createQueenPool } from '../../lib/db/queen-pool' import { logger } from '../../lib/logger' import { + liveWorkerCapacity, type WorkerCapacityBreakdown, workerCapacityBreakdown, } from '../services/queen-dispatch' @@ -225,10 +227,15 @@ export function createQueenPublicResearchRoute( const url = databaseUrl() let runtime: { status: 'live' | 'offline' } = { status: 'offline' } let busyIndices: number[] = [] + let breakdown = capacityBreakdown() + let capacityAvailable = !!deps.workerCapacityBreakdown if (url) { const pool = createPool(url) try { + if (!deps.workerCapacityBreakdown) + breakdown = await liveWorkerCapacity(pool as Pool) + capacityAvailable = true const active = await pool.query( `SELECT key_index FROM queen_dispatch @@ -239,6 +246,8 @@ export function createQueenPublicResearchRoute( busyIndices = active.rows.map((row) => Number(row.key_index)) runtime = { status: 'live' } } catch (error) { + if (!capacityAvailable) + return c.json({ error: 'Worker capacity is unavailable' }, 503) logger.warn('Queen public research telemetry query failed', { error: error instanceof Error ? error.message : String(error), }) @@ -254,7 +263,6 @@ export function createQueenPublicResearchRoute( // One authority, one number: the projection's capacity IS the breakdown's // effective capacity, so an operator reading "4" and the factors below it // can never see two totals that disagree about the same configuration. - const breakdown = capacityBreakdown() return c.json({ ...graph, runtime, 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 25c57410c2..c639c94d6f 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 @@ -40,10 +40,14 @@ */ import { Hono } from 'hono' +import type { Pool } from 'pg' import { createQueenPool } from '../../lib/db/queen-pool' import { logger } from '../../lib/logger' import { workerModelRanking } from '../../lib/model-ranking' -import { configuredWorkerCapacity } from '../services/queen-dispatch' +import { + configuredWorkerCapacity, + liveWorkerCapacity, +} from '../services/queen-dispatch' interface QueryResult { rowCount: number | null @@ -635,7 +639,9 @@ export function createQueenPublicStatusRoute(deps: QueenPublicStatusDeps = {}) { // slot, and each reading keeps its own meaning. const running = asCount(countRow.running) const startedUnfinished = asCount(countRow.started_running) - const capacity = workerCapacity() + const capacity = deps.workerCapacity + ? workerCapacity() + : (await liveWorkerCapacity(pool as Pool)).effectiveCapacity 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 diff --git a/trios/agent-server/apps/server/src/api/server.ts b/trios/agent-server/apps/server/src/api/server.ts index 7959e7361f..c591e6373f 100644 --- a/trios/agent-server/apps/server/src/api/server.ts +++ b/trios/agent-server/apps/server/src/api/server.ts @@ -40,6 +40,7 @@ import { createMonitoringRoutes } from './routes/monitoring' import { createOAuthRoutes } from './routes/oauth' import { createOpenClawRoutes } from './routes/openclaw' import { createProviderRoutes } from './routes/provider' +import { createQueenContributorKeysRoute } from './routes/queen-contributor-keys' import { createQueenDashboardRoute } from './routes/queen-dashboard' import { createQueenExportRoute } from './routes/queen-export' import { @@ -390,6 +391,8 @@ export async function createHttpServer(config: HttpServerConfig) { .route('/queen/public-board', createQueenPublicBoardRoute()) // Who lent a lane and what it did; no titles, no worker text, no key. .route('/queen/public-leaderboard', createQueenPublicLeaderboardRoute()) + // Separate server-to-server capability; never a public-read or operator-token route. + .route('/queen/contributor-keys', createQueenContributorKeysRoute()) .route('/queen/registry', queenRegistryRoutes) // The shell only. It holds no state and no token; every byte of data it // shows comes from /queen/lease, which stays guarded. See the route header diff --git a/trios/agent-server/apps/server/src/api/services/queen-contributor-keys.ts b/trios/agent-server/apps/server/src/api/services/queen-contributor-keys.ts new file mode 100644 index 0000000000..f320608b50 --- /dev/null +++ b/trios/agent-server/apps/server/src/api/services/queen-contributor-keys.ts @@ -0,0 +1,656 @@ +import { + createCipheriv, + createDecipheriv, + createHash, + randomBytes, + timingSafeEqual, +} from 'node:crypto' +import type { Pool } from 'pg' +import { logger } from '../../lib/logger' +import { CONTRIBUTOR_POLICY as P } from './queen-contributor-policy' + +export type ContributorProvider = 'nvidia' | 'zai' +export type ProbeStatus = + | 'ok' + | 'invalid' + | 'rate_limited' + | 'unavailable' + | 'not_checked' +export interface KeyProbe { + status: ProbeStatus + checkedAt: string | null + latencyMs: number | null +} +export interface EnvironmentKey { + id: number + apiKey: string + provider: ContributorProvider + model: string + baseUrl: string + contextWindow?: number +} +export interface ContributorKey { + id: number + source: 'managed' | 'environment' + provider: ContributorProvider + model: string + label: string + fingerprint: string + enabled: boolean + createdAt: string | null + lastProbe: KeyProbe +} +interface KeyRow { + key_index: number + source: 'managed' | 'environment' + fingerprint: string + owner_subject: string + owner_name: string + provider: ContributorProvider + model: string + label: string + sealed: string | null + enabled: boolean + revision: number + probe_status: ProbeStatus + probe_at: Date | string | null + probe_ms: number | null + created_at: Date | string +} +export class ContributorError extends Error { + constructor( + public readonly code: string, + public readonly status = 400, + ) { + super(code) + } +} +export const CONTRIBUTOR_PROVIDERS = [ + { + id: 'nvidia' as const, + label: 'NVIDIA NIM', + model: P.NVIDIA_MODEL, + baseUrl: P.NVIDIA_URL, + }, + { id: 'zai' as const, label: 'Z.ai', model: P.ZAI_MODEL, baseUrl: P.ZAI_URL }, +] +const ready = new WeakMap>() +export function contributorsEnabled(): boolean { + return ( + (process.env.QUEEN_CONTRIBUTOR_PROXY_TOKEN?.length ?? 0) >= + P.PROXY_TOKEN_MIN_BYTES + ) +} +export function trustedContributor( + authorization: string | undefined, + subject: string | undefined, +): string { + const expected = process.env.QUEEN_CONTRIBUTOR_PROXY_TOKEN ?? '' + if (!contributorsEnabled()) + throw new ContributorError('contributor_keys_unavailable', 503) + const provided = authorization?.startsWith('Bearer ') + ? authorization.slice(7) + : '' + const a = createHash('sha256').update(provided).digest() + const b = createHash('sha256').update(expected).digest() + if (!provided || !timingSafeEqual(a, b)) + throw new ContributorError('forbidden', 403) + if ( + !subject || + subject !== subject.trim() || + !/^telegram:[1-9][0-9]{0,19}$/.test(subject) + ) { + throw new ContributorError('invalid_contributor', 400) + } + return subject +} +export function contributorGithub(subject: string): string | undefined { + try { + const identities = JSON.parse( + process.env.QUEEN_CONTRIBUTOR_IDENTITIES ?? '{}', + ) + const login = Object.hasOwn(identities, subject) + ? identities[subject] + : undefined + return typeof login === 'string' && + login === login.trim() && + /^[a-zA-Z\d](?:[a-zA-Z\d]|-(?=[a-zA-Z\d])){0,38}$/.test(login) + ? login + : undefined + } catch { + return undefined + } +} +export function contributorName(subject: string): string { + const github = contributorGithub(subject) + // Never publish a Telegram subject. The pseudonym is stable but not an ID. + return github + ? `@${github}` + : `contributor ${fingerprint(subject).slice(0, 8)}` +} +export function fingerprint(key: string): string { + return createHash('sha256').update(key.trim()).digest('hex') +} +function masterKey(): Buffer { + const encoded = process.env.QUEEN_CONTRIBUTOR_ENCRYPTION_KEY ?? '' + const key = Buffer.from(encoded, 'base64') + if (key.length !== 32 || key.toString('base64') !== encoded) { + throw new ContributorError('key_storage_unavailable', 503) + } + return key +} +export function sealCredential( + secret: string, + subject: string, + digest: string, +): string { + const nonce = randomBytes(12) + const cipher = createCipheriv('aes-256-gcm', masterKey(), nonce) + cipher.setAAD(Buffer.from(`${subject}\n${digest}`)) + const ciphertext = Buffer.concat([ + cipher.update(secret, 'utf8'), + cipher.final(), + ]) + return [nonce, cipher.getAuthTag(), ciphertext] + .map((part) => part.toString('base64')) + .join('.') +} +export function openCredential( + sealed: string, + subject: string, + digest: string, +): string { + const parts = sealed.split('.') + if (parts.length !== 3) + throw new ContributorError('key_storage_unavailable', 503) + try { + const [nonce, tag, ciphertext] = parts.map((part) => + Buffer.from(part, 'base64'), + ) + const cipher = createDecipheriv('aes-256-gcm', masterKey(), nonce) + cipher.setAAD(Buffer.from(`${subject}\n${digest}`)) + cipher.setAuthTag(tag) + return Buffer.concat([cipher.update(ciphertext), cipher.final()]).toString( + 'utf8', + ) + } catch { + throw new ContributorError('key_storage_unavailable', 503) + } +} +export async function ensureContributorKeys(pool: Pool): Promise { + let promise = ready.get(pool) + if (!promise) { + promise = pool + .query(` + CREATE SEQUENCE IF NOT EXISTS queen_contributor_key_index AS integer + INCREMENT BY -1 MINVALUE -2147483647 MAXVALUE -1 START WITH -1 NO CYCLE; + CREATE TABLE IF NOT EXISTS queen_contributor_keys ( + key_index integer PRIMARY KEY DEFAULT nextval('queen_contributor_key_index'), + source text NOT NULL CHECK (source IN ('managed','environment')), + fingerprint text NOT NULL UNIQUE, + owner_subject text NOT NULL, + owner_name text NOT NULL, + provider text NOT NULL CHECK (provider IN ('nvidia','zai')), + model text NOT NULL, + label text NOT NULL, + sealed text, + enabled boolean NOT NULL DEFAULT false, + revision integer NOT NULL DEFAULT 0, + probe_status text NOT NULL DEFAULT 'not_checked' + CHECK (probe_status IN ('ok','invalid','rate_limited','unavailable','not_checked')), + probe_at timestamptz, + probe_ms integer, + created_at timestamptz NOT NULL DEFAULT now(), + CHECK ((source='managed' AND key_index<0 AND sealed IS NOT NULL) + OR (source='environment' AND key_index>=0 AND sealed IS NULL)) + ); + `) + .then(() => undefined) + .catch((error) => { + ready.delete(pool) + throw error + }) + ready.set(pool, promise) + } + return promise +} +const date = (value: Date | string | null): string | null => + value ? new Date(value).toISOString() : null +function publicKey(row: KeyRow): ContributorKey { + return { + id: row.key_index, + source: row.source, + provider: row.provider, + model: row.model, + label: row.label, + fingerprint: row.fingerprint.slice(0, 12), + enabled: row.enabled, + createdAt: date(row.created_at), + lastProbe: { + status: row.probe_status, + checkedAt: date(row.probe_at), + latencyMs: row.probe_ms, + }, + } +} +function environmentOwned( + key: EnvironmentKey, + subject: string, + owners: Record, +): boolean { + const github = contributorGithub(subject) + return ( + !!github && owners[key.id]?.toLowerCase() === `@${github}`.toLowerCase() + ) +} +export async function listContributorKeys( + pool: Pool, + subject: string, + environment: EnvironmentKey[], + owners: Record, +): Promise { + await ensureContributorKeys(pool) + const { rows } = await pool.query( + 'SELECT * FROM queen_contributor_keys WHERE owner_subject=$1 ORDER BY key_index', + [subject], + ) + const bindings = await pool.query< + Pick + >( + "SELECT key_index,owner_subject,fingerprint FROM queen_contributor_keys WHERE source='environment'", + ) + const out = rows.filter((row) => row.source === 'managed').map(publicKey) + for (const key of environment) { + const digest = fingerprint(key.apiKey) + const bound = bindings.rows.find((row) => row.key_index === key.id) + if ( + bound + ? bound.owner_subject !== subject || bound.fingerprint !== digest + : !environmentOwned(key, subject, owners) + ) + continue + const stored = rows.find( + (row) => + row.source === 'environment' && + row.key_index === key.id && + row.fingerprint === digest, + ) + out.push( + stored + ? publicKey(stored) + : { + id: key.id, + source: 'environment', + provider: key.provider, + model: key.model, + label: `${key.provider} #${key.id + 1}`, + fingerprint: digest.slice(0, 12), + enabled: true, + createdAt: null, + lastProbe: { + status: 'not_checked', + checkedAt: null, + latencyMs: null, + }, + }, + ) + } + return out +} + +/** No response bodies, redirects, arbitrary URLs, retries or provider messages escape. */ +let probesInFlight = 0 +export async function probeCredential( + key: Pick, + fetcher: typeof fetch = fetch, +): Promise { + const started = Date.now() + const provider = CONTRIBUTOR_PROVIDERS.find((p) => p.id === key.provider) + if (!provider) throw new ContributorError('unsupported_provider') + if (probesInFlight >= P.MAX_PROBES_IN_FLIGHT) + throw new ContributorError('probe_rate_limited', 429) + probesInFlight++ + let status: ProbeStatus = 'unavailable' + try { + const response = await fetcher(`${provider.baseUrl}/chat/completions`, { + method: 'POST', + redirect: 'error', + signal: AbortSignal.timeout(P.PROBE_TIMEOUT_MS), + headers: { + Authorization: `Bearer ${key.apiKey}`, + 'Content-Type': 'application/json', + }, + body: JSON.stringify({ + model: key.model, + messages: [{ role: 'user', content: 'Reply OK.' }], + max_tokens: P.PROBE_MAX_TOKENS, + stream: false, + }), + }) + if (response.status === 401 || response.status === 403) status = 'invalid' + else if (response.status === 429) status = 'rate_limited' + else if (response.ok) { + // A login/proxy HTML response or an empty 200 is not usable inference. + const body = (await response.json()) as { + choices?: Array<{ + message?: { content?: unknown; reasoning_content?: unknown } + }> + } + const message = body.choices?.[0]?.message + if (typeof message?.content === 'string' && message.content.trim()) + status = 'ok' + else if ( + typeof message?.reasoning_content === 'string' && + message.reasoning_content.trim() + ) + status = 'ok' + } + await response.body?.cancel().catch(() => {}) + } catch { + /* The closed result intentionally excludes errors carrying credentials. */ + } finally { + probesInFlight-- + } + return { + status, + checkedAt: new Date().toISOString(), + latencyMs: Date.now() - started, + } +} + +export async function addContributorKey( + pool: Pool, + subject: string, + input: unknown, + environment: EnvironmentKey[], + fetcher: typeof fetch = fetch, +): Promise { + if (!input || typeof input !== 'object' || Array.isArray(input)) + throw new ContributorError('invalid_key') + const value = input as Record + if ( + Object.keys(value).some( + (key) => !['provider', 'apiKey', 'label'].includes(key), + ) + ) + throw new ContributorError('invalid_key') + const provider = CONTRIBUTOR_PROVIDERS.find((p) => p.id === value.provider) + if (!provider) throw new ContributorError('unsupported_provider') + if (typeof value.apiKey !== 'string') + throw new ContributorError('invalid_key') + const secret = value.apiKey + if ( + !secret || + Buffer.byteLength(secret) > P.MAX_KEY_BYTES || + /[^!-~]/.test(secret) + ) + throw new ContributorError('invalid_key') + if (value.label !== undefined && typeof value.label !== 'string') + throw new ContributorError('invalid_label') + const label = + typeof value.label === 'string' ? value.label.trim() : provider.label + if ( + label.length > P.LABEL_LIMIT || + [...label].some( + (char) => char.charCodeAt(0) < 32 || char.charCodeAt(0) === 127, + ) + ) + throw new ContributorError('invalid_label') + const digest = fingerprint(secret) + if (environment.some((key) => fingerprint(key.apiKey) === digest)) + throw new ContributorError('key_already_connected', 409) + const sealed = sealCredential(secret, subject, digest) + await ensureContributorKeys(pool) + const duplicate = await pool.query( + 'SELECT * FROM queen_contributor_keys WHERE fingerprint=$1', + [digest], + ) + if (duplicate.rows[0]) { + if (duplicate.rows[0].owner_subject !== subject) + throw new ContributorError('key_already_connected', 409) + return publicKey(duplicate.rows[0]) + } + const checked = await probeCredential( + { provider: provider.id, model: provider.model, apiKey: secret }, + fetcher, + ) + const client = await pool.connect() + try { + await client.query('BEGIN') + await client.query('SELECT pg_advisory_xact_lock(hashtext($1))', [ + `contributor:${subject}`, + ]) + const count = await client.query( + 'SELECT count(*)::integer AS n FROM queen_contributor_keys WHERE owner_subject=$1', + [subject], + ) + if (Number(count.rows[0]?.n) >= P.MAX_KEYS_PER_OWNER) + throw new ContributorError('key_limit', 409) + const { rows } = await client.query( + `INSERT INTO queen_contributor_keys + (source,fingerprint,owner_subject,owner_name,provider,model,label,sealed,enabled,probe_status,probe_at,probe_ms) + VALUES ('managed',$1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11) RETURNING *`, + [ + digest, + subject, + contributorName(subject), + provider.id, + provider.model, + label, + sealed, + checked.status === 'ok', + checked.status, + checked.checkedAt, + checked.latencyMs, + ], + ) + await client.query('COMMIT') + return publicKey(rows[0]) + } catch (error) { + await client.query('ROLLBACK') + if ((error as { code?: string }).code === '23505') + throw new ContributorError('key_already_connected', 409) + throw error + } finally { + client.release() + } +} + +async function ownedRow( + pool: Pool, + subject: string, + id: number, + environment: EnvironmentKey[], + owners: Record, + includeSecret: boolean, +): Promise<{ row: KeyRow; secret: string }> { + await ensureContributorKeys(pool) + if (id >= 0) { + const current = environment.find((key) => key.id === id) + if (!current) throw new ContributorError('key_not_found', 404) + const digest = fingerprint(current.apiKey) + const existing = await pool.query( + 'SELECT * FROM queen_contributor_keys WHERE key_index=$1', + [id], + ) + if (existing.rows[0]) { + const row = existing.rows[0] + if (row.owner_subject !== subject) + throw new ContributorError('key_not_found', 404) + if (row.fingerprint !== digest) + throw new ContributorError('key_binding_conflict', 409) + return { row, secret: current.apiKey } + } + if (!environmentOwned(current, subject, owners)) + throw new ContributorError('key_not_found', 404) + await pool.query( + `INSERT INTO queen_contributor_keys + (key_index,source,fingerprint,owner_subject,owner_name,provider,model,label,enabled) + VALUES ($1,'environment',$2,$3,$4,$5,$6,$7,true) ON CONFLICT DO NOTHING`, + [ + id, + digest, + subject, + contributorName(subject), + current.provider, + current.model, + `${current.provider} #${id + 1}`, + ], + ) + const { rows } = await pool.query( + 'SELECT * FROM queen_contributor_keys WHERE key_index=$1 AND owner_subject=$2 AND fingerprint=$3', + [id, subject, digest], + ) + if (!rows[0]) throw new ContributorError('key_binding_conflict', 409) + return { row: rows[0], secret: current.apiKey } + } + const { rows } = await pool.query( + 'SELECT * FROM queen_contributor_keys WHERE key_index=$1 AND owner_subject=$2', + [id, subject], + ) + const row = rows[0] + if (!row || row.source !== 'managed' || !row.sealed) + throw new ContributorError('key_not_found', 404) + return { + row, + secret: includeSecret + ? openCredential(row.sealed, subject, row.fingerprint) + : '', + } +} + +export async function changeContributorKey( + pool: Pool, + subject: string, + id: number, + action: 'probe' | 'enable' | 'disable', + environment: EnvironmentKey[], + owners: Record, + fetcher: typeof fetch = fetch, +): Promise { + const { row, secret } = await ownedRow( + pool, + subject, + id, + environment, + owners, + action !== 'disable', + ) + if (action === 'disable') { + const { rows } = await pool.query( + 'UPDATE queen_contributor_keys SET enabled=false,revision=revision+1 WHERE key_index=$1 AND owner_subject=$2 RETURNING *', + [id, subject], + ) + return publicKey(rows[0]) + } + // Database claim, so multiple tabs/processes cannot fan out probes on one key. + const claimed = await pool.query( + `UPDATE queen_contributor_keys SET probe_at=now() + WHERE key_index=$1 AND owner_subject=$2 + AND (probe_at IS NULL OR probe_at < now() - ($3::integer * interval '1 second')) + RETURNING key_index,revision`, + [id, subject, P.PROBE_COOLDOWN_SECONDS], + ) + if (!claimed.rows.length) + throw new ContributorError('probe_rate_limited', 429) + const checked = await probeCredential( + { provider: row.provider, model: row.model, apiKey: secret }, + fetcher, + ) + const { rows } = await pool.query( + `UPDATE queen_contributor_keys + SET probe_status=$3,probe_at=$4,probe_ms=$5, + enabled=CASE WHEN $3='invalid' THEN false WHEN $6 AND $3='ok' AND revision=$7 THEN true ELSE enabled END + WHERE key_index=$1 AND owner_subject=$2 RETURNING *`, + [ + id, + subject, + checked.status, + checked.checkedAt, + checked.latencyMs, + action === 'enable', + claimed.rows[0].revision, + ], + ) + return publicKey(rows[0]) +} + +export interface ContributorRuntime { + managed: EnvironmentKey[] + disabled: number[] +} +async function contributorRegistryAvailable(pool: Pool): Promise { + if (contributorsEnabled()) { + await ensureContributorKeys(pool) + return true + } + // The management capability gates requests, never the already-persisted + // consent ledger. Revoking that capability cannot revive a disabled key. + const result = await pool.query<{ registry: string | null }>( + "SELECT to_regclass('queen_contributor_keys')::text AS registry", + ) + return !!result.rows[0]?.registry +} +export async function contributorRuntime( + pool: Pool, + environment: EnvironmentKey[], +): Promise { + if (!(await contributorRegistryAvailable(pool))) + return { managed: [], disabled: [] } + const { rows } = await pool.query( + 'SELECT * FROM queen_contributor_keys', + ) + const disabled: number[] = [] + const managed: EnvironmentKey[] = [] + const digests = new Set(environment.map((key) => fingerprint(key.apiKey))) + for (const row of rows) { + if (row.source === 'environment') { + const current = environment.find((key) => key.id === row.key_index) + if ( + !row.enabled && + current && + fingerprint(current.apiKey) === row.fingerprint + ) + disabled.push(row.key_index) + } else if ( + row.enabled && + row.probe_status !== 'invalid' && + row.probe_status !== 'not_checked' && + row.sealed && + !digests.has(row.fingerprint) + ) { + const provider = CONTRIBUTOR_PROVIDERS.find((p) => p.id === row.provider) + if (!provider) continue + let secret: string + try { + secret = openCredential(row.sealed, row.owner_subject, row.fingerprint) + } catch { + // One unreadable managed key must not revive disabled keys or stop + // independent environment credentials. Never log the cipher/error. + logger.warn('Contributor credential unavailable for allocation', { + keyIndex: row.key_index, + }) + continue + } + managed.push({ + id: row.key_index, + provider: row.provider, + model: row.model, + apiKey: secret, + baseUrl: provider.baseUrl, + contextWindow: 65536, + }) + digests.add(row.fingerprint) + } + } + return { managed, disabled } +} +export async function contributorOwnerNames( + pool: Pool, +): Promise> { + if (!(await contributorRegistryAvailable(pool))) return {} + const { rows } = await pool.query>( + 'SELECT key_index,owner_name FROM queen_contributor_keys', + ) + return Object.fromEntries(rows.map((row) => [row.key_index, row.owner_name])) +} diff --git a/trios/agent-server/apps/server/src/api/services/queen-contributor-policy.gen.ts b/trios/agent-server/apps/server/src/api/services/queen-contributor-policy.gen.ts new file mode 100644 index 0000000000..6e522e1f02 --- /dev/null +++ b/trios/agent-server/apps/server/src/api/services/queen-contributor-policy.gen.ts @@ -0,0 +1,59 @@ +// Generated by `t27c gen-ts` from queen-contributor-keys.t27. Do not edit. +// +// Every name, value and type here comes from that spec. Change the spec and +// regenerate; an edit made here is lost the next time anyone builds. + +export const KIND = "automation" satisfies string; +export const ID = "queen-contributor-keys" satisfies string; +export const REPO = "BrowserOS" satisfies string; +export const VERSION = 1 satisfies number; +export const PROBE_TIMEOUT_MS = 15000 satisfies number; +export const PROBE_MAX_TOKENS = 16 satisfies number; +export const MAX_PROBES_IN_FLIGHT = 4 satisfies number; +export const PROBE_COOLDOWN_SECONDS = 60 satisfies number; +export const MAX_KEYS_PER_OWNER = 100 satisfies number; +export const MAX_KEY_BYTES = 4096 satisfies number; +export const MAX_BODY_BYTES = 8192 satisfies number; +export const LABEL_LIMIT = 80 satisfies number; +export const PROXY_TOKEN_MIN_BYTES = 32 satisfies number; +export const OWNER_REASSIGNMENT = false satisfies boolean; +export const CREDIT_WALLET = false satisfies boolean; +export const ENVIRONMENT_INDICES_COMPACTED = false satisfies boolean; +export const NVIDIA_URL = "https://integrate.api.nvidia.com/v1" satisfies string; +export const NVIDIA_MODEL = "nvidia/nemotron-3-ultra-550b-a55b" satisfies string; +export const ZAI_URL = "https://api.z.ai/api/paas/v4" satisfies string; +export const ZAI_MODEL = "glm-4.5-flash" satisfies string; +// t27c gen-ts: fn trusted_proxy was not emitted -- this backend lowers declarations, not bodies. +// t27c gen-ts: fn source_matches_index was not emitted -- this backend lowers declarations, not bodies. +// t27c gen-ts: fn may_add_key was not emitted -- this backend lowers declarations, not bodies. +// t27c gen-ts: fn may_start_probe was not emitted -- this backend lowers declarations, not bodies. +// t27c gen-ts: fn cooldown_passed was not emitted -- this backend lowers declarations, not bodies. +// t27c gen-ts: fn credential_fits was not emitted -- this backend lowers declarations, not bodies. +// t27c gen-ts: fn same_binding was not emitted -- this backend lowers declarations, not bodies. +// t27c gen-ts: fn enabled_after_probe was not emitted -- this backend lowers declarations, not bodies. +// t27c gen-ts: fn wallet_credit was not emitted -- this backend lowers declarations, not bodies. +// t27c gen-ts: a TestBlock was not emitted -- it is checked by the compiler, not by the artifact. +// t27c gen-ts: a TestBlock was not emitted -- it is checked by the compiler, not by the artifact. +// t27c gen-ts: a TestBlock was not emitted -- it is checked by the compiler, not by the artifact. +// t27c gen-ts: a TestBlock was not emitted -- it is checked by the compiler, not by the artifact. +// t27c gen-ts: a TestBlock was not emitted -- it is checked by the compiler, not by the artifact. +// t27c gen-ts: a TestBlock was not emitted -- it is checked by the compiler, not by the artifact. +// t27c gen-ts: a TestBlock was not emitted -- it is checked by the compiler, not by the artifact. +// t27c gen-ts: a TestBlock was not emitted -- it is checked by the compiler, not by the artifact. +// t27c gen-ts: a TestBlock was not emitted -- it is checked by the compiler, not by the artifact. + +// Declaration order, which the spec's own laws depend on. +export const __STRUCT_ORDER__ = [] as const; +export const __DECL_ORDER__ = ["KIND", "ID", "REPO", "VERSION", "PROBE_TIMEOUT_MS", "PROBE_MAX_TOKENS", "MAX_PROBES_IN_FLIGHT", "PROBE_COOLDOWN_SECONDS", "MAX_KEYS_PER_OWNER", "MAX_KEY_BYTES", "MAX_BODY_BYTES", "LABEL_LIMIT", "PROXY_TOKEN_MIN_BYTES", "OWNER_REASSIGNMENT", "CREDIT_WALLET", "ENVIRONMENT_INDICES_COMPACTED", "NVIDIA_URL", "NVIDIA_MODEL", "ZAI_URL", "ZAI_MODEL"] as const; + +// What this spec holds and this backend did not print. Empty is the whole +// story most of the time; an entry here is a promise the artifact does not +// keep, and reading it is how a tool tells a partial module from a complete +// one without parsing comments. +// +// A LIST, not a map keyed by name. `specs/numeric/formats.t27` declares four +// separate consts called `result`, one per test block, and under an object +// literal three of the four omissions vanished into the fourth -- a record of +// what went missing that itself went missing. A spec is free to reuse a +// name; this file is not free to lose the second one. +export const __NOT_EMITTED__ = Object.freeze([] as const); diff --git a/trios/agent-server/apps/server/src/api/services/queen-contributor-policy.ts b/trios/agent-server/apps/server/src/api/services/queen-contributor-policy.ts new file mode 100644 index 0000000000..b63aea53be --- /dev/null +++ b/trios/agent-server/apps/server/src/api/services/queen-contributor-policy.ts @@ -0,0 +1,4 @@ +// Numeric rules come from the canonical t27 specification, generated by t27c. +import * as policy from './queen-contributor-policy.gen' + +export const CONTRIBUTOR_POLICY = policy 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..a6e65d7b43 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 @@ -29,6 +29,11 @@ import { statSync } from 'node:fs' import type { Pool } from 'pg' import { logger } from '../../lib/logger' import { shellArgv } from '../../tools/filesystem/bash' +import { + type ContributorRuntime, + contributorRuntime, + type EnvironmentKey, +} from './queen-contributor-keys' import { describeReading, diskLineUsedPercent, @@ -414,7 +419,20 @@ export function reviewExtraLanesPerCredential( return Math.min(parsed, 2) } -export function workerCapacityBreakdown(): WorkerCapacityBreakdown { +export function workerCapacityBreakdown( + runtime?: ContributorRuntime, +): WorkerCapacityBreakdown { + if (runtime && (runtime.managed.length || runtime.disabled.length)) { + const lanes = contributorWorkerCandidates(runtime) + return { + connectedCredentials: lanes.length, + lanesPerCredential: configuredRemoteLanesPerCredential(), + effectiveCapacity: Math.min( + lanes.reduce((sum, lane) => sum + (lane.laneCount ?? 1), 0), + queenWorkerLimit(), + ), + } + } const endpoint = configuredWorkerBaseUrl() if (endpoint) { // An explicitly configured Ollama is one measured inference server even @@ -478,6 +496,14 @@ export function configuredWorkerCapacity(): number { return workerCapacityBreakdown().effectiveCapacity } +export async function liveWorkerCapacity( + pool: Pool, +): Promise { + return workerCapacityBreakdown( + await contributorRuntime(pool, environmentContributorKeys()), + ) +} + /** * A worker endpoint named by URL. It can be a keyless local Ollama or a remote * OpenAI-compatible API backed by the generic worker-key pool. @@ -808,13 +834,21 @@ function endpointPoolProvider( * The first pool's URL and keys, for the model probes (lib/model-ranking.ts). * Undefined for a local Ollama: it serves one model and has nothing to rank. */ -export function workerProbeEndpoint(): - | { baseUrl: string; keys: string[] } - | undefined { +export function workerProbeEndpoint( + runtime?: ContributorRuntime, +): { baseUrl: string; keys: string[] } | undefined { const pool = configuredEndpointPools()[0] if (!pool || pool.provider === 'ollama' || pool.keys.length === 0) return undefined - return { baseUrl: pool.baseUrl, keys: pool.keys } + const keys = pool.keys.filter( + (_, index) => !runtime?.disabled.includes(index), + ) + keys.push( + ...(runtime?.managed + .filter((key) => key.baseUrl === pool.baseUrl) + .map((key) => key.apiKey) ?? []), + ) + return keys.length ? { baseUrl: pool.baseUrl, keys } : undefined } function configuredEndpointProvider( @@ -901,7 +935,7 @@ function configuredEndpointProvider( * least-loaded credential. The credential index is stored with the dispatch, * so repeated 429s remain attributable without publishing a secret. */ -export function resolveWorkerProvider( +function resolveEnvironmentWorkerProvider( takenKeyIndices: number[] = [], afterKeyIndex?: number, ): WorkerProvider | null { @@ -971,6 +1005,122 @@ export function resolveWorkerProvider( return null } +/** The environment's existing durable positions are assigned BEFORE filtering. */ +function environmentWorkerCandidates(): WorkerProvider[] { + if (configuredWorkerBaseUrl()) { + return configuredEndpointPools().flatMap((pool) => + (pool.provider === 'ollama' ? [pool.keys[0] || 'local'] : pool.keys).map( + (apiKey, position) => ({ + provider: pool.provider, + model: pool.model, + baseUrl: pool.baseUrl, + apiKey, + keyIndex: (pool.number - 1) * POOL_KEY_STRIDE + position, + keyCount: pool.keys.length, + poolNumber: pool.number, + contextWindow: pool.contextWindow, + laneCount: + pool.provider === 'ollama' + ? 1 + : configuredRemoteLanesPerCredential(), + }), + ), + ) + } + for (const provider of WORKER_PROVIDERS) { + const keys = keysFor(provider.envVar) + if (!keys.length) continue + return keys.map((apiKey, keyIndex) => ({ + provider: provider.provider, + model: process.env.TRIOS_QUEEN_WORKER_MODEL || provider.model, + apiKey, + keyIndex, + keyCount: keys.length, + laneCount: workerLanesFor(provider.provider), + })) + } + return [] +} + +export function environmentContributorKeys(): EnvironmentKey[] { + return environmentWorkerCandidates().flatMap((key) => { + const provider = + key.baseUrl === 'https://integrate.api.nvidia.com/v1' + ? 'nvidia' + : key.baseUrl === 'https://api.z.ai/api/paas/v4' || + key.provider === 'zai' + ? 'zai' + : undefined + if (!provider || key.keyIndex === undefined || !key.apiKey) return [] + return [ + { + id: key.keyIndex, + provider, + apiKey: key.apiKey, + model: key.model, + baseUrl: key.baseUrl ?? 'https://api.z.ai/api/paas/v4', + contextWindow: key.contextWindow, + }, + ] + }) +} + +function contributorWorkerCandidates( + runtime: ContributorRuntime, +): WorkerProvider[] { + const disabled = new Set(runtime.disabled) + return [ + ...environmentWorkerCandidates().filter( + (key) => key.keyIndex !== undefined && !disabled.has(key.keyIndex), + ), + ...runtime.managed.map((key) => ({ + provider: key.provider === 'zai' ? 'zai' : 'openai-compatible', + model: key.model, + baseUrl: key.baseUrl, + apiKey: key.apiKey, + keyIndex: key.id, + keyCount: runtime.managed.length, + contextWindow: key.contextWindow, + laneCount: configuredRemoteLanesPerCredential(), + })), + ] +} + +export function resolveWorkerProvider( + takenKeyIndices: number[] = [], + afterKeyIndex?: number, + runtime?: ContributorRuntime, +): WorkerProvider | null { + if (!runtime || (!runtime.managed.length && !runtime.disabled.length)) { + return resolveEnvironmentWorkerProvider(takenKeyIndices, afterKeyIndex) + } + const candidates = contributorWorkerCandidates(runtime) + if (!candidates.length) return null + const cursor = candidates.findIndex((key) => key.keyIndex === afterKeyIndex) + let chosen: WorkerProvider | undefined + let smallest = Number.POSITIVE_INFINITY + for (let step = 1; step <= candidates.length; step++) { + const candidate = candidates[(cursor + step) % candidates.length] + const occupancy = takenKeyIndices.filter( + (index) => index === candidate.keyIndex, + ).length + if (occupancy < (candidate.laneCount ?? 1) && occupancy < smallest) { + chosen = { ...candidate, laneIndex: occupancy } + smallest = occupancy + } + } + return ( + chosen ?? { + provider: candidates[0].provider, + model: candidates[0].model, + exhausted: candidates.reduce( + (sum, lane) => sum + (lane.laneCount ?? 1), + 0, + ), + } + ) +} + /** * Every credential lane a one-shot REVIEW may use right now, across every * connected pool and every legacy provider that holds a key. @@ -992,7 +1142,34 @@ export function resolveWorkerProvider( */ export function reviewLaneCandidates( takenKeyIndices: number[] = [], + runtime?: ContributorRuntime, ): WorkerProvider[] { + if (runtime && (runtime.managed.length || runtime.disabled.length)) { + const managedSecrets = new Set(runtime.managed.map((key) => key.apiKey)) + const environment = reviewLaneCandidates(takenKeyIndices).filter( + (candidate) => + !managedSecrets.has(candidate.apiKey ?? '') && + (candidate.keyIndex === undefined || + !runtime.disabled.includes(candidate.keyIndex)), + ) + const managed = contributorWorkerCandidates(runtime) + .filter( + (candidate) => + candidate.keyIndex !== undefined && candidate.keyIndex < 0, + ) + .flatMap((candidate) => { + const busy = takenKeyIndices.filter( + (id) => id === candidate.keyIndex, + ).length + const laneCount = + (candidate.laneCount ?? 1) + + reviewExtraLanesPerCredential(candidate.provider) + return busy < laneCount + ? [{ ...candidate, laneIndex: busy, laneCount }] + : [] + }) + return [...environment, ...managed] + } const busy = (index: number) => takenKeyIndices.filter((taken) => taken === index).length const out: WorkerProvider[] = [] @@ -4206,7 +4383,8 @@ export async function dispatchBee( ): Promise { const branch = `queen-${issue}` - const chosen = resolveWorkerProvider(takenKeyIndices, afterKeyIndex) + const runtime = await contributorRuntime(pool, environmentContributorKeys()) + const chosen = resolveWorkerProvider(takenKeyIndices, afterKeyIndex, runtime) if (chosen?.exhausted !== undefined) { // Not a missing credential: every key this deployment has is already // carrying a bee. Named separately because the fix is different - one more @@ -4322,7 +4500,9 @@ export async function dispatchBee( ? ` pool ${chosen.poolNumber}` : '') + (chosen.keyCount && chosen.keyCount > 1 - ? ` key ${((chosen.keyIndex ?? 0) % POOL_KEY_STRIDE) + 1}/${chosen.keyCount}` + ? (chosen.keyIndex ?? 0) < 0 + ? ` contributor key ${chosen.keyIndex}` + : ` key ${((chosen.keyIndex ?? 0) % POOL_KEY_STRIDE) + 1}/${chosen.keyCount}` : '') + (chosen.laneCount && chosen.laneCount > 1 ? ` lane ${(chosen.laneIndex ?? 0) + 1}/${chosen.laneCount}` diff --git a/trios/agent-server/apps/server/src/api/services/queen-leaderboard.ts b/trios/agent-server/apps/server/src/api/services/queen-leaderboard.ts index 665a03ca0f..9d0da933d7 100644 --- a/trios/agent-server/apps/server/src/api/services/queen-leaderboard.ts +++ b/trios/agent-server/apps/server/src/api/services/queen-leaderboard.ts @@ -39,6 +39,7 @@ * leaves the environment, and only its index appears in the database. */ import type { Pool } from 'pg' +import { contributorOwnerNames } from './queen-contributor-keys' /** An issue the Queen accepted, on this key. */ export const ACCEPTED_XP = 100 @@ -203,7 +204,7 @@ export async function keyWork( (snapshot->>'finished_at')::timestamptz, coalesce(snapshot->'owned_paths', '[]'::jsonb) FROM queen_dispatch_history - WHERE snapshot->>'key_index' ~ '^[0-9]+$' + WHERE snapshot->>'key_index' ~ '^-?[0-9]+$' ${windowed ? "AND archived_at > now() - ($1::integer * interval '1 day')" : ''} ) SELECT key_index, @@ -259,6 +260,9 @@ export async function leaderboard( days, measuredAt: new Date().toISOString(), scoring: { acceptedXp: ACCEPTED_XP, specXp: SPEC_XP, hourXp: HOUR_XP }, - contributors: rank(work, parseOwners(process.env.TRIOS_KEY_OWNERS)), + contributors: rank(work, { + ...parseOwners(process.env.TRIOS_KEY_OWNERS), + ...(await contributorOwnerNames(pool)), + }), } } diff --git a/trios/agent-server/apps/server/src/api/services/queen-reviewer.ts b/trios/agent-server/apps/server/src/api/services/queen-reviewer.ts index 712ab20a0f..5217fd79f1 100644 --- a/trios/agent-server/apps/server/src/api/services/queen-reviewer.ts +++ b/trios/agent-server/apps/server/src/api/services/queen-reviewer.ts @@ -778,7 +778,9 @@ export interface ReviewDeps { criteria: string[], baseSha?: string | null, ) => Promise - laneCandidates: (takenKeyIndices: number[]) => WorkerProvider[] + laneCandidates: ( + takenKeyIndices: number[], + ) => WorkerProvider[] | Promise llm: ( lane: WorkerProvider, system: string, diff --git a/trios/agent-server/apps/server/src/api/services/queen-tick.ts b/trios/agent-server/apps/server/src/api/services/queen-tick.ts index d543a4cceb..2fc10388f8 100644 --- a/trios/agent-server/apps/server/src/api/services/queen-tick.ts +++ b/trios/agent-server/apps/server/src/api/services/queen-tick.ts @@ -43,6 +43,7 @@ import { logger } from '../../lib/logger' import { startModelProbes, workerModelRanking } from '../../lib/model-ranking' import { outstandingEscalations } from '../routes/queen-needs-you' import { githubCiDeps, takeBackRefusedAcceptances } from './queen-ci-verdict' +import { contributorRuntime } from './queen-contributor-keys' import { type CriterionRun, criteriaCounts, @@ -55,8 +56,10 @@ import { baseRef, DISPATCH_OUTCOME_LABELS, dispatchBee, + environmentContributorKeys, reapDispatchesFromPreviousBoot, reapStalledDispatches, + reviewLaneCandidates, setDurableCloseListener, type Witness, type WorkerProvider, @@ -122,7 +125,7 @@ export function latestProviderKeyIndex( if ( typeof index === 'number' && Number.isInteger(index) && - index >= 0 && + index >= -2147483647 && Number.isFinite(at) && at > latestAt ) { @@ -2622,7 +2625,15 @@ export async function reviewFinishedDispatches( pool: Pool, overrides: Partial = {}, ): Promise { - const deps: ReviewDeps = { ...defaultReviewDeps(), ...overrides } + const deps: ReviewDeps = { + ...defaultReviewDeps(), + laneCandidates: async (taken) => + reviewLaneCandidates( + taken, + await contributorRuntime(pool, environmentContributorKeys()), + ), + ...overrides, + } // TWO THINGS ABOUT THIS QUERY, BOTH MEASURED ON 2026-09-03. // // The say rows are joined with NOTHING between them, not a newline. The @@ -3197,15 +3208,14 @@ export async function reviewFinishedDispatches( // review, unlike `markReviewerLaneFailed`, which is a half-hour // backoff and belongs only to a lane that is broken rather than busy. const triedLanes = new Set() - const pick = () => + const pick = async () => chooseReviewerLane( - deps - .laneCandidates(taken) + (await deps.laneCandidates(taken)) .filter((lane) => !reviewerLaneBackedOff(lane)) .filter((lane) => !triedLanes.has(reviewerLaneKey(lane))), bee, ) - let choice = pick() + let choice = await pick() if (!choice) { reviewerSkipped = 'no reviewer lane is free' } else { @@ -3282,7 +3292,7 @@ export async function reviewFinishedDispatches( // would turn a provider's bad minute into the Queen's bad // half-hour. if (!answer.transient) markReviewerLaneFailed(lane) - choice = pick() + choice = await pick() continue } // THE COMPILER'S LINES ONLY, as `reviewerMessage` is given @@ -4394,6 +4404,10 @@ export function startQueenTick(): void { const probeEndpoint = workerProbeEndpoint() if (ranking && probeEndpoint) { startModelProbes(ranking, probeEndpoint, { + resolveEndpoint: async () => + workerProbeEndpoint( + await contributorRuntime(pool, environmentContributorKeys()), + ), onRound: (snapshot, chosen) => logger.info('Queen worker models ranked', { chosen, diff --git a/trios/agent-server/apps/server/src/lib/model-ranking.ts b/trios/agent-server/apps/server/src/lib/model-ranking.ts index 6dc88deb9f..c362d3fd13 100644 --- a/trios/agent-server/apps/server/src/lib/model-ranking.ts +++ b/trios/agent-server/apps/server/src/lib/model-ranking.ts @@ -516,6 +516,10 @@ export function startModelProbes( options: { intervalMs?: number fetchImpl?: typeof fetch + /** Re-read consent before each new probe; disabling a key survives this timer. */ + resolveEndpoint?: () => Promise< + { baseUrl: string; keys: string[] } | undefined + > onRound?: (snapshot: RankedModel[], chosen: string) => void } = {}, ): () => void { @@ -529,11 +533,15 @@ export function startModelProbes( let probed = 0 for (const model of ranking.candidates) { if (!ranking.dueForProbe(model)) continue + const current = options.resolveEndpoint + ? await options.resolveEndpoint().catch(() => undefined) + : endpoint + if (!current?.keys.length) continue probed++ - const key = endpoint.keys[keyCursor++ % endpoint.keys.length] + const key = current.keys[keyCursor++ % current.keys.length] ranking.recordProbe( model, - await probeModel(endpoint.baseUrl, key, model, { + await probeModel(current.baseUrl, key, model, { fetchImpl: options.fetchImpl, }), ) diff --git a/trios/agent-server/apps/server/tests/api/queen-contributor-keys.test.ts b/trios/agent-server/apps/server/tests/api/queen-contributor-keys.test.ts new file mode 100644 index 0000000000..dca1e6c967 --- /dev/null +++ b/trios/agent-server/apps/server/tests/api/queen-contributor-keys.test.ts @@ -0,0 +1,357 @@ +import { afterEach, beforeEach, describe, expect, it } from 'bun:test' +import { readFile } from 'node:fs/promises' +import { Hono } from 'hono' +import type { Pool } from 'pg' +import { createQueenContributorKeysRoute } from '../../src/api/routes/queen-contributor-keys' +import { + addContributorKey, + contributorGithub, + fingerprint, + openCredential, + probeCredential, + sealCredential, + trustedContributor, +} from '../../src/api/services/queen-contributor-keys' +import { CONTRIBUTOR_POLICY } from '../../src/api/services/queen-contributor-policy' +import { + environmentContributorKeys, + resolveWorkerProvider, + reviewLaneCandidates, + workerCapacityBreakdown, + workerProbeEndpoint, +} from '../../src/api/services/queen-dispatch' +import { latestProviderKeyIndex } from '../../src/api/services/queen-tick' +import { loadCompiler } from '../../src/inngest/t27-consts' +import { ModelRanking, startModelProbes } from '../../src/lib/model-ranking' + +const subject = 'telegram:111222333' +const token = 'test-proxy-capability-at-least-32-bytes' +const names = [ + 'QUEEN_CONTRIBUTOR_PROXY_TOKEN', + 'QUEEN_CONTRIBUTOR_ENCRYPTION_KEY', + 'QUEEN_CONTRIBUTOR_IDENTITIES', + 'TRIOS_QUEEN_WORKER_BASE_URL', + 'TRIOS_QUEEN_WORKER_PROVIDER', + 'TRIOS_QUEEN_WORKER_MODEL', + 'TRIOS_QUEEN_WORKER_API_KEY', + 'TRIOS_QUEEN_WORKER_API_KEY_2', + 'TRIOS_QUEEN_WORKER_API_KEY_3', + 'TRIOS_QUEEN_REMOTE_CONCURRENCY_PER_KEY', + 'ZAI_API_KEY', + 'ANTHROPIC_API_KEY', +] +const before = new Map() +beforeEach(() => { + for (const name of names) before.set(name, process.env[name]) + process.env.QUEEN_CONTRIBUTOR_PROXY_TOKEN = token + process.env.QUEEN_CONTRIBUTOR_ENCRYPTION_KEY = Buffer.alloc(32, 37).toString( + 'base64', + ) + process.env.QUEEN_CONTRIBUTOR_IDENTITIES = JSON.stringify({ + [subject]: 'dmitrii-f-t27', + }) +}) +afterEach(() => { + for (const name of names) { + const value = before.get(name) + if (value === undefined) delete process.env[name] + else process.env[name] = value + } +}) + +describe('contributor capability and vault', () => { + it('refuses raw secret whitespace, client ownership claims and malformed labels before any database access', async () => { + const pool = {} as Pool + for (const input of [ + { provider: 'zai', apiKey: 'secret\n' }, + { provider: 'zai', apiKey: 'secret', github: 'dmitrii-f-t27' }, + { provider: 'zai', apiKey: 'secret', label: 42 }, + { provider: 'zai', apiKey: 'secret', label: 'a'.repeat(81) }, + ]) + await expect( + addContributorKey(pool, subject, input, []), + ).rejects.toThrow() + }) + it('requires its separate capability, a verified subject and configured service', () => { + expect(trustedContributor(`Bearer ${token}`, subject)).toBe(subject) + expect(() => + trustedContributor('Bearer operator-or-app-token', subject), + ).toThrow('forbidden') + expect(() => trustedContributor(undefined, subject)).toThrow('forbidden') + expect(() => + trustedContributor(`Bearer ${token}`, '@dmitrii-f-t27'), + ).toThrow('invalid_contributor') + expect(() => trustedContributor(`Bearer ${token}`, `${subject}\n`)).toThrow( + 'invalid_contributor', + ) + delete process.env.QUEEN_CONTRIBUTOR_PROXY_TOKEN + expect(() => trustedContributor(`Bearer ${token}`, subject)).toThrow( + 'unavailable', + ) + }) + it('only the operator map can assert a GitHub login', () => { + expect(contributorGithub(subject)).toBe('dmitrii-f-t27') + expect(contributorGithub('telegram:999')).toBeUndefined() + process.env.QUEEN_CONTRIBUTOR_IDENTITIES = '{bad' + expect(contributorGithub(subject)).toBeUndefined() + }) + it('uses randomized authenticated encryption bound to owner and fingerprint', () => { + const secret = 'provider-test-secret' + const digest = fingerprint(secret) + const first = sealCredential(secret, subject, digest) + expect(first).not.toContain(secret) + expect(sealCredential(secret, subject, digest)).not.toBe(first) + expect(openCredential(first, subject, digest)).toBe(secret) + expect(() => openCredential(first, 'telegram:999', digest)).toThrow( + 'unavailable', + ) + expect(() => + openCredential(first, subject, fingerprint('different')), + ).toThrow('unavailable') + const tampered = first.split('.') + tampered[2] = Buffer.from('tampered').toString('base64') + expect(() => openCredential(tampered.join('.'), subject, digest)).toThrow( + 'unavailable', + ) + process.env.QUEEN_CONTRIBUTOR_ENCRYPTION_KEY = 'short' + expect(() => sealCredential(secret, subject, digest)).toThrow('unavailable') + }) + it('rejects unauthenticated and malformed requests before touching storage', async () => { + let queries = 0 + const app = new Hono().route( + '/queen/contributor-keys', + createQueenContributorKeysRoute({ + pool: () => { + queries++ + throw new Error('must not read') + }, + }), + ) + expect((await app.request('/queen/contributor-keys')).status).toBe(403) + const headers = { + Authorization: `Bearer ${token}`, + 'X-Queen-Contributor-Id': subject, + } + expect( + ( + await app.request('/queen/contributor-keys/-1/reassign', { + method: 'POST', + headers, + }) + ).status, + ).toBe(400) + expect( + ( + await app.request('/queen/contributor-keys/1e3/probe', { + method: 'POST', + headers, + }) + ).status, + ).toBe(400) + expect( + ( + await app.request('/queen/contributor-keys', { + method: 'POST', + headers, + body: 'a'.repeat(9000), + }) + ).status, + ).toBe(413) + expect(queries).toBe(0) + }) +}) + +describe('bounded provider evidence', () => { + it('the ranking timer rereads consent and never reuses a disabled startup credential', async () => { + const calls: string[] = [] + let reads = 0 + let completed!: () => void + const done = new Promise((resolve) => { + completed = resolve + }) + const stop = startModelProbes( + new ModelRanking(['first', 'second']), + { baseUrl: 'https://example.invalid', keys: ['disabled-startup-key'] }, + { + resolveEndpoint: async () => + ++reads === 1 + ? { baseUrl: 'https://example.invalid', keys: ['enabled-key'] } + : undefined, + fetchImpl: (async (_url, init) => { + calls.push(new Headers(init?.headers).get('authorization') ?? '') + return Response.json({ + usage: { completion_tokens: 10 }, + choices: [{ message: { content: 'OK', tool_calls: [{}] } }], + }) + }) as typeof fetch, + onRound: () => completed(), + }, + ) + try { + await done + } finally { + stop() + } + expect(reads).toBe(2) + expect(calls.length).toBeGreaterThan(0) + expect(calls.every((value) => value === 'Bearer enabled-key')).toBe(true) + }) + const key = { + provider: 'nvidia' as const, + model: 'test-model', + apiKey: 'never-reflect-this-secret', + } + it('uses a fixed endpoint, timeout, bounded generation and refuses redirects', async () => { + const result = await probeCredential(key, (async (url, init) => { + expect(url).toBe(`${CONTRIBUTOR_POLICY.NVIDIA_URL}/chat/completions`) + expect(init?.redirect).toBe('error') + expect(init?.signal).toBeInstanceOf(AbortSignal) + expect(JSON.parse(String(init?.body)).max_tokens).toBe( + CONTRIBUTOR_POLICY.PROBE_MAX_TOKENS, + ) + return Response.json({ choices: [{ message: { content: 'OK' } }] }) + }) as typeof fetch) + expect(result.status).toBe('ok') + expect(result.checkedAt).not.toBeNull() + }) + it('distinguishes invalid, throttled, unavailable and unusable 200 without leaking errors', async () => { + for (const [status, expected] of [ + [401, 'invalid'], + [403, 'invalid'], + [429, 'rate_limited'], + [502, 'unavailable'], + ] as const) { + const result = await probeCredential( + key, + (async () => new Response(key.apiKey, { status })) as typeof fetch, + ) + expect(result.status).toBe(expected) + expect(JSON.stringify(result)).not.toContain(key.apiKey) + } + expect( + ( + await probeCredential(key, (async () => + Response.json({})) as typeof fetch) + ).status, + ).toBe('unavailable') + const failed = await probeCredential(key, (async () => { + throw new Error(key.apiKey) + }) as typeof fetch) + expect(JSON.stringify(failed)).not.toContain(key.apiKey) + }) +}) + +describe('stable allocation and attribution', () => { + it('keeps keyless Ollama and legacy review fallback when managed keys are connected', () => { + const runtime = { + disabled: [], + managed: [ + { + id: -1, + provider: 'nvidia' as const, + model: 'managed-model', + apiKey: 'managed', + baseUrl: CONTRIBUTOR_POLICY.NVIDIA_URL, + }, + ], + } + process.env.TRIOS_QUEEN_WORKER_BASE_URL = 'http://127.0.0.1:11434' + process.env.TRIOS_QUEEN_WORKER_PROVIDER = 'ollama' + delete process.env.TRIOS_QUEEN_WORKER_API_KEY + delete process.env.TRIOS_QUEEN_WORKER_API_KEY_2 + delete process.env.TRIOS_QUEEN_WORKER_API_KEY_3 + expect(resolveWorkerProvider([], undefined, runtime)).toMatchObject({ + provider: 'ollama', + keyIndex: 0, + apiKey: 'local', + }) + expect(workerCapacityBreakdown(runtime).connectedCredentials).toBe(2) + expect( + reviewLaneCandidates([], runtime).map((key) => key.provider), + ).toEqual(['ollama', 'openai-compatible']) + delete process.env.TRIOS_QUEEN_WORKER_BASE_URL + delete process.env.TRIOS_QUEEN_WORKER_PROVIDER + process.env.ZAI_API_KEY = 'zai-test' + process.env.ANTHROPIC_API_KEY = 'anthropic-test' + expect( + reviewLaneCandidates([], { ...runtime, disabled: [0] }).map((key) => [ + key.provider, + key.keyIndex, + ]), + ).toEqual([ + ['anthropic', undefined], + ['openai-compatible', -1], + ]) + process.env.ANTHROPIC_API_KEY = 'managed' + expect( + reviewLaneCandidates([-1, -1], { ...runtime, disabled: [0] }), + ).toEqual([]) + }) + it('disabling an environment key preserves the remaining historical indices', () => { + process.env.TRIOS_QUEEN_WORKER_BASE_URL = CONTRIBUTOR_POLICY.NVIDIA_URL + process.env.TRIOS_QUEEN_WORKER_PROVIDER = 'openai-compatible' + process.env.TRIOS_QUEEN_WORKER_MODEL = 'measured-model' + process.env.TRIOS_QUEEN_WORKER_API_KEY = 'key-a' + process.env.TRIOS_QUEEN_WORKER_API_KEY_2 = 'key-b' + process.env.TRIOS_QUEEN_WORKER_API_KEY_3 = 'key-c' + const roster = environmentContributorKeys() + expect(roster.map((key) => key.id)).toEqual([0, 1, 2]) + const runtime = { managed: [], disabled: [0, 1] } + expect(resolveWorkerProvider([], undefined, runtime)).toMatchObject({ + keyIndex: 2, + apiKey: 'key-c', + }) + expect(workerCapacityBreakdown(runtime).connectedCredentials).toBe(1) + expect(environmentContributorKeys().map((key) => key.id)).toEqual([0, 1, 2]) + expect( + workerProbeEndpoint({ managed: [], disabled: [0, 1, 2] }), + ).toBeUndefined() + }) + it('managed IDs remain negative in round-robin and the next persisted cursor', () => { + process.env.TRIOS_QUEEN_WORKER_BASE_URL = CONTRIBUTOR_POLICY.NVIDIA_URL + process.env.TRIOS_QUEEN_WORKER_PROVIDER = 'openai-compatible' + process.env.TRIOS_QUEEN_WORKER_API_KEY = 'key-a' + delete process.env.TRIOS_QUEEN_WORKER_API_KEY_2 + delete process.env.TRIOS_QUEEN_WORKER_API_KEY_3 + const runtime = { + disabled: [], + managed: [ + { + id: -9, + provider: 'nvidia' as const, + model: 'model', + apiKey: 'managed', + baseUrl: CONTRIBUTOR_POLICY.NVIDIA_URL, + }, + ], + } + expect(resolveWorkerProvider([], 0, runtime)?.keyIndex).toBe(-9) + expect(resolveWorkerProvider([], -9, runtime)?.keyIndex).toBe(0) + expect( + latestProviderKeyIndex([ + { key_index: 0, dispatched_at: '2026-10-01T00:00:00Z' }, + { key_index: -9, dispatched_at: '2026-10-01T00:01:00Z' }, + ]), + ).toBe(-9) + }) + it('the real t27 compiler verifies the host policy bindings', async () => { + const analyze = await loadCompiler( + await readFile( + new URL('../../../../specs/t27_compiler.wasm', import.meta.url), + ), + ) + const source = await readFile( + new URL( + '../../../../specs/automation/queen-contributor-keys.t27', + import.meta.url, + ), + 'utf8', + ) + const result = analyze(source) + expect(result.typecheckOk).toBe(true) + expect(result.discarded).toBe(0) + expect(CONTRIBUTOR_POLICY.__NOT_EMITTED__).toEqual([]) + for (const name of CONTRIBUTOR_POLICY.__DECL_ORDER__) + expect(result.consts[name]?.value).toEqual(CONTRIBUTOR_POLICY[name]) + }) +}) diff --git a/trios/agent-server/apps/server/tests/pglive/queen-contributor-keys-live.test.ts b/trios/agent-server/apps/server/tests/pglive/queen-contributor-keys-live.test.ts new file mode 100644 index 0000000000..5647be61bc --- /dev/null +++ b/trios/agent-server/apps/server/tests/pglive/queen-contributor-keys-live.test.ts @@ -0,0 +1,488 @@ +import { afterAll, beforeAll, beforeEach, describe, expect, it } from 'bun:test' +import { randomUUID } from 'node:crypto' +import { Hono } from 'hono' +import { Pool } from 'pg' +import { createQueenContributorKeysRoute } from '../../src/api/routes/queen-contributor-keys' +import { + addContributorKey, + changeContributorKey, + contributorOwnerNames, + contributorRuntime, + ensureContributorKeys, + listContributorKeys, +} from '../../src/api/services/queen-contributor-keys' +import { CONTRIBUTOR_POLICY } from '../../src/api/services/queen-contributor-policy' +import { keyWork, leaderboard } from '../../src/api/services/queen-leaderboard' + +// Deliberate, isolated PostgreSQL. A skip is not a live database pass. +const url = + process.env.QUEEN_CONTRIBUTOR_TEST_DATABASE_URL ?? + process.env.TRIOS_PG_TEST_URL +const live = url ? describe : describe.skip +live('contributor ownership and allocation against real PostgreSQL', () => { + const subject = 'telegram:111222333' + const other = 'telegram:999888777' + const token = 'test-only-proxy-capability-32-bytes-minimum' + const schema = `contributor_${randomUUID().replaceAll('-', '')}` + let pool: Pool + let admin: Pool + const saved = new Map() + const names = [ + 'QUEEN_CONTRIBUTOR_PROXY_TOKEN', + 'QUEEN_CONTRIBUTOR_ENCRYPTION_KEY', + 'QUEEN_CONTRIBUTOR_IDENTITIES', + 'TRIOS_KEY_OWNERS', + ] + const fetchOk = (async () => + Response.json({ + choices: [{ message: { content: 'OK' } }], + })) as typeof fetch + const environment = [ + { + id: 10000, + provider: 'nvidia' as const, + model: 'production-model', + apiKey: 'existing-private-credential', + baseUrl: CONTRIBUTOR_POLICY.NVIDIA_URL, + }, + ] + const owners = { 10000: '@dmitrii-f-t27' } + beforeAll(async () => { + for (const name of names) saved.set(name, process.env[name]) + process.env.QUEEN_CONTRIBUTOR_PROXY_TOKEN = token + process.env.QUEEN_CONTRIBUTOR_ENCRYPTION_KEY = Buffer.alloc( + 32, + 91, + ).toString('base64') + process.env.QUEEN_CONTRIBUTOR_IDENTITIES = JSON.stringify({ + [subject]: 'dmitrii-f-t27', + [other]: 'other-person', + }) + process.env.TRIOS_KEY_OWNERS = '10000=@dmitrii-f-t27' + admin = new Pool({ connectionString: url }) + await admin.query(`CREATE SCHEMA ${schema}`) + pool = new Pool({ + connectionString: url, + options: `-c search_path=${schema}`, + max: 8, + }) + await ensureContributorKeys(pool) + await pool.query(`CREATE TABLE queen_dispatch (key_index integer,review_state text,dispatched_at timestamptz,finished_at timestamptz,owned_paths jsonb); + CREATE TABLE queen_dispatch_history (snapshot jsonb,archived_at timestamptz);`) + }) + beforeEach(async () => { + await pool.query( + 'TRUNCATE queen_contributor_keys,queen_dispatch,queen_dispatch_history', + ) + await pool.query( + 'ALTER SEQUENCE queen_contributor_key_index RESTART WITH -1', + ) + }) + afterAll(async () => { + await pool?.end() + await admin?.query(`DROP SCHEMA IF EXISTS ${schema} CASCADE`) + await admin?.end() + for (const name of names) { + const value = saved.get(name) + if (value === undefined) delete process.env[name] + else process.env[name] = value + } + }) + it('keeps plaintext out of SQL rows and out of every own-account JSON response', async () => { + const secret = 'a-private-managed-test-credential' + const added = await addContributorKey( + pool, + subject, + { provider: 'nvidia', apiKey: secret }, + environment, + fetchOk, + ) + expect(added.id).toBeLessThan(0) + expect(added.enabled).toBe(true) + const stored = await pool.query('SELECT * FROM queen_contributor_keys') + expect(JSON.stringify(stored.rows)).not.toContain(secret) + expect(stored.rows[0].owner_name).toBe('@dmitrii-f-t27') + const listed = await listContributorKeys(pool, subject, environment, owners) + expect(listed.map((key) => key.id)).toEqual([added.id, 10000]) + expect(JSON.stringify(listed)).not.toContain(secret) + const runtime = await contributorRuntime(pool, environment) + expect(runtime.managed).toHaveLength(1) + expect(runtime.managed[0].apiKey).toBe(secret) + }) + it('rejects IDOR at the actual route for every mutation and excludes others from GET', async () => { + const added = await addContributorKey( + pool, + subject, + { provider: 'zai', apiKey: 'key-owned-by-first' }, + environment, + fetchOk, + ) + const app = new Hono().route( + '/queen/contributor-keys', + createQueenContributorKeysRoute({ + pool: () => pool, + environment: () => environment, + owners: () => owners, + fetcher: fetchOk, + }), + ) + const headers = { + Authorization: `Bearer ${token}`, + 'X-Queen-Contributor-Id': other, + } + const list = await app.request('/queen/contributor-keys', { headers }) + expect(list.status).toBe(200) + expect((await list.json()).keys).toEqual([]) + for (const id of [added.id, 10000]) + for (const action of ['disable', 'enable', 'probe']) { + const denied = await app.request( + `/queen/contributor-keys/${id}/${action}`, + { method: 'POST', headers }, + ) + expect(denied.status).toBe(404) + expect(await denied.json()).toEqual({ error: 'key_not_found' }) + } + expect((await contributorRuntime(pool, environment)).managed).toHaveLength( + 1, + ) + }) + it('excludes unreadable managed keys and allows their owner to disable without the master key', async () => { + const added = await addContributorKey( + pool, + subject, + { provider: 'nvidia', apiKey: 'corrupt-test' }, + environment, + fetchOk, + ) + const healthy = await addContributorKey( + pool, + subject, + { provider: 'nvidia', apiKey: 'healthy-test' }, + environment, + fetchOk, + ) + await changeContributorKey( + pool, + subject, + 10000, + 'disable', + environment, + owners, + fetchOk, + ) + await pool.query( + 'UPDATE queen_contributor_keys SET sealed=$2 WHERE key_index=$1', + [added.id, 'corrupt.ciphertext.value'], + ) + expect(await contributorRuntime(pool, environment)).toMatchObject({ + disabled: [10000], + managed: [{ id: healthy.id }], + }) + const master = process.env.QUEEN_CONTRIBUTOR_ENCRYPTION_KEY + delete process.env.QUEEN_CONTRIBUTOR_ENCRYPTION_KEY + try { + await expect( + changeContributorKey( + pool, + other, + added.id, + 'disable', + environment, + owners, + fetchOk, + ), + ).rejects.toMatchObject({ code: 'key_not_found' }) + const disabled = await changeContributorKey( + pool, + subject, + added.id, + 'disable', + environment, + owners, + fetchOk, + ) + expect(disabled.enabled).toBe(false) + expect(await contributorRuntime(pool, environment)).toEqual({ + disabled: [10000], + managed: [], + }) + } finally { + process.env.QUEEN_CONTRIBUTOR_ENCRYPTION_KEY = master + } + }) + it('retains persisted consent and ownership when the management capability is removed', async () => { + const added = await addContributorKey( + pool, + subject, + { provider: 'zai', apiKey: 'persisted-consent' }, + environment, + fetchOk, + ) + await changeContributorKey( + pool, + subject, + 10000, + 'disable', + environment, + owners, + fetchOk, + ) + const capability = process.env.QUEEN_CONTRIBUTOR_PROXY_TOKEN + delete process.env.QUEEN_CONTRIBUTOR_PROXY_TOKEN + try { + expect(await contributorRuntime(pool, environment)).toMatchObject({ + disabled: [10000], + managed: [{ id: added.id }], + }) + expect(await contributorOwnerNames(pool)).toEqual({ + [added.id]: '@dmitrii-f-t27', + 10000: '@dmitrii-f-t27', + }) + await changeContributorKey( + pool, + subject, + added.id, + 'disable', + environment, + owners, + fetchOk, + ) + expect(await contributorRuntime(pool, environment)).toEqual({ + disabled: [10000], + managed: [], + }) + } finally { + process.env.QUEEN_CONTRIBUTOR_PROXY_TOKEN = capability + } + }) + it('rejects duplicates across both owners and environment, including concurrent additions', async () => { + await expect( + addContributorKey( + pool, + subject, + { provider: 'nvidia', apiKey: environment[0].apiKey }, + environment, + fetchOk, + ), + ).rejects.toThrow('key_already_connected') + const results = await Promise.allSettled( + [subject, other].map((owner) => + addContributorKey( + pool, + owner, + { provider: 'zai', apiKey: 'one-shared-test-credential' }, + environment, + fetchOk, + ), + ), + ) + expect( + results.filter((result) => result.status === 'fulfilled'), + ).toHaveLength(1) + expect( + results.filter((result) => result.status === 'rejected'), + ).toHaveLength(1) + expect( + (await pool.query('SELECT * FROM queen_contributor_keys')).rows, + ).toHaveLength(1) + }) + it('disables the matching environment fingerprint without storing or reassigning it', async () => { + await changeContributorKey( + pool, + subject, + 10000, + 'disable', + environment, + owners, + fetchOk, + ) + expect((await contributorRuntime(pool, environment)).disabled).toEqual([ + 10000, + ]) + const row = (await pool.query('SELECT * FROM queen_contributor_keys')) + .rows[0] + expect(row.sealed).toBeNull() + expect(JSON.stringify(row)).not.toContain(environment[0].apiKey) + await expect( + changeContributorKey( + pool, + other, + 10000, + 'enable', + environment, + { 10000: '@other-person' }, + fetchOk, + ), + ).rejects.toThrow('key_not_found') + expect( + await listContributorKeys(pool, other, environment, { + 10000: '@other-person', + }), + ).toEqual([]) + expect( + ( + await contributorRuntime(pool, [ + { ...environment[0], apiKey: 'replaced-secret' }, + ]) + ).disabled, + ).toEqual([]) + await expect( + changeContributorKey( + pool, + subject, + 10000, + 'probe', + [{ ...environment[0], apiKey: 'replaced-secret' }], + owners, + fetchOk, + ), + ).rejects.toThrow('key_binding_conflict') + }) + it('does not let a late enable probe undo a later disable', async () => { + const added = await addContributorKey( + pool, + subject, + { provider: 'nvidia', apiKey: 'test-race-key' }, + environment, + fetchOk, + ) + await pool.query( + "UPDATE queen_contributor_keys SET probe_at=now()-interval '2 minutes'", + ) + let release!: () => void + let started!: () => void + const entered = new Promise((resolve) => { + started = resolve + }) + const pending = new Promise((resolve) => { + release = resolve + }) + const enable = changeContributorKey( + pool, + subject, + added.id, + 'enable', + environment, + owners, + (async () => { + started() + await pending + return Response.json({ choices: [{ message: { content: 'OK' } }] }) + }) as typeof fetch, + ) + await entered + await changeContributorKey( + pool, + subject, + added.id, + 'disable', + environment, + owners, + fetchOk, + ) + release() + expect((await enable).enabled).toBe(false) + expect((await contributorRuntime(pool, environment)).managed).toEqual([]) + }) + it('records failures, bounds repeated probes and requires usable inference before enabling', async () => { + const bad = await addContributorKey( + pool, + subject, + { provider: 'zai', apiKey: 'bad-test-key' }, + environment, + (async () => + new Response('secret-containing upstream error', { + status: 401, + })) as typeof fetch, + ) + expect(bad.enabled).toBe(false) + expect(bad.lastProbe.status).toBe('invalid') + await expect( + changeContributorKey( + pool, + subject, + bad.id, + 'probe', + environment, + owners, + fetchOk, + ), + ).rejects.toThrow('probe_rate_limited') + await pool.query( + "UPDATE queen_contributor_keys SET probe_at=now()-interval '2 minutes'", + ) + expect( + ( + await changeContributorKey( + pool, + subject, + bad.id, + 'enable', + environment, + owners, + fetchOk, + ) + ).enabled, + ).toBe(true) + }) + it('counts negative IDs in archived work and retains existing environment XP under the same owner', async () => { + const added = await addContributorKey( + pool, + subject, + { provider: 'nvidia', apiKey: 'worked-key' }, + environment, + fetchOk, + ) + await pool.query( + `INSERT INTO queen_dispatch VALUES (10000,'accept','2026-09-30T00:00:00Z','2026-09-30T01:00:00Z','["specs/existing.t27"]')`, + ) + await pool.query(`INSERT INTO queen_dispatch_history VALUES ($1,now())`, [ + JSON.stringify({ + key_index: added.id, + review_state: 'accept', + dispatched_at: '2026-09-30T00:00:00Z', + finished_at: '2026-09-30T02:00:00Z', + owned_paths: [], + }), + ]) + expect((await keyWork(pool)).map((row) => row.keyIndex)).toContain(added.id) + const board = await leaderboard(pool) + expect(board.contributors).toHaveLength(1) + expect(board.contributors[0]).toMatchObject({ + github: 'dmitrii-f-t27', + keys: [added.id, 10000], + xp: 430, + accepted: 2, + specs: 1, + hours: 3, + }) + expect((await contributorOwnerNames(pool))[added.id]).toBe('@dmitrii-f-t27') + const app = new Hono().route( + '/queen/contributor-keys', + createQueenContributorKeysRoute({ + pool: () => pool, + environment: () => environment, + owners: () => owners, + fetcher: fetchOk, + }), + ) + const response = await app.request( + '/queen/contributor-keys/10000/disable', + { + method: 'POST', + headers: { + Authorization: `Bearer ${token}`, + 'X-Queen-Contributor-Id': subject, + }, + }, + ) + expect(response.status).toBe(200) + expect((await response.json()).key.contribution).toEqual({ + xp: 310, + accepted: 1, + specs: 1, + finished: 1, + hours: 1, + }) + }) +}) diff --git a/trios/agent-server/biome.json b/trios/agent-server/biome.json index d414afc9b2..927ab04577 100644 --- a/trios/agent-server/biome.json +++ b/trios/agent-server/biome.json @@ -7,7 +7,11 @@ }, "files": { "ignoreUnknown": false, - "includes": ["**", "!**/apps/eval/src/dashboard/index.html"] + "includes": [ + "**", + "!**/apps/eval/src/dashboard/index.html", + "!**/apps/server/src/api/services/queen-contributor-policy.gen.ts" + ] }, "formatter": { "enabled": true, diff --git a/trios/agent-server/docs/contributor-keys.md b/trios/agent-server/docs/contributor-keys.md new file mode 100644 index 0000000000..b1fade5cea --- /dev/null +++ b/trios/agent-server/docs/contributor-keys.md @@ -0,0 +1,117 @@ +# Contributor keys and XP + +Tracks gHashTag/999-multibots-telegraf#3251 and gHashTag/t27#5472. +The contract is `specs/automation/queen-contributor-keys.t27`, mirrored from +the canonical spec in gHashTag/t27#5473. Regenerate the host policy with +`t27c gen-ts specs/automation/queen-contributor-keys.t27` and save its output +as `apps/server/src/api/services/queen-contributor-policy.gen.ts`; never edit +the generated module. SQL, cryptography and HTTP are explicit host adapters. +The formatter excludes this one generated file so commit hooks preserve the +compiler's bytes. TypeScript and the compiler-AST parity test still check it. + +The personal account can list, check, add, enable and disable its own Queen +provider keys. XP comes from the existing dispatch and dispatch-history rows. +This feature never writes to the application's token wallet. + +## Deployment + +Deploy the reviewed commit explicitly to the actual Queen service. The PR +base is `feat/queen-supervisor`; at the 2026-10-01 rollout preparation, Railway +was configured to build `fix/queen-worker-provider-and-prompt-size`. Merging +the PR base alone does not establish that the live service runs this change. + +The management API is off until `QUEEN_CONTRIBUTOR_PROXY_TOKEN` contains at least +32 bytes. Configure a separate random service capability in the Queen and +the app render proxy. Never send it to the browser. It is deliberately +different from `TRIOS_API_TOKEN` and from a user's app session. + +The render proxy verifies its normal app session, then sets +`X-Queen-Contributor-Id: telegram:` on the server-to-server +request. It must discard any subject/header the browser supplied. + +Set `QUEEN_CONTRIBUTOR_IDENTITIES` to an operator-verified JSON map from +these subjects to GitHub logins. A browser cannot assert a GitHub login. +An unmapped account receives a stable pseudonym in the public XP table. + +For existing environment credentials, `TRIOS_KEY_OWNERS` maps the Queen's +actual durable indices to `@github-login`. Verify the active pool ordering +before changing that map; environment suffix numbers and dashboard row +numbers are not interchangeable. The primary pool starts at zero and pool +two starts at 10000. The map is the explicit authority for the existing +history. No automatic ownership migration or foreign key claiming occurs. + +Adding new keys also requires `QUEEN_CONTRIBUTOR_ENCRYPTION_KEY`: canonical +base64 encoding of 32 cryptographically random bytes. Preserve this key +across releases and database backups. Losing it makes stored managed keys +unreadable. Existing environment-key listing and checks do not require it. + +The Queen database role needs permission to create the optional +`queen_contributor_keys` table and `queen_contributor_key_index` sequence in +the configured Queen schema. Setup runs on first use, and before scheduling +when the feature is enabled. Storage/auth failures return a closed error; +they do not silently substitute empty data or expose provider errors. +An unreadable managed credential is excluded individually from allocation and +logged by index only. Its owner can still disable it without the master key; +other credentials and existing disable directives remain effective. +Once the registry exists, persisted consent and owner snapshots remain in +force even if the management capability is removed. An old deployment with +no registry keeps its legacy allocation; a registry lookup failure is closed. + +## API + +All routes require `Authorization: Bearer ` and the +verified contributor header. Responses use `Cache-Control: no-store`. + +- `GET /queen/contributor-keys`: own keys, provider options, per-key and + total XP, verified attribution, and the existing public XP formula. +- `POST /queen/contributor-keys`: `{provider, apiKey, label?}`. Providers + are `nvidia` and `zai`; URLs and models are server-controlled. +- `POST /queen/contributor-keys/:id/probe`: check an owned key. +- `POST /queen/contributor-keys/:id/enable`: check and enable on success. +- `POST /queen/contributor-keys/:id/disable`: stop assigning new work. + +Successful mutations return `{key}` with the same metadata and actual +contribution fields as GET. Secrets, ciphertext, provider response bodies, +upstream error messages and operator configuration are never returned. +Errors are `{error: }`. Foreign/absent keys both return 404. + +A probe performs one bounded chat-completion request. An HTTP 200 with no +model output is not a successful check. Invalid credentials are disabled; +temporary rate limits or network failures do not revoke prior consent. +Each key allows one explicit check per minute, and the process caps concurrent +explicit checks at four. A provider may charge for these small requests. + +Disabling stops future worker/review/model-probe assignment; a turn already +running keeps the credential it was given. A late enable response cannot +undo a later disable. Model-probe scheduling rereads consent before each +new model probe. + +## Stable identity + +Existing environment keys retain their current positive/zero indices. +Overrides bind index, SHA-256 fingerprint and immutable owner. Filtering +happens after legacy indexing, so disabling key zero does not rename key one. +An environment replacement at a bound index is a conflict, not a transfer. + +New keys use immutable negative database indices, unique fingerprints and +AES-256-GCM ciphertext authenticated with owner and fingerprint. There is +no ownership-update or hard-delete endpoint. Disabled keys keep their XP. +Public attribution reads the saved owner name for registered keys; changing +an environment mapping cannot reassign a registered credential. + +## Verification + +Run the API tests with Bun 1.3.6: + +```text +bun test tests/api/queen-contributor-keys.test.ts +``` + +From `apps/server`, set `QUEEN_CONTRIBUTOR_TEST_DATABASE_URL` to a disposable +PostgreSQL database and run `tests/pglive/queen-contributor-keys-live.test.ts`. +The existing CI's `TRIOS_PG_TEST_URL` is also accepted; production +`DATABASE_URL` is never used by these tests. +It creates and drops an isolated schema. These checks cover real SQL owner +filters, duplicate races, encryption at rest, disable/enable races, malformed +requests, provider failure classification, negative-ID history and unchanged +legacy XP. Without the database variable the live tests explicitly skip. diff --git a/trios/agent-server/specs/automation/queen-contributor-keys.controls.json b/trios/agent-server/specs/automation/queen-contributor-keys.controls.json new file mode 100644 index 0000000000..40eee31cfa --- /dev/null +++ b/trios/agent-server/specs/automation/queen-contributor-keys.controls.json @@ -0,0 +1,47 @@ +[ + { + "name": "a short proxy token is sufficient", + "from": "configured_bytes >= PROXY_TOKEN_MIN_BYTES", + "to": "configured_bytes >= PROXY_TOKEN_MIN_BYTES - 1" + }, + { + "name": "managed keys collide with environment zero", + "from": "managed && index < 0", + "to": "managed && index <= 0" + }, + { + "name": "owner key limit admits one more key", + "from": "owned < MAX_KEYS_PER_OWNER", + "to": "owned <= MAX_KEYS_PER_OWNER" + }, + { + "name": "probe fanout admits one more request", + "from": "active < MAX_PROBES_IN_FLIGHT", + "to": "active <= MAX_PROBES_IN_FLIGHT" + }, + { + "name": "cooldown permits the excluded boundary", + "from": "elapsed_seconds > PROBE_COOLDOWN_SECONDS", + "to": "elapsed_seconds >= PROBE_COOLDOWN_SECONDS" + }, + { + "name": "empty credentials are accepted", + "from": "bytes > 0 && bytes <= MAX_KEY_BYTES", + "to": "bytes <= MAX_KEY_BYTES" + }, + { + "name": "either owner or fingerprint match is enough", + "from": "return subject_matches && fingerprint_matches;", + "to": "return subject_matches || fingerprint_matches;" + }, + { + "name": "a probe can override a concurrent disable", + "from": "success && request_enable && same_revision", + "to": "success && (request_enable || same_revision)" + }, + { + "name": "adding keys credits a wallet", + "from": "pub const CREDIT_WALLET : bool = false;", + "to": "pub const CREDIT_WALLET : bool = true;" + } +] diff --git a/trios/agent-server/specs/automation/queen-contributor-keys.t27 b/trios/agent-server/specs/automation/queen-contributor-keys.t27 new file mode 100644 index 0000000000..6649fef86c --- /dev/null +++ b/trios/agent-server/specs/automation/queen-contributor-keys.t27 @@ -0,0 +1,138 @@ +// SPDX-License-Identifier: Apache-2.0 +// specs/automation/queen-contributor-keys.t27 -- private contributor key lifecycle. +// Host: BrowserOS trios/agent-server/apps/server/src/api/services/queen-contributor-keys.ts. +// Trace: gHashTag/t27#5472; gHashTag/999-multibots-telegraf#3251. +// Generated: t27c gen-ts supplies the host policy constants. +// Claim status: executable contract; PostgreSQL and live provider evidence are separate. +// XP follows work recorded by the Queen. Adding credentials never credits a wallet. +// Negative managed indices and nonnegative environment indices share one history. +// v1 (2026-10-01): trusted proxy, persistent ownership and bounded provider probes. +// phi^2 + 1/phi^2 = 3 | TRINITY + +module automation::queen_contributor_keys { + pub const KIND : str = "automation"; + pub const ID : str = "queen-contributor-keys"; + pub const REPO : str = "BrowserOS"; + pub const VERSION : u8 = 1; + + pub const PROBE_TIMEOUT_MS : u32 = 15000; + pub const PROBE_MAX_TOKENS : u16 = 16; + pub const MAX_PROBES_IN_FLIGHT : u16 = 4; + pub const PROBE_COOLDOWN_SECONDS : u16 = 60; + pub const MAX_KEYS_PER_OWNER : u16 = 100; + pub const MAX_KEY_BYTES : u16 = 4096; + pub const MAX_BODY_BYTES : u16 = 8192; + pub const LABEL_LIMIT : u16 = 80; + pub const PROXY_TOKEN_MIN_BYTES : u16 = 32; + pub const OWNER_REASSIGNMENT : bool = false; + pub const CREDIT_WALLET : bool = false; + pub const ENVIRONMENT_INDICES_COMPACTED : bool = false; + pub const NVIDIA_URL : str = "https://integrate.api.nvidia.com/v1"; + pub const NVIDIA_MODEL : str = "nvidia/nemotron-3-ultra-550b-a55b"; + pub const ZAI_URL : str = "https://api.z.ai/api/paas/v4"; + pub const ZAI_MODEL : str = "glm-4.5-flash"; + + fn trusted_proxy(configured_bytes: u32, token_matches: bool, subject_valid: bool) -> bool { + return configured_bytes >= PROXY_TOKEN_MIN_BYTES && token_matches && subject_valid; + } + + fn source_matches_index(managed: bool, index: i32) -> bool { + return (managed && index < 0) || (!managed && index >= 0); + } + + fn may_add_key(owned: u32) -> bool { + return owned < MAX_KEYS_PER_OWNER; + } + + fn may_start_probe(active: u32) -> bool { + return active < MAX_PROBES_IN_FLIGHT; + } + + fn cooldown_passed(has_probe: bool, elapsed_seconds: u32) -> bool { + return !has_probe || elapsed_seconds > PROBE_COOLDOWN_SECONDS; + } + + fn credential_fits(bytes: u32, printable_without_space: bool) -> bool { + return bytes > 0 && bytes <= MAX_KEY_BYTES && printable_without_space; + } + + fn same_binding(subject_matches: bool, fingerprint_matches: bool) -> bool { + return subject_matches && fingerprint_matches; + } + + fn enabled_after_probe(was_enabled: bool, invalid: bool, success: bool, request_enable: bool, same_revision: bool) -> bool { + if invalid { return false; } + if success && request_enable && same_revision { return true; } + return was_enabled; + } + + fn wallet_credit(added: u32) -> u32 { + if CREDIT_WALLET { return added; } + return 0; + } + + test "only a configured matching proxy and valid person subject are trusted" { + assert(trusted_proxy(32, true, true)); + assert(!trusted_proxy(31, true, true)); + assert(!trusted_proxy(32, false, true)); + assert(!trusted_proxy(32, true, false)); + } + + test "managed indices are negative and environment history keeps its nonnegative indices" { + assert(source_matches_index(true, -1)); + assert(!source_matches_index(true, 0)); + assert(!source_matches_index(false, -1)); + assert(source_matches_index(false, 0)); + assert(source_matches_index(false, 10000)); + assert(!ENVIRONMENT_INDICES_COMPACTED); + } + + test "the owner key limit is checked before another key is stored" { + assert(may_add_key(0)); + assert(may_add_key(99)); + assert(!may_add_key(100)); + } + + test "probe fanout and per-key cooldown are bounded" { + assert(may_start_probe(3)); + assert(!may_start_probe(4)); + assert(cooldown_passed(false, 0)); + assert(!cooldown_passed(true, 59)); + assert(!cooldown_passed(true, 60)); + assert(cooldown_passed(true, 61)); + } + + test "credentials are bounded nonempty printable values" { + assert(credential_fits(1, true)); + assert(credential_fits(4096, true)); + assert(!credential_fits(0, true)); + assert(!credential_fits(4097, true)); + assert(!credential_fits(32, false)); + } + + test "an existing binding cannot change its owner or credential" { + assert(same_binding(true, true)); + assert(!same_binding(false, true)); + assert(!same_binding(true, false)); + assert(!OWNER_REASSIGNMENT); + } + + test "a failed probe disables an invalid key and a concurrent disable wins" { + assert(!enabled_after_probe(true, true, false, false, true)); + assert(enabled_after_probe(false, false, true, true, true)); + assert(!enabled_after_probe(false, false, true, true, false)); + assert(!enabled_after_probe(false, false, true, false, true)); + assert(enabled_after_probe(true, false, false, false, true)); + } + + test "adding credentials never credits a wallet" { + assert(wallet_credit(1) == 0); + assert(wallet_credit(100) == 0); + } + + test "VERSION says what the header's last note says" { + assert(VERSION == 1); + } +} + +// phi^2 + 1/phi^2 = 3 | TRINITY