Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -2,16 +2,22 @@
* Public native research graph for t27.ai.
*
* The canonical graph is the evidence-backed file used by /queen/tree. This
* route adds directionally-correct prerequisites/unlocks and a secret-free
* view of paid worker-slot utilisation. It deliberately keeps graph state and
* worker activity separate: "partial" means the repository has incomplete
* route adds directionally-correct prerequisites/unlocks, a secret-free
* view of paid worker-slot utilisation, and - since #1308 - the anonymous
* capacity factors behind that utilisation: how many credentials are
* connected and how many lanes each carries, read from the same dispatch
* authority that allocates against them. It deliberately keeps graph state
* and worker activity separate: "partial" means the repository has incomplete
* evidence, not that a model is currently spending tokens on it.
*/

import { Hono } from 'hono'
import { Pool } from 'pg'
import { logger } from '../../lib/logger'
import { configuredWorkerCapacity } from '../services/queen-dispatch'
import {
type WorkerCapacityBreakdown,
workerCapacityBreakdown,
} from '../services/queen-dispatch'
import {
isTreeLoadFailure,
loadTree as loadCanonicalTree,
Expand All @@ -33,7 +39,8 @@ interface QueenPublicResearchDeps {
loadTree?: () => Promise<Tree | TreeLoadFailure | null>
databaseUrl?: () => string | undefined
createPool?: (url: string) => ResearchPool
workerCapacity?: () => number
/** The closed capacity authority; defaults to dispatch's own breakdown. */
workerCapacityBreakdown?: () => WorkerCapacityBreakdown
publicOrigin?: (requestUrl: string) => string
}

Expand Down Expand Up @@ -119,6 +126,9 @@ function projectTree(tree: Tree) {
}

function workerProjection(capacity: number, busyIndices: number[]) {
// The capacity arrives from the closed breakdown authority (#1308); the
// anonymous factor fields join it in the response without changing any of
// the contracts below.
const safeCapacity = Math.max(0, Math.floor(capacity))
// key_index identifies a credential, not a logical lane. With an explicit
// multi-lane plan two rows may legitimately carry the same index; counting
Expand Down Expand Up @@ -149,7 +159,8 @@ export function createQueenPublicResearchRoute(
const createPool =
deps.createPool ??
((url: string) => new Pool({ connectionString: url }) as ResearchPool)
const workerCapacity = deps.workerCapacity ?? configuredWorkerCapacity
const capacityBreakdown =
deps.workerCapacityBreakdown ?? workerCapacityBreakdown
const publicOrigin = deps.publicOrigin ?? configuredPublicOrigin

return new Hono().get('/', async (c) => {
Expand Down Expand Up @@ -194,10 +205,17 @@ export function createQueenPublicResearchRoute(
// container. Build copyable A2A links from Railway's trusted public domain
// rather than leaking the internal scheme into the bootstrap contract.
const origin = publicOrigin(c.req.url)
// 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,
workers: workerProjection(workerCapacity(), busyIndices),
workers: {
...workerProjection(breakdown.effectiveCapacity, busyIndices),
...breakdown,
},
agentBootstrap: {
version: 'trinity-research-a2a/v1',
mode: 'public-read-only',
Expand Down
72 changes: 62 additions & 10 deletions trios/agent-server/apps/server/src/api/services/queen-dispatch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -203,14 +203,22 @@ export interface WorkerProvider {
* one function - the count a dashboard shows and the index a bee takes are the
* same list, never two different stories about one secret. First occurrence
* wins, so the unsuffixed variable stays index 0 in every ordering.
*
* Values are TRIMMED before they are judged (#1308). ' key' and 'key' pasted
* into two boxes are one credential wearing its whitespace differently, and a
* value that is nothing but whitespace is the empty box one paste later. The
* trimmed value is also the one stored: a key that authenticates never needed
* its padding, and handing the trimmed form out keeps the count and the
* selection - which both read this list - from ever disagreeing.
*/
function keysFor(envVar: string): string[] {
const keys: string[] = []
const seen = new Set<string>()
const admit = (value: string | undefined) => {
if (!value || value.length === 0 || seen.has(value)) return
seen.add(value)
keys.push(value)
const trimmed = (value ?? '').trim()
if (trimmed.length === 0 || seen.has(trimmed)) return
seen.add(trimmed)
keys.push(trimmed)
}
admit(process.env[envVar])
for (let i = 2; i <= 16; i++) {
Expand Down Expand Up @@ -241,16 +249,60 @@ function workerLanesFor(provider: string): number {
}

/**
* Number of genuinely independent worker credentials available to the first
* configured provider. The values never leave this module; the public research
* projection uses only the count to show whether paid capacity is idle.
* The closed, anonymous capacity breakdown every capacity number is made of
* (#1308).
*
* `workers.capacity` answering 4 does not say WHICH 4: two subscriptions at a
* lane each and one subscription at two lanes each are the same total with
* completely different operator implications - the first hides a disconnected
* paid subscription, the second promises parallelism a single rate limit
* cannot back. This is the ONE authority both `configuredWorkerCapacity` and
* the public research telemetry read, so the factorisation a dashboard shows
* and the ceiling dispatch allocates against are the same statement, never two
* different stories about one configuration.
*
* CLOSED means three integers and nothing else. No hashes, no key suffixes, no
* slot indexes, no provider variable names, no values: anything shaped like a
* credential is a disclosure, and a count cannot be inverted into one.
*/
export function configuredWorkerCapacity(): number {
export interface WorkerCapacityBreakdown {
connectedCredentials: number
lanesPerCredential: number
effectiveCapacity: number
}

export function workerCapacityBreakdown(): WorkerCapacityBreakdown {
for (const candidate of WORKER_PROVIDERS) {
const count = keysFor(candidate.envVar).length
if (count > 0) return count * workerLanesFor(candidate.provider)
const keys = keysFor(candidate.envVar)
if (keys.length > 0) {
const lanesPerCredential = workerLanesFor(candidate.provider)
return {
connectedCredentials: keys.length,
lanesPerCredential,
effectiveCapacity: keys.length * lanesPerCredential,
}
}
}
return 0
// Nothing is connected. The lanes factor keeps its safe default rather than
// zeroing, because it describes the bound the NEXT connected credential
// would run under; the total is still zero, from zero credentials alone.
return {
connectedCredentials: 0,
lanesPerCredential: configuredWorkerLanesPerCredential(),
effectiveCapacity: 0,
}
}

/**
* Number of genuinely independent worker credentials available to the first
* configured provider, multiplied by that provider's lanes. The values never
* leave this module; the public research projection uses only the count to
* show whether paid capacity is idle. Delegates to the breakdown authority so
* the number allocated against and the number explained publicly can never
* diverge (#1308).
*/
export function configuredWorkerCapacity(): number {
return workerCapacityBreakdown().effectiveCapacity
}

/**
Expand Down
181 changes: 181 additions & 0 deletions trios/agent-server/apps/server/tests/api/queen-dispatch.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import {
recordDispatch,
resolveWorkerProvider,
setDurableCloseListener,
workerCapacityBreakdown,
workspaceRoot,
} from '../../src/api/services/queen-dispatch'
import { logger } from '../../src/lib/logger'
Expand All @@ -30,6 +31,11 @@ const KEYS = [
'OPENROUTER_API_KEY',
'MOONSHOT_API_KEY',
'OPENAI_API_KEY',
// #1308's factorisation table configures a second OpenAI slot. Bun runs the
// api test files in ONE process, so a suffixed name missing from this list
// is not a tidiness problem: it survives into the next FILE and makes a
// "nothing is connected" case read one connected credential.
'OPENAI_API_KEY_2',
'TRIOS_QUEEN_WORKER_MODEL',
'TRIOS_ZAI_CONCURRENCY_PER_KEY',
]
Expand Down Expand Up @@ -253,6 +259,181 @@ describe('queen dispatch precheck', () => {
})
})

/**
* #1308. `workers.capacity` answers a number; this breakdown answers what the
* number is MADE of. An operator seeing capacity 4 cannot act on it without
* knowing whether it is two subscriptions at a lane each - one of which may be
* quietly disconnected - or one subscription at two lanes each, and a total
* alone keeps that a guess.
*/
describe('worker capacity breakdown', () => {
// Scenario 1 of the issue: two distinct configured Z.ai credentials and two
// lanes per credential. The response is closed - three integers, no trace of
// WHICH credentials produced them.
it('factors capacity into connected credentials and lanes per credential', () => {
process.env.ZAI_API_KEY = 'planted-secret-a'
process.env.ZAI_API_KEY_2 = 'planted-secret-b'
process.env.TRIOS_ZAI_CONCURRENCY_PER_KEY = '2'
expect(workerCapacityBreakdown()).toEqual({
connectedCredentials: 2,
lanesPerCredential: 2,
effectiveCapacity: 4,
})
// The same authority dispatch allocates against, not a second story.
expect(configuredWorkerCapacity()).toBe(4)
expect(resolveWorkerProvider([0, 1, 0, 1])?.exhausted).toBe(4)
})

// FR-002/FR-003: closed and anonymous. Anything beyond these three fields -
// a hash, a suffix, an index, a variable name, a value - is a disclosure.
it('is three numeric fields and nothing else', () => {
process.env.ZAI_API_KEY = 'planted-secret-a'
process.env.ZAI_API_KEY_2 = 'planted-secret-b'
const breakdown = workerCapacityBreakdown() as unknown as Record<
string,
unknown
>
expect(Object.keys(breakdown).sort()).toEqual([
'connectedCredentials',
'effectiveCapacity',
'lanesPerCredential',
])
for (const value of Object.values(breakdown)) {
expect(typeof value).toBe('number')
expect(Number.isInteger(value)).toBe(true)
}
const serialized = JSON.stringify(breakdown)
expect(serialized).not.toContain('planted-secret')
expect(serialized).not.toContain('ZAI_API_KEY')
expect(serialized).not.toContain('ANTHROPIC_API_KEY')
})

// Scenario 2: a credential duplicated across slots is one account with one
// rate limit (#1293). Neither factor may be inflated by it.
it('counts a duplicated credential once so nothing is inflated', () => {
process.env.ZAI_API_KEY = 'planted-secret-a'
process.env.ZAI_API_KEY_2 = 'planted-secret-a'
process.env.ZAI_API_KEY_3 = 'planted-secret-b'
process.env.TRIOS_ZAI_CONCURRENCY_PER_KEY = '2'
expect(workerCapacityBreakdown()).toEqual({
connectedCredentials: 2,
lanesPerCredential: 2,
effectiveCapacity: 4,
})
expect(configuredWorkerCapacity()).toBe(4)
})

// FR-003: counted after TRIMMING. ' key' and 'key' in two boxes are one
// credential wearing its whitespace differently, and the count must say so
// before a second slot is handed a secret the first is already spending.
it('trims values before counting, so padded duplicates are one credential', () => {
process.env.ZAI_API_KEY = ' planted-secret-a '
process.env.ZAI_API_KEY_2 = 'planted-secret-a'
process.env.ZAI_API_KEY_3 = ' planted-secret-b '
expect(workerCapacityBreakdown().connectedCredentials).toBe(2)
// Selection reads the same trimmed list, so the two can never disagree.
expect(resolveWorkerProvider([])?.keyCount).toBe(2)
})

it('treats a whitespace-only value as the empty box it supplies nothing from', () => {
process.env.ZAI_API_KEY = ' '
expect(workerCapacityBreakdown().connectedCredentials).toBe(0)
expect(configuredWorkerCapacity()).toBe(0)
})

// Scenario 3: no supported provider credentials. Every factor is zero or
// its safe default, and nothing about a secret leaves with it.
it('reports zeros and the safe lane default when nothing is connected', () => {
expect(workerCapacityBreakdown()).toEqual({
connectedCredentials: 0,
lanesPerCredential: 1,
effectiveCapacity: 0,
})
expect(configuredWorkerCapacity()).toBe(0)
})

// FR-005: the lane factor keeps its existing safe default and bound.
it('keeps the safe default of one lane and the bound of four', () => {
process.env.ZAI_API_KEY = 'planted-secret-a'
process.env.TRIOS_ZAI_CONCURRENCY_PER_KEY = '0'
expect(workerCapacityBreakdown()).toEqual({
connectedCredentials: 1,
lanesPerCredential: 1,
effectiveCapacity: 1,
})
process.env.TRIOS_ZAI_CONCURRENCY_PER_KEY = '99'
expect(workerCapacityBreakdown()).toEqual({
connectedCredentials: 1,
lanesPerCredential: 4,
effectiveCapacity: 4,
})
})

// The lane override belongs to Z.ai's tiered plans; another provider's
// capacity stays one credential times one lane, exactly as before.
it('does not apply the Z.ai lane factor to another provider', () => {
process.env.ANTHROPIC_API_KEY = 'planted-anthropic-secret'
process.env.TRIOS_ZAI_CONCURRENCY_PER_KEY = '3'
expect(workerCapacityBreakdown()).toEqual({
connectedCredentials: 1,
lanesPerCredential: 1,
effectiveCapacity: 1,
})
expect(configuredWorkerCapacity()).toBe(1)
})

it('factors only the first configured provider, in preference order', () => {
process.env.ZAI_API_KEY = 'planted-secret-a'
process.env.ANTHROPIC_API_KEY = 'planted-anthropic-secret'
process.env.OPENAI_API_KEY = 'planted-openai-secret'
expect(workerCapacityBreakdown().connectedCredentials).toBe(1)
})

// FR-004: effective capacity is the number dispatch allocates against, so
// the two must be one number in every configuration, not two that happen to
// agree today. Each fixture starts from a cleared environment.
it('equals configuredWorkerCapacity in every tested configuration', () => {
const configurations: Array<() => void> = [
() => undefined,
() => {
process.env.ZAI_API_KEY = 'a'
},
() => {
process.env.ZAI_API_KEY = 'a'
process.env.ZAI_API_KEY_2 = 'b'
process.env.TRIOS_ZAI_CONCURRENCY_PER_KEY = '2'
},
() => {
process.env.ZAI_API_KEY = 'a'
process.env.ZAI_API_KEY_2 = 'a'
process.env.TRIOS_ZAI_CONCURRENCY_PER_KEY = '4'
},
() => {
process.env.ANTHROPIC_API_KEY = 'anthropic-a'
process.env.TRIOS_ZAI_CONCURRENCY_PER_KEY = '2'
},
() => {
process.env.OPENAI_API_KEY = 'openai-a'
process.env.OPENAI_API_KEY_2 = 'openai-b'
},
() => {
process.env.ZAI_API_KEY = ' a '
process.env.ZAI_API_KEY_2 = 'a'
process.env.ZAI_API_KEY_3 = ' b '
},
]
for (const configure of configurations) {
for (const key of KEYS) delete process.env[key]
configure()
const breakdown = workerCapacityBreakdown()
expect(breakdown.effectiveCapacity).toBe(configuredWorkerCapacity())
expect(breakdown.effectiveCapacity).toBe(
breakdown.connectedCredentials * breakdown.lanesPerCredential,
)
}
})
})

/** Every statement a call made, with the values it bound. */
function recordingPool(
answer: (sql: string, attempt: number) => unknown = () => ({
Expand Down
Loading
Loading