import type { SiteSettings } from "../lib/settings"; import { LANES, type AutoQueueKind } from "../lib/autoQueueTypes"; import { allOperations, operationLabel, operationsForLane, } from "../lib/operations"; import type { WorkerSummary } from "../jobs/workerPool"; import type { ChannelSnapshot } from "../controller/channelSnapshot"; import { workerMatches } from "../lib/workers"; import { buildOperationBands, type OperationBand, } from "./pipeline/buildBands"; // WHAT THE CONSOLE KNOWS ABOUT A LANE THAT IS NOT ITS RUNNER'S STATE. // // This file used to be the two SWEEP-fed lanes: a state, a reason, a scope, a // plan and an arm button apiece, for the two pipelines the auto-queue runner // could not dispatch. It can dispatch them now (slice 1.2) and the sweeps are // gone (1.3), so a lane's state, its work list and its controls all come from // `/api/auto-queue/status` like the other two lanes'. // // TWO THINGS ARE LEFT, and both are properties of the REGISTRY rather than of a // runner, which is why they are still built here: // // * the BANDS — one per catalogued operation, its five populations off the // channel snapshots. The rail draws them, and they are per OPERATION where // everything else on this page is per LANE. // * WHICH OPERATIONS EACH LANE DISPATCHES, with their labels and the delegate // workers that can take them. Settings-dependent (`digest.remoteEnabled` // moves the digest operation between two queue keys that are both the // digest lane), so it is a server read, and the console is a client // component. export type LaneOperationView = { id: string; label: string; }; export type LaneWorker = { // The configured worker's id (slot-expansion suffixes folded back together). id: string; name: string; kind: "remote" | "llm"; slots: number; // Slots currently enabled and not degraded — 0 across every delegate is the // "configured but nothing can take it right now" warning state. available: number; }; export type AutoQueueLanesPayload = { // One band per pipeline, in rail order. See components/pipelines/buildBands.ts. bands: OperationBand[]; // The operations each lane dispatches, in dependency order. Empty for the two // bucket lanes, whose work list is a snapshot bucket rather than an operation. operations: Record; // WHERE A LANE'S WORK CAN RUN besides this machine, folded to one row per // configured worker with its slot count. READ-ONLY and never part of any // policy — delegation is a worker-pool concern decided per item at dispatch, // and the console only reports it. Derived through the SAME workerMatches rule // the dispatchers grant by, so this row and the actual routing cannot // disagree. runsOn: Record; }; export type AutoQueueLanesInputs = { settings: SiteSettings; // Only the snapshots are read — the bands are the five populations each // channel's report already carries. Typed as the half that is used so a // caller need not have a full brief to ask. briefs: ReadonlyArray<{ snapshot: ChannelSnapshot | null }>; // `pool.summary()`, taken once by the shell rather than per lane. workerSummary: readonly WorkerSummary[]; }; export function buildAutoQueueLanes( inputs: AutoQueueLanesInputs, ): AutoQueueLanesPayload { const { settings, briefs, workerSummary } = inputs; // allOperations, not one lane's members: the rail is the CATALOG view — every // operation that is switched on, whatever queue it runs on. const kinds = allOperations(settings); const bands = buildOperationBands({ snapshots: briefs.map((c) => c.snapshot ?? null), operationIds: kinds.map((k) => k.id), }); const contendsForOf = new Map( kinds.map((k) => [k.id, (k.laneFor?.(settings) ?? k.lane).contendsFor]), ); // An llm worker serves an operation its tags (or the llm default set) name; a // remote worker takes units only when TAGGED (the unit-dispatch opt-in) and // its tags intersect [operation, contended resource]. Local workers are not // listed — the lane always runs here, and saying so would be noise. // ONE PASS PER WORKER SLOT, not per (operation, slot) pair. A worker serving // two of the lane's operations is one row with one slot count — the slots are // the machine's, not the operation's, and counting them twice would advertise // capacity that does not exist. const runsOnFor = (ids: readonly string[]): LaneWorker[] => { const byBase = new Map(); for (const w of workerSummary) { const takesAny = ids.some((id) => w.kind === "llm" ? workerMatches({ kind: "llm", tags: w.tags }, [id]) : w.kind === "remote" && (w.tags?.length ?? 0) > 0 && workerMatches({ kind: "remote", tags: w.tags }, [ id, contendsForOf.get(id) ?? "cpu", ]), ); if (!takesAny) continue; const base = w.id.split("#")[0]; const entry = byBase.get(base) ?? { id: base, name: w.name.replace(/ #\d+$/, ""), kind: w.kind as "remote" | "llm", slots: 0, available: 0, }; entry.slots++; if (w.state === "enabled" && !w.degraded) entry.available++; byBase.set(base, entry); } return [...byBase.values()]; }; const operations = {} as Record; const runsOn = {} as Record; for (const lane of LANES) { const ids = operationsForLane(lane, settings).map((op) => op.id); operations[lane] = ids.map((id) => ({ id, label: operationLabel(id) })); runsOn[lane] = runsOnFor(ids); } return { bands, operations, runsOn }; }