From 0c3fed31d26894bdae193fafc78cdb9c6e4da549 Mon Sep 17 00:00:00 2001 From: BakerSean168 Date: Sat, 3 Oct 2026 16:25:54 +0000 Subject: [PATCH 1/5] feat: add Antigravity agent bridge --- README.md | 14 +++ bin/agy-subagent.mjs | 64 ++++++++++++ docs/architecture.md | 14 +++ extension/antigravity.js | 100 +++++++++++++++++++ extension/index.js | 6 ++ package.json | 4 +- skills/forgeflow/SKILL.md | 7 ++ tests/antigravity.test.js | 147 ++++++++++++++++++++++++++++ tests/architecture-boundary.test.js | 2 +- tests/extension.test.js | 22 +++++ 10 files changed, 377 insertions(+), 3 deletions(-) create mode 100644 bin/agy-subagent.mjs create mode 100644 extension/antigravity.js create mode 100644 tests/antigravity.test.js diff --git a/README.md b/README.md index 94d2118..8e66ce0 100644 --- a/README.md +++ b/README.md @@ -30,6 +30,20 @@ Provider infrastructure remains below Pi model selection. ForgeFlow logical role may resolve to physical models, while endpoint selection, credentials, channel health, weights, quotas, and transport belong to the provider layer. +## Antigravity delegation + +When the installed `pi-subagents` owner is present, ForgeFlow also registers two +external agents backed by the locally authenticated Antigravity CLI (`agy`): + +- `antigravity` (alias `agy`) runs Antigravity in read-only `plan` mode; +- `antigravity-writer` (alias `agy-writer`) runs in `accept-edits` mode. + +`pi-subagents` still owns child lifecycle, status, timeout, and stop. ForgeFlow only +bridges its stdin handoff to `agy --print`; Antigravity keeps its own authentication, +model selection, quota, and execution runtime. The bridge is local-only and requires +`agy` on `PATH`. It does not turn Antigravity subscription quota into a Pi model +provider. + See [`docs/model-policy.md`](./docs/model-policy.md) for logical-model policy and [`docs/architecture.md`](./docs/architecture.md) for the ownership boundary. diff --git a/bin/agy-subagent.mjs b/bin/agy-subagent.mjs new file mode 100644 index 0000000..943f022 --- /dev/null +++ b/bin/agy-subagent.mjs @@ -0,0 +1,64 @@ +#!/usr/bin/env node + +import { spawn } from "node:child_process"; + +const VALID_MODES = new Set(["plan", "accept-edits"]); +const mode = process.argv[2]; + +if (!VALID_MODES.has(mode)) { + process.stderr.write("Usage: agy-subagent.mjs \n"); + process.exit(64); +} + +let prompt = ""; +process.stdin.setEncoding("utf8"); +for await (const chunk of process.stdin) { + prompt += chunk; +} + +if (!prompt.trim()) { + process.stderr.write("Antigravity subagent prompt must not be empty.\n"); + process.exit(64); +} + +const args = [ + "--mode", + mode, + "--output-format", + "text", + "--dangerously-skip-permissions", + "--print-timeout", + "0s", + `--print=${prompt}` +]; + +const child = spawn("agy", args, { + env: process.env, + stdio: ["ignore", "pipe", "pipe"] +}); + +child.stdout.pipe(process.stdout); +child.stderr.pipe(process.stderr); + +for (const signal of ["SIGINT", "SIGTERM"]) { + process.on(signal, () => { + if (!child.killed) child.kill(signal); + }); +} + +child.on("error", (error) => { + process.stderr.write( + `Failed to launch Antigravity CLI (agy): ${error.message}\n` + ); + process.exitCode = 127; +}); + +child.on("exit", (code, signal) => { + if (signal) { + process.stderr.write(`Antigravity CLI terminated by ${signal}.\n`); + process.exitCode = 1; + return; + } + + process.exitCode = code ?? 1; +}); diff --git a/docs/architecture.md b/docs/architecture.md index a469b4a..d5bbf73 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -30,6 +30,20 @@ workflow state, provider/channel routing, worktrees, or resume, it belongs in Pi an existing plugin. If it is an engineering method or delivery procedure, it belongs in a Skill. Provider/channel routing remains in the provider layer. +### External coding agents + +ForgeFlow may register a thin transport adapter when an installed coding agent cannot +consume the `pi-subagents` stdin contract directly. The Antigravity integration is +one example: `pi-subagents` remains the orchestration owner, while a small bridge +converts the assembled stdin prompt into one `agy --print=` argument. + +The bridge owns no sessions, retries, workflow state, model routing, credentials, or +quota. `agy` remains the authority for Antigravity authentication, model entitlement, +and subscription usage. Two runtime agents are exposed when the `pi-subagents` +registration owner is present: `antigravity`/`agy` for plan-mode analysis and +`antigravity-writer`/`agy-writer` for explicit workspace mutation. Generic external +CLI runners are local-only under the current `pi-subagents` contract. + ## Model policy boundary ForgeFlow registers stable logical roles as Pi virtual models. Their physical model diff --git a/extension/antigravity.js b/extension/antigravity.js new file mode 100644 index 0000000..8520f0a --- /dev/null +++ b/extension/antigravity.js @@ -0,0 +1,100 @@ +import { fileURLToPath } from "node:url"; + +import { registerAgentViaEvents } from "pi-subagents/agents"; + +const AGY_BRIDGE_PATH = fileURLToPath( + new URL("../bin/agy-subagent.mjs", import.meta.url) +); + +const RUNTIME_OWNER_UNAVAILABLE = + /pi-subagents is not installed, not ready, or does not support runtime agent event registration/; + +function definition(mode, options) { + return { + description: options.description, + aliases: options.aliases, + systemPrompt: options.systemPrompt, + systemPromptMode: "replace", + inheritProjectContext: true, + inheritGlobalContext: false, + inheritSkills: false, + defaultAsync: true, + acceptanceRole: options.acceptanceRole, + runner: { + type: "external-cli", + command: process.execPath, + args: [AGY_BRIDGE_PATH, mode], + promptDelivery: "stdin" + } + }; +} + +export const ANTIGRAVITY_AGENTS = Object.freeze([ + Object.freeze({ + name: "antigravity", + definition: Object.freeze( + definition("plan", { + aliases: ["agy"], + acceptanceRole: "read-only", + description: + "Read-only Antigravity coding agent using the locally authenticated agy CLI and its own subscription quota", + systemPrompt: + "Analyze the requested engineering task using Antigravity in plan mode. Inspect as needed, do not modify the workspace, and return concise findings with concrete evidence." + }) + ) + }), + Object.freeze({ + name: "antigravity-writer", + definition: Object.freeze( + definition("accept-edits", { + aliases: ["agy-writer"], + acceptanceRole: "writer", + description: + "Workspace-writing Antigravity coding agent using the locally authenticated agy CLI and its own subscription quota", + systemPrompt: + "Implement the requested engineering task in the current workspace using Antigravity. Keep changes scoped, run relevant validation, and return a concise completion report with changed files and test evidence." + }) + ) + }) +]); + +function disposeAll(registrations) { + for (const registration of registrations.reverse()) { + registration.dispose(); + } +} + +export function registerAntigravityAgents(pi) { + const registrations = []; + + try { + for (const agent of ANTIGRAVITY_AGENTS) { + registrations.push( + registerAgentViaEvents({ + pi, + name: agent.name, + definition: agent.definition + }) + ); + } + } catch (error) { + disposeAll(registrations); + + if ( + error instanceof Error && + RUNTIME_OWNER_UNAVAILABLE.test(error.message) + ) { + return undefined; + } + + throw error; + } + + return { + dispose() { + disposeAll(registrations); + } + }; +} + +export { AGY_BRIDGE_PATH }; diff --git a/extension/index.js b/extension/index.js index 49d9feb..898b075 100644 --- a/extension/index.js +++ b/extension/index.js @@ -1,6 +1,7 @@ import { fileURLToPath } from "node:url"; import { registerRequiredChildExtensions } from "pi-subagents/required-child-extensions"; +import { registerAntigravityAgents } from "./antigravity.js"; import { renderPreflight } from "./invariants.js"; import { registerForgeFlowVirtualModels } from "./model-policy.js"; @@ -23,6 +24,7 @@ export default function registerForgeFlow(pi) { registerForgeFlowVirtualModels(pi); let requiredChildRegistration; + let antigravityRegistration; pi.on("before_agent_start", (event) => { event.systemPromptOptions.sections.forgeflow_policy = buildForgeFlowPromptSection(event.prompt); @@ -30,14 +32,18 @@ export default function registerForgeFlow(pi) { pi.on("session_start", (_event, ctx) => { requiredChildRegistration?.dispose(); + antigravityRegistration?.dispose(); requiredChildRegistration = registerRequiredChildExtensions({ sessionId: ctx.sessionManager.getSessionId(), extensions: [{ id: "forgeflow", path: FORGEFLOW_EXTENSION_PATH }] }); + antigravityRegistration = registerAntigravityAgents(pi); }); pi.on("session_shutdown", () => { + antigravityRegistration?.dispose(); + antigravityRegistration = undefined; requiredChildRegistration?.dispose(); requiredChildRegistration = undefined; }); diff --git a/package.json b/package.json index 0dfb822..d7f9454 100644 --- a/package.json +++ b/package.json @@ -4,7 +4,7 @@ "description": "Thin software-engineering governance for Pi Agent and pi-subagents.", "type": "module", "license": "MIT", - "files": ["extension", "skills", "docs/architecture.md", "docs/model-policy.md", "README.md", "LICENSE"], + "files": ["bin", "extension", "skills", "docs/architecture.md", "docs/model-policy.md", "README.md", "LICENSE"], "keywords": [ "pi-package", "pi", @@ -35,6 +35,6 @@ }, "scripts": { "test": "node --test tests/*.test.js", - "check": "npm test && node --check extension/index.js && node --check extension/model-policy.js && node --check extension/invariants.js" + "check": "npm test && node --check bin/agy-subagent.mjs && node --check extension/antigravity.js && node --check extension/index.js && node --check extension/model-policy.js && node --check extension/invariants.js" } } diff --git a/skills/forgeflow/SKILL.md b/skills/forgeflow/SKILL.md index 826569c..433a9a7 100644 --- a/skills/forgeflow/SKILL.md +++ b/skills/forgeflow/SKILL.md @@ -39,6 +39,13 @@ ForgeFlow. Use `pi-subagents` for child execution, reviewer runs, review loops, typed gates, and runtime acceptance evidence. Treat it as the single orchestration owner. +When the operator explicitly asks to delegate work to Antigravity, use +`antigravity`/`agy` for read-only analysis or `antigravity-writer`/`agy-writer` for +workspace changes. Those agents consume the local Antigravity CLI's own +subscription quota; they are external agents, not ForgeFlow physical model roles. +Never run an Antigravity writer concurrently with another writer in the same +worktree. + Use installed engineering Skills such as `test-driven-development`, `spec-driven-development`, `pr-gate`, and `delivery-verification` for reusable workflow methods. Exact-head CI, test execution, artifact identity, and deployment diff --git a/tests/antigravity.test.js b/tests/antigravity.test.js new file mode 100644 index 0000000..64ab568 --- /dev/null +++ b/tests/antigravity.test.js @@ -0,0 +1,147 @@ +import assert from "node:assert/strict"; +import { spawnSync } from "node:child_process"; +import { + chmodSync, + mkdtempSync, + readFileSync, + writeFileSync +} from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import test from "node:test"; + +import { registerAgent } from "pi-subagents/agents"; +import { + AGY_BRIDGE_PATH, + ANTIGRAVITY_AGENTS, + registerAntigravityAgents +} from "../extension/antigravity.js"; + +test("Antigravity runtime agents expose read-only and writer roles", () => { + assert.deepEqual( + ANTIGRAVITY_AGENTS.map(({ name }) => name), + ["antigravity", "antigravity-writer"] + ); + + const [reader, writer] = ANTIGRAVITY_AGENTS; + + assert.equal(reader.definition.acceptanceRole, "read-only"); + assert.deepEqual(reader.definition.aliases, ["agy"]); + assert.equal(reader.definition.runner.type, "external-cli"); + assert.equal(reader.definition.runner.command, process.execPath); + assert.deepEqual(reader.definition.runner.args, [AGY_BRIDGE_PATH, "plan"]); + + assert.equal(writer.definition.acceptanceRole, "writer"); + assert.deepEqual(writer.definition.aliases, ["agy-writer"]); + assert.deepEqual(writer.definition.runner.args, [ + AGY_BRIDGE_PATH, + "accept-edits" + ]); +}); + +test("Antigravity definitions satisfy the pi-subagents runtime validator", () => { + const pi = { + on() {}, + registerTool() {} + }; + const registrations = ANTIGRAVITY_AGENTS.map((agent) => + registerAgent({ + pi, + name: agent.name, + definition: agent.definition + }) + ); + + for (const registration of registrations.reverse()) { + registration.dispose(); + } +}); + +test("runtime registration is owned by pi-subagents and disposed together", () => { + const requests = []; + const disposed = []; + + const pi = { + events: { + emit(event, request) { + assert.equal(event, "pi-subagents:runtime-agent-register:v1"); + requests.push(request); + request.result = { + ok: true, + registration: { + dispose() { + disposed.push(request.name); + } + } + }; + } + } + }; + + const registration = registerAntigravityAgents(pi); + assert.ok(registration); + assert.deepEqual( + requests.map(({ name }) => name), + ["antigravity", "antigravity-writer"] + ); + + registration.dispose(); + assert.deepEqual(disposed, ["antigravity-writer", "antigravity"]); +}); + +test("missing pi-subagents owner leaves ForgeFlow usable", () => { + const pi = { + events: { + emit() {} + } + }; + + assert.equal(registerAntigravityAgents(pi), undefined); +}); + +test("agy bridge converts stdin into one safe --print argument", () => { + const dir = mkdtempSync(path.join(os.tmpdir(), "forgeflow-agy-test-")); + const capture = path.join(dir, "capture.json"); + const fakeAgy = path.join(dir, process.platform === "win32" ? "agy.cmd" : "agy"); + + if (process.platform === "win32") { + writeFileSync( + fakeAgy, + `@echo off\r\nnode -e "require('fs').writeFileSync(process.env.AGY_CAPTURE, JSON.stringify(process.argv.slice(1)))" %*\r\n` + ); + } else { + writeFileSync( + fakeAgy, + `#!/usr/bin/env node +const fs = require("node:fs"); +fs.writeFileSync(process.env.AGY_CAPTURE, JSON.stringify(process.argv.slice(2))); +` + ); + chmodSync(fakeAgy, 0o755); + } + + const prompt = "Inspect this; echo $(danger) && keep newlines\nsecond line"; + const result = spawnSync(process.execPath, [AGY_BRIDGE_PATH, "plan"], { + input: prompt, + encoding: "utf8", + env: { + ...process.env, + AGY_CAPTURE: capture, + PATH: `${dir}${path.delimiter}${process.env.PATH ?? ""}` + } + }); + + assert.equal(result.status, 0, result.stderr); + + const args = JSON.parse(readFileSync(capture, "utf8")); + assert.deepEqual(args.slice(0, 7), [ + "--mode", + "plan", + "--output-format", + "text", + "--dangerously-skip-permissions", + "--print-timeout", + "0s" + ]); + assert.equal(args[7], `--print=${prompt}`); +}); diff --git a/tests/architecture-boundary.test.js b/tests/architecture-boundary.test.js index f1f3da4..ec62adb 100644 --- a/tests/architecture-boundary.test.js +++ b/tests/architecture-boundary.test.js @@ -71,6 +71,6 @@ test("Pi package has one explicit runtime entrypoint and one plugin dependency", assert.equal(packageJson.peerDependencies["@earendil-works/pi-coding-agent"], ">=0.99.0"); assert.deepEqual( new Set(packageJson.files), - new Set(["extension", "skills", "docs/architecture.md", "docs/model-policy.md", "README.md", "LICENSE"]) + new Set(["bin", "extension", "skills", "docs/architecture.md", "docs/model-policy.md", "README.md", "LICENSE"]) ); }); diff --git a/tests/extension.test.js b/tests/extension.test.js index aa18889..1ab914d 100644 --- a/tests/extension.test.js +++ b/tests/extension.test.js @@ -7,13 +7,30 @@ import registerForgeFlow from "../extension/index.js"; test("extension registers Pi lifecycle hooks and injects policy without a model call", () => { const handlers = new Map(); const virtualModels = []; + const runtimeAgentRequests = []; + const disposedRuntimeAgents = []; const pi = { on(event, handler) { handlers.set(event, handler); return () => handlers.delete(event); }, + registerTool() {}, registerVirtualModel(definition) { virtualModels.push(definition); + }, + events: { + emit(event, request) { + if (event !== "pi-subagents:runtime-agent-register:v1") return; + runtimeAgentRequests.push(request); + request.result = { + ok: true, + registration: { + dispose() { + disposedRuntimeAgents.push(request.name); + } + } + }; + } } }; @@ -56,8 +73,13 @@ test("extension registers Pi lifecycle hooks and injects policy without a model () => registerRequiredChildExtensions({ sessionId, extensions: [] }), /already registered/ ); + assert.deepEqual( + runtimeAgentRequests.map(({ name }) => name), + ["antigravity", "antigravity-writer"] + ); assert.doesNotThrow(() => handlers.get("session_shutdown")()); + assert.deepEqual(disposedRuntimeAgents, ["antigravity-writer", "antigravity"]); const afterShutdown = registerRequiredChildExtensions({ sessionId, extensions: [] }); afterShutdown.dispose(); }); From 944ee8f8d4670727f910d63e5c88fd852f90d256 Mon Sep 17 00:00:00 2001 From: BakerSean168 Date: Sun, 4 Oct 2026 07:50:03 +0000 Subject: [PATCH 2/5] feat(model-policy): add adaptive thinking envelopes --- docs/model-policy.md | 43 +++++++++++++++------ extension/model-policy.js | 77 ++++++++++++++++++++++++++++++++++---- tests/model-policy.test.js | 62 +++++++++++++++++++++++++++++- 3 files changed, 162 insertions(+), 20 deletions(-) diff --git a/docs/model-policy.md b/docs/model-policy.md index 8209de3..7f7a86e 100644 --- a/docs/model-policy.md +++ b/docs/model-policy.md @@ -22,11 +22,15 @@ This keeps fast-moving physical model names out of ForgeFlow workflows and agent ForgeFlow registers these Pi virtual models: -- `forgeflow/planner` -- `forgeflow/worker` -- `forgeflow/reviewer` -- `forgeflow/scout` -- `forgeflow/oracle` +| Role | Purpose | Typical effort envelope | +| --- | --- | --- | +| `forgeflow/planner` | decompose work, architecture, execution planning | `high` by default, up to `xhigh` | +| `forgeflow/worker` | implementation, focused debugging, routine code changes | `low` by default, up to `medium` | +| `forgeflow/reviewer` | diff review, acceptance reasoning, risk checks | `high` by default, up to `xhigh` | +| `forgeflow/scout` | repository exploration, cheap search, fact gathering | `low` by default, up to `medium` | +| `forgeflow/oracle` | expensive expert escalation for ambiguous, cross-system, or hard root-cause questions | `xhigh` by default and capped at `xhigh` | + +`oracle` is an engineering consultation role, not Oracle Cloud or the Oracle2 host. It should be invoked sparingly when the planner/reviewer needs a higher-cost second opinion or a difficult decision resolved; it is not the default implementation worker. The names are stable contracts. The physical model behind each role is operator policy. @@ -67,23 +71,33 @@ Policy schema v1: "roles": { "planner": { "model": "gateway/frontier-reasoning", - "thinkingLevel": "high" + "defaultThinkingLevel": "high", + "minThinkingLevel": "medium", + "maxThinkingLevel": "xhigh" }, "worker": { "model": "gateway/fast-coder", - "thinkingLevel": "low" + "defaultThinkingLevel": "low", + "minThinkingLevel": "minimal", + "maxThinkingLevel": "medium" }, "reviewer": { "model": "gateway/frontier-review", - "thinkingLevel": "high" + "defaultThinkingLevel": "high", + "minThinkingLevel": "high", + "maxThinkingLevel": "xhigh" }, "scout": { "model": "gateway/fast-general", - "thinkingLevel": "low" + "defaultThinkingLevel": "low", + "minThinkingLevel": "off", + "maxThinkingLevel": "medium" }, "oracle": { "model": "gateway/frontier-reasoning", - "thinkingLevel": "high" + "defaultThinkingLevel": "xhigh", + "minThinkingLevel": "high", + "maxThinkingLevel": "xhigh" } } } @@ -91,7 +105,14 @@ Policy schema v1: `model` must be a fully qualified **physical** Pi model, `provider/model`. Model IDs may contain additional slashes. A role may not point at another `forgeflow/*` virtual model. -`thinkingLevel` is optional. When omitted, the selected virtual thinking level is passed through to the physical model. +Thinking effort is part of the role policy, not an afterthought. The adaptive fields are: + +- `defaultThinkingLevel`: used when the caller does not request a level; +- `minThinkingLevel`: floor for the role, preventing an underpowered call; +- `maxThinkingLevel`: ceiling for the role, preventing routine work from consuming frontier effort; +- legacy `thinkingLevel`: an exact fixed pin retained for backwards compatibility. It cannot be combined with the adaptive fields. + +When the caller explicitly selects a virtual thinking level, ForgeFlow preserves that request inside the configured envelope. For example a planner configured as `medium..xhigh` can run at `high` or escalate to `xhigh`, while a worker capped at `medium` cannot accidentally consume `xhigh`. `max` remains available in the schema, but operator policy should only enable it for a physical model whose registry metadata explicitly supports that level. Continuations and retries keep the already-selected physical model and effort for turn stability. The policy file is read when a new user/direct request is routed, so changing a mapping does not require changing ForgeFlow code. Pi still records the actual physical model on every assistant response. diff --git a/extension/model-policy.js b/extension/model-policy.js index ef7965f..a0208e8 100644 --- a/extension/model-policy.js +++ b/extension/model-policy.js @@ -69,17 +69,56 @@ function validateRolePolicy(role, value, source) { ); } - let thinkingLevel; - if (value.thinkingLevel !== undefined) { - if (typeof value.thinkingLevel !== "string" || !THINKING_LEVELS.has(value.thinkingLevel)) { + function optionalThinkingLevel(field) { + if (value[field] === undefined) return undefined; + if (typeof value[field] !== "string" || !THINKING_LEVELS.has(value[field])) { throw new Error( - `ForgeFlow model policy '${source}' role '${role}' has invalid 'thinkingLevel'.` + `ForgeFlow model policy '${source}' role '${role}' has invalid '${field}'.` ); } - thinkingLevel = value.thinkingLevel; + return value[field]; } - return { provider, id, thinkingLevel }; + const thinkingLevel = optionalThinkingLevel("thinkingLevel"); + const defaultThinkingLevel = optionalThinkingLevel("defaultThinkingLevel"); + const minThinkingLevel = optionalThinkingLevel("minThinkingLevel"); + const maxThinkingLevel = optionalThinkingLevel("maxThinkingLevel"); + + const adaptiveFields = [defaultThinkingLevel, minThinkingLevel, maxThinkingLevel]; + if (thinkingLevel !== undefined && adaptiveFields.some((entry) => entry !== undefined)) { + throw new Error( + `ForgeFlow model policy '${source}' role '${role}' cannot combine legacy 'thinkingLevel' with adaptive thinking fields.` + ); + } + + const rank = (level) => level === undefined ? undefined : [...THINKING_LEVELS].indexOf(level); + const minRank = rank(minThinkingLevel); + const maxRank = rank(maxThinkingLevel); + const defaultRank = rank(defaultThinkingLevel); + if (minRank !== undefined && maxRank !== undefined && minRank > maxRank) { + throw new Error( + `ForgeFlow model policy '${source}' role '${role}' has minThinkingLevel above maxThinkingLevel.` + ); + } + if (defaultRank !== undefined && minRank !== undefined && defaultRank < minRank) { + throw new Error( + `ForgeFlow model policy '${source}' role '${role}' has defaultThinkingLevel below minThinkingLevel.` + ); + } + if (defaultRank !== undefined && maxRank !== undefined && defaultRank > maxRank) { + throw new Error( + `ForgeFlow model policy '${source}' role '${role}' has defaultThinkingLevel above maxThinkingLevel.` + ); + } + + return { + provider, + id, + ...(thinkingLevel !== undefined ? { thinkingLevel } : {}), + ...(defaultThinkingLevel !== undefined ? { defaultThinkingLevel } : {}), + ...(minThinkingLevel !== undefined ? { minThinkingLevel } : {}), + ...(maxThinkingLevel !== undefined ? { maxThinkingLevel } : {}) + }; } export function parseModelPolicy(text, source = "") { @@ -133,6 +172,30 @@ export function loadModelPolicy(cwd, env = process.env, home = homedir(), option }; } +function resolveThinkingLevel(requested, target) { + if (target.thinkingLevel !== undefined) { + return target.thinkingLevel; + } + + let level = requested ?? target.defaultThinkingLevel; + if (level === undefined) return undefined; + + const ordered = [...THINKING_LEVELS]; + const rank = ordered.indexOf(level); + if (rank === -1) return target.defaultThinkingLevel; + + const minRank = target.minThinkingLevel === undefined + ? undefined + : ordered.indexOf(target.minThinkingLevel); + const maxRank = target.maxThinkingLevel === undefined + ? undefined + : ordered.indexOf(target.maxThinkingLevel); + + if (minRank !== undefined && rank < minRank) level = target.minThinkingLevel; + if (maxRank !== undefined && ordered.indexOf(level) > maxRank) level = target.maxThinkingLevel; + return level; +} + function stickyRoute(request) { if (request.reason === "retry" && request.failed) { return { @@ -181,7 +244,7 @@ export function createRoleRouter(role, options = {}) { return { model, - thinkingLevel: target.thinkingLevel ?? request.thinkingLevel + thinkingLevel: resolveThinkingLevel(request.thinkingLevel, target) }; }; } diff --git a/tests/model-policy.test.js b/tests/model-policy.test.js index b281347..f5381ab 100644 --- a/tests/model-policy.test.js +++ b/tests/model-policy.test.js @@ -18,7 +18,12 @@ test("model policy keeps physical model ids outside ForgeFlow code paths", () => version: 1, roles: { worker: { model: "litellm/gpt-fast", thinkingLevel: "low" }, - reviewer: { model: "newapi/reasoner/v2", thinkingLevel: "high" } + reviewer: { + model: "newapi/reasoner/v2", + defaultThinkingLevel: "high", + minThinkingLevel: "medium", + maxThinkingLevel: "xhigh" + } } })); @@ -30,7 +35,9 @@ test("model policy keeps physical model ids outside ForgeFlow code paths", () => assert.deepEqual(policy.roles.reviewer, { provider: "newapi", id: "reasoner/v2", - thinkingLevel: "high" + defaultThinkingLevel: "high", + minThinkingLevel: "medium", + maxThinkingLevel: "xhigh" }); }); @@ -60,6 +67,57 @@ test("model policy fails closed on virtual recursion, unknown roles, and unquali ); }); +test("adaptive thinking policy applies role floors, ceilings, and defaults", () => { + const physical = { provider: "litellm", id: "reasoner-next" }; + const route = createRoleRouter("planner", { + loadPolicy() { + return { + path: "/policy.json", + policy: { + version: 1, + roles: { + planner: { + provider: "litellm", + id: "reasoner-next", + defaultThinkingLevel: "high", + minThinkingLevel: "medium", + maxThinkingLevel: "xhigh" + } + } + } + }; + } + }); + const ctx = { + cwd: "/repo", + isProjectTrusted: () => true, + modelRegistry: { find: () => physical } + }; + + assert.deepEqual(route({ reason: "user" }, ctx), { model: physical, thinkingLevel: "high" }); + assert.deepEqual(route({ reason: "user", thinkingLevel: "low" }, ctx), { model: physical, thinkingLevel: "medium" }); + assert.deepEqual(route({ reason: "user", thinkingLevel: "max" }, ctx), { model: physical, thinkingLevel: "xhigh" }); + assert.deepEqual(route({ reason: "user", thinkingLevel: "high" }, ctx), { model: physical, thinkingLevel: "high" }); +}); + +test("model policy rejects invalid adaptive thinking envelopes", () => { + assert.throws( + () => parseModelPolicy(JSON.stringify({ + version: 1, + roles: { worker: { model: "litellm/model", minThinkingLevel: "high", maxThinkingLevel: "low" } } + }), "range.json"), + /minThinkingLevel above maxThinkingLevel/ + ); + + assert.throws( + () => parseModelPolicy(JSON.stringify({ + version: 1, + roles: { worker: { model: "litellm/model", thinkingLevel: "low", maxThinkingLevel: "high" } } + }), "mixed.json"), + /cannot combine legacy 'thinkingLevel'/ + ); +}); + test("role router resolves a fresh physical model for a user turn", () => { const physical = { provider: "litellm", id: "coder-next" }; const route = createRoleRouter("worker", { From 14b9073ab3e45134c6314ffe4b101f26ca4ef34f Mon Sep 17 00:00:00 2001 From: BakerSean168 Date: Sun, 4 Oct 2026 08:21:50 +0000 Subject: [PATCH 3/5] feat(model-policy): add deterministic multi-model routing --- docs/architecture.md | 14 +- docs/model-policy.md | 114 ++++++++----- extension/index.js | 1 + extension/model-policy.js | 336 ++++++++++++++++++++++++++++--------- tests/model-policy.test.js | 322 +++++++++++++++++++++++++++++++++++ 5 files changed, 668 insertions(+), 119 deletions(-) diff --git a/docs/architecture.md b/docs/architecture.md index d5bbf73..c9c8537 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -56,9 +56,17 @@ ambient-extension discovery, so this keeps the `forgeflow/*` roles available in foreground, detached, nested, and recovery child sessions without hard-coding an installation path in operator profile settings. -New user/direct requests resolve the current role mapping. Continuation/retry -requests stay on the physical model already handling the turn to preserve cache and -reasoning-signature continuity. +New user/direct requests resolve the current role mapping. Policy v2 may expose several +physical candidates for one role; ForgeFlow deterministically selects exactly one from +explicit task class plus thinking-effort policy before provider execution. Pi then owns +the request lifecycle, while LiteLLM may choose among channels for that already-selected +physical model. Continuation/retry requests stay on the physical model already handling +the turn to preserve cache and reasoning-signature continuity. + +The router records its v2 decision in Pi's native virtual-model state rather than a +ForgeFlow database. Explicit task classes use `[[forgeflow:task=]]` in the delegated +prompt; ForgeFlow deliberately does not add a hidden LLM classifier or per-turn semantic +router. Project-local `.pi/forgeflow-models.json` is considered only when Pi reports the project trusted. User-level policy under `~/.pi/forgeflow-models.json` remains diff --git a/docs/model-policy.md b/docs/model-policy.md index 7f7a86e..70f4462 100644 --- a/docs/model-policy.md +++ b/docs/model-policy.md @@ -63,58 +63,96 @@ ForgeFlow resolves the physical mapping in this order: An explicit `FORGEFLOW_MODEL_POLICY` path is authoritative, must be absolute, and fails closed if it does not exist. Requiring an absolute operator path prevents a globally inherited relative environment value from resolving to a file supplied by an untrusted checkout. Project-local policy is ignored for untrusted projects so a checked-out repository cannot silently redirect prompts to another provider. -Policy schema v1: +### Policy schema v2 + +Version 2 makes a role a **candidate-route policy** instead of a permanent alias for one physical model. Each request still resolves to exactly one physical model before provider execution. ```json { - "version": 1, + "version": 2, "roles": { - "planner": { - "model": "gateway/frontier-reasoning", - "defaultThinkingLevel": "high", - "minThinkingLevel": "medium", - "maxThinkingLevel": "xhigh" - }, "worker": { - "model": "gateway/fast-coder", - "defaultThinkingLevel": "low", - "minThinkingLevel": "minimal", - "maxThinkingLevel": "medium" - }, - "reviewer": { - "model": "gateway/frontier-review", - "defaultThinkingLevel": "high", - "minThinkingLevel": "high", - "maxThinkingLevel": "xhigh" - }, - "scout": { - "model": "gateway/fast-general", - "defaultThinkingLevel": "low", - "minThinkingLevel": "off", - "maxThinkingLevel": "medium" - }, - "oracle": { - "model": "gateway/frontier-reasoning", - "defaultThinkingLevel": "xhigh", - "minThinkingLevel": "high", - "maxThinkingLevel": "xhigh" + "defaultRoute": "fast", + "defaultTaskClass": "mechanical", + "routes": [ + { + "id": "fast", + "model": "litellm/deepseek-v4-flash", + "taskClasses": ["mechanical", "implementation"], + "defaultThinkingLevel": "low", + "minThinkingLevel": "minimal", + "maxThinkingLevel": "low" + }, + { + "id": "standard", + "model": "litellm/gpt-6-astra", + "taskClasses": ["implementation", "debug"], + "defaultThinkingLevel": "medium", + "minThinkingLevel": "medium", + "maxThinkingLevel": "medium" + } + ] } } } ``` -`model` must be a fully qualified **physical** Pi model, `provider/model`. Model IDs may contain additional slashes. A role may not point at another `forgeflow/*` virtual model. +`model` must be a fully qualified **physical** Pi model, `provider/model`. A route may not point at another `forgeflow/*` virtual model. Route ids are stable operator-facing names, not model ids. + +Supported task classes are: + +- `recon` — repository search, inventory, fact gathering; +- `mechanical` — small, well-scoped, mechanically verifiable edits; +- `implementation` — ordinary feature implementation; +- `debug` — diagnosis and repair of a concrete failure; +- `review` — diff/acceptance/risk review; +- `architecture` — system design and cross-cutting tradeoffs; +- `product-judgment` — UX, intent, taste, and ambiguous product tradeoffs; +- `root-cause` — difficult cross-system diagnosis or causal analysis. + +A parent can explicitly classify a delegated request by putting this marker in the delegated user prompt: + +```text +[[forgeflow:task=debug]] +``` + +The marker is deterministic routing metadata. ForgeFlow does **not** use an LLM or keyword classifier to guess a task class. If the marker is absent, routes are ranked by how closely their thinking envelope matches the selected virtual thinking level; `defaultRoute` breaks equal-distance ties. `defaultTaskClass` is only the recorded semantic default for unclassified requests. + +When a task class is explicit, ForgeFlow only considers routes that declare that class (or a generic route with no `taskClasses` if no exact class route exists). It does not silently cross a semantic boundary such as `product-judgment` → `root-cause` merely because another model is available. + +If the best unclassified route names a model that is not present in Pi's physical model registry, ForgeFlow deterministically tries the next compatible route. Runtime provider/channel failures are **not** model fallback signals: once a launch has selected a physical model, LiteLLM owns channel failover for that model and Pi surfaces an exhausted model failure back to the parent. -Thinking effort is part of the role policy, not an afterthought. The adaptive fields are: +### Thinking effort + +Thinking effort is part of each route policy: + +- `defaultThinkingLevel`: used when a caller does not provide a level; +- `minThinkingLevel`: floor for the route; +- `maxThinkingLevel`: ceiling for the route; +- legacy `thinkingLevel`: an exact fixed pin retained for version-1 compatibility and single-route policies. + +The supported levels are `off`, `minimal`, `low`, `medium`, `high`, `xhigh`, and `max`. `max` should only be enabled when the physical model's registry metadata explicitly supports it. + +The selected virtual level participates in route choice. A `worker` at `low` can therefore resolve to a cheap implementation model while the same `forgeflow/worker` selected at `medium` resolves to a stronger implementation model. After route selection, the level is clamped to that route's configured envelope. + +### Decision telemetry + +Version 2 returns a small JSON-serializable routing decision as Pi virtual-model state. Pi stores that state on the session branch as its native `pi.virtual-model-state` entry, so ForgeFlow does not create a second telemetry database. A decision records: + +```text +role +routeId +taskClass +physical model +effective thinking level +selection basis (explicit task class or effort envelope) +``` -- `defaultThinkingLevel`: used when the caller does not request a level; -- `minThinkingLevel`: floor for the role, preventing an underpowered call; -- `maxThinkingLevel`: ceiling for the role, preventing routine work from consuming frontier effort; -- legacy `thinkingLevel`: an exact fixed pin retained for backwards compatibility. It cannot be combined with the adaptive fields. +This is enough to audit why a launch used a model and to build routing data later without introducing a learned router now. -When the caller explicitly selects a virtual thinking level, ForgeFlow preserves that request inside the configured envelope. For example a planner configured as `medium..xhigh` can run at `high` or escalate to `xhigh`, while a worker capped at `medium` cannot accidentally consume `xhigh`. `max` remains available in the schema, but operator policy should only enable it for a physical model whose registry metadata explicitly supports that level. Continuations and retries keep the already-selected physical model and effort for turn stability. +### Version 1 compatibility -The policy file is read when a new user/direct request is routed, so changing a mapping does not require changing ForgeFlow code. Pi still records the actual physical model on every assistant response. +Version 1 single-model policies remain valid. They keep their previous semantics and can be migrated gradually. Version 2 is preferred for operator policy because it separates a stable role from the physical model pool behind that role. ## Turn stickiness diff --git a/extension/index.js b/extension/index.js index 898b075..77034c1 100644 --- a/extension/index.js +++ b/extension/index.js @@ -10,6 +10,7 @@ const FORGEFLOW_EXTENSION_PATH = fileURLToPath(import.meta.url); const CORE_POLICY = [ "ForgeFlow is a thin engineering-governance extension for Pi; Pi and installed plugins own execution, sessions, delegation, review loops, acceptance gates, worktrees, missions, schedules, and resume.", "ForgeFlow may define stable logical model roles, but Pi owns virtual-model dispatch and the provider layer owns channel, credential, quota, and transport routing.", + "When delegating through forgeflow/* roles, classify the task explicitly when useful by prefixing the delegated prompt with [[forgeflow:task=]], where class is recon, mechanical, implementation, debug, review, architecture, product-judgment, or root-cause. Omit the marker when the role's effort-based default routing is desired.", "Do not create a second agent runtime, workflow database, provider gateway, reviewer runtime, PR gate, or duplicate subagent scheduler inside ForgeFlow.", "Keep one mutation writer per working tree. Independent reviewers must not mutate the candidate under review.", "Agent completion is evidence, not authority. For PR delivery, final acceptance must be tied to the authoritative current head and stale evidence must be invalidated after every new push.", diff --git a/extension/model-policy.js b/extension/model-policy.js index a0208e8..f01e0c6 100644 --- a/extension/model-policy.js +++ b/extension/model-policy.js @@ -11,8 +11,18 @@ export const FORGEFLOW_MODEL_ROLES = Object.freeze([ "scout", "oracle" ]); +export const FORGEFLOW_TASK_CLASSES = Object.freeze([ + "recon", + "mechanical", + "implementation", + "debug", + "review", + "architecture", + "product-judgment", + "root-cause" +]); -const THINKING_LEVELS = new Set([ +const THINKING_LEVELS = Object.freeze([ "off", "minimal", "low", @@ -21,6 +31,10 @@ const THINKING_LEVELS = new Set([ "xhigh", "max" ]); +const THINKING_LEVEL_SET = new Set(THINKING_LEVELS); +const TASK_CLASS_SET = new Set(FORGEFLOW_TASK_CLASSES); +const ROUTE_ID_RE = /^[a-z0-9][a-z0-9._-]{0,62}$/; +const TASK_CLASS_MARKER_RE = /\[\[forgeflow:task=([a-z0-9-]+)\]\]/i; function modelPolicyCandidates(cwd, env = process.env, home = homedir(), options = {}) { const explicit = env[FORGEFLOW_MODEL_POLICY_ENV]?.trim(); @@ -44,76 +58,59 @@ export function resolveModelPolicyPath(cwd, env = process.env, home = homedir(), return modelPolicyCandidates(cwd, env, home, options).find((path) => existsSync(path)); } -function validateRolePolicy(role, value, source) { - if (!value || typeof value !== "object" || Array.isArray(value)) { - throw new Error(`ForgeFlow model policy '${source}' has invalid role '${role}'; expected an object.`); - } - - if (typeof value.model !== "string" || !value.model.trim()) { - throw new Error(`ForgeFlow model policy '${source}' role '${role}' requires a non-empty 'model'.`); +function parseQualifiedPhysicalModel(rawModel, subject) { + if (typeof rawModel !== "string" || !rawModel.trim()) { + throw new Error(`${subject} requires a non-empty 'model'.`); } - const model = value.model.trim(); + const model = rawModel.trim(); const separator = model.indexOf("/"); if (separator <= 0 || separator === model.length - 1) { - throw new Error( - `ForgeFlow model policy '${source}' role '${role}' must use a qualified physical model 'provider/model'.` - ); + throw new Error(`${subject} must use a qualified physical model 'provider/model'.`); } const provider = model.slice(0, separator); const id = model.slice(separator + 1); if (provider === FORGEFLOW_MODEL_PROVIDER) { - throw new Error( - `ForgeFlow model policy '${source}' role '${role}' cannot route to another ForgeFlow virtual model.` - ); + throw new Error(`${subject} cannot route to another ForgeFlow virtual model.`); } + return { provider, id }; +} - function optionalThinkingLevel(field) { - if (value[field] === undefined) return undefined; - if (typeof value[field] !== "string" || !THINKING_LEVELS.has(value[field])) { - throw new Error( - `ForgeFlow model policy '${source}' role '${role}' has invalid '${field}'.` - ); - } - return value[field]; +function optionalThinkingLevel(value, field, subject) { + if (value[field] === undefined) return undefined; + if (typeof value[field] !== "string" || !THINKING_LEVEL_SET.has(value[field])) { + throw new Error(`${subject} has invalid '${field}'.`); } + return value[field]; +} - const thinkingLevel = optionalThinkingLevel("thinkingLevel"); - const defaultThinkingLevel = optionalThinkingLevel("defaultThinkingLevel"); - const minThinkingLevel = optionalThinkingLevel("minThinkingLevel"); - const maxThinkingLevel = optionalThinkingLevel("maxThinkingLevel"); +function validateThinkingPolicy(value, subject) { + const thinkingLevel = optionalThinkingLevel(value, "thinkingLevel", subject); + const defaultThinkingLevel = optionalThinkingLevel(value, "defaultThinkingLevel", subject); + const minThinkingLevel = optionalThinkingLevel(value, "minThinkingLevel", subject); + const maxThinkingLevel = optionalThinkingLevel(value, "maxThinkingLevel", subject); const adaptiveFields = [defaultThinkingLevel, minThinkingLevel, maxThinkingLevel]; if (thinkingLevel !== undefined && adaptiveFields.some((entry) => entry !== undefined)) { - throw new Error( - `ForgeFlow model policy '${source}' role '${role}' cannot combine legacy 'thinkingLevel' with adaptive thinking fields.` - ); + throw new Error(`${subject} cannot combine legacy 'thinkingLevel' with adaptive thinking fields.`); } - const rank = (level) => level === undefined ? undefined : [...THINKING_LEVELS].indexOf(level); + const rank = (level) => level === undefined ? undefined : THINKING_LEVELS.indexOf(level); const minRank = rank(minThinkingLevel); const maxRank = rank(maxThinkingLevel); const defaultRank = rank(defaultThinkingLevel); if (minRank !== undefined && maxRank !== undefined && minRank > maxRank) { - throw new Error( - `ForgeFlow model policy '${source}' role '${role}' has minThinkingLevel above maxThinkingLevel.` - ); + throw new Error(`${subject} has minThinkingLevel above maxThinkingLevel.`); } if (defaultRank !== undefined && minRank !== undefined && defaultRank < minRank) { - throw new Error( - `ForgeFlow model policy '${source}' role '${role}' has defaultThinkingLevel below minThinkingLevel.` - ); + throw new Error(`${subject} has defaultThinkingLevel below minThinkingLevel.`); } if (defaultRank !== undefined && maxRank !== undefined && defaultRank > maxRank) { - throw new Error( - `ForgeFlow model policy '${source}' role '${role}' has defaultThinkingLevel above maxThinkingLevel.` - ); + throw new Error(`${subject} has defaultThinkingLevel above maxThinkingLevel.`); } return { - provider, - id, ...(thinkingLevel !== undefined ? { thinkingLevel } : {}), ...(defaultThinkingLevel !== undefined ? { defaultThinkingLevel } : {}), ...(minThinkingLevel !== undefined ? { minThinkingLevel } : {}), @@ -121,6 +118,81 @@ function validateRolePolicy(role, value, source) { }; } +function validateV1RolePolicy(role, value, source) { + const subject = `ForgeFlow model policy '${source}' role '${role}'`; + if (!value || typeof value !== "object" || Array.isArray(value)) { + throw new Error(`${subject} must be an object.`); + } + return { + ...parseQualifiedPhysicalModel(value.model, subject), + ...validateThinkingPolicy(value, subject) + }; +} + +function validateTaskClasses(value, subject) { + if (value === undefined) return []; + if (!Array.isArray(value) || value.length === 0) { + throw new Error(`${subject} 'taskClasses' must be a non-empty array when provided.`); + } + const classes = []; + for (const taskClass of value) { + if (typeof taskClass !== "string" || !TASK_CLASS_SET.has(taskClass)) { + throw new Error(`${subject} has invalid task class '${String(taskClass)}'.`); + } + if (!classes.includes(taskClass)) classes.push(taskClass); + } + return classes; +} + +function validateV2RolePolicy(role, value, source) { + const subject = `ForgeFlow model policy '${source}' role '${role}'`; + if (!value || typeof value !== "object" || Array.isArray(value)) { + throw new Error(`${subject} must be an object.`); + } + if (!Array.isArray(value.routes) || value.routes.length === 0) { + throw new Error(`${subject} requires a non-empty 'routes' array.`); + } + + const routeIds = new Set(); + const routes = value.routes.map((route, index) => { + const routeSubject = `${subject} route #${index + 1}`; + if (!route || typeof route !== "object" || Array.isArray(route)) { + throw new Error(`${routeSubject} must be an object.`); + } + if (typeof route.id !== "string" || !ROUTE_ID_RE.test(route.id)) { + throw new Error(`${routeSubject} requires a safe non-empty 'id'.`); + } + if (routeIds.has(route.id)) { + throw new Error(`${subject} has duplicate route id '${route.id}'.`); + } + routeIds.add(route.id); + return { + routeId: route.id, + ...parseQualifiedPhysicalModel(route.model, `${subject} route '${route.id}'`), + taskClasses: validateTaskClasses(route.taskClasses, `${subject} route '${route.id}'`), + ...validateThinkingPolicy(route, `${subject} route '${route.id}'`) + }; + }); + + if (typeof value.defaultRoute !== "string" || !routeIds.has(value.defaultRoute)) { + throw new Error(`${subject} 'defaultRoute' must name one configured route.`); + } + + let defaultTaskClass; + if (value.defaultTaskClass !== undefined) { + if (typeof value.defaultTaskClass !== "string" || !TASK_CLASS_SET.has(value.defaultTaskClass)) { + throw new Error(`${subject} has invalid 'defaultTaskClass'.`); + } + defaultTaskClass = value.defaultTaskClass; + } + + return { + routes, + defaultRoute: value.defaultRoute, + ...(defaultTaskClass !== undefined ? { defaultTaskClass } : {}) + }; +} + export function parseModelPolicy(text, source = "") { let raw; try { @@ -134,8 +206,8 @@ export function parseModelPolicy(text, source = "") { if (!raw || typeof raw !== "object" || Array.isArray(raw)) { throw new Error(`ForgeFlow model policy '${source}' must be a JSON object.`); } - if (raw.version !== 1) { - throw new Error(`ForgeFlow model policy '${source}' requires version 1.`); + if (raw.version !== 1 && raw.version !== 2) { + throw new Error(`ForgeFlow model policy '${source}' requires version 1 or 2.`); } if (!raw.roles || typeof raw.roles !== "object" || Array.isArray(raw.roles)) { throw new Error(`ForgeFlow model policy '${source}' requires a 'roles' object.`); @@ -143,9 +215,10 @@ export function parseModelPolicy(text, source = "") { const roles = {}; for (const role of FORGEFLOW_MODEL_ROLES) { - if (raw.roles[role] !== undefined) { - roles[role] = validateRolePolicy(role, raw.roles[role], source); - } + if (raw.roles[role] === undefined) continue; + roles[role] = raw.version === 1 + ? validateV1RolePolicy(role, raw.roles[role], source) + : validateV2RolePolicy(role, raw.roles[role], source); } for (const role of Object.keys(raw.roles)) { @@ -154,7 +227,7 @@ export function parseModelPolicy(text, source = "") { } } - return { version: 1, roles }; + return { version: raw.version, roles }; } export function loadModelPolicy(cwd, env = process.env, home = homedir(), options = {}) { @@ -173,53 +246,173 @@ export function loadModelPolicy(cwd, env = process.env, home = homedir(), option } function resolveThinkingLevel(requested, target) { - if (target.thinkingLevel !== undefined) { - return target.thinkingLevel; - } + if (target.thinkingLevel !== undefined) return target.thinkingLevel; let level = requested ?? target.defaultThinkingLevel; if (level === undefined) return undefined; - const ordered = [...THINKING_LEVELS]; - const rank = ordered.indexOf(level); + const rank = THINKING_LEVELS.indexOf(level); if (rank === -1) return target.defaultThinkingLevel; const minRank = target.minThinkingLevel === undefined ? undefined - : ordered.indexOf(target.minThinkingLevel); + : THINKING_LEVELS.indexOf(target.minThinkingLevel); const maxRank = target.maxThinkingLevel === undefined ? undefined - : ordered.indexOf(target.maxThinkingLevel); + : THINKING_LEVELS.indexOf(target.maxThinkingLevel); if (minRank !== undefined && rank < minRank) level = target.minThinkingLevel; - if (maxRank !== undefined && ordered.indexOf(level) > maxRank) level = target.maxThinkingLevel; + if (maxRank !== undefined && THINKING_LEVELS.indexOf(level) > maxRank) level = target.maxThinkingLevel; return level; } +function thinkingDistance(requested, route) { + if (!requested || !THINKING_LEVEL_SET.has(requested)) return 0; + const effective = resolveThinkingLevel(requested, route) ?? requested; + return Math.abs(THINKING_LEVELS.indexOf(effective) - THINKING_LEVELS.indexOf(requested)); +} + +function messageText(message) { + if (!message || message.role !== "user") return ""; + if (typeof message.content === "string") return message.content; + if (!Array.isArray(message.content)) return ""; + return message.content + .map((part) => { + if (!part || typeof part !== "object") return ""; + if (typeof part.text === "string") return part.text; + if (typeof part.content === "string") return part.content; + return ""; + }) + .join("\n"); +} + +function explicitTaskClass(messages = []) { + for (let index = messages.length - 1; index >= 0; index -= 1) { + const message = messages[index]; + if (!message || message.role !== "user") continue; + const text = messageText(message); + const match = text.match(TASK_CLASS_MARKER_RE); + if (!match) return undefined; + const taskClass = match[1].toLowerCase(); + if (!TASK_CLASS_SET.has(taskClass)) { + throw new Error( + `ForgeFlow task marker names unknown class '${taskClass}'. Expected one of: ${FORGEFLOW_TASK_CLASSES.join(", ")}.` + ); + } + return taskClass; + } + return undefined; +} + +function routeCandidates(rolePolicy, requestedThinking, taskClass) { + let pool = rolePolicy.routes; + if (taskClass) { + const exact = pool.filter((route) => route.taskClasses.includes(taskClass)); + if (exact.length > 0) { + pool = exact; + } else { + const generic = pool.filter((route) => route.taskClasses.length === 0); + if (generic.length === 0) { + throw new Error(`No ForgeFlow route handles explicit task class '${taskClass}'.`); + } + pool = generic; + } + } + + const indexed = pool.map((route) => ({ + route, + distance: thinkingDistance(requestedThinking, route), + isDefault: route.routeId === rolePolicy.defaultRoute, + index: rolePolicy.routes.indexOf(route) + })); + indexed.sort((left, right) => { + if (left.distance !== right.distance) return left.distance - right.distance; + if (left.isDefault !== right.isDefault) return left.isDefault ? -1 : 1; + return left.index - right.index; + }); + return indexed.map((entry) => entry.route); +} + function stickyRoute(request) { if (request.reason === "retry" && request.failed) { return { model: request.failed.model, - thinkingLevel: request.failed.thinkingLevel ?? request.thinkingLevel + thinkingLevel: request.failed.thinkingLevel ?? request.thinkingLevel, + ...(request.state !== undefined ? { state: request.state } : {}) }; } if (request.reason === "continuation" && request.previous) { return { model: request.previous.model, - thinkingLevel: request.previous.thinkingLevel ?? request.thinkingLevel + thinkingLevel: request.previous.thinkingLevel ?? request.thinkingLevel, + ...(request.state !== undefined ? { state: request.state } : {}) }; } return undefined; } +function resolvePhysicalModel(ctx, route) { + const model = ctx.modelRegistry.find(route.provider, route.id); + if (!model || model.api === "pi-virtual") return undefined; + return model; +} + +function routeV1(role, request, ctx, path, target) { + const model = ctx.modelRegistry.find(target.provider, target.id); + if (!model) { + throw new Error( + `ForgeFlow role '${role}' targets unknown physical model '${target.provider}/${target.id}' from '${path}'.` + ); + } + if (model.api === "pi-virtual") { + throw new Error( + `ForgeFlow role '${role}' must resolve directly to a physical model; '${target.provider}/${target.id}' is virtual.` + ); + } + return { + model, + thinkingLevel: resolveThinkingLevel(request.thinkingLevel, target) + }; +} + +function routeV2(role, request, ctx, path, target) { + const taskClass = explicitTaskClass(request.messages); + const candidates = routeCandidates(target, request.thinkingLevel, taskClass); + const attempted = []; + + for (const route of candidates) { + attempted.push(`${route.provider}/${route.id}`); + const model = resolvePhysicalModel(ctx, route); + if (!model) continue; + + const thinkingLevel = resolveThinkingLevel(request.thinkingLevel, route); + return { + model, + thinkingLevel, + state: { + decisionVersion: 1, + policyVersion: 2, + role, + routeId: route.routeId, + taskClass: taskClass ?? target.defaultTaskClass ?? null, + model: `${route.provider}/${route.id}`, + thinkingLevel, + basis: taskClass ? "explicit-task-class" : "effort-envelope" + } + }; + } + + throw new Error( + `ForgeFlow role '${role}' has no available physical route in '${path}'. Tried: ${attempted.join(", ")}.` + ); +} + export function createRoleRouter(role, options = {}) { const loadPolicy = options.loadPolicy ?? loadModelPolicy; return function route(request, ctx) { const sticky = stickyRoute(request); - if (sticky) { - return sticky; - } + if (sticky) return sticky; const allowProject = typeof ctx.isProjectTrusted === "function" && ctx.isProjectTrusted(); const { path, policy } = loadPolicy(ctx.cwd, process.env, homedir(), { allowProject }); @@ -230,22 +423,9 @@ export function createRoleRouter(role, options = {}) { ); } - const model = ctx.modelRegistry.find(target.provider, target.id); - if (!model) { - throw new Error( - `ForgeFlow role '${role}' targets unknown physical model '${target.provider}/${target.id}' from '${path}'.` - ); - } - if (model.api === "pi-virtual") { - throw new Error( - `ForgeFlow role '${role}' must resolve directly to a physical model; '${target.provider}/${target.id}' is virtual.` - ); - } - - return { - model, - thinkingLevel: resolveThinkingLevel(request.thinkingLevel, target) - }; + return policy.version === 1 + ? routeV1(role, request, ctx, path, target) + : routeV2(role, request, ctx, path, target); }; } @@ -259,7 +439,7 @@ export function registerForgeFlowVirtualModels(pi, options = {}) { provider: FORGEFLOW_MODEL_PROVIDER, id: role, name: `ForgeFlow ${role}`, - thinkingLevels: ["off", "minimal", "low", "medium", "high", "xhigh", "max"], + thinkingLevels: THINKING_LEVELS, route: createRoleRouter(role, options) }); } diff --git a/tests/model-policy.test.js b/tests/model-policy.test.js index f5381ab..888831c 100644 --- a/tests/model-policy.test.js +++ b/tests/model-policy.test.js @@ -283,3 +283,325 @@ test("ForgeFlow registers stable logical roles as Pi virtual models", () => { assert.equal(typeof registration.route, "function"); } }); + +test("v2 model policy parses deterministic multi-model routes", () => { + const policy = parseModelPolicy(JSON.stringify({ + version: 2, + roles: { + worker: { + defaultRoute: "fast", + defaultTaskClass: "mechanical", + routes: [ + { + id: "fast", + model: "litellm/deepseek-v4-flash", + taskClasses: ["mechanical", "implementation"], + defaultThinkingLevel: "low", + minThinkingLevel: "minimal", + maxThinkingLevel: "low" + }, + { + id: "standard", + model: "litellm/gpt-6-astra", + taskClasses: ["implementation", "debug"], + defaultThinkingLevel: "medium", + minThinkingLevel: "medium", + maxThinkingLevel: "medium" + } + ] + } + } + }), "v2.json"); + + assert.equal(policy.version, 2); + assert.equal(policy.roles.worker.defaultRoute, "fast"); + assert.equal(policy.roles.worker.defaultTaskClass, "mechanical"); + assert.deepEqual( + policy.roles.worker.routes.map((route) => route.routeId), + ["fast", "standard"] + ); + assert.equal(policy.roles.worker.routes[1].provider, "litellm"); + assert.equal(policy.roles.worker.routes[1].id, "gpt-6-astra"); +}); + +test("v2 router selects a physical model by effort and records the decision in router state", () => { + const deepseek = { provider: "litellm", id: "deepseek-v4-flash", api: "openai-completions" }; + const astra = { provider: "litellm", id: "gpt-6-astra", api: "openai-responses" }; + const models = new Map([ + ["litellm/deepseek-v4-flash", deepseek], + ["litellm/gpt-6-astra", astra] + ]); + const route = createRoleRouter("worker", { + loadPolicy() { + return { + path: "/policy-v2.json", + policy: parseModelPolicy(JSON.stringify({ + version: 2, + roles: { + worker: { + defaultRoute: "fast", + defaultTaskClass: "mechanical", + routes: [ + { + id: "fast", + model: "litellm/deepseek-v4-flash", + taskClasses: ["mechanical", "implementation"], + minThinkingLevel: "minimal", + maxThinkingLevel: "low" + }, + { + id: "standard", + model: "litellm/gpt-6-astra", + taskClasses: ["implementation", "debug"], + minThinkingLevel: "medium", + maxThinkingLevel: "medium" + } + ] + } + } + })) + }; + } + }); + const ctx = { + cwd: "/repo", + isProjectTrusted: () => true, + modelRegistry: { + find(provider, id) { + return models.get(`${provider}/${id}`); + } + } + }; + + const low = route({ reason: "user", thinkingLevel: "low", messages: [] }, ctx); + assert.equal(low.model, deepseek); + assert.equal(low.thinkingLevel, "low"); + assert.deepEqual(low.state, { + decisionVersion: 1, + policyVersion: 2, + role: "worker", + routeId: "fast", + taskClass: "mechanical", + model: "litellm/deepseek-v4-flash", + thinkingLevel: "low", + basis: "effort-envelope" + }); + + const medium = route({ reason: "user", thinkingLevel: "medium", messages: [] }, ctx); + assert.equal(medium.model, astra); + assert.equal(medium.thinkingLevel, "medium"); + assert.equal(medium.state.routeId, "standard"); +}); + +test("explicit task-class marker overrides the cheap route and clamps effort inside the selected route", () => { + const deepseek = { provider: "litellm", id: "deepseek-v4-flash", api: "openai-completions" }; + const astra = { provider: "litellm", id: "gpt-6-astra", api: "openai-responses" }; + const route = createRoleRouter("worker", { + loadPolicy() { + return { + path: "/policy-v2.json", + policy: parseModelPolicy(JSON.stringify({ + version: 2, + roles: { + worker: { + defaultRoute: "fast", + routes: [ + { + id: "fast", + model: "litellm/deepseek-v4-flash", + taskClasses: ["mechanical", "implementation"], + minThinkingLevel: "minimal", + maxThinkingLevel: "low" + }, + { + id: "standard", + model: "litellm/gpt-6-astra", + taskClasses: ["implementation", "debug"], + minThinkingLevel: "medium", + maxThinkingLevel: "medium" + } + ] + } + } + })) + }; + } + }); + const ctx = { + cwd: "/repo", + isProjectTrusted: () => true, + modelRegistry: { + find(provider, id) { + if (`${provider}/${id}` === "litellm/deepseek-v4-flash") return deepseek; + if (`${provider}/${id}` === "litellm/gpt-6-astra") return astra; + } + } + }; + + const result = route({ + reason: "user", + thinkingLevel: "low", + messages: [{ role: "user", content: [{ type: "text", text: "[[forgeflow:task=debug]] fix the race" }] }] + }, ctx); + + assert.equal(result.model, astra); + assert.equal(result.thinkingLevel, "medium"); + assert.equal(result.state.taskClass, "debug"); + assert.equal(result.state.basis, "explicit-task-class"); +}); + +test("v2 router can select a different Oracle model for product judgement", () => { + const sol = { provider: "litellm", id: "gpt-6.1-sol", api: "openai-completions" }; + const opus = { provider: "litellm", id: "claude-opus-5-5", api: "openai-responses" }; + const route = createRoleRouter("oracle", { + loadPolicy() { + return { + path: "/policy-v2.json", + policy: parseModelPolicy(JSON.stringify({ + version: 2, + roles: { + oracle: { + defaultRoute: "reasoning", + defaultTaskClass: "root-cause", + routes: [ + { + id: "reasoning", + model: "litellm/gpt-6.1-sol", + taskClasses: ["root-cause", "architecture"], + minThinkingLevel: "high", + maxThinkingLevel: "xhigh" + }, + { + id: "product", + model: "litellm/claude-opus-5-5", + taskClasses: ["product-judgment"], + minThinkingLevel: "high", + maxThinkingLevel: "xhigh" + } + ] + } + } + })) + }; + } + }); + const ctx = { + cwd: "/repo", + isProjectTrusted: () => true, + modelRegistry: { + find(provider, id) { + if (`${provider}/${id}` === "litellm/gpt-6.1-sol") return sol; + if (`${provider}/${id}` === "litellm/claude-opus-5-5") return opus; + } + } + }; + + const result = route({ + reason: "user", + thinkingLevel: "high", + messages: [{ role: "user", content: "[[forgeflow:task=product-judgment]] choose the better UX tradeoff" }] + }, ctx); + assert.equal(result.model, opus); + assert.equal(result.state.routeId, "product"); +}); + +test("v2 route selection skips an unavailable candidate but does not cross an explicit task-class boundary", () => { + const fallback = { provider: "litellm", id: "available", api: "openai-completions" }; + const policy = parseModelPolicy(JSON.stringify({ + version: 2, + roles: { + worker: { + defaultRoute: "preferred", + routes: [ + { id: "preferred", model: "litellm/missing", minThinkingLevel: "low", maxThinkingLevel: "low" }, + { id: "available", model: "litellm/available", minThinkingLevel: "medium", maxThinkingLevel: "medium" } + ] + }, + oracle: { + defaultRoute: "reasoning", + routes: [ + { id: "reasoning", model: "litellm/available", taskClasses: ["root-cause"], minThinkingLevel: "high", maxThinkingLevel: "xhigh" }, + { id: "product", model: "litellm/missing", taskClasses: ["product-judgment"], minThinkingLevel: "high", maxThinkingLevel: "xhigh" } + ] + } + } + })); + const loadPolicy = () => ({ path: "/policy-v2.json", policy }); + const ctx = { + cwd: "/repo", + isProjectTrusted: () => true, + modelRegistry: { find: (_provider, id) => id === "available" ? fallback : undefined } + }; + + const worker = createRoleRouter("worker", { loadPolicy }); + const result = worker({ reason: "user", thinkingLevel: "low", messages: [] }, ctx); + assert.equal(result.model, fallback); + assert.equal(result.state.routeId, "available"); + assert.equal(result.thinkingLevel, "medium"); + + const oracle = createRoleRouter("oracle", { loadPolicy }); + assert.throws( + () => oracle({ + reason: "user", + thinkingLevel: "high", + messages: [{ role: "user", content: "[[forgeflow:task=product-judgment]] decide" }] + }, ctx), + /no available physical route/ + ); +}); + +test("v2 policy rejects bad route identities, task classes, and default routes", () => { + assert.throws( + () => parseModelPolicy(JSON.stringify({ + version: 2, + roles: { worker: { defaultRoute: "missing", routes: [{ id: "fast", model: "litellm/a" }] } } + }), "missing-default.json"), + /defaultRoute/ + ); + assert.throws( + () => parseModelPolicy(JSON.stringify({ + version: 2, + roles: { + worker: { + defaultRoute: "fast", + routes: [{ id: "fast", model: "litellm/a", taskClasses: ["mystery"] }] + } + } + }), "bad-class.json"), + /invalid task class/ + ); +}); + +test("task-class markers do not leak from an older user turn", () => { + const cheap = { provider: "litellm", id: "cheap", api: "openai-completions" }; + const expensive = { provider: "litellm", id: "expensive", api: "openai-completions" }; + const policy = parseModelPolicy(JSON.stringify({ + version: 2, + roles: { + worker: { + defaultRoute: "cheap", + defaultTaskClass: "mechanical", + routes: [ + { id: "cheap", model: "litellm/cheap", minThinkingLevel: "low", maxThinkingLevel: "low" }, + { id: "expensive", model: "litellm/expensive", taskClasses: ["debug"], minThinkingLevel: "medium", maxThinkingLevel: "medium" } + ] + } + } + })); + const route = createRoleRouter("worker", { loadPolicy: () => ({ path: "/policy.json", policy }) }); + const result = route({ + reason: "user", + thinkingLevel: "low", + messages: [ + { role: "user", content: "[[forgeflow:task=debug]] old turn" }, + { role: "assistant", content: [{ type: "text", text: "done" }] }, + { role: "user", content: "new unclassified turn" } + ] + }, { + cwd: "/repo", + isProjectTrusted: () => true, + modelRegistry: { find: (_provider, id) => id === "cheap" ? cheap : expensive } + }); + assert.equal(result.model, cheap); + assert.equal(result.state.basis, "effort-envelope"); +}); From 93729e532245b71793172648e913eec7fecd7fd8 Mon Sep 17 00:00:00 2001 From: BakerSean168 Date: Sun, 4 Oct 2026 08:26:17 +0000 Subject: [PATCH 4/5] feat(worker): prefer Codex Team GPT-6.1 Sol --- docs/model-policy.md | 4 ++- tests/model-policy.test.js | 73 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 76 insertions(+), 1 deletion(-) diff --git a/docs/model-policy.md b/docs/model-policy.md index 70f4462..30b7972 100644 --- a/docs/model-policy.md +++ b/docs/model-policy.md @@ -25,7 +25,7 @@ ForgeFlow registers these Pi virtual models: | Role | Purpose | Typical effort envelope | | --- | --- | --- | | `forgeflow/planner` | decompose work, architecture, execution planning | `high` by default, up to `xhigh` | -| `forgeflow/worker` | implementation, focused debugging, routine code changes | `low` by default, up to `medium` | +| `forgeflow/worker` | implementation, focused debugging, routine code changes | prefers Codex Team `gpt-6.1-sol` at `medium`; falls back by policy | | `forgeflow/reviewer` | diff review, acceptance reasoning, risk checks | `high` by default, up to `xhigh` | | `forgeflow/scout` | repository exploration, cheap search, fact gathering | `low` by default, up to `medium` | | `forgeflow/oracle` | expensive expert escalation for ambiguous, cross-system, or hard root-cause questions | `xhigh` by default and capped at `xhigh` | @@ -137,6 +137,8 @@ The selected virtual level participates in route choice. A `worker` at `low` can ### Decision telemetry +For the current operator policy, `forgeflow/worker` prefers the native `openai-codex/gpt-6.1-sol` route. This is the Codex Team/Media execution lane; `Media` is the product/quota lane name, while the physical Pi model id remains `openai-codex/gpt-6.1-sol`. If that physical model is not present in Pi's registry, deterministic routing falls through to the configured LiteLLM worker routes without changing provider/channel ownership. + Version 2 returns a small JSON-serializable routing decision as Pi virtual-model state. Pi stores that state on the session branch as its native `pi.virtual-model-state` entry, so ForgeFlow does not create a second telemetry database. A decision records: ```text diff --git a/tests/model-policy.test.js b/tests/model-policy.test.js index 888831c..e3cb442 100644 --- a/tests/model-policy.test.js +++ b/tests/model-policy.test.js @@ -605,3 +605,76 @@ test("task-class markers do not leak from an older user turn", () => { assert.equal(result.model, cheap); assert.equal(result.state.basis, "effort-envelope"); }); + +test("worker policy can prefer native Codex Team GPT-6.1 Sol and fall back when unavailable", () => { + const codex = { provider: "openai-codex", id: "gpt-6.1-sol", api: "openai-codex-responses" }; + const astra = { provider: "litellm", id: "gpt-6-astra", api: "openai-responses" }; + const deepseek = { provider: "litellm", id: "deepseek-v4-flash", api: "openai-completions" }; + const policy = parseModelPolicy(JSON.stringify({ + version: 2, + roles: { + worker: { + defaultRoute: "codex-team", + defaultTaskClass: "implementation", + routes: [ + { + id: "codex-team", + model: "openai-codex/gpt-6.1-sol", + taskClasses: ["mechanical", "implementation", "debug"], + defaultThinkingLevel: "medium", + minThinkingLevel: "low", + maxThinkingLevel: "high" + }, + { + id: "standard", + model: "litellm/gpt-6-astra", + taskClasses: ["implementation", "debug"], + minThinkingLevel: "medium", + maxThinkingLevel: "medium" + }, + { + id: "fast", + model: "litellm/deepseek-v4-flash", + taskClasses: ["mechanical", "implementation"], + minThinkingLevel: "minimal", + maxThinkingLevel: "low" + } + ] + } + } + })); + const route = createRoleRouter("worker", { loadPolicy: () => ({ path: "/policy.json", policy }) }); + + const withCodex = { + cwd: "/repo", + isProjectTrusted: () => true, + modelRegistry: { + find(provider, id) { + if (`${provider}/${id}` === "openai-codex/gpt-6.1-sol") return codex; + if (`${provider}/${id}` === "litellm/gpt-6-astra") return astra; + if (`${provider}/${id}` === "litellm/deepseek-v4-flash") return deepseek; + } + } + }; + const preferred = route({ reason: "user", thinkingLevel: "low", messages: [] }, withCodex); + assert.equal(preferred.model, codex); + assert.equal(preferred.thinkingLevel, "low"); + assert.equal(preferred.state.routeId, "codex-team"); + + const withoutCodex = { + ...withCodex, + modelRegistry: { + find(provider, id) { + if (`${provider}/${id}` === "litellm/gpt-6-astra") return astra; + if (`${provider}/${id}` === "litellm/deepseek-v4-flash") return deepseek; + } + } + }; + const fallbackLow = route({ reason: "user", thinkingLevel: "low", messages: [] }, withoutCodex); + assert.equal(fallbackLow.model, deepseek); + assert.equal(fallbackLow.state.routeId, "fast"); + + const fallbackMedium = route({ reason: "user", thinkingLevel: "medium", messages: [] }, withoutCodex); + assert.equal(fallbackMedium.model, astra); + assert.equal(fallbackMedium.state.routeId, "standard"); +}); From 214b1a4b1c326f90ccdabebe8a4abdb18b99eabb Mon Sep 17 00:00:00 2001 From: BakerSean168 Date: Sun, 4 Oct 2026 12:56:18 +0000 Subject: [PATCH 5/5] feat(model-policy): add supply priority and usage telemetry --- README.md | 10 +- bin/forgeflow-usage.mjs | 70 +++++++ docs/architecture.md | 30 +-- docs/model-policy.md | 83 ++++++-- extension/index.js | 12 +- extension/model-policy.js | 390 ++++++++++++++++++++++++++++++++++--- extension/usage.js | 159 +++++++++++++++ package-lock.json | 3 + package.json | 43 +++- tests/extension.test.js | 1 + tests/model-policy.test.js | 185 ++++++++++++++++++ tests/usage.test.js | 114 +++++++++++ 12 files changed, 1031 insertions(+), 69 deletions(-) create mode 100755 bin/forgeflow-usage.mjs create mode 100644 extension/usage.js create mode 100644 tests/usage.test.js diff --git a/README.md b/README.md index 8e66ce0..4307478 100644 --- a/README.md +++ b/README.md @@ -27,8 +27,14 @@ as `test-driven-development`, `spec-driven-development`, `pr-gate`, and does not wrap or duplicate those surfaces. Provider infrastructure remains below Pi model selection. ForgeFlow logical roles -may resolve to physical models, while endpoint selection, credentials, channel -health, weights, quotas, and transport belong to the provider layer. +choose model capability and thinking effort; policy v3 can also order equivalent quota +sources (for example Business Team before a commercial relay). Endpoint selection, +credentials, channel health, and transport remain provider-layer concerns. LiteLLM still +owns commercial-channel selection after ForgeFlow has selected the commercial physical +model. + +ForgeFlow records credential-free local model/supply usage in `~/.pi/forgeflow-usage.jsonl`; +`forgeflow-usage` summarizes it. LiteLLM SpendLogs remain authoritative for relay-channel spend. ## Antigravity delegation diff --git a/bin/forgeflow-usage.mjs b/bin/forgeflow-usage.mjs new file mode 100755 index 0000000..0ee4bcc --- /dev/null +++ b/bin/forgeflow-usage.mjs @@ -0,0 +1,70 @@ +#!/usr/bin/env node +import { existsSync } from "node:fs"; +import { readUsageEvents, resolveUsageLogPath } from "../extension/usage.js"; + +const args = new Set(process.argv.slice(2)); +const json = args.has("--json"); +const path = resolveUsageLogPath(); +if (!existsSync(path)) { + console.error(`ForgeFlow usage log not found: ${path}`); + process.exitCode = 1; +} else { + const events = readUsageEvents(path); + const usageEvents = events.filter((event) => event.event === "model_usage"); + const decisions = events.filter((event) => event.event === "route_decision"); + const groups = new Map(); + + for (const event of usageEvents) { + const supply = event.supplyId ?? event.provider ?? "unknown"; + const key = `${supply}\u0000${event.model ?? "unknown"}`; + const row = groups.get(key) ?? { + supply, + model: event.model ?? "unknown", + requests: 0, + failures: 0, + tokens: 0, + input: 0, + output: 0, + cacheRead: 0, + reasoning: 0, + cost: 0 + }; + row.requests += 1; + if (event.status !== "success") row.failures += 1; + row.tokens += Number(event.usage?.totalTokens ?? 0); + row.input += Number(event.usage?.input ?? 0); + row.output += Number(event.usage?.output ?? 0); + row.cacheRead += Number(event.usage?.cacheRead ?? 0); + row.reasoning += Number(event.usage?.reasoning ?? 0); + row.cost += Number(event.usage?.cost?.total ?? 0); + groups.set(key, row); + } + + const failovers = new Map(); + for (const event of decisions) { + if (event.basis !== "supply-failover") continue; + const reason = event.failoverReason ?? "unknown"; + failovers.set(reason, (failovers.get(reason) ?? 0) + 1); + } + + const result = { + path, + usage: [...groups.values()].sort((a, b) => b.tokens - a.tokens || b.requests - a.requests), + failovers: [...failovers.entries()].map(([reason, count]) => ({ reason, count })).sort((a, b) => b.count - a.count) + }; + + if (json) { + console.log(JSON.stringify(result, null, 2)); + } else { + console.log(`ForgeFlow usage: ${path}`); + console.log("SUPPLY\tMODEL\tREQUESTS\tFAIL\tTOKENS\tCOST"); + for (const row of result.usage) { + console.log(`${row.supply}\t${row.model}\t${row.requests}\t${row.failures}\t${row.tokens}\t${row.cost.toFixed(6)}`); + } + if (result.failovers.length > 0) { + console.log("\nFAILOVER_REASON\tCOUNT"); + for (const row of result.failovers) console.log(`${row.reason}\t${row.count}`); + } + } +} + diff --git a/docs/architecture.md b/docs/architecture.md index c9c8537..e50a734 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -56,17 +56,23 @@ ambient-extension discovery, so this keeps the `forgeflow/*` roles available in foreground, detached, nested, and recovery child sessions without hard-coding an installation path in operator profile settings. -New user/direct requests resolve the current role mapping. Policy v2 may expose several -physical candidates for one role; ForgeFlow deterministically selects exactly one from -explicit task class plus thinking-effort policy before provider execution. Pi then owns -the request lifecycle, while LiteLLM may choose among channels for that already-selected -physical model. Continuation/retry requests stay on the physical model already handling -the turn to preserve cache and reasoning-signature continuity. - -The router records its v2 decision in Pi's native virtual-model state rather than a -ForgeFlow database. Explicit task classes use `[[forgeflow:task=]]` in the delegated -prompt; ForgeFlow deliberately does not add a hidden LLM classifier or per-turn semantic -router. +New user/direct requests resolve the current role mapping. Policy v3 separates model +capability from supply: a role chooses a logical model/effort, then an ordered supply +group chooses the physical Pi model that pays for it. The current GPT-6.1 Sol supply +order is Business Team (`openai-codex`) before the commercial relay (`litellm`). +LiteLLM then chooses only among channels for the already-selected commercial physical +model. + +Continuations remain sticky. Retries remain sticky unless the failed response is +classified as a narrow supply failure (quota/rate/capacity/upstream/transport/auth/model +availability); only then may v3 move to the next source in the same supply group. Context, +request, policy, and tool/schema failures do not consume the next paid source. + +The router records its decision in Pi's native virtual-model state. ForgeFlow additionally +maintains a credential-free local usage JSONL projection for subscription/native-provider +traffic, while LiteLLM SpendLogs remain authoritative for commercial relay/channel spend. +Explicit task classes use `[[forgeflow:task=]]`; ForgeFlow deliberately does not add +a hidden LLM classifier or per-turn semantic router. Project-local `.pi/forgeflow-models.json` is considered only when Pi reports the project trusted. User-level policy under `~/.pi/forgeflow-models.json` remains @@ -77,7 +83,7 @@ Provider gateways such as LiteLLM sit below this boundary: ForgeFlow chooses a logical role and Pi resolves its physical model; the provider plane chooses the endpoint/channel/key for that already-selected physical model. -See `docs/model-policy.md` for the v1 schema and role mapping. +See `docs/model-policy.md` for policy schemas, supply priority, failure classification, and usage projection. ## Invariant preflight diff --git a/docs/model-policy.md b/docs/model-policy.md index 30b7972..23dfeb6 100644 --- a/docs/model-policy.md +++ b/docs/model-policy.md @@ -7,13 +7,15 @@ The model path is: ```text ForgeFlow role policy ↓ -Pi virtual model +model capability + thinking effort ↓ -physical model +supply policy (subscription/team before commercial relay) + ↓ +Pi physical model ↓ provider gateway / native provider ↓ -channel, credential, quota, transport +channel, credential, transport ``` This keeps fast-moving physical model names out of ForgeFlow workflows and agent definitions. A model upgrade should normally be a policy-file change, not a ForgeFlow code change. @@ -137,7 +139,7 @@ The selected virtual level participates in route choice. A `worker` at `low` can ### Decision telemetry -For the current operator policy, `forgeflow/worker` prefers the native `openai-codex/gpt-6.1-sol` route. This is the Codex Team/Media execution lane; `Media` is the product/quota lane name, while the physical Pi model id remains `openai-codex/gpt-6.1-sol`. If that physical model is not present in Pi's registry, deterministic routing falls through to the configured LiteLLM worker routes without changing provider/channel ownership. +For the current operator policy, GPT-6.1 Sol roles use the `gpt-6.1-sol` supply group. Its first source is native `openai-codex/gpt-6.1-sol` (Business Team / Codex Team; `Media` is the product/quota lane label), followed by `litellm/gpt-6.1-sol` as the commercial relay. If the native source is absent from Pi's registry, a new request selects the commercial source immediately; if it fails at runtime with a qualifying supply error, retry advances to the commercial source. Version 2 returns a small JSON-serializable routing decision as Pi virtual-model state. Pi stores that state on the session branch as its native `pi.virtual-model-state` entry, so ForgeFlow does not create a second telemetry database. A decision records: @@ -152,23 +154,76 @@ selection basis (explicit task class or effort envelope) This is enough to audit why a launch used a model and to build routing data later without introducing a learned router now. -### Version 1 compatibility +### Policy schema v3: supply priority -Version 1 single-model policies remain valid. They keep their previous semantics and can be migrated gradually. Version 2 is preferred for operator policy because it separates a stable role from the physical model pool behind that role. +Version 3 separates **what model capability the role needs** from **which quota source should pay for that model**. A supply group is an ordered list of physical Pi models that represent the same logical model capability: -## Turn stickiness +```json +{ + "version": 3, + "supplies": { + "gpt-6.1-sol": { + "sources": [ + { + "id": "business-team", + "model": "openai-codex/gpt-6.1-sol", + "priority": 10 + }, + { + "id": "commercial-relay", + "model": "litellm/gpt-6.1-sol", + "priority": 20 + } + ] + } + }, + "roles": { + "worker": { + "defaultRoute": "frontier", + "defaultTaskClass": "implementation", + "routes": [ + { + "id": "frontier", + "supplyGroup": "gpt-6.1-sol", + "taskClasses": ["mechanical", "implementation", "debug"], + "defaultThinkingLevel": "medium", + "minThinkingLevel": "low", + "maxThinkingLevel": "high" + } + ] + } + } +} +``` + +For a new request ForgeFlow tries the lowest numeric supply priority whose physical model exists in Pi's model registry. With the operator policy above, `openai-codex/gpt-6.1-sol` consumes Business Team/Codex Team quota first; `litellm/gpt-6.1-sol` is the commercial-relay source. The `Media` label is a product/quota lane name rather than a distinct physical model id. + +A supply failover is allowed only on a narrow set of source-specific failures: rate limiting, exhausted quota/credits, capacity/overload, 502/503/504-style upstream failures, transport failures, authentication expiry, or source-specific model unavailability. Context overflow, malformed/bad requests, content-policy failures, tool/schema errors, and other request-semantic failures stay on the current source. This prevents a bad request from silently burning a second paid channel. -Pi documents that continuation requests should normally remain on the model that handled the turn, and retries should remain on the failed request's physical model unless a router deliberately performs a recovery switch. +Continuation requests stay on the last successful physical model. On a retry caused by a qualifying supply failure, ForgeFlow moves only to the next source in that logical supply group and records the reason in virtual-model state. Once the commercial relay is selected, LiteLLM remains responsible for endpoint/channel routing inside that physical model group; ForgeFlow does not select ACS versus 4Router. -ForgeFlow v1 follows that conservative rule: +### Usage projection -- `continuation` → reuse `previous`; -- `retry` → reuse `failed`; -- new `user` or `direct` request → resolve the current role policy again. +ForgeFlow writes a compact local JSONL usage projection to `~/.pi/forgeflow-usage.jsonl` by default. Override it with the absolute `FORGEFLOW_USAGE_LOG` path. The projection contains only routing metadata and Pi's normalized assistant usage; it never records credentials or request payloads. + +Each v3 route decision records the role, route, logical model, supply id/priority, physical model, effective thinking level, and any supply-failover reason. Each assistant `message_end` records provider/model, stop status, normalized input/output/cache/reasoning token counts, and Pi's normalized cost fields. Use: + +```bash +forgeflow-usage +forgeflow-usage --json +``` + +Business Team usage is therefore visible from Pi/ForgeFlow even though it bypasses LiteLLM. Commercial API truth remains in LiteLLM's `LiteLLM_SpendLogs`, which additionally identifies the actual relay endpoint/deployment such as ACS or 4Router. The two logs intentionally preserve their own authority rather than pretending a subscription-backed Codex request passed through LiteLLM. + +### Version 1/2 compatibility + +Version 1 single-model policies and version 2 multi-route policies remain valid. Version 3 is preferred for operator policy when one logical model can be supplied by multiple quota sources. + +## Turn stickiness -This preserves provider prompt caches and reasoning signatures inside a turn while still allowing model policy to change between user turns. +Pi documents that continuation requests should normally remain on the model that handled the turn. ForgeFlow keeps that rule. Version 1/2 retries also stay on the failed physical model. Version 3 permits one explicit exception: a retry whose failure is classified as supply-specific may move to the next source in the same logical supply group. -Cross-model retry escalation can be added later as an explicit policy version rather than hidden fallback behavior. +This preserves provider prompt caches and reasoning signatures during healthy execution while making quota/capacity exhaustion recoverable without re-launching the whole child. The switch is source failover, not semantic model replacement. ## Provider boundary diff --git a/extension/index.js b/extension/index.js index 77034c1..3f756f7 100644 --- a/extension/index.js +++ b/extension/index.js @@ -4,6 +4,7 @@ import { registerRequiredChildExtensions } from "pi-subagents/required-child-ext import { registerAntigravityAgents } from "./antigravity.js"; import { renderPreflight } from "./invariants.js"; import { registerForgeFlowVirtualModels } from "./model-policy.js"; +import { createUsageRecorder } from "./usage.js"; const FORGEFLOW_EXTENSION_PATH = fileURLToPath(import.meta.url); @@ -22,7 +23,8 @@ export function buildForgeFlowPromptSection(prompt) { } export default function registerForgeFlow(pi) { - registerForgeFlowVirtualModels(pi); + const usageRecorder = createUsageRecorder(); + registerForgeFlowVirtualModels(pi, { onDecision: usageRecorder.recordDecision }); let requiredChildRegistration; let antigravityRegistration; @@ -31,6 +33,14 @@ export default function registerForgeFlow(pi) { event.systemPromptOptions.sections.forgeflow_policy = buildForgeFlowPromptSection(event.prompt); }); + pi.on("message_end", (event, ctx) => { + try { + usageRecorder.recordMessage(event, ctx); + } catch { + // Observability is best-effort and must never break an agent turn. + } + }); + pi.on("session_start", (_event, ctx) => { requiredChildRegistration?.dispose(); antigravityRegistration?.dispose(); diff --git a/extension/model-policy.js b/extension/model-policy.js index f01e0c6..d36437b 100644 --- a/extension/model-policy.js +++ b/extension/model-policy.js @@ -34,6 +34,8 @@ const THINKING_LEVELS = Object.freeze([ const THINKING_LEVEL_SET = new Set(THINKING_LEVELS); const TASK_CLASS_SET = new Set(FORGEFLOW_TASK_CLASSES); const ROUTE_ID_RE = /^[a-z0-9][a-z0-9._-]{0,62}$/; +const SUPPLY_ID_RE = ROUTE_ID_RE; +const SUPPLY_GROUP_RE = /^[a-z0-9][a-z0-9._-]{0,95}$/; const TASK_CLASS_MARKER_RE = /\[\[forgeflow:task=([a-z0-9-]+)\]\]/i; function modelPolicyCandidates(cwd, env = process.env, home = homedir(), options = {}) { @@ -193,6 +195,121 @@ function validateV2RolePolicy(role, value, source) { }; } +function validateSupplyGroups(rawSupplies, source) { + if (rawSupplies === undefined) return {}; + if (!rawSupplies || typeof rawSupplies !== "object" || Array.isArray(rawSupplies)) { + throw new Error(`ForgeFlow model policy '${source}' 'supplies' must be an object.`); + } + + const supplies = {}; + for (const [groupId, value] of Object.entries(rawSupplies)) { + const subject = `ForgeFlow model policy '${source}' supply group '${groupId}'`; + if (!SUPPLY_GROUP_RE.test(groupId)) { + throw new Error(`${subject} has an invalid group id.`); + } + if (!value || typeof value !== "object" || Array.isArray(value)) { + throw new Error(`${subject} must be an object.`); + } + if (!Array.isArray(value.sources) || value.sources.length === 0) { + throw new Error(`${subject} requires a non-empty 'sources' array.`); + } + + const sourceIds = new Set(); + const sources = value.sources.map((entry, index) => { + const sourceSubject = `${subject} source #${index + 1}`; + if (!entry || typeof entry !== "object" || Array.isArray(entry)) { + throw new Error(`${sourceSubject} must be an object.`); + } + if (typeof entry.id !== "string" || !SUPPLY_ID_RE.test(entry.id)) { + throw new Error(`${sourceSubject} requires a safe non-empty 'id'.`); + } + if (sourceIds.has(entry.id)) { + throw new Error(`${subject} has duplicate source id '${entry.id}'.`); + } + sourceIds.add(entry.id); + const priority = entry.priority ?? (index + 1) * 10; + if (!Number.isSafeInteger(priority) || priority < 0) { + throw new Error(`${sourceSubject} requires a non-negative integer 'priority'.`); + } + return { + supplyId: entry.id, + priority, + ...parseQualifiedPhysicalModel(entry.model, `${subject} source '${entry.id}'`) + }; + }); + + sources.sort((left, right) => left.priority - right.priority || left.supplyId.localeCompare(right.supplyId)); + supplies[groupId] = { groupId, sources }; + } + return supplies; +} + +function validateV3RolePolicy(role, value, source, supplies) { + const subject = `ForgeFlow model policy '${source}' role '${role}'`; + if (!value || typeof value !== "object" || Array.isArray(value)) { + throw new Error(`${subject} must be an object.`); + } + if (!Array.isArray(value.routes) || value.routes.length === 0) { + throw new Error(`${subject} requires a non-empty 'routes' array.`); + } + + const routeIds = new Set(); + const routes = value.routes.map((route, index) => { + const routeSubject = `${subject} route #${index + 1}`; + if (!route || typeof route !== "object" || Array.isArray(route)) { + throw new Error(`${routeSubject} must be an object.`); + } + if (typeof route.id !== "string" || !ROUTE_ID_RE.test(route.id)) { + throw new Error(`${routeSubject} requires a safe non-empty 'id'.`); + } + if (routeIds.has(route.id)) { + throw new Error(`${subject} has duplicate route id '${route.id}'.`); + } + routeIds.add(route.id); + + const hasModel = route.model !== undefined; + const hasSupplyGroup = route.supplyGroup !== undefined; + if (hasModel === hasSupplyGroup) { + throw new Error(`${subject} route '${route.id}' requires exactly one of 'model' or 'supplyGroup'.`); + } + + let target; + if (hasSupplyGroup) { + if (typeof route.supplyGroup !== "string" || !Object.hasOwn(supplies, route.supplyGroup)) { + throw new Error(`${subject} route '${route.id}' references unknown supply group '${String(route.supplyGroup)}'.`); + } + target = { supplyGroup: route.supplyGroup }; + } else { + target = parseQualifiedPhysicalModel(route.model, `${subject} route '${route.id}'`); + } + + return { + routeId: route.id, + ...target, + taskClasses: validateTaskClasses(route.taskClasses, `${subject} route '${route.id}'`), + ...validateThinkingPolicy(route, `${subject} route '${route.id}'`) + }; + }); + + if (typeof value.defaultRoute !== "string" || !routeIds.has(value.defaultRoute)) { + throw new Error(`${subject} 'defaultRoute' must name one configured route.`); + } + + let defaultTaskClass; + if (value.defaultTaskClass !== undefined) { + if (typeof value.defaultTaskClass !== "string" || !TASK_CLASS_SET.has(value.defaultTaskClass)) { + throw new Error(`${subject} has invalid 'defaultTaskClass'.`); + } + defaultTaskClass = value.defaultTaskClass; + } + + return { + routes, + defaultRoute: value.defaultRoute, + ...(defaultTaskClass !== undefined ? { defaultTaskClass } : {}) + }; +} + export function parseModelPolicy(text, source = "") { let raw; try { @@ -206,19 +323,22 @@ export function parseModelPolicy(text, source = "") { if (!raw || typeof raw !== "object" || Array.isArray(raw)) { throw new Error(`ForgeFlow model policy '${source}' must be a JSON object.`); } - if (raw.version !== 1 && raw.version !== 2) { - throw new Error(`ForgeFlow model policy '${source}' requires version 1 or 2.`); + if (raw.version !== 1 && raw.version !== 2 && raw.version !== 3) { + throw new Error(`ForgeFlow model policy '${source}' requires version 1, 2, or 3.`); } if (!raw.roles || typeof raw.roles !== "object" || Array.isArray(raw.roles)) { throw new Error(`ForgeFlow model policy '${source}' requires a 'roles' object.`); } + const supplies = raw.version === 3 ? validateSupplyGroups(raw.supplies, source) : {}; const roles = {}; for (const role of FORGEFLOW_MODEL_ROLES) { if (raw.roles[role] === undefined) continue; roles[role] = raw.version === 1 ? validateV1RolePolicy(role, raw.roles[role], source) - : validateV2RolePolicy(role, raw.roles[role], source); + : raw.version === 2 + ? validateV2RolePolicy(role, raw.roles[role], source) + : validateV3RolePolicy(role, raw.roles[role], source, supplies); } for (const role of Object.keys(raw.roles)) { @@ -227,7 +347,7 @@ export function parseModelPolicy(text, source = "") { } } - return { version: raw.version, roles }; + return { version: raw.version, roles, ...(raw.version === 3 ? { supplies } : {}) }; } export function loadModelPolicy(cwd, env = process.env, home = homedir(), options = {}) { @@ -333,14 +453,7 @@ function routeCandidates(rolePolicy, requestedThinking, taskClass) { return indexed.map((entry) => entry.route); } -function stickyRoute(request) { - if (request.reason === "retry" && request.failed) { - return { - model: request.failed.model, - thinkingLevel: request.failed.thinkingLevel ?? request.thinkingLevel, - ...(request.state !== undefined ? { state: request.state } : {}) - }; - } +function continuationRoute(request) { if (request.reason === "continuation" && request.previous) { return { model: request.previous.model, @@ -351,8 +464,76 @@ function stickyRoute(request) { return undefined; } -function resolvePhysicalModel(ctx, route) { - const model = ctx.modelRegistry.find(route.provider, route.id); +function failedStickyRoute(request) { + if (request.reason !== "retry" || !request.failed) return undefined; + return { + model: request.failed.model, + thinkingLevel: request.failed.thinkingLevel ?? request.thinkingLevel, + ...(request.state !== undefined ? { state: request.state } : {}) + }; +} + +function errorText(messageOrText) { + if (typeof messageOrText === "string") return messageOrText; + if (!messageOrText || typeof messageOrText !== "object") return ""; + const values = [ + messageOrText.errorMessage, + messageOrText.rawStopReason, + ...(Array.isArray(messageOrText.diagnostics) + ? messageOrText.diagnostics.flatMap((entry) => [entry?.error?.message, entry?.error?.code]) + : []) + ]; + return values.filter((value) => typeof value === "string" || typeof value === "number").join(" "); +} + +/** + * Classify only failures that are plausibly specific to the current supply source. + * Semantic/request failures deliberately return undefined so retries stay sticky. + */ +export function classifySupplyFailure(messageOrText) { + const text = errorText(messageOrText).toLowerCase(); + if (!text) return undefined; + + const requestFailures = [ + /context.{0,24}(length|window|limit|overflow|too long)/, + /maximum context/, + /too many (input )?tokens/, + /content.?policy/, + /safety (policy|filter|violation)/, + /invalid[_ -]?request/, + /bad request/, + /malformed/, + /invalid (tool|schema|json|parameter|argument)/, + /tool.{0,24}(error|invalid|schema)/ + ]; + if (requestFailures.some((pattern) => pattern.test(text))) return undefined; + + if (/(?:http\s*)?429\b|too many requests|rate[_ -]?limit|rate limit|throttl/.test(text)) { + return "rate_limited"; + } + if (/insufficient[_ -]?quota|quota.{0,32}(exceed|exhaust|limit|deplet)|usage.{0,24}(limit|exhaust)|credit.{0,24}(exhaust|deplet|insufficient)|plan.{0,24}(limit|exhaust)/.test(text)) { + return "quota_exhausted"; + } + if (/overloaded|over capacity|capacity.{0,24}(exceed|unavailable|full)|server busy|temporar(?:y|ily) unavailable/.test(text)) { + return "capacity_unavailable"; + } + if (/(?:http\s*)?(502|503|504)\b|bad gateway|service unavailable|gateway timeout|upstream.{0,32}(error|fail|timeout|unavailable)/.test(text)) { + return "transient_upstream"; + } + if (/econnreset|econnrefused|etimedout|connection (?:reset|refused|closed|failed)|network (?:error|failure)|socket hang up|dns.{0,16}(fail|error)|fetch failed/.test(text)) { + return "transport_unavailable"; + } + if (/(?:http\s*)?401\b|unauthori[sz]ed|authentication.{0,24}(fail|expired|required)|invalid api key|api key.{0,24}(invalid|expired)|token.{0,24}(expired|invalid)/.test(text)) { + return "auth_unavailable"; + } + if (/model[_ -]?not[_ -]?found|model not found|model.{0,24}(not available|unavailable|unsupported)|not authorized.{0,24}model|not authorised.{0,24}model/.test(text)) { + return "model_unavailable"; + } + return undefined; +} + +function resolvePhysicalModel(ctx, target) { + const model = ctx.modelRegistry.find(target.provider, target.id); if (!model || model.api === "pi-virtual") return undefined; return model; } @@ -375,7 +556,16 @@ function routeV1(role, request, ctx, path, target) { }; } -function routeV2(role, request, ctx, path, target) { +function emitDecision(options, decision, ctx) { + if (typeof options.onDecision !== "function") return; + try { + options.onDecision(decision, ctx); + } catch { + // Usage/telemetry must never break routing. + } +} + +function routeV2(role, request, ctx, path, target, options) { const taskClass = explicitTaskClass(request.messages); const candidates = routeCandidates(target, request.thinkingLevel, taskClass); const attempted = []; @@ -386,20 +576,77 @@ function routeV2(role, request, ctx, path, target) { if (!model) continue; const thinkingLevel = resolveThinkingLevel(request.thinkingLevel, route); - return { - model, + const state = { + decisionVersion: 1, + policyVersion: 2, + role, + routeId: route.routeId, + taskClass: taskClass ?? target.defaultTaskClass ?? null, + model: `${route.provider}/${route.id}`, thinkingLevel, - state: { - decisionVersion: 1, - policyVersion: 2, + basis: taskClass ? "explicit-task-class" : "effort-envelope" + }; + emitDecision(options, state, ctx); + return { model, thinkingLevel, state }; + } + + throw new Error( + `ForgeFlow role '${role}' has no available physical route in '${path}'. Tried: ${attempted.join(", ")}.` + ); +} + +function supplyTargets(policy, route) { + if (!route.supplyGroup) return [{ ...route, supplyGroup: null, supplyId: null, priority: null }]; + const group = policy.supplies?.[route.supplyGroup]; + if (!group) return []; + return group.sources.map((source) => ({ ...source, supplyGroup: route.supplyGroup })); +} + +function v3DecisionState({ role, target, route, taskClass, source, thinkingLevel, basis, failoverReason, failedSupplyId, priorState }) { + const physical = `${source.provider}/${source.id}`; + return { + decisionVersion: 2, + policyVersion: 3, + role, + routeId: route.routeId, + taskClass: taskClass ?? target.defaultTaskClass ?? priorState?.taskClass ?? null, + logicalModel: route.supplyGroup ?? physical, + supplyGroup: source.supplyGroup ?? null, + supplyId: source.supplyId ?? null, + supplyPriority: source.priority ?? null, + model: physical, + thinkingLevel, + basis, + ...(failoverReason ? { failoverReason } : {}), + ...(failedSupplyId ? { failedSupplyId } : {}), + ...(basis === "supply-failover" ? { failoverCount: (priorState?.failoverCount ?? 0) + 1 } : {}) + }; +} + +function routeV3(role, request, ctx, path, policy, target, options) { + const taskClass = explicitTaskClass(request.messages); + const routes = routeCandidates(target, request.thinkingLevel, taskClass); + const attempted = []; + + for (const route of routes) { + for (const source of supplyTargets(policy, route)) { + attempted.push(`${source.provider}/${source.id}`); + const model = resolvePhysicalModel(ctx, source); + if (!model) continue; + + const thinkingLevel = resolveThinkingLevel(request.thinkingLevel, route); + const state = v3DecisionState({ role, - routeId: route.routeId, - taskClass: taskClass ?? target.defaultTaskClass ?? null, - model: `${route.provider}/${route.id}`, + target, + route, + taskClass, + source, thinkingLevel, - basis: taskClass ? "explicit-task-class" : "effort-envelope" - } - }; + basis: taskClass ? "explicit-task-class" : route.supplyGroup ? "supply-priority" : "effort-envelope" + }); + emitDecision(options, state, ctx); + return { model, thinkingLevel, state }; + } } throw new Error( @@ -407,12 +654,86 @@ function routeV2(role, request, ctx, path, target) { ); } +function sourceMatchesModel(source, model) { + return Boolean(model) && source.provider === model.provider && source.id === model.id; +} + +function findRetryRoute(target, state, supplyGroup) { + if (state?.routeId) { + const exact = target.routes.find((route) => route.routeId === state.routeId && route.supplyGroup === supplyGroup); + if (exact) return exact; + } + return target.routes.find((route) => route.supplyGroup === supplyGroup); +} + +function findSupplyGroupForFailedModel(policy, failedModel) { + for (const group of Object.values(policy.supplies ?? {})) { + const source = group.sources.find((candidate) => sourceMatchesModel(candidate, failedModel)); + if (source) return { group, source }; + } + return undefined; +} + +function retryV3(role, request, ctx, path, policy, target, options) { + if (request.reason !== "retry" || !request.failed) return undefined; + const failureClass = classifySupplyFailure(request.failed.message); + if (!failureClass) return failedStickyRoute(request); + + const priorState = request.state && typeof request.state === "object" ? request.state : undefined; + let groupId = typeof priorState?.supplyGroup === "string" ? priorState.supplyGroup : undefined; + let group = groupId ? policy.supplies?.[groupId] : undefined; + let failedSource = group?.sources.find((source) => sourceMatchesModel(source, request.failed.model)); + + if (!group || !failedSource) { + const found = findSupplyGroupForFailedModel(policy, request.failed.model); + group = found?.group; + failedSource = found?.source; + groupId = group?.groupId; + } + if (!group || !failedSource || !groupId) return failedStickyRoute(request); + + const route = findRetryRoute(target, priorState, groupId); + if (!route) return failedStickyRoute(request); + + const failedIndex = group.sources.findIndex((source) => source.supplyId === failedSource.supplyId); + const nextSources = group.sources.slice(failedIndex + 1); + for (const source of nextSources) { + const model = resolvePhysicalModel(ctx, source); + if (!model) continue; + const requestedThinking = request.failed.thinkingLevel ?? request.thinkingLevel; + const thinkingLevel = resolveThinkingLevel(requestedThinking, route); + const state = v3DecisionState({ + role, + target, + route, + taskClass: priorState?.taskClass ?? explicitTaskClass(request.messages), + source: { ...source, supplyGroup: groupId }, + thinkingLevel, + basis: "supply-failover", + failoverReason: failureClass, + failedSupplyId: failedSource.supplyId, + priorState + }); + emitDecision(options, state, ctx); + return { model, thinkingLevel, state }; + } + + return failedStickyRoute(request); +} + export function createRoleRouter(role, options = {}) { const loadPolicy = options.loadPolicy ?? loadModelPolicy; return function route(request, ctx) { - const sticky = stickyRoute(request); - if (sticky) return sticky; + const continuation = continuationRoute(request); + if (continuation) return continuation; + if ( + request.reason === "retry" + && request.failed + && (!request.state || request.state.policyVersion !== 3) + ) { + return failedStickyRoute(request); + } const allowProject = typeof ctx.isProjectTrusted === "function" && ctx.isProjectTrusted(); const { path, policy } = loadPolicy(ctx.cwd, process.env, homedir(), { allowProject }); @@ -423,9 +744,16 @@ export function createRoleRouter(role, options = {}) { ); } - return policy.version === 1 - ? routeV1(role, request, ctx, path, target) - : routeV2(role, request, ctx, path, target); + if (request.reason === "retry" && request.failed) { + if (policy.version === 3) { + return retryV3(role, request, ctx, path, policy, target, options); + } + return failedStickyRoute(request); + } + + if (policy.version === 1) return routeV1(role, request, ctx, path, target); + if (policy.version === 2) return routeV2(role, request, ctx, path, target, options); + return routeV3(role, request, ctx, path, policy, target, options); }; } diff --git a/extension/usage.js b/extension/usage.js new file mode 100644 index 0000000..5956b1f --- /dev/null +++ b/extension/usage.js @@ -0,0 +1,159 @@ +import { appendFileSync, chmodSync, mkdirSync, readFileSync } from "node:fs"; +import { homedir } from "node:os"; +import { dirname, isAbsolute, join } from "node:path"; + +import { classifySupplyFailure } from "./model-policy.js"; + +export const FORGEFLOW_USAGE_LOG_ENV = "FORGEFLOW_USAGE_LOG"; +export const FORGEFLOW_USAGE_EVENT_VERSION = 1; +const VIRTUAL_MODEL_STATE_ENTRY = "pi.virtual-model-state"; +const FORGEFLOW_PROVIDER = "forgeflow"; + +export function resolveUsageLogPath(env = process.env, home = homedir()) { + const explicit = env[FORGEFLOW_USAGE_LOG_ENV]?.trim(); + if (explicit) { + if (!isAbsolute(explicit)) { + throw new Error(`${FORGEFLOW_USAGE_LOG_ENV} must be an absolute path.`); + } + return explicit; + } + return join(home, ".pi", "forgeflow-usage.jsonl"); +} + +function appendJsonLine(path, value) { + mkdirSync(dirname(path), { recursive: true, mode: 0o700 }); + appendFileSync(path, `${JSON.stringify(value)}\n`, { encoding: "utf8", mode: 0o600 }); + try { + chmodSync(path, 0o600); + } catch { + // Permissions can be owned by an external filesystem policy; logging is best-effort. + } +} + +function latestForgeFlowState(ctx) { + let branch; + try { + branch = ctx?.sessionManager?.getBranch?.(); + } catch { + return undefined; + } + if (!Array.isArray(branch)) return undefined; + + for (let index = branch.length - 1; index >= 0; index -= 1) { + const entry = branch[index]; + if (!entry || entry.type !== "custom" || entry.customType !== VIRTUAL_MODEL_STATE_ENTRY) continue; + const data = entry.data; + if (!data || data.provider !== FORGEFLOW_PROVIDER || typeof data.modelId !== "string") continue; + const state = data.state && typeof data.state === "object" ? data.state : undefined; + return { role: data.modelId, state }; + } + return undefined; +} + +function usageNumbers(usage) { + if (!usage || typeof usage !== "object") return undefined; + const cost = usage.cost && typeof usage.cost === "object" ? usage.cost : {}; + return { + input: Number(usage.input ?? 0), + output: Number(usage.output ?? 0), + cacheRead: Number(usage.cacheRead ?? 0), + cacheWrite: Number(usage.cacheWrite ?? 0), + reasoning: usage.reasoning === undefined ? null : Number(usage.reasoning), + totalTokens: Number(usage.totalTokens ?? 0), + cost: { + input: Number(cost.input ?? 0), + output: Number(cost.output ?? 0), + cacheRead: Number(cost.cacheRead ?? 0), + cacheWrite: Number(cost.cacheWrite ?? 0), + total: Number(cost.total ?? 0) + } + }; +} + +function sessionId(ctx) { + try { + return ctx?.sessionManager?.getSessionId?.() ?? null; + } catch { + return null; + } +} + +function baseEvent(type, ctx) { + return { + version: FORGEFLOW_USAGE_EVENT_VERSION, + event: type, + timestamp: new Date().toISOString(), + sessionId: sessionId(ctx) + }; +} + +export function createUsageRecorder(options = {}) { + const path = options.path ?? resolveUsageLogPath(options.env ?? process.env, options.home ?? homedir()); + const append = options.append ?? ((value) => appendJsonLine(path, value)); + + function recordDecision(decision, ctx) { + if (!decision || typeof decision !== "object") return; + append({ + ...baseEvent("route_decision", ctx), + role: decision.role ?? null, + routeId: decision.routeId ?? null, + taskClass: decision.taskClass ?? null, + logicalModel: decision.logicalModel ?? decision.model ?? null, + supplyGroup: decision.supplyGroup ?? null, + supplyId: decision.supplyId ?? null, + supplyPriority: decision.supplyPriority ?? null, + physicalModel: decision.model ?? null, + thinkingLevel: decision.thinkingLevel ?? null, + basis: decision.basis ?? null, + failoverReason: decision.failoverReason ?? null, + failedSupplyId: decision.failedSupplyId ?? null, + failoverCount: decision.failoverCount ?? 0 + }); + } + + function recordMessage(event, ctx) { + const message = event?.message; + if (!message || message.role !== "assistant") return; + const virtual = latestForgeFlowState(ctx); + const state = virtual?.state; + const physicalModel = `${message.provider}/${message.model}`; + const stateMatches = state?.model === physicalModel; + const usage = usageNumbers(message.usage); + + append({ + ...baseEvent("model_usage", ctx), + role: stateMatches ? (state.role ?? virtual?.role ?? null) : null, + routeId: stateMatches ? (state.routeId ?? null) : null, + taskClass: stateMatches ? (state.taskClass ?? null) : null, + logicalModel: stateMatches ? (state.logicalModel ?? message.model) : message.model, + supplyGroup: stateMatches ? (state.supplyGroup ?? null) : null, + supplyId: stateMatches ? (state.supplyId ?? null) : null, + provider: message.provider, + model: message.model, + responseModel: message.responseModel ?? null, + api: message.api ?? null, + thinkingLevel: message.thinkingLevel ?? null, + providerThinkingLevel: message.providerThinkingLevel ?? null, + stopReason: message.stopReason ?? null, + status: message.stopReason === "error" ? "failure" : message.stopReason === "aborted" ? "aborted" : "success", + supplyFailureClass: message.stopReason === "error" ? (classifySupplyFailure(message) ?? null) : null, + usage + }); + } + + return { path, recordDecision, recordMessage }; +} + +export function readUsageEvents(path) { + const text = readFileSync(path, "utf8"); + return text + .split(/\r?\n/) + .filter(Boolean) + .map((line, index) => { + try { + return JSON.parse(line); + } catch (error) { + throw new Error(`Invalid ForgeFlow usage JSONL at line ${index + 1}: ${error instanceof Error ? error.message : String(error)}`); + } + }); +} diff --git a/package-lock.json b/package-lock.json index 0004848..20727bd 100644 --- a/package-lock.json +++ b/package-lock.json @@ -11,6 +11,9 @@ "dependencies": { "pi-subagents": "0.75.0" }, + "bin": { + "forgeflow-usage": "bin/forgeflow-usage.mjs" + }, "peerDependencies": { "@earendil-works/pi-agent-core": "*", "@earendil-works/pi-ai": "*", diff --git a/package.json b/package.json index d7f9454..fcff048 100644 --- a/package.json +++ b/package.json @@ -4,7 +4,15 @@ "description": "Thin software-engineering governance for Pi Agent and pi-subagents.", "type": "module", "license": "MIT", - "files": ["bin", "extension", "skills", "docs/architecture.md", "docs/model-policy.md", "README.md", "LICENSE"], + "files": [ + "bin", + "extension", + "skills", + "docs/architecture.md", + "docs/model-policy.md", + "README.md", + "LICENSE" + ], "keywords": [ "pi-package", "pi", @@ -23,18 +31,35 @@ "typebox": "*" }, "peerDependenciesMeta": { - "@earendil-works/pi-agent-core": { "optional": true }, - "@earendil-works/pi-ai": { "optional": true }, - "@earendil-works/pi-coding-agent": { "optional": true }, - "@earendil-works/pi-tui": { "optional": true }, - "typebox": { "optional": true } + "@earendil-works/pi-agent-core": { + "optional": true + }, + "@earendil-works/pi-ai": { + "optional": true + }, + "@earendil-works/pi-coding-agent": { + "optional": true + }, + "@earendil-works/pi-tui": { + "optional": true + }, + "typebox": { + "optional": true + } }, "pi": { - "extensions": ["./extension/index.js"], - "skills": ["./skills"] + "extensions": [ + "./extension/index.js" + ], + "skills": [ + "./skills" + ] }, "scripts": { "test": "node --test tests/*.test.js", - "check": "npm test && node --check bin/agy-subagent.mjs && node --check extension/antigravity.js && node --check extension/index.js && node --check extension/model-policy.js && node --check extension/invariants.js" + "check": "npm test && node --check bin/agy-subagent.mjs && node --check bin/forgeflow-usage.mjs && node --check extension/antigravity.js && node --check extension/index.js && node --check extension/model-policy.js && node --check extension/usage.js && node --check extension/invariants.js" + }, + "bin": { + "forgeflow-usage": "./bin/forgeflow-usage.mjs" } } diff --git a/tests/extension.test.js b/tests/extension.test.js index 1ab914d..8492d9b 100644 --- a/tests/extension.test.js +++ b/tests/extension.test.js @@ -45,6 +45,7 @@ test("extension registers Pi lifecycle hooks and injects policy without a model "forgeflow/oracle" ]); assert.equal(typeof handlers.get("before_agent_start"), "function"); + assert.equal(typeof handlers.get("message_end"), "function"); assert.equal(typeof handlers.get("session_start"), "function"); assert.equal(typeof handlers.get("session_shutdown"), "function"); diff --git a/tests/model-policy.test.js b/tests/model-policy.test.js index e3cb442..deff9a8 100644 --- a/tests/model-policy.test.js +++ b/tests/model-policy.test.js @@ -7,6 +7,7 @@ import test from "node:test"; import { FORGEFLOW_MODEL_PROVIDER, FORGEFLOW_MODEL_ROLES, + classifySupplyFailure, createRoleRouter, parseModelPolicy, registerForgeFlowVirtualModels, @@ -678,3 +679,187 @@ test("worker policy can prefer native Codex Team GPT-6.1 Sol and fall back when assert.equal(fallbackMedium.model, astra); assert.equal(fallbackMedium.state.routeId, "standard"); }); + +test("v3 supply groups separate logical model capability from quota source", () => { + const policy = parseModelPolicy(JSON.stringify({ + version: 3, + supplies: { + "gpt-6.1-sol": { + sources: [ + { id: "commercial-relay", model: "litellm/gpt-6.1-sol", priority: 20 }, + { id: "business-team", model: "openai-codex/gpt-6.1-sol", priority: 10 } + ] + } + }, + roles: { + worker: { + defaultRoute: "frontier", + defaultTaskClass: "implementation", + routes: [{ + id: "frontier", + supplyGroup: "gpt-6.1-sol", + taskClasses: ["implementation", "debug"], + defaultThinkingLevel: "medium", + minThinkingLevel: "low", + maxThinkingLevel: "high" + }] + } + } + }), "v3.json"); + + assert.equal(policy.version, 3); + assert.deepEqual( + policy.supplies["gpt-6.1-sol"].sources.map(({ supplyId, priority, provider, id }) => ({ supplyId, priority, provider, id })), + [ + { supplyId: "business-team", priority: 10, provider: "openai-codex", id: "gpt-6.1-sol" }, + { supplyId: "commercial-relay", priority: 20, provider: "litellm", id: "gpt-6.1-sol" } + ] + ); + assert.equal(policy.roles.worker.routes[0].supplyGroup, "gpt-6.1-sol"); +}); + +test("v3 route prefers Business Team supply and only fails over on supply failures", () => { + const team = { provider: "openai-codex", id: "gpt-6.1-sol", api: "openai-codex-responses" }; + const relay = { provider: "litellm", id: "gpt-6.1-sol", api: "openai-completions" }; + const policy = parseModelPolicy(JSON.stringify({ + version: 3, + supplies: { + "gpt-6.1-sol": { + sources: [ + { id: "business-team", model: "openai-codex/gpt-6.1-sol", priority: 10 }, + { id: "commercial-relay", model: "litellm/gpt-6.1-sol", priority: 20 } + ] + } + }, + roles: { + worker: { + defaultRoute: "frontier", + defaultTaskClass: "implementation", + routes: [{ + id: "frontier", + supplyGroup: "gpt-6.1-sol", + taskClasses: ["implementation", "debug"], + defaultThinkingLevel: "medium", + minThinkingLevel: "low", + maxThinkingLevel: "high" + }] + } + } + })); + const decisions = []; + const route = createRoleRouter("worker", { + loadPolicy: () => ({ path: "/policy-v3.json", policy }), + onDecision: (decision) => decisions.push(decision) + }); + const ctx = { + cwd: "/repo", + isProjectTrusted: () => true, + modelRegistry: { + find(provider, id) { + if (`${provider}/${id}` === "openai-codex/gpt-6.1-sol") return team; + if (`${provider}/${id}` === "litellm/gpt-6.1-sol") return relay; + } + } + }; + + const initial = route({ reason: "user", thinkingLevel: "medium", messages: [] }, ctx); + assert.equal(initial.model, team); + assert.equal(initial.state.supplyId, "business-team"); + assert.equal(initial.state.logicalModel, "gpt-6.1-sol"); + assert.equal(initial.state.basis, "supply-priority"); + + const retry = route({ + reason: "retry", + thinkingLevel: "medium", + state: initial.state, + messages: [], + failed: { + model: team, + thinkingLevel: "medium", + message: { role: "assistant", stopReason: "error", errorMessage: "HTTP 429 rate limit exceeded" } + } + }, ctx); + assert.equal(retry.model, relay); + assert.equal(retry.state.supplyId, "commercial-relay"); + assert.equal(retry.state.failedSupplyId, "business-team"); + assert.equal(retry.state.failoverReason, "rate_limited"); + assert.equal(retry.state.basis, "supply-failover"); + assert.equal(retry.state.failoverCount, 1); + assert.equal(decisions.length, 2); +}); + +test("v3 retry stays on Business Team for request and context failures", () => { + const team = { provider: "openai-codex", id: "gpt-6.1-sol", api: "openai-codex-responses" }; + const relay = { provider: "litellm", id: "gpt-6.1-sol", api: "openai-completions" }; + const policy = parseModelPolicy(JSON.stringify({ + version: 3, + supplies: { + "gpt-6.1-sol": { sources: [ + { id: "business-team", model: "openai-codex/gpt-6.1-sol", priority: 10 }, + { id: "commercial-relay", model: "litellm/gpt-6.1-sol", priority: 20 } + ] } + }, + roles: { + worker: { + defaultRoute: "frontier", + routes: [{ id: "frontier", supplyGroup: "gpt-6.1-sol", minThinkingLevel: "low", maxThinkingLevel: "high" }] + } + } + })); + const route = createRoleRouter("worker", { loadPolicy: () => ({ path: "/policy.json", policy }) }); + const ctx = { + cwd: "/repo", + isProjectTrusted: () => true, + modelRegistry: { + find(provider, id) { + if (provider === "openai-codex" && id === "gpt-6.1-sol") return team; + if (provider === "litellm" && id === "gpt-6.1-sol") return relay; + } + } + }; + const initial = route({ reason: "user", thinkingLevel: "medium", messages: [] }, ctx); + const retry = route({ + reason: "retry", + thinkingLevel: "medium", + state: initial.state, + messages: [], + failed: { + model: team, + thinkingLevel: "medium", + message: { role: "assistant", stopReason: "error", errorMessage: "maximum context length exceeded" } + } + }, ctx); + assert.equal(retry.model, team); + assert.equal(retry.state, initial.state); +}); + +test("supply failure classifier is intentionally narrow", () => { + assert.equal(classifySupplyFailure("429 Too Many Requests"), "rate_limited"); + assert.equal(classifySupplyFailure("insufficient_quota: usage limit reached"), "quota_exhausted"); + assert.equal(classifySupplyFailure("503 Service Unavailable"), "transient_upstream"); + assert.equal(classifySupplyFailure("ECONNRESET upstream connection reset"), "transport_unavailable"); + assert.equal(classifySupplyFailure("401 unauthorized token expired"), "auth_unavailable"); + assert.equal(classifySupplyFailure("model_not_found"), "model_unavailable"); + assert.equal(classifySupplyFailure("context window exceeded"), undefined); + assert.equal(classifySupplyFailure("400 bad request invalid parameter"), undefined); + assert.equal(classifySupplyFailure("content policy violation"), undefined); +}); + +test("v3 policy fails closed on invalid or unknown supply references", () => { + assert.throws( + () => parseModelPolicy(JSON.stringify({ + version: 3, + supplies: {}, + roles: { worker: { defaultRoute: "x", routes: [{ id: "x", supplyGroup: "missing" }] } } + }), "missing-supply.json"), + /unknown supply group/ + ); + assert.throws( + () => parseModelPolicy(JSON.stringify({ + version: 3, + supplies: { x: { sources: [{ id: "a", model: "litellm/a", priority: -1 }] } }, + roles: { worker: { defaultRoute: "x", routes: [{ id: "x", supplyGroup: "x" }] } } + }), "bad-priority.json"), + /non-negative integer 'priority'/ + ); +}); diff --git a/tests/usage.test.js b/tests/usage.test.js new file mode 100644 index 0000000..cd92f5f --- /dev/null +++ b/tests/usage.test.js @@ -0,0 +1,114 @@ +import assert from "node:assert/strict"; +import { mkdtempSync, readFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; + +import { createUsageRecorder, readUsageEvents, resolveUsageLogPath } from "../extension/usage.js"; + +test("usage recorder stores route decisions and assistant usage without secrets", () => { + const dir = mkdtempSync(join(tmpdir(), "forgeflow-usage-")); + const path = join(dir, "usage.jsonl"); + const recorder = createUsageRecorder({ path }); + const state = { + decisionVersion: 2, + policyVersion: 3, + role: "worker", + routeId: "frontier", + taskClass: "implementation", + logicalModel: "gpt-6.1-sol", + supplyGroup: "gpt-6.1-sol", + supplyId: "business-team", + supplyPriority: 10, + model: "openai-codex/gpt-6.1-sol", + thinkingLevel: "medium", + basis: "supply-priority" + }; + const ctx = { + sessionManager: { + getSessionId: () => "session-1", + getBranch: () => [{ + type: "custom", + customType: "pi.virtual-model-state", + data: { provider: "forgeflow", modelId: "worker", state } + }] + } + }; + + recorder.recordDecision(state, ctx); + recorder.recordMessage({ + message: { + role: "assistant", + provider: "openai-codex", + model: "gpt-6.1-sol", + api: "openai-codex-responses", + thinkingLevel: "medium", + providerThinkingLevel: "medium", + stopReason: "stop", + usage: { + input: 100, + output: 20, + cacheRead: 40, + cacheWrite: 0, + reasoning: 8, + totalTokens: 120, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 } + } + } + }, ctx); + + const events = readUsageEvents(path); + assert.equal(events.length, 2); + assert.equal(events[0].event, "route_decision"); + assert.equal(events[0].supplyId, "business-team"); + assert.equal(events[1].event, "model_usage"); + assert.equal(events[1].role, "worker"); + assert.equal(events[1].supplyId, "business-team"); + assert.equal(events[1].usage.totalTokens, 120); + assert.equal(events[1].usage.reasoning, 8); + assert.equal(events[1].status, "success"); + assert.doesNotMatch(readFileSync(path, "utf8"), /api[_-]?key|authorization|bearer|errorMessage/i); +}); + +test("usage recorder does not attribute stale virtual state to a different physical model", () => { + const dir = mkdtempSync(join(tmpdir(), "forgeflow-usage-")); + const path = join(dir, "usage.jsonl"); + const recorder = createUsageRecorder({ path }); + const ctx = { + sessionManager: { + getSessionId: () => "session-2", + getBranch: () => [{ + type: "custom", + customType: "pi.virtual-model-state", + data: { + provider: "forgeflow", + modelId: "worker", + state: { role: "worker", model: "openai-codex/gpt-6.1-sol", supplyId: "business-team" } + } + }] + } + }; + + recorder.recordMessage({ + message: { + role: "assistant", + provider: "litellm", + model: "gpt-6-astra", + api: "openai-responses", + stopReason: "stop", + usage: { input: 1, output: 1, cacheRead: 0, cacheWrite: 0, totalTokens: 2, cost: { total: 0 } } + } + }, ctx); + + const [event] = readUsageEvents(path); + assert.equal(event.role, null); + assert.equal(event.supplyId, null); + assert.equal(event.model, "gpt-6-astra"); +}); + +test("usage log override must be absolute", () => { + assert.throws( + () => resolveUsageLogPath({ FORGEFLOW_USAGE_LOG: "relative.jsonl" }, "/home/test"), + /must be an absolute path/ + ); +});