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
164 changes: 164 additions & 0 deletions trios/agent-server/apps/server/src/api/routes/queen-runners.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,164 @@
/**
* @license
* Copyright 2025 BrowserOS
* SPDX-License-Identifier: AGPL-3.0-or-later
*
* THE RUNNER CABINET, AND THE DOOR A RUNNER KNOCKS ON.
*
* Two routes, two credentials that never meet:
*
* /queen/me/runners the PERSON, with the bearer token app.t27.ai issued
* them (verified by asking its issuer, see
* queen-app-identity.ts). Lists, mints and revokes that
* person's runner tokens - nobody else's.
*
* /queen/runner the RUNNER, with the token minted above. Today it can
* say it is alive and learn its lane; taking work and
* handing it back are the next stage, and the answer says
* so instead of pretending to offer work.
*
* Neither route accepts, stores or returns a provider key. The key stays on
* the runner's machine; that is the whole point of a runner.
*/
import { Hono } from 'hono'
import type { Pool } from 'pg'
import { createQueenPool } from '../../lib/db/queen-pool'
import { logger } from '../../lib/logger'
import {
bearerOf,
createAppIdentity,
type Identify,
IdentityUnavailableError,
} from '../services/queen-app-identity'
import {
cleanLabel,
createRunner,
heartbeatRunner,
laneOf,
listRunners,
MAX_RUNNERS_PER_PERSON,
revokeRunner,
viewOf,
} from '../services/queen-runners'

/** What a runner speaks. Bumped when the claim/complete stage lands. */
export const RUNNER_PROTOCOL = 1

type Queryable = Pick<Pool, 'query'>

export interface RunnerRouteDeps {
/** The database, or null when none is configured (503). */
pool: () => Queryable | null
identify: Identify
}

/**
* One pool for the life of the process, created on first use. A pool per
* request that is never ended is a connection leak with a delay on it.
*/
let sharedPool: Pool | undefined
function defaultPool(): Queryable | null {
const url = process.env.DATABASE_URL || process.env.RAILWAY_SSOT_URL
if (!url) return null
if (!sharedPool) sharedPool = createQueenPool(url)
return sharedPool
}

function defaults(deps: Partial<RunnerRouteDeps>): RunnerRouteDeps {
return {
pool: deps.pool ?? defaultPool,
identify: deps.identify ?? createAppIdentity(),
}
}

const NO_DATABASE = { error: 'No database configured' } as const

export function createQueenCabinetRoute(given: Partial<RunnerRouteDeps> = {}) {
const deps = defaults(given)
return new Hono<{
Variables: { person: { telegramId: string; name: string } }
}>()
.use('/*', async (c, next) => {
const bearer = bearerOf(c.req.header('authorization'))
if (!bearer) return c.json({ error: 'Sign in to app.t27.ai first' }, 401)
try {
const person = await deps.identify(bearer)
if (!person) return c.json({ error: 'The session was refused' }, 401)
c.set('person', person)
} catch (error) {
if (!(error instanceof IdentityUnavailableError)) throw error
logger.warn('Runner cabinet could not verify a session', {
error: error.message,
})
return c.json({ error: 'Sign-in service did not answer' }, 503)
}
await next()
return
})
.get('/', async (c) => {
const pool = deps.pool()
if (!pool) return c.json(NO_DATABASE, 503)
const rows = await listRunners(pool, c.get('person').telegramId)
return c.json(
{
runners: rows.map((r) => viewOf(r)),
limit: MAX_RUNNERS_PER_PERSON,
protocol: RUNNER_PROTOCOL,
},
200,
{ 'Cache-Control': 'no-store' },
)
})
.post('/', async (c) => {
const pool = deps.pool()
if (!pool) return c.json(NO_DATABASE, 503)
const body = await c.req.json().catch(() => null)
const label = cleanLabel(body?.label)
if (!label) return c.json({ error: 'A runner needs a name' }, 400)
const made = await createRunner(pool, c.get('person'), label)
if (!made.ok) {
return c.json(
{
error: `At most ${MAX_RUNNERS_PER_PERSON} runners; revoke one first`,
},
409,
)
}
// The only answer that ever carries the token. The page shows it once.
return c.json({ runner: viewOf(made.runner), token: made.token }, 201, {
'Cache-Control': 'no-store',
})
})
.delete('/:id', async (c) => {
const pool = deps.pool()
if (!pool) return c.json(NO_DATABASE, 503)
const id = Number(c.req.param('id'))
if (!Number.isSafeInteger(id) || id <= 0)
return c.json({ error: 'No such runner' }, 404)
const done = await revokeRunner(pool, c.get('person').telegramId, id)
// Someone else's runner and no runner at all are the same answer: the
// cabinet does not confirm which ids exist.
return done ? c.body(null, 204) : c.json({ error: 'No such runner' }, 404)
})
}

export function createQueenRunnerRoute(given: Partial<RunnerRouteDeps> = {}) {
const deps = defaults(given)
return new Hono().post('/heartbeat', async (c) => {
const pool = deps.pool()
if (!pool) return c.json(NO_DATABASE, 503)
const token = bearerOf(c.req.header('authorization'))
const runner = token ? await heartbeatRunner(pool, token) : null
if (!runner) return c.json({ error: 'Unknown or revoked runner' }, 401)
return c.json(
{
runner: { id: runner.id, label: runner.label, lane: laneOf(runner.id) },
protocol: RUNNER_PROTOCOL,
work: null,
note: 'Registered. Handing tasks to runners is the next stage; until it lands there is nothing to take.',
},
200,
{ 'Cache-Control': 'no-store' },
)
})
}
16 changes: 15 additions & 1 deletion trios/agent-server/apps/server/src/api/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,10 @@ import {
createQueenRoadmapDataRoute,
createQueenRoadmapRoute,
} from './routes/queen-roadmap'
import {
createQueenCabinetRoute,
createQueenRunnerRoute,
} from './routes/queen-runners'
import { createQueenTreeRoute } from './routes/queen-tree'
import { createRefinePromptRoutes } from './routes/refine-prompt'
import { createShutdownRoute } from './routes/shutdown'
Expand All @@ -89,7 +93,11 @@ import { convertOpenClawHistoryToAgentHistory } from './services/openclaw/histor
import { getOpenClawService } from './services/openclaw/openclaw-service'
import { TaskQueueService } from './services/task-queue-service'
import type { Env, HttpServerConfig } from './types'
import { publicReadCorsMiddleware, trustedCorsMiddleware } from './utils/cors'
import {
appCabinetCorsMiddleware,
publicReadCorsMiddleware,
trustedCorsMiddleware,
} from './utils/cors'
import { requireTrustedAppOrigin } from './utils/request-auth'

async function assertPortAvailable(port: number): Promise<void> {
Expand Down Expand Up @@ -386,6 +394,8 @@ export async function createHttpServer(config: HttpServerConfig) {
.use('/queen/public-credits', publicReadCorsMiddleware())
.use('/queen/public-earnings', publicReadCorsMiddleware())
.use('/queen/scheduler', publicReadCorsMiddleware())
// The runner cabinet: exactly https://app.t27.ai, bearer only, no cookies.
.use('/queen/me/*', appCabinetCorsMiddleware())
.use('/*', trustedCorsMiddleware())
// The Inngest server registers and invokes functions here; each request
// is signed with INNGEST_SIGNING_KEY and verified by the SDK, so this sits
Expand All @@ -403,6 +413,10 @@ 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())
// A person's runner tokens (their bearer, verified by its issuer), and the
// door their runner knocks on (the runner token). Neither takes a key.
.route('/queen/me/runners', createQueenCabinetRoute())
.route('/queen/runner', createQueenRunnerRoute())
.route('/queen/public-credits', createQueenPublicCreditsRoute())
// What accepted spec work earned; recorded, not withdrawable. Repository,
// issue, commit and declared .t27 paths only - no titles, notes or keys.
Expand Down
185 changes: 185 additions & 0 deletions trios/agent-server/apps/server/src/api/services/queen-app-identity.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,185 @@
/**
* @license
* Copyright 2025 BrowserOS
* SPDX-License-Identifier: AGPL-3.0-or-later
*
* WHO IS ASKING, FOR A PERSON SIGNED IN TO https://app.t27.ai.
*
* This server has never known a person. Its only credentials are the
* operator's TRIOS_API_TOKEN and the loopback local-auth token; Telegram
* sessions belong to the player's own service (vibee-render, in
* 999-multibots-telegraf), which issues the access token the board at
* app.t27.ai/queen/ already holds in memory.
*
* So the answer to "who is this" is asked of the service that issued the
* token, with the token: `tools/call whoami` on its /mcp, the same call the
* board itself makes to put a name beside the avatar. Its answer carries
* `telegram_id`, and that id is the only identity the runner cabinet keys on.
*
* WHAT THIS FILE DOES WITH THE TOKEN. It forwards it once to the issuer and
* nowhere else. It is never logged, never stored and never echoed; the cache
* below is keyed by its SHA-256, so a heap dump holds a hash, not a credential.
*
* THREE ANSWERS, NOT TWO. `null` means the issuer refused the token (sign in
* again); a thrown IdentityUnavailableError means the issuer did not answer
* (try again later). Collapsing them would tell a person whose session is fine
* that they are signed out every time vibee-render redeploys.
*/
import { createHash } from 'node:crypto'

export const DEFAULT_APP_IDENTITY_URL =
'https://vibee-render-production.up.railway.app'

/** How long a verified answer is trusted before the issuer is asked again. */
export const IDENTITY_CACHE_MS = 60_000
/** Bounded: a flood of distinct tokens must not grow the heap without limit. */
const IDENTITY_CACHE_MAX = 500
/** The issuer is abandoned after this long; the caller says "try again". */
export const IDENTITY_TIMEOUT_MS = 5_000

export interface AppPerson {
/** Digits only: Telegram ids are integers, and the column is text. */
telegramId: string
/** Display name from the issuer's profile, for the person's own runners. */
name: string
}

export class IdentityUnavailableError extends Error {
constructor(reason: string) {
super(`app identity unavailable: ${reason}`)
this.name = 'IdentityUnavailableError'
}
}

type FetchLike = (
url: string,
init: {
method: string
headers: Record<string, string>
body: string
signal: AbortSignal
},
) => Promise<{ ok: boolean; status: number; json(): Promise<unknown> }>

const isRecord = (value: unknown): value is Record<string, unknown> =>
!!value && typeof value === 'object' && !Array.isArray(value)

/**
* The payload of one MCP `tools/call` result: structured content when present,
* otherwise the first text block parsed as JSON. Same rule as the board's
* mcpAnswer.ts, because the same server is answering both.
*/
export function mcpPayload(result: unknown): unknown {
if (!isRecord(result)) return undefined
if (result.structuredContent !== undefined) return result.structuredContent
if (!Array.isArray(result.content)) return undefined
const text = result.content.find(
(part): part is { type: string; text: string } =>
isRecord(part) && part.type === 'text' && typeof part.text === 'string',
)?.text
if (text === undefined) return undefined
try {
return JSON.parse(text)
} catch {
return text
}
}

/** The person in a whoami answer, or null when it names nobody. */
export function personFromWhoami(body: unknown): AppPerson | null {
if (!isRecord(body) || body.error !== undefined) return null
const result = body.result
if (isRecord(result) && result.isError === true) return null
const data = mcpPayload(result)
if (!isRecord(data)) return null
const raw = data.telegram_id
const telegramId =
typeof raw === 'number' && Number.isSafeInteger(raw) && raw > 0
? String(raw)
: typeof raw === 'string' && /^[1-9]\d{0,19}$/.test(raw)
? raw
: null
if (!telegramId) return null
const profile = isRecord(data['профиль']) ? data['профиль'] : {}
const name =
[profile.display_name, profile.first_name, profile.username]
.find((v): v is string => typeof v === 'string' && v.trim() !== '')
?.trim()
.slice(0, 40) ?? `tg ${telegramId.slice(-4)}`
return { telegramId, name }
}

/** `Authorization: Bearer x` -> 'x'. Anything else is no credential. */
export function bearerOf(header: string | undefined): string | null {
if (!header) return null
const match = /^Bearer ([A-Za-z0-9._~+/=-]{16,4096})$/.exec(header.trim())
return match ? match[1] : null
}

export interface AppIdentityOptions {
baseUrl?: string
fetch?: FetchLike
now?: () => number
}

export type Identify = (bearer: string) => Promise<AppPerson | null>

export function createAppIdentity(options: AppIdentityOptions = {}): Identify {
const baseUrl = (
options.baseUrl ??
process.env.TRIOS_APP_IDENTITY_URL ??
DEFAULT_APP_IDENTITY_URL
).replace(/\/+$/, '')
const doFetch: FetchLike = options.fetch ?? (fetch as unknown as FetchLike)
const now = options.now ?? Date.now
const cache = new Map<string, { person: AppPerson | null; until: number }>()

return async (bearer) => {
const key = createHash('sha256').update(bearer).digest('hex')
const hit = cache.get(key)
if (hit && hit.until > now()) return hit.person

const abort = new AbortController()
const timer = setTimeout(() => abort.abort(), IDENTITY_TIMEOUT_MS)
let status: number
let body: unknown
try {
const res = await doFetch(`${baseUrl}/mcp`, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
Accept: 'application/json',
Authorization: `Bearer ${bearer}`,
},
body: JSON.stringify({
jsonrpc: '2.0',
id: 1,
method: 'tools/call',
params: { name: 'whoami', arguments: {} },
}),
signal: abort.signal,
})
status = res.status
body = res.ok ? await res.json() : null
} catch (error) {
throw new IdentityUnavailableError(
error instanceof Error ? error.name : 'fetch failed',
)
} finally {
clearTimeout(timer)
}

let person: AppPerson | null
if (status === 401 || status === 403) person = null
else if (status >= 200 && status < 300) person = personFromWhoami(body)
else throw new IdentityUnavailableError(`http ${status}`)

if (cache.size >= IDENTITY_CACHE_MAX) {
// Oldest first: Map iterates in insertion order.
const oldest = cache.keys().next().value
if (oldest !== undefined) cache.delete(oldest)
}
cache.set(key, { person, until: now() + IDENTITY_CACHE_MS })
return person
}
}
Loading
Loading