From 0b6ae49ed21abee31db2abdbc186615de01ea77c Mon Sep 17 00:00:00 2001 From: shuntianyifang <2016860013@qq.com> Date: Wed, 7 Oct 2026 15:41:07 +0800 Subject: [PATCH 1/3] fix(api): account for worker-interrupted usage and block unconfirmed spend --- README.md | 8 +- apps/docs/docs/faq.md | 15 + apps/docs/docs/reference/api.md | 7 + .../projects/[projectId]/insights/page.tsx | 14 +- apps/web/lib/overview-presentation.ts | 25 +- apps/web/test/insights-page.test.tsx | 20 +- apps/web/test/overview-presentation.test.ts | 9 +- packages/sdk/openapi.json | 12 +- packages/sdk/src/schema.d.ts | 8 +- services/api/src/insights/costs.ts | 136 +++-- .../api/src/routes/v1/insights-pipeline.ts | 4 +- .../api/src/routes/v1/project-overview.ts | 2 +- services/api/src/stories/service.ts | 176 ++++++- services/api/src/stories/titles.ts | 2 + services/api/src/turns/engines.ts | 18 +- services/api/src/turns/usage-journal.ts | 131 +++++ services/api/src/worker.ts | 29 +- .../api/test/budget-accounting-policy.test.ts | 50 ++ .../interrupted-usage.integration.test.ts | 465 ++++++++++++++++++ .../test/project-overview.integration.test.ts | 3 +- services/api/test/usage-journal.test.ts | 119 +++++ 21 files changed, 1172 insertions(+), 81 deletions(-) create mode 100644 services/api/src/turns/usage-journal.ts create mode 100644 services/api/test/budget-accounting-policy.test.ts create mode 100644 services/api/test/interrupted-usage.integration.test.ts create mode 100644 services/api/test/usage-journal.test.ts diff --git a/README.md b/README.md index 17f44fbd..ddaf748a 100644 --- a/README.md +++ b/README.md @@ -131,8 +131,12 @@ Self-hosting gives you control of the Facility services, database, and workspace Code and Codex communicate with the model service you configure for the engine. Choose providers and credentials that fit your infrastructure requirements. -Project budgets are checked before new provider calls. Usage is recorded afterwards, so an -in-flight call can take spending beyond the monthly limit. Retained workspaces also need an +Project budgets are checked before new provider calls. An interrupted worker's usage is +recovered from retained per-turn journals where possible. Missing +or incomplete usage pauses new model calls while budget enforcement is enabled; see the +[budget FAQ](apps/docs/docs/faq.md#does-a-project-budget-delete-or-stop-a-workspace). +Usage is recorded afterwards, so an in-flight call can take spending beyond the monthly limit. +Retained workspaces also need an explicit storage and deletion policy. The [hardening guide](apps/docs/docs/reference/hardening.md) covers isolation, credentials, diff --git a/apps/docs/docs/faq.md b/apps/docs/docs/faq.md index 1e8b943f..2cc598ff 100644 --- a/apps/docs/docs/faq.md +++ b/apps/docs/docs/faq.md @@ -57,6 +57,21 @@ afterwards; later turns are blocked once the monthly limit is reached. The workt and engine sessions remain. Unknown model pricing is rejected while budget enforcement is enabled, and unavailable workspace pricing is not reported as zero. +If a worker dies, recovery reads the turn's retained usage journal. Final counters are charged +once; partial counters are a lower bound. Missing, unreadable, or incomplete usage is marked +unpriced and the budget state becomes `unconfirmed`: new agent and title-generation model calls +are blocked while enforcement is enabled. Raising the limit, resolving attention, or starting a +new month does not confirm the bill. Recovery retries retained journals without waking compute +or rerunning the agent. A finalized journal clears the block after reconciliation. +The same recovery pass marks missing bills on worker-interrupted turns failed by older +Facility versions as unconfirmed; it cannot invent usage that was never retained. + +When no final journal can be recovered, inspect the provider bill and retained session files; +the true charge may exceed the recorded lower bound. There is currently no manual billing +reconciliation endpoint. An authorized budget administrator can explicitly disable enforcement +to continue, accepting that spend remains incomplete. Deleting the workspace loses evidence +and does not clear the accounting block. + To use a newly supported model, update the agent's `model` field and ensure the deployed Facility version contains its price-book entry. See [cost and budget API behavior](reference/api.md) for the supported additions and diff --git a/apps/docs/docs/reference/api.md b/apps/docs/docs/reference/api.md index 9f10bbd4..d86f350e 100644 --- a/apps/docs/docs/reference/api.md +++ b/apps/docs/docs/reference/api.md @@ -189,6 +189,13 @@ verified on 2026-09-29. These entries use standard global API pricing and 5-minute cache writes, not fast mode, batch, regional premiums, or 1-hour cache writes. A valid engine-reported cost takes precedence over this fallback. An enabled budget still blocks an unpriced model or a project over its limit. +Interrupted turns with incomplete usage return `state: "unconfirmed"` in budget responses. +Their usage rows have `priced: false`, `source: "unpriced"`, and a nullable cost or a measured +lower bound. New agent and title-generation model calls are blocked until complete retained +usage is reconciled or an authorized administrator explicitly disables budget enforcement. +Raising the limit, resolving attention, and a new calendar month do not clear uncertainty. +Recovered usage is counted when it is settled, as with live-worker accounting; reconciliation does not double-charge +already settled rows. The periodic worker retries available journals without waking compute. Updating the catalog does not reprice persisted turns or change agent defaults. ## Authentication and authorization diff --git a/apps/web/app/(app)/projects/[projectId]/insights/page.tsx b/apps/web/app/(app)/projects/[projectId]/insights/page.tsx index 14b3dad0..d1db41dd 100644 --- a/apps/web/app/(app)/projects/[projectId]/insights/page.tsx +++ b/apps/web/app/(app)/projects/[projectId]/insights/page.tsx @@ -114,9 +114,11 @@ export default async function InsightsPage({ params }: { params: Promise<{ proje {budget.data.state.replaceAll("_", " ")}

- {money(budget.data.spent_cents)} spent this month. New turns are blocked once the limit is - reached; a provider call already in progress is allowed to finish and is accounted - afterwards. + {budget.data.state === "unconfirmed" ? "At least " : ""} + {money(budget.data.spent_cents)} recorded this month. New turns are blocked once the limit + is reached; a provider call already in progress is allowed to finish and is accounted + afterwards. If an interrupted turn's full usage cannot be recovered, spending is shown as + a lower bound and new model calls are blocked while the budget is enabled.

@@ -167,6 +169,7 @@ function UsageTable({ name: string; turns: number; costCents: number; + unpricedTurns?: number; inputTokens: number; outputTokens: number; cacheReadTokens: number; @@ -187,7 +190,10 @@ function UsageTable({ > {row.name} {row.turns} turns - {money(row.costCents)} + + {(row.unpricedTurns ?? 0) > 0 ? "≥ " : ""} + {money(row.costCents)} + )) )} diff --git a/apps/web/lib/overview-presentation.ts b/apps/web/lib/overview-presentation.ts index 34b2a80d..9b31dd68 100644 --- a/apps/web/lib/overview-presentation.ts +++ b/apps/web/lib/overview-presentation.ts @@ -223,20 +223,27 @@ export function attentionSignals( }); } const budget = overview.spend.budget; - if (budget.available && (budget.state === "exceeded" || budget.state === "warning")) { + if ( + budget.available && + (budget.state === "exceeded" || budget.state === "warning" || budget.state === "unconfirmed") + ) { entries.push({ key: "budget", - tone: budget.state === "exceeded" ? "bad" : "human", + tone: budget.state === "warning" ? "human" : "bad", title: - budget.state === "exceeded" - ? "Monthly budget exhausted: new agent turns are blocked" - : "Monthly budget warning", + budget.state === "unconfirmed" + ? "Spending unconfirmed: new model calls are blocked" + : budget.state === "exceeded" + ? "Monthly budget exhausted: new agent turns are blocked" + : "Monthly budget warning", storyId: null, storyTitle: null, summary: `${money(budget.spentCents)} of ${money(budget.monthlyLimitCents ?? 0)} used this month.${ - budget.state === "exceeded" - ? " Raise the limit or wait for the next month to continue." - : " New turns are blocked once the limit is reached." + budget.state === "unconfirmed" + ? " Interrupted usage could not be fully recovered. Raising the limit does not clear this block." + : budget.state === "exceeded" + ? " Raise the limit or wait for the next month to continue." + : " New turns are blocked once the limit is reached." }`, at: null, action: { @@ -335,6 +342,8 @@ export function budgetReading(budget: ProjectOverview["spend"]["budget"]) { return { label: "Budget disabled", tone: "machine" as Tone }; case "exceeded": return { label: "Budget exhausted", tone: "bad" as Tone }; + case "unconfirmed": + return { label: "Spending unconfirmed", tone: "bad" as Tone }; case "warning": return { label: "Budget warning", tone: "human" as Tone }; default: diff --git a/apps/web/test/insights-page.test.tsx b/apps/web/test/insights-page.test.tsx index 7311a49b..797f3021 100644 --- a/apps/web/test/insights-page.test.tsx +++ b/apps/web/test/insights-page.test.tsx @@ -3,7 +3,12 @@ import { renderToStaticMarkup } from "react-dom/server"; import { describe, expect, it, vi } from "vitest"; import InsightsPage from "../app/(app)/projects/[projectId]/insights/page"; -const mocks = vi.hoisted(() => ({ denied: false, measured: 1, failed: 2 })); +const mocks = vi.hoisted(() => ({ + denied: false, + measured: 1, + failed: 2, + budgetUnconfirmed: false, +})); vi.mock("next/navigation", () => ({ useRouter: () => ({ refresh() {} }) })); vi.mock("../components/insights/budget-form", () => ({ BudgetForm: () => null })); vi.mock("../lib/api", () => ({ @@ -51,7 +56,10 @@ vi.mock("../lib/api", () => ({ recentAudit: [], }, }, - projectBudget: async () => ({ ok: true, data: { state: "not_configured", spent_cents: 0 } }), + projectBudget: async () => ({ + ok: true, + data: { state: mocks.budgetUnconfirmed ? "unconfirmed" : "not_configured", spent_cents: 0 }, + }), }, })); @@ -62,6 +70,14 @@ async function page() { } describe("Insights spend rendering", () => { + it("shows interrupted budget spending as a lower bound and explains the model-call block", async () => { + mocks.budgetUnconfirmed = true; + const html = await page(); + expect(html).toContain("unconfirmed"); + expect(html).toContain("At least"); + expect(html).toContain("new model calls are blocked while the budget is enabled"); + mocks.budgetUnconfirmed = false; + }); it("shows partial cost and the missing usage in the actual page", async () => { mocks.denied = false; mocks.measured = 1; diff --git a/apps/web/test/overview-presentation.test.ts b/apps/web/test/overview-presentation.test.ts index 663ba1ef..ea32b983 100644 --- a/apps/web/test/overview-presentation.test.ts +++ b/apps/web/test/overview-presentation.test.ts @@ -347,12 +347,17 @@ describe("spend readings", () => { expect(budgetReading({ available: false, reason: "permission" }).label).toBe( "Not visible for your role", ); - const budget = (state: "not_configured" | "disabled" | "ok" | "warning" | "exceeded") => - ({ available: true, state }) as ProjectOverview["spend"]["budget"]; + const budget = ( + state: "not_configured" | "disabled" | "ok" | "warning" | "exceeded" | "unconfirmed", + ) => ({ available: true, state }) as ProjectOverview["spend"]["budget"]; expect(budgetReading(budget("not_configured")).label).toBe("No monthly budget"); expect(budgetReading(budget("ok"))).toEqual({ label: "Within budget", tone: "ok" }); expect(budgetReading(budget("warning")).tone).toBe("human"); expect(budgetReading(budget("exceeded")).tone).toBe("bad"); + expect(budgetReading(budget("unconfirmed"))).toEqual({ + label: "Spending unconfirmed", + tone: "bad", + }); }); it("reports recorded workspace states without claiming live inspection", () => { expect( diff --git a/packages/sdk/openapi.json b/packages/sdk/openapi.json index 3e976a4b..dcbbb9ce 100644 --- a/packages/sdk/openapi.json +++ b/packages/sdk/openapi.json @@ -6074,7 +6074,8 @@ "disabled", "ok", "warning", - "exceeded" + "exceeded", + "unconfirmed" ] }, "enforcement": { @@ -6299,7 +6300,8 @@ "disabled", "ok", "warning", - "exceeded" + "exceeded", + "unconfirmed" ] }, "enforcement": { @@ -6606,7 +6608,8 @@ "disabled", "ok", "warning", - "exceeded" + "exceeded", + "unconfirmed" ] }, "enforcement": { @@ -8999,7 +9002,8 @@ "disabled", "ok", "warning", - "exceeded" + "exceeded", + "unconfirmed" ] }, "enabled": { diff --git a/packages/sdk/src/schema.d.ts b/packages/sdk/src/schema.d.ts index 8e0c75a5..943b444f 100644 --- a/packages/sdk/src/schema.d.ts +++ b/packages/sdk/src/schema.d.ts @@ -4878,7 +4878,7 @@ export interface operations { remaining_cents: number | null; percent_used: number | null; /** @enum {string} */ - state: "not_configured" | "disabled" | "ok" | "warning" | "exceeded"; + state: "not_configured" | "disabled" | "ok" | "warning" | "exceeded" | "unconfirmed"; /** @enum {string} */ enforcement: "block_new_turns"; }; @@ -5001,7 +5001,7 @@ export interface operations { remaining_cents: number | null; percent_used: number | null; /** @enum {string} */ - state: "not_configured" | "disabled" | "ok" | "warning" | "exceeded"; + state: "not_configured" | "disabled" | "ok" | "warning" | "exceeded" | "unconfirmed"; /** @enum {string} */ enforcement: "block_new_turns"; }; @@ -5142,7 +5142,7 @@ export interface operations { remainingCents: number | null; percentUsed: number | null; /** @enum {string} */ - state: "not_configured" | "disabled" | "ok" | "warning" | "exceeded"; + state: "not_configured" | "disabled" | "ok" | "warning" | "exceeded" | "unconfirmed"; /** @enum {string} */ enforcement: "block_new_turns"; }; @@ -5952,7 +5952,7 @@ export interface operations { /** @enum {boolean} */ available: true; /** @enum {string} */ - state: "not_configured" | "disabled" | "ok" | "warning" | "exceeded"; + state: "not_configured" | "disabled" | "ok" | "warning" | "exceeded" | "unconfirmed"; enabled: boolean; monthlyLimitCents: number | null; warningPercent: number | null; diff --git a/services/api/src/insights/costs.ts b/services/api/src/insights/costs.ts index 0e0e6f0a..d28e07e7 100644 --- a/services/api/src/insights/costs.ts +++ b/services/api/src/insights/costs.ts @@ -5,7 +5,7 @@ import type { AgentTurnUsage } from "../turns/engines.js"; export class BudgetPolicyError extends Error { constructor( - readonly code: "budget_exceeded" | "budget_model_unpriced", + readonly code: "budget_exceeded" | "budget_model_unpriced" | "budget_usage_unconfirmed", message: string, ) { super(message); @@ -20,7 +20,7 @@ export type BudgetState = { spentCents: number; remainingCents: number | null; percentUsed: number | null; - state: "not_configured" | "disabled" | "ok" | "warning" | "exceeded"; + state: "not_configured" | "disabled" | "ok" | "warning" | "exceeded" | "unconfirmed"; }; export class CostBudgetService { @@ -29,6 +29,12 @@ export class CostBudgetService { async assertTurnAllowed(orgId: string, projectId: string, model: string, now = new Date()) { const state = await this.budgetState(orgId, projectId, now); if (!state.budget?.enabled) return state; + if (state.state === "unconfirmed") { + throw new BudgetPolicyError( + "budget_usage_unconfirmed", + "Project spending is unconfirmed after an interrupted turn; new model calls are blocked while the budget is enabled.", + ); + } if (!normalizeModel(model)) { throw new BudgetPolicyError( "budget_model_unpriced", @@ -55,53 +61,77 @@ export class CostBudgetService { usage?: AgentTurnUsage; durationMs: number; status: "succeeded" | "failed"; + /** False means these counters are only a lower bound, never a zero-cost completion. */ + complete?: boolean; }) { - if (!input.usage) return null; - const calculated = costCents({ - model: input.model, - inputTokens: input.usage.inputTokens, - outputTokens: input.usage.outputTokens, - cacheReadTokens: input.usage.cacheReadTokens, - cacheWriteTokens: input.usage.cacheWriteTokens, - }); - const providerCost = finiteNonNegative(input.usage.reportedCostCents); - const cost = providerCost ?? calculated; - const row = ( - await this.db - .insert(turnUsage) - .values({ - id: newId("evt"), - orgId: input.orgId, - projectId: input.projectId, - storyId: input.storyId, - turnId: input.turnId, - agentName: input.agentName, - engine: input.engine, + if (!input.usage && input.complete !== false) return null; + const usage = input.usage ?? { + inputTokens: 0, + outputTokens: 0, + cacheReadTokens: 0, + cacheWriteTokens: 0, + }; + const calculated = input.usage + ? costCents({ model: input.model, - inputTokens: input.usage.inputTokens, - outputTokens: input.usage.outputTokens, - cacheReadTokens: input.usage.cacheReadTokens, - cacheWriteTokens: input.usage.cacheWriteTokens, - costCents: cost, - priced: cost !== null, - source: - providerCost !== undefined - ? "provider" - : calculated !== null - ? "price_book" - : "unpriced", - durationMs: Math.max(0, Math.round(input.durationMs)), - status: input.status, + inputTokens: usage.inputTokens, + outputTokens: usage.outputTokens, + cacheReadTokens: usage.cacheReadTokens, + cacheWriteTokens: usage.cacheWriteTokens, }) - .onConflictDoNothing({ target: turnUsage.turnId }) - .returning() - )[0]; + : null; + const providerCost = finiteNonNegative(usage.reportedCostCents); + const cost = providerCost ?? calculated; + const values: typeof turnUsage.$inferInsert = { + id: newId("evt"), + orgId: input.orgId, + projectId: input.projectId, + storyId: input.storyId, + turnId: input.turnId, + agentName: input.agentName, + engine: input.engine, + model: input.model, + inputTokens: usage.inputTokens, + outputTokens: usage.outputTokens, + cacheReadTokens: usage.cacheReadTokens, + cacheWriteTokens: usage.cacheWriteTokens, + costCents: cost, + priced: input.complete !== false && cost !== null, + source: + input.complete === false + ? "unpriced" + : providerCost !== undefined + ? "provider" + : calculated !== null + ? "price_book" + : "unpriced", + durationMs: Math.min(2_147_483_647, Math.max(0, Math.round(input.durationMs))), + status: input.status, + }; + const insert = this.db.insert(turnUsage).values(values); + const { id: _id, ...settled } = values; + // A late live dispatcher can also settle the recovery placeholder. Never + // overwrite an already priced charge. Confirmed usage enters the budget + // when it is accounted, just like an ordinary live-worker completion. + const write = input.usage + ? insert.onConflictDoUpdate({ + target: turnUsage.turnId, + set: { ...settled, ...(values.priced ? { createdAt: new Date() } : {}) }, + setWhere: and( + eq(turnUsage.orgId, input.orgId), + eq(turnUsage.projectId, input.projectId), + eq(turnUsage.storyId, input.storyId), + eq(turnUsage.priced, false), + ), + }) + : insert.onConflictDoNothing({ target: turnUsage.turnId }); + const row = (await write.returning())[0]; return row ?? null; } async budgetState(orgId: string, projectId: string, now = new Date()): Promise { const [windowStart, windowEnd] = monthWindow(now); - const [budget, totals] = await Promise.all([ + const [budget, totals, unconfirmed] = await Promise.all([ this.db .select() .from(projectBudgets) @@ -122,6 +152,18 @@ export class CostBudgetService { ), ) .then((rows) => rows[0]), + // A calendar rollover does not establish the missing provider bill. + this.db + .select({ id: turnUsage.id }) + .from(turnUsage) + .where( + and( + eq(turnUsage.orgId, orgId), + eq(turnUsage.projectId, projectId), + eq(turnUsage.priced, false), + ), + ) + .limit(1), ]); const spentCents = totals?.spentCents ?? 0; if (!budget) { @@ -144,11 +186,13 @@ export class CostBudgetService { : (spentCents / budget.monthlyLimitCents) * 100; const state = !budget.enabled ? "disabled" - : spentCents >= budget.monthlyLimitCents - ? "exceeded" - : percentUsed >= budget.warningPercent - ? "warning" - : "ok"; + : unconfirmed.length > 0 + ? "unconfirmed" + : spentCents >= budget.monthlyLimitCents + ? "exceeded" + : percentUsed >= budget.warningPercent + ? "warning" + : "ok"; return { budget, windowStart, windowEnd, spentCents, remainingCents, percentUsed, state }; } diff --git a/services/api/src/routes/v1/insights-pipeline.ts b/services/api/src/routes/v1/insights-pipeline.ts index fb05a568..63c05231 100644 --- a/services/api/src/routes/v1/insights-pipeline.ts +++ b/services/api/src/routes/v1/insights-pipeline.ts @@ -29,7 +29,7 @@ const BudgetResponse = z.object({ spent_cents: z.number(), remaining_cents: z.number().nullable(), percent_used: z.number().nullable(), - state: z.enum(["not_configured", "disabled", "ok", "warning", "exceeded"]), + state: z.enum(["not_configured", "disabled", "ok", "warning", "exceeded", "unconfirmed"]), enforcement: z.literal("block_new_turns"), }); const UsageSummary = z.object({ @@ -102,7 +102,7 @@ const ObservabilityResponse = z.object({ spentCents: z.number(), remainingCents: z.number().nullable(), percentUsed: z.number().nullable(), - state: z.enum(["not_configured", "disabled", "ok", "warning", "exceeded"]), + state: z.enum(["not_configured", "disabled", "ok", "warning", "exceeded", "unconfirmed"]), enforcement: z.literal("block_new_turns"), }), workspaces: z.object({ diff --git a/services/api/src/routes/v1/project-overview.ts b/services/api/src/routes/v1/project-overview.ts index 8ee110d8..264aa50e 100644 --- a/services/api/src/routes/v1/project-overview.ts +++ b/services/api/src/routes/v1/project-overview.ts @@ -126,7 +126,7 @@ const AgentSpend = z.object({ }); const BudgetSpend = z.object({ available: z.literal(true), - state: z.enum(["not_configured", "disabled", "ok", "warning", "exceeded"]), + state: z.enum(["not_configured", "disabled", "ok", "warning", "exceeded", "unconfirmed"]), enabled: z.boolean(), monthlyLimitCents: z.number().nullable(), warningPercent: z.number().nullable(), diff --git a/services/api/src/stories/service.ts b/services/api/src/stories/service.ts index dc83b3d2..526a635f 100644 --- a/services/api/src/stories/service.ts +++ b/services/api/src/stories/service.ts @@ -13,15 +13,19 @@ import { storyEvidenceEvents, storyMessages, turnEvents, + turnGitEvidence, turns, + turnUsage, userIdentities, users, workspaces, } from "@facility/db"; import { and, asc, desc, eq, inArray, isNull, lt, lte, notInArray, sql } from "drizzle-orm"; +import { CostBudgetService } from "../insights/costs.js"; import { ACTIVITY_NOISE_TYPES, presentTurnEvent } from "../turns/activity.js"; import { stopInterruptedEngineProcess } from "../turns/engines.js"; import { appendTurnEvent } from "../turns/events.js"; +import { readUsageJournal } from "../turns/usage-journal.js"; import { appendWorkspaceEvent } from "../workspaces/events.js"; import { shouldSuspendFailedWorkspace } from "../workspaces/failed-turn-policy.js"; import type { @@ -645,6 +649,33 @@ export class StoryWorkspaceService { .returning() )[0]; if (!updated) return undefined; + // Book uncertainty before committing failure or admitting a successor. A + // legacy turn without phase evidence may already have charged its provider. + const phases = await tx + .select({ data: turnEvents.data }) + .from(turnEvents) + .where( + and( + eq(turnEvents.orgId, input.orgId), + eq(turnEvents.projectId, input.projectId), + eq(turnEvents.turnId, turn.id), + eq(turnEvents.type, "turn.phase"), + ), + ); + const mayHaveSpent = providerMayHaveRun(phases); + if (mayHaveSpent) + await new CostBudgetService(tx).record({ + orgId: input.orgId, + projectId: input.projectId, + storyId: turn.storyId, + turnId: turn.id, + agentName: turn.agentName, + engine: turn.engine, + model: turn.model, + durationMs: turn.startedAt ? now.getTime() - turn.startedAt.getTime() : 0, + status: "failed", + complete: false, + }); await tx .update(stories) .set({ @@ -662,14 +693,27 @@ export class StoryWorkspaceService { kind: "worker_interrupted", title: `${turn.agentName} was interrupted`, detail: - "The worker heartbeat expired. The message, worktree, and native session files were retained; retry this attention item after inspecting any partial changes.", + "The worker heartbeat expired. Retained usage will be reconciled where available. Unconfirmed spending blocks new model calls while the project budget is enabled; inspect the usage and partial work before retrying. The message, worktree, and native session files were retained.", }); - return updated; + return { ...updated, mayHaveSpent }; }); if (!recovered) return false; let processCleanup = "workspace unavailable"; try { + const evidence = ( + await this.db + .select({ workspaceId: turnGitEvidence.workspaceId }) + .from(turnGitEvidence) + .where( + and( + eq(turnGitEvidence.orgId, input.orgId), + eq(turnGitEvidence.projectId, input.projectId), + eq(turnGitEvidence.turnId, recovered.id), + ), + ) + .limit(1) + )[0]; const workspace = ( await this.db .select() @@ -679,6 +723,7 @@ export class StoryWorkspaceService { eq(workspaces.orgId, input.orgId), eq(workspaces.projectId, input.projectId), eq(workspaces.storyId, recovered.storyId), + ...(evidence?.workspaceId ? [eq(workspaces.id, evidence.workspaceId)] : []), ), ) .orderBy(desc(workspaces.createdAt)) @@ -690,6 +735,29 @@ export class StoryWorkspaceService { locatorFromRow(workspace), recovered.id, ); + if (recovered.mayHaveSpent) { + const measured = await readUsageJournal( + this.runtime, + locatorFromRow(workspace), + recovered.id, + recovered.engine, + ); + await new CostBudgetService(this.db).record({ + orgId: input.orgId, + projectId: input.projectId, + storyId: recovered.storyId, + turnId: recovered.id, + agentName: recovered.agentName, + engine: recovered.engine, + model: recovered.model, + ...measured, + durationMs: + recovered.startedAt && recovered.endedAt + ? recovered.endedAt.getTime() - recovered.startedAt.getTime() + : 0, + status: "failed", + }); + } } } catch (error) { processCleanup = @@ -706,6 +774,100 @@ export class StoryWorkspaceService { return true; } + /** Retry accounting without waking compute or rerunning a provider call. */ + async reconcileInterruptedUsage(input: { orgId: string; projectId: string; turnId: string }) { + const turn = await scopedTurn(this.db, input.orgId, input.projectId, input.turnId); + if ( + turn.state !== "failed" || + turn.error !== "Worker heartbeat expired before the agent turn completed." + ) + return false; + const settled = await this.db + .select({ priced: turnUsage.priced }) + .from(turnUsage) + .where( + and( + eq(turnUsage.orgId, input.orgId), + eq(turnUsage.projectId, input.projectId), + eq(turnUsage.turnId, turn.id), + ), + ) + .limit(1); + if (settled[0]?.priced) return false; + if (!settled[0]) { + // Upgrade recovery: old workers marked failure without writing any usage row. + const phases = await this.db + .select({ data: turnEvents.data }) + .from(turnEvents) + .where( + and( + eq(turnEvents.orgId, input.orgId), + eq(turnEvents.projectId, input.projectId), + eq(turnEvents.turnId, turn.id), + eq(turnEvents.type, "turn.phase"), + ), + ); + if (!providerMayHaveRun(phases)) return false; + await new CostBudgetService(this.db).record({ + ...input, + storyId: turn.storyId, + agentName: turn.agentName, + engine: turn.engine, + model: turn.model, + durationMs: + turn.startedAt && turn.endedAt ? turn.endedAt.getTime() - turn.startedAt.getTime() : 0, + status: "failed", + complete: false, + }); + } + const evidence = await this.db + .select({ workspaceId: turnGitEvidence.workspaceId }) + .from(turnGitEvidence) + .where( + and( + eq(turnGitEvidence.orgId, input.orgId), + eq(turnGitEvidence.projectId, input.projectId), + eq(turnGitEvidence.turnId, turn.id), + ), + ) + .limit(1); + const workspace = ( + await this.db + .select() + .from(workspaces) + .where( + and( + eq(workspaces.orgId, input.orgId), + eq(workspaces.projectId, input.projectId), + eq(workspaces.storyId, turn.storyId), + ...(evidence[0]?.workspaceId ? [eq(workspaces.id, evidence[0].workspaceId)] : []), + ), + ) + .orderBy(desc(workspaces.createdAt)) + .limit(1) + )[0]; + if (!workspace?.externalRef) return false; + const measured = await readUsageJournal( + this.runtime, + locatorFromRow(workspace), + turn.id, + turn.engine, + ); + if (!measured.usage) return false; + const recorded = await new CostBudgetService(this.db).record({ + ...input, + storyId: turn.storyId, + agentName: turn.agentName, + engine: turn.engine, + model: turn.model, + ...measured, + durationMs: + turn.startedAt && turn.endedAt ? turn.endedAt.getTime() - turn.startedAt.getTime() : 0, + status: "failed", + }); + return recorded?.priced === true; + } + async flagAttention(input: { orgId: string; projectId: string; @@ -2296,6 +2458,16 @@ function workspaceConfiguration(input: Omit) { }; } +function providerMayHaveRun(phases: Array<{ data: unknown }>) { + if (phases.length === 0) return true; + return phases.some(({ data }) => { + const phase = + data && typeof data === "object" ? (data as { phase?: unknown }).phase : undefined; + // Only explicit pre-engine phases establish that no provider call began. + return !["credentials", "workspace", "environment"].includes(String(phase)); + }); +} + function workspaceInput(row: typeof workspaces.$inferSelect): Omit { const value = row.environment as { image?: unknown; diff --git a/services/api/src/stories/titles.ts b/services/api/src/stories/titles.ts index 35cdfe48..0d53bfae 100644 --- a/services/api/src/stories/titles.ts +++ b/services/api/src/stories/titles.ts @@ -143,6 +143,8 @@ export class StoryTitleService { const budget = await this.deps.budget.budgetState(input.orgId, input.projectId); if (budget.state === "exceeded") return this.settle(input, attempt, provider, "budget_exceeded"); + if (budget.state === "unconfirmed") + return this.settle(input, attempt, provider, "budget_usage_unconfirmed"); const model = this.deps.models?.[provider] ?? TITLE_MODELS[provider]; const apiKey = credentials[provider]; diff --git a/services/api/src/turns/engines.ts b/services/api/src/turns/engines.ts index cb9639b4..8979a143 100644 --- a/services/api/src/turns/engines.ts +++ b/services/api/src/turns/engines.ts @@ -4,6 +4,7 @@ import type { WorkspaceLocator, WorkspaceRuntime, } from "../workspaces/runtime.js"; +import { ENGINE_USAGE_PROCESS } from "./usage-journal.js"; export type AgentTurnRequest = { turnId: string; @@ -81,9 +82,22 @@ abstract class CliAgentEngine implements AgentEngine { try { result = await this.runtime.exec(request.workspace, { command: "sh", - args: ["-c", ENGINE_PROCESS_WRAPPER, "facility-engine", command, ...args], + args: [ + "-c", + ENGINE_PROCESS_WRAPPER, + "facility-engine", + "node", + "-e", + ENGINE_USAGE_PROCESS, + command, + ...args, + ], cwd: request.cwd, - env: { ...(request.environment ?? {}), FACILITY_TURN_ID: request.turnId }, + env: { + ...(request.environment ?? {}), + FACILITY_TURN_ID: request.turnId, + FACILITY_ENGINE: this.name, + }, timeoutMs: request.timeoutMs ?? 24 * 60 * 60 * 1_000, signal: request.signal, onObservation: (data) => diff --git a/services/api/src/turns/usage-journal.ts b/services/api/src/turns/usage-journal.ts new file mode 100644 index 00000000..032e3b98 --- /dev/null +++ b/services/api/src/turns/usage-journal.ts @@ -0,0 +1,131 @@ +import type { WorkspaceLocator, WorkspaceRuntime } from "../workspaces/runtime.js"; +import type { AgentTurnUsage } from "./engines.js"; + +/** Runs inside the durable workspace, independently of the observing worker. + * Persist only counters: transcripts, prompts and credentials never enter this journal. + */ +export const ENGINE_USAGE_PROCESS = String.raw` +const fs = require('node:fs'); +const path = require('node:path'); +const { spawn } = require('node:child_process'); +const turnId = process.env.FACILITY_TURN_ID; +const engine = process.env.FACILITY_ENGINE; +if (!/^turn_[A-Za-z0-9_-]+$/.test(turnId) || !['codex', 'claude_code'].includes(engine)) process.exit(64); +const dir = path.join(path.dirname(process.env.HOME), 'engine-usage'); +fs.mkdirSync(dir, { recursive: true, mode: 0o700 }); +const file = path.join(dir, turnId + '.json'); +let usage; +let complete = false; +let invalid = false; +let buffer = ''; +const messages = new Map(); +const keys = ['inputTokens', 'outputTokens', 'cacheReadTokens', 'cacheWriteTokens']; +const integer = n => Number.isSafeInteger(n) && n >= 0; +function counters(value) { + if (!value || typeof value !== 'object') return; + const values = [value.input_tokens, value.output_tokens, + value.cache_read_input_tokens ?? value.cached_input_tokens ?? 0, + value.cache_creation_input_tokens ?? 0]; + if (!values.every(integer)) return; + return Object.fromEntries(keys.map((key, i) => [key, values[i]])); +} +function save() { + fs.writeFileSync(file + '.tmp', JSON.stringify({ version: 1, turnId, engine, complete: complete && !invalid, usage }), { mode: 0o600 }); + fs.renameSync(file + '.tmp', file); +} +function accept(line) { + if (!line.trim()) return; + let event; + try { event = JSON.parse(line); } catch { invalid = true; save(); return; } + if (!event || typeof event !== 'object' || Array.isArray(event)) { invalid = true; save(); return; } + if (engine === 'claude_code' && event.type === 'assistant') { + const message = event.message; + const measured = counters(message?.usage); + if (measured && typeof message.id === 'string' && message.id.length <= 200) { + messages.set(message.id, measured); + usage = Object.fromEntries(keys.map(key => [key, [...messages.values()].reduce((sum, item) => sum + item[key], 0)])); + save(); + } + } + if ((engine === 'claude_code' && event.type === 'result') || (engine === 'codex' && event.type === 'turn.completed')) { + const measured = counters(event.usage); + if (measured) { + usage = measured; + if (engine === 'claude_code' && Number.isFinite(event.total_cost_usd) && event.total_cost_usd >= 0) + usage.reportedCostCents = event.total_cost_usd * 100; + complete = true; + } else invalid = true; + save(); + } +} +save(); +const child = spawn(process.argv[1], process.argv.slice(2), { stdio: ['inherit', 'pipe', 'inherit'] }); +// A dead observer closes stdout. Keep draining the CLI and journaling its usage. +process.stdout.on('error', () => {}); +for (const signal of ['SIGTERM', 'SIGINT', 'SIGHUP']) process.on(signal, () => child.kill(signal)); +child.stdout.setEncoding('utf8'); +child.stdout.on('data', chunk => { + // Journal first; losing the API worker must not lose the provider's final usage. + buffer += chunk; + let index; + while ((index = buffer.indexOf('\n')) >= 0) { accept(buffer.slice(0, index)); buffer = buffer.slice(index + 1); } + if (buffer.length > 4 * 1024 * 1024) { invalid = true; buffer = ''; save(); } + process.stdout.write(chunk); +}); +child.on('error', () => { invalid = true; save(); process.exitCode = 127; }); +child.on('close', code => { if (buffer) accept(buffer); save(); process.exitCode = code ?? 1; }); +`; + +export type RecoveredUsage = { usage?: AgentTurnUsage; complete: boolean }; + +export function parseUsageJournal(raw: string, turnId: string, engine: string): RecoveredUsage { + const unknown = { complete: false }; + if (raw.length > 4096) return unknown; + let value: Record; + try { + value = JSON.parse(raw); + } catch { + return unknown; + } + if (value?.version !== 1 || value.turnId !== turnId || value.engine !== engine) return unknown; + const usage = value.usage as AgentTurnUsage | undefined; + if ( + !usage || + [usage.inputTokens, usage.outputTokens, usage.cacheReadTokens, usage.cacheWriteTokens].some( + (n) => !Number.isSafeInteger(n) || n < 0, + ) + ) + return unknown; + if ( + usage.reportedCostCents !== undefined && + (!Number.isFinite(usage.reportedCostCents) || usage.reportedCostCents < 0) + ) + return unknown; + return { usage, complete: value.complete === true }; +} + +export async function readUsageJournal( + runtime: WorkspaceRuntime, + workspace: WorkspaceLocator, + turnId: string, + engine: string, +): Promise { + if (!/^turn_[A-Za-z0-9_-]+$/.test(turnId)) return { complete: false }; + try { + const result = await runtime.exec(workspace, { + resume: false, + command: "node", + args: [ + "-e", + "const fs=require('node:fs'),path=require('node:path');const file=path.join(path.dirname(process.env.HOME),'engine-usage',process.argv[1]+'.json');const fd=fs.openSync(file,'r');const buf=Buffer.alloc(4097);const size=fs.readSync(fd,buf,0,buf.length,0);fs.closeSync(fd);process.stdout.write(buf.subarray(0,size));", + turnId, + ], + timeoutMs: 15_000, + }); + return result.exitCode === 0 + ? parseUsageJournal(result.stdout, turnId, engine) + : { complete: false }; + } catch { + return { complete: false }; + } +} diff --git a/services/api/src/worker.ts b/services/api/src/worker.ts index 3dfbdfb5..9111710c 100644 --- a/services/api/src/worker.ts +++ b/services/api/src/worker.ts @@ -1,4 +1,4 @@ -import { createDb, type FacilityDb, turns } from "@facility/db"; +import { createDb, type FacilityDb, turns, turnUsage } from "@facility/db"; import { and, asc, eq, isNull, lte, or } from "drizzle-orm"; import PgBoss from "pg-boss"; import pino from "pino"; @@ -217,6 +217,33 @@ export async function recoverInterruptedTurns( }) => Promise, ) { const staleBefore = new Date(now.getTime() - leaseTimeoutMs); + const unsettled = await db + .select({ orgId: turns.orgId, projectId: turns.projectId, turnId: turns.id }) + .from(turns) + .leftJoin( + turnUsage, + and( + eq(turns.id, turnUsage.turnId), + eq(turns.orgId, turnUsage.orgId), + eq(turns.projectId, turnUsage.projectId), + ), + ) + .where( + and( + or(isNull(turnUsage.id), eq(turnUsage.priced, false)), + eq(turns.state, "failed"), + eq(turns.error, "Worker heartbeat expired before the agent turn completed."), + ), + ) + .orderBy(asc(turns.createdAt)) + .limit(1_000); + for (const turn of unsettled) + await stories.reconcileInterruptedUsage(turn).catch(() => { + // Preserve the accounting block and still recover other dead workers. + console.error( + JSON.stringify({ event: "turn.usage_reconciliation_failed", turnId: turn.turnId }), + ); + }); const running = await db .select({ id: turns.id, diff --git a/services/api/test/budget-accounting-policy.test.ts b/services/api/test/budget-accounting-policy.test.ts new file mode 100644 index 00000000..c4bbc034 --- /dev/null +++ b/services/api/test/budget-accounting-policy.test.ts @@ -0,0 +1,50 @@ +import type { FacilityDb } from "@facility/db"; +import { describe, expect, it, vi } from "vitest"; +import { type BudgetState, CostBudgetService } from "../src/insights/costs.js"; + +function service(state: BudgetState["state"], enabled = true) { + const costs = new CostBudgetService({} as FacilityDb); + const reading = { + state, + budget: + state === "not_configured" + ? null + : { + enabled, + monthlyLimitCents: 100_000, + }, + spentCents: state === "exceeded" ? 100_000 : 0, + remainingCents: 100_000, + } as BudgetState; + vi.spyOn(costs, "budgetState").mockResolvedValue(reading); + return costs; +} + +describe("budget accounting admission", () => { + it.each([ + "gpt-5.5", + "private-unpriced-model", + ])("refuses %s while spending is unconfirmed even with available budget", async (model) => { + await expect( + service("unconfirmed").assertTurnAllowed("org_test", "proj_test", model), + ).rejects.toMatchObject({ code: "budget_usage_unconfirmed" }); + }); + it("keeps the explicit disabled-budget escape under administrator control", async () => { + await expect( + service("disabled", false).assertTurnAllowed("org_test", "proj_test", "gpt-5.5"), + ).resolves.toMatchObject({ state: "disabled" }); + }); + it("allows priced work under a confirmed budget", async () => { + await expect( + service("ok").assertTurnAllowed("org_test", "proj_test", "gpt-5.5"), + ).resolves.toMatchObject({ state: "ok" }); + }); + it("preserves denial for unknown model pricing and an exhausted budget", async () => { + await expect( + service("ok").assertTurnAllowed("org_test", "proj_test", "private-unpriced-model"), + ).rejects.toMatchObject({ code: "budget_model_unpriced" }); + await expect( + service("exceeded").assertTurnAllowed("org_test", "proj_test", "gpt-5.5"), + ).rejects.toMatchObject({ code: "budget_exceeded" }); + }); +}); diff --git a/services/api/test/interrupted-usage.integration.test.ts b/services/api/test/interrupted-usage.integration.test.ts new file mode 100644 index 00000000..f6f086eb --- /dev/null +++ b/services/api/test/interrupted-usage.integration.test.ts @@ -0,0 +1,465 @@ +import { randomUUID } from "node:crypto"; +import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { newId } from "@facility/core"; +import { + createDb, + migrate, + orgs, + projectBudgets, + projects, + stories, + storyConversations, + turnGitEvidence, + turns, + turnUsage, + workspaces, +} from "@facility/db"; +import { eq } from "drizzle-orm"; +import { afterAll, beforeAll, describe, expect, it, vi } from "vitest"; +import { CostBudgetService } from "../src/insights/costs.js"; +import { StoryWorkspaceService } from "../src/stories/service.js"; +import { StoryTitleService } from "../src/stories/titles.js"; +import { appendTurnEvent } from "../src/turns/events.js"; +import { ENGINE_USAGE_PROCESS } from "../src/turns/usage-journal.js"; +import { recoverInterruptedTurns } from "../src/worker.js"; +import { FakeWorkspaceRuntime } from "../src/workspaces/fake.js"; + +const databaseUrl = + process.env.DATABASE_URL ?? "postgres://facility:facility@localhost:5461/facility_test"; +const { db, client } = createDb(databaseUrl); +const orgId = newId("org"); +const otherOrgId = newId("org"); +let root: string; +let runtime: FakeWorkspaceRuntime; +let service: StoryWorkspaceService; +const costs = new CostBudgetService(db); +const now = new Date(); +const stale = new Date(now.getTime() - 300_000); +const counters = { inputTokens: 100, outputTokens: 20, cacheReadTokens: 0, cacheWriteTokens: 0 }; + +beforeAll(async () => { + await migrate(databaseUrl); // A missing local database fails; these financial regressions never skip. + root = await mkdtemp(join(tmpdir(), "facility-interrupted-usage-")); + runtime = new FakeWorkspaceRuntime(root); + service = new StoryWorkspaceService(db, runtime); + await db.insert(orgs).values([ + { id: orgId, name: "Recovery usage", slug: `recovery-${randomUUID()}`, settings: {} }, + { id: otherOrgId, name: "Other", slug: `recovery-other-${randomUUID()}`, settings: {} }, + ]); +}); +afterAll(async () => { + await client.end(); + if (root) await rm(root, { recursive: true, force: true }); +}); + +async function seed(engine: "codex" | "claude_code" = "codex", startedAt = stale) { + const projectId = newId("proj"), + storyId = newId("story"), + conversationId = newId("sess"), + turnId = newId("turn"), + workspaceId = newId("ws"); + await db + .insert(projects) + .values({ id: projectId, orgId, name: "Recovery", slug: `usage-${projectId}`, settings: {} }); + await db + .insert(projectBudgets) + .values({ id: newId("bud"), orgId, projectId, monthlyLimitCents: 100_000, enabled: true }); + await db.insert(stories).values({ + id: storyId, + orgId, + projectId, + provider: "manual", + externalId: storyId, + title: "Interrupted", + titleSource: "pending", + status: "working", + createdBy: { type: "user", id: "test" }, + }); + await db.insert(storyConversations).values({ id: conversationId, orgId, projectId, storyId }); + const workspace = await runtime.create({ id: workspaceId, image: "fake" }); + await db.insert(workspaces).values({ + id: workspaceId, + orgId, + projectId, + storyId, + provider: "fake", + volumeRef: workspace.volumeRef, + externalRef: workspace.externalRef, + environment: { image: "fake", ports: [] }, + state: "running", + }); + await db.insert(turns).values({ + id: turnId, + orgId, + projectId, + storyId, + conversationId, + agentName: "builder", + manifestHash: "hash", + manifest: {}, + engine, + model: "gpt-5.5", + state: "running", + triggerType: "manual", + createdBy: { type: "user", id: "test" }, + startedAt, + updatedAt: startedAt, + }); + return { + orgId, + projectId, + storyId, + turnId, + workspace, + engine, + staleBefore: new Date(now.getTime() - 60_000), + }; +} +async function journal(input: Awaited>, data: unknown) { + const dir = join(input.workspace.volumeRef, ".facility", "engine-usage"); + await mkdir(dir, { recursive: true }); + await writeFile( + join(dir, `${input.turnId}.json`), + typeof data === "string" ? data : JSON.stringify(data), + ); +} +function evidence(input: Awaited>, complete = true) { + return { version: 1, turnId: input.turnId, engine: input.engine, usage: counters, complete }; +} + +describe("dead worker accounting", () => { + it("recovers a native CLI's durable final usage into the monthly budget exactly once", async () => { + const input = await seed(); + // No live provider: a real workspace process emits deterministic native CLI JSONL. + await runtime.exec(input.workspace, { + command: process.execPath, + args: [ + "-e", + ENGINE_USAGE_PROCESS, + process.execPath, + "-e", + `console.log(JSON.stringify({type:'turn.completed',usage:{input_tokens:100,output_tokens:20}}))`, + ], + env: { FACILITY_TURN_ID: input.turnId, FACILITY_ENGINE: input.engine }, + }); + expect(await service.recoverInterruptedTurn(input)).toBe(true); + const recorded = await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId)); + expect(recorded).toEqual([ + expect.objectContaining({ ...counters, status: "failed", priced: true }), + ]); + expect((await costs.budgetState(orgId, input.projectId)).spentCents).toBeGreaterThan(0); + expect(await service.recoverInterruptedTurn(input)).toBe(false); + expect(await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId))).toEqual( + recorded, + ); + await expect(costs.assertTurnAllowed(orgId, input.projectId, "gpt-5.5")).resolves.toMatchObject( + { state: "ok" }, + ); + await db + .update(projectBudgets) + .set({ monthlyLimitCents: 0 }) + .where(eq(projectBudgets.projectId, input.projectId)); + await expect(costs.assertTurnAllowed(orgId, input.projectId, "gpt-5.5")).rejects.toMatchObject({ + code: "budget_exceeded", + }); + }); + it("books a partial lower bound but blocks new calls even after raising the budget", async () => { + const input = await seed(); + await journal(input, evidence(input, false)); + await service.recoverInterruptedTurn(input); + expect(await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId))).toEqual([ + expect.objectContaining({ ...counters, priced: false, source: "unpriced" }), + ]); + expect((await costs.budgetState(orgId, input.projectId)).spentCents).toBeGreaterThan(0); + await db + .update(projectBudgets) + .set({ monthlyLimitCents: 2_000_000_000 }) + .where(eq(projectBudgets.projectId, input.projectId)); + await expect(costs.assertTurnAllowed(orgId, input.projectId, "gpt-5.5")).rejects.toMatchObject({ + code: "budget_usage_unconfirmed", + }); + expect(await costs.budgetState(otherOrgId, input.projectId)).toMatchObject({ + spentCents: 0, + state: "not_configured", + }); + await db + .update(projectBudgets) + .set({ enabled: false }) + .where(eq(projectBudgets.projectId, input.projectId)); + await expect(costs.assertTurnAllowed(orgId, input.projectId, "gpt-5.5")).resolves.toMatchObject( + { state: "disabled" }, + ); + }); + it("does not consume a newer workspace's usage when the recorded workspace was destroyed", async () => { + const input = await seed(); + await journal(input, evidence(input)); + await db.insert(turnGitEvidence).values({ + orgId, + projectId: input.projectId, + storyId: input.storyId, + turnId: input.turnId, + workspaceId: input.workspace.id, + engineSessionId: newId("esess"), + initialSha: "a".repeat(40), + }); + await runtime.destroy(input.workspace); + await db + .update(workspaces) + .set({ state: "destroyed", destroyedAt: new Date() }) + .where(eq(workspaces.id, input.workspace.id)); + const newer = await runtime.create({ id: newId("ws"), image: "fake" }); + await db.insert(workspaces).values({ + id: newer.id, + orgId, + projectId: input.projectId, + storyId: input.storyId, + provider: "fake", + externalRef: newer.externalRef, + volumeRef: newer.volumeRef, + environment: { image: "fake", ports: [] }, + state: "running", + createdAt: new Date(now.getTime() + 1000), + }); + await journal( + { ...input, workspace: newer }, + { ...evidence(input), usage: { ...counters, inputTokens: 9999 } }, + ); + await service.recoverInterruptedTurn(input); + expect(await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId))).toEqual([ + expect.objectContaining({ inputTokens: 0, costCents: null, priced: false }), + ]); + }); + it.each([ + "absent", + "malformed", + "wrong-turn", + "wrong-engine", + "negative", + "oversized", + "unavailable", + ])("keeps %s evidence unknown and fails closed", async (kind) => { + const input = await seed(); + if (kind === "malformed") await journal(input, "{"); + if (kind === "wrong-turn") await journal(input, { ...evidence(input), turnId: "turn_other" }); + if (kind === "wrong-engine") + await journal(input, { ...evidence(input), engine: "claude_code" }); + if (kind === "negative") + await journal(input, { ...evidence(input), usage: { ...counters, inputTokens: -1 } }); + if (kind === "oversized") await journal(input, "x".repeat(5000)); + if (kind === "unavailable") await runtime.suspend(input.workspace); + await service.recoverInterruptedTurn(input); + expect(await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId))).toEqual([ + expect.objectContaining({ costCents: null, priced: false }), + ]); + await expect(costs.assertTurnAllowed(orgId, input.projectId, "gpt-5.5")).rejects.toMatchObject({ + code: "budget_usage_unconfirmed", + }); + }); + it("rejects cross-tenant and cross-project recovery before touching a journal", async () => { + const input = await seed(); + await journal(input, evidence(input)); + await expect( + service.recoverInterruptedTurn({ ...input, orgId: otherOrgId }), + ).rejects.toMatchObject({ code: "turn_not_found" }); + const other = await seed(); + await expect( + service.recoverInterruptedTurn({ ...input, projectId: other.projectId }), + ).rejects.toMatchObject({ code: "turn_not_found" }); + expect(await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId))).toEqual([]); + await db.update(turns).set({ state: "failed" }).where(eq(turns.id, other.turnId)); + }); + it("leaves a live lease alone and does not book an old or replayed journal", async () => { + const input = await seed("codex", now); + await journal(input, evidence(input)); + expect(await service.recoverInterruptedTurn(input)).toBe(false); + expect(await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId))).toEqual([]); + }); + it("does not block for an interrupted preparation phase that never reached the engine", async () => { + const input = await seed(); + await appendTurnEvent(db, { + orgId, + projectId: input.projectId, + storyId: input.storyId, + turnId: input.turnId, + type: "turn.phase", + data: { phase: "environment" }, + }); + await service.recoverInterruptedTurn(input); + expect(await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId))).toEqual([]); + expect((await costs.budgetState(orgId, input.projectId)).state).toBe("ok"); + }); + it.each([ + {}, + { phase: "unknown-stage" }, + ])("does not infer zero spend from malformed or unknown phase evidence %#", async (data) => { + const input = await seed(); + await appendTurnEvent(db, { + orgId, + projectId: input.projectId, + storyId: input.storyId, + turnId: input.turnId, + type: "turn.phase", + data, + }); + await service.recoverInterruptedTurn(input); + await expect(costs.assertTurnAllowed(orgId, input.projectId, "gpt-5.5")).rejects.toMatchObject({ + code: "budget_usage_unconfirmed", + }); + }); + it("keeps an already settled charge when recovery sees an incomplete journal", async () => { + const input = await seed(); + await costs.record({ + ...input, + agentName: "builder", + model: "gpt-5.5", + usage: counters, + durationMs: 10, + status: "failed", + }); + const before = await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId)); + await journal(input, evidence(input, false)); + await service.recoverInterruptedTurn(input); + expect(await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId))).toEqual( + before, + ); + }); + it("bookmarks uncertainty before the worker activates a queued successor", async () => { + const input = await seed(); + const callback = vi.fn(async ({ projectId }: { projectId: string }) => { + if (projectId === input.projectId) + await expect(costs.assertTurnAllowed(orgId, projectId, "gpt-5.5")).rejects.toMatchObject({ + code: "budget_usage_unconfirmed", + }); + }); + await recoverInterruptedTurns(db, service, now, 60_000, callback); + expect(callback).toHaveBeenCalledWith(expect.objectContaining({ turnId: input.turnId })); + }); + it("prevents title generation from making a model call while spending is unconfirmed", async () => { + const input = await seed(); + await service.recoverInterruptedTurn(input); + const complete = vi.fn(async () => ({ + title: "Must not call", + inputTokens: 1, + outputTokens: 1, + })); + const titleService = new StoryTitleService(db, { + credentials: async () => ({ openai: "local-fake" }), + budget: costs, + complete, + }); + // The title service must see a stored request before consulting the budget. + const { storyMessages } = await import("@facility/db"); + const turn = (await db.select().from(turns).where(eq(turns.id, input.turnId)))[0]; + if (!turn) throw new Error("fixture turn missing"); + await db.insert(storyMessages).values({ + id: newId("msg"), + orgId, + projectId: input.projectId, + storyId: input.storyId, + conversationId: turn.conversationId, + seq: 1, + role: "user", + body: "Generate a title", + actor: { type: "user", id: "test" }, + }); + expect(await titleService.generate(input)).toMatchObject({ + reason: "budget_usage_unconfirmed", + }); + expect(complete).not.toHaveBeenCalled(); + }); + it("does not clear unknown spending when the month rolls over", async () => { + const previous = new Date(Date.UTC(now.getUTCFullYear(), now.getUTCMonth(), 0, 12)); + const input = await seed("codex", previous); + await service.recoverInterruptedTurn(input); + expect((await costs.budgetState(orgId, input.projectId, now)).spentCents).toBe(0); + await expect( + costs.assertTurnAllowed( + orgId, + input.projectId, + "gpt-5.5", + new Date(Date.UTC(now.getUTCFullYear(), now.getUTCMonth() + 1, 1)), + ), + ).rejects.toMatchObject({ code: "budget_usage_unconfirmed" }); + }); + it("settles later durable evidence on a recovery pass without rerunning or waking compute", async () => { + const input = await seed(); + await service.recoverInterruptedTurn(input); + await journal(input, evidence(input)); + const wake = vi.spyOn(runtime, "wake"); + await recoverInterruptedTurns(db, service, now, 60_000); + expect(wake).not.toHaveBeenCalled(); + wake.mockRestore(); + await expect(costs.assertTurnAllowed(orgId, input.projectId, "gpt-5.5")).resolves.toMatchObject( + { state: "ok" }, + ); + const before = await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId)); + await journal(input, evidence(input, false)); + expect(await service.reconcileInterruptedUsage(input)).toBe(false); + expect(await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId))).toEqual( + before, + ); + }); + it("marks a pre-upgrade failed worker's missing bill unknown during recovery", async () => { + const input = await seed(); + await db + .update(turns) + .set({ + state: "failed", + endedAt: now, + error: "Worker heartbeat expired before the agent turn completed.", + }) + .where(eq(turns.id, input.turnId)); + await recoverInterruptedTurns(db, service, now, 60_000); + expect(await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId))).toEqual([ + expect.objectContaining({ priced: false, costCents: null }), + ]); + await expect(costs.assertTurnAllowed(orgId, input.projectId, "gpt-5.5")).rejects.toMatchObject({ + code: "budget_usage_unconfirmed", + }); + }); + it("allows a late live dispatcher to settle an unknown row exactly once", async () => { + const input = await seed(); + await service.recoverInterruptedTurn(input); + const request = { + ...input, + agentName: "builder", + model: "gpt-5.5", + usage: counters, + durationMs: 10, + status: "failed" as const, + }; + expect(await costs.record(request)).toMatchObject({ priced: true }); + expect(await costs.record(request)).toBeNull(); + expect((await costs.budgetState(orgId, input.projectId)).state).toBe("ok"); + }); + it("counts a previous-month worker's recovered charge when it is actually settled", async () => { + const previous = new Date(Date.UTC(now.getUTCFullYear(), now.getUTCMonth(), 0, 12)); + const input = await seed("codex", previous); + await journal(input, evidence(input)); + await service.recoverInterruptedTurn(input); + expect((await costs.budgetState(orgId, input.projectId, previous)).spentCents).toBe(0); + expect((await costs.budgetState(orgId, input.projectId, now)).spentCents).toBeGreaterThan(0); + }); + it("charges this month's budget when an old unknown placeholder is finally settled", async () => { + const previous = new Date(Date.UTC(now.getUTCFullYear(), now.getUTCMonth(), 0, 12)); + const input = await seed("codex", previous); + await service.recoverInterruptedTurn(input); + await db + .update(turnUsage) + .set({ createdAt: previous }) + .where(eq(turnUsage.turnId, input.turnId)); + await journal(input, evidence(input)); + await service.reconcileInterruptedUsage(input); + expect((await costs.budgetState(orgId, input.projectId, now)).spentCents).toBeGreaterThan(0); + expect((await costs.budgetState(orgId, input.projectId, previous)).spentCents).toBe(0); + await db + .update(projectBudgets) + .set({ monthlyLimitCents: 0 }) + .where(eq(projectBudgets.projectId, input.projectId)); + await expect( + costs.assertTurnAllowed(orgId, input.projectId, "gpt-5.5", now), + ).rejects.toMatchObject({ code: "budget_exceeded" }); + }); +}); diff --git a/services/api/test/project-overview.integration.test.ts b/services/api/test/project-overview.integration.test.ts index aaf267d7..5b3f5af6 100644 --- a/services/api/test/project-overview.integration.test.ts +++ b/services/api/test/project-overview.integration.test.ts @@ -756,7 +756,8 @@ describe("project overview", async () => { ]); if (!overview.spend.budget.available) throw new Error("budget must be available"); expect(overview.spend.budget).toMatchObject({ - state: "warning", + // An unpriced turn leaves the true bill unknown even below the warning ceiling. + state: "unconfirmed", monthlyLimitCents: 1_000, spentCents: 300, }); diff --git a/services/api/test/usage-journal.test.ts b/services/api/test/usage-journal.test.ts new file mode 100644 index 00000000..542d08eb --- /dev/null +++ b/services/api/test/usage-journal.test.ts @@ -0,0 +1,119 @@ +import { spawn } from "node:child_process"; +import { mkdir, mkdtemp, readFile, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { afterEach, describe, expect, it } from "vitest"; +import { ENGINE_USAGE_PROCESS, parseUsageJournal } from "../src/turns/usage-journal.js"; + +const roots: string[] = []; +afterEach(async () => { + for (const root of roots.splice(0)) await rm(root, { recursive: true, force: true }); +}); +const usage = { inputTokens: 100, outputTokens: 20, cacheReadTokens: 3, cacheWriteTokens: 2 }; +const native = { + input_tokens: 100, + output_tokens: 20, + cache_read_input_tokens: 3, + cache_creation_input_tokens: 2, +}; +function journal(overrides: Record = {}) { + return JSON.stringify({ + version: 1, + turnId: "turn_test", + engine: "codex", + complete: true, + usage, + ...overrides, + }); +} +async function run(engine: string, events: unknown[], loseObserver = false) { + const root = await mkdtemp(join(tmpdir(), "facility-usage-unit-")); + roots.push(root); + const home = join(root, "home"); + await mkdir(home); + const fake = `const events=${JSON.stringify(events)};for(const event of events)process.stdout.write(JSON.stringify(event)+'\\n');`; + const child = spawn( + process.execPath, + ["-e", ENGINE_USAGE_PROCESS, process.execPath, "-e", fake], + { + env: { ...process.env, HOME: home, FACILITY_TURN_ID: "turn_test", FACILITY_ENGINE: engine }, + stdio: ["ignore", "pipe", "pipe"], + }, + ); + if (loseObserver) child.stdout.destroy(); + else child.stdout.resume(); + const errors: Buffer[] = []; + child.stderr.on("data", (data) => errors.push(data)); + const code = await new Promise((resolve, reject) => { + child.on("error", reject); + child.on("close", resolve); + }); + expect(Buffer.concat(errors).toString()).toBe(""); + expect(code).toBe(0); + return readFile(join(root, "engine-usage", "turn_test.json"), "utf8"); +} + +describe("durable per-turn usage", () => { + it.each([ + false, + true, + ])("books Codex terminal counters even with lost observer=%s", async (lost) => { + const raw = await run("codex", [{ type: "turn.completed", usage: native }], lost); + expect(parseUsageJournal(raw, "turn_test", "codex")).toEqual({ complete: true, usage }); + }); + it("deduplicates Claude message counters and keeps partial usage unconfirmed", async () => { + const message = { + type: "assistant", + message: { id: "msg-1", usage: native, content: [{ text: "private-output" }] }, + }; + const raw = await run("claude_code", [message, message]); + expect(parseUsageJournal(raw, "turn_test", "claude_code")).toEqual({ complete: false, usage }); + expect(raw).not.toContain("private-output"); + expect(raw).not.toContain("msg-1"); + }); + it("uses Claude's final provider total instead of adding it to partial counters", async () => { + const raw = await run("claude_code", [ + { type: "assistant", message: { id: "msg-1", usage: native } }, + { type: "result", usage: native, total_cost_usd: 0.03 }, + ]); + expect(parseUsageJournal(raw, "turn_test", "claude_code")).toEqual({ + complete: true, + usage: { ...usage, reportedCostCents: 3 }, + }); + }); + it("does not claim zero spending when a CLI exits without usage", async () => { + expect( + parseUsageJournal( + await run("codex", [{ type: "thread.started", thread_id: "thread" }]), + "turn_test", + "codex", + ), + ).toEqual({ complete: false }); + }); + it.each([ + "", + "null", + "[]", + "{", + "x".repeat(4097), + journal({ version: 2 }), + journal({ turnId: "turn_other" }), + journal({ engine: "claude_code" }), + journal({ usage: { ...usage, inputTokens: -1 } }), + journal({ usage: { ...usage, outputTokens: 1.5 } }), + journal({ usage: { ...usage, reportedCostCents: -1 } }), + ])("rejects malformed or unrelated evidence %#", (raw) => { + expect(parseUsageJournal(raw, "turn_test", "codex")).toEqual({ complete: false }); + }); + it("rejects malformed terminal counters rather than treating them as a paid zero", async () => { + expect( + parseUsageJournal( + await run("codex", [ + { type: "turn.completed", usage: { input_tokens: -1, output_tokens: 0 } }, + ]), + "turn_test", + "codex", + ), + ).toEqual({ complete: false }); + }); +}); From ea4da90fc1716075b51b878c0451d42d394a54ad Mon Sep 17 00:00:00 2001 From: shuntianyifang <2016860013@qq.com> Date: Wed, 7 Oct 2026 16:26:14 +0800 Subject: [PATCH 2/3] test(api): verify durable usage in production workspaces --- .../workspace-runtime.integration.test.ts | 79 +++++++++++++++++++ 1 file changed, 79 insertions(+) diff --git a/services/api/test/workspace-runtime.integration.test.ts b/services/api/test/workspace-runtime.integration.test.ts index 969af72a..7c804d3d 100644 --- a/services/api/test/workspace-runtime.integration.test.ts +++ b/services/api/test/workspace-runtime.integration.test.ts @@ -3,6 +3,7 @@ import { readFile } from "node:fs/promises"; import Docker from "dockerode"; import { describe, expect, it } from "vitest"; import { stopInterruptedEngineProcess } from "../src/turns/engines.js"; +import { ENGINE_USAGE_PROCESS, readUsageJournal } from "../src/turns/usage-journal.js"; import { DockerWorkspaceRuntime } from "../src/workspaces/docker.js"; import { exportWorkspaceBackup, @@ -13,6 +14,84 @@ import { const enabled = process.env.FACILITY_E2E_DOCKER === "1"; describe.skipIf(!enabled)("DockerWorkspaceRuntime integration", () => { + it("retains per-turn usage after the observer exits and compute is replaced", async () => { + const runtime = new DockerWorkspaceRuntime(new Docker()); + const workspace = await runtime.create({ + id: `ws_${randomBytes(12).toString("hex")}`, + image: process.env.FACILITY_WORKSPACE_TEST_IMAGE ?? "facility-runner:serialized", + }); + const usage = { inputTokens: 100, outputTokens: 20, cacheReadTokens: 3, cacheWriteTokens: 2 }; + const native = { + input_tokens: 100, + output_tokens: 20, + cache_read_input_tokens: 3, + cache_creation_input_tokens: 2, + }; + try { + for (const engine of ["codex", "claude_code"] as const) { + const turnId = `turn_docker_${engine}`; + const event = + engine === "codex" + ? { type: "turn.completed", usage: native } + : { type: "result", usage: native, total_cost_usd: 0.03, result: "private-output" }; + const launched = await runtime.exec(workspace, { + command: "sh", + args: [ + "-lc", + 'nohup node -e "$FACILITY_USAGE_SCRIPT" node -e "$FACILITY_FAKE_CLI" /dev/null 2>&1 &', + ], + env: { + FACILITY_TURN_ID: turnId, + FACILITY_ENGINE: engine, + FACILITY_USAGE_SCRIPT: ENGINE_USAGE_PROCESS, + FACILITY_FAKE_CLI: `setTimeout(() => console.log(JSON.stringify(${JSON.stringify(event)})), 100);`, + }, + }); + expect(launched.exitCode, launched.stderr).toBe(0); + await waitUntil( + async () => (await readUsageJournal(runtime, workspace, turnId, engine)).complete, + ); + expect(await readUsageJournal(runtime, workspace, turnId, engine)).toEqual({ + complete: true, + usage: engine === "codex" ? usage : { ...usage, reportedCostCents: 3 }, + }); + const contents = await runtime.exec(workspace, { + command: "sh", + args: ["-lc", 'cat "$(dirname "$HOME")/engine-usage/$FACILITY_TURN_ID.json"'], + env: { FACILITY_TURN_ID: turnId }, + }); + expect(contents.stdout).not.toContain("private-output"); + } + const partial = await runtime.exec(workspace, { + command: "node", + args: [ + "-e", + ENGINE_USAGE_PROCESS, + "node", + "-e", + `console.log(JSON.stringify({type:'assistant',message:{id:'msg-partial',usage:${JSON.stringify(native)}}}))`, + ], + env: { FACILITY_TURN_ID: "turn_docker_partial", FACILITY_ENGINE: "claude_code" }, + }); + expect(partial.exitCode, partial.stderr).toBe(0); + expect( + await readUsageJournal(runtime, workspace, "turn_docker_partial", "claude_code"), + ).toEqual({ complete: false, usage }); + await runtime.replaceCompute(workspace); + expect(await readUsageJournal(runtime, workspace, "turn_docker_codex", "codex")).toEqual({ + complete: false, + }); + await expect(runtime.inspect(workspace)).resolves.toMatchObject({ state: "sleeping" }); + await runtime.wake(workspace); + expect(await readUsageJournal(runtime, workspace, "turn_docker_codex", "codex")).toEqual({ + complete: true, + usage, + }); + } finally { + await runtime.destroy(workspace); + } + }, 180_000); + it("reattaches the same named volume after its compute is removed", async () => { const id = `ws_${randomBytes(12).toString("hex")}`; const gatewayToken = randomBytes(32).toString("base64url"); From 45dc057b9075e083f59c596d031ca06847f90ff3 Mon Sep 17 00:00:00 2001 From: shuntianyifang <2016860013@qq.com> Date: Wed, 7 Oct 2026 18:39:35 +0800 Subject: [PATCH 3/3] test(api): protect interrupted accounting during concurrent settlement --- .../interrupted-usage.integration.test.ts | 53 +++++++++++++++++++ 1 file changed, 53 insertions(+) diff --git a/services/api/test/interrupted-usage.integration.test.ts b/services/api/test/interrupted-usage.integration.test.ts index f6f086eb..1b9a8357 100644 --- a/services/api/test/interrupted-usage.integration.test.ts +++ b/services/api/test/interrupted-usage.integration.test.ts @@ -434,6 +434,59 @@ describe("dead worker accounting", () => { expect(await costs.record(request)).toBeNull(); expect((await costs.budgetState(orgId, input.projectId)).state).toBe("ok"); }); + it("keeps the project blocked until every interrupted turn is settled", async () => { + const input = await seed(); + const secondId = newId("turn"); + const original = (await db.select().from(turns).where(eq(turns.id, input.turnId)))[0]; + if (!original) throw new Error("fixture turn missing"); + await service.recoverInterruptedTurn(input); + await db.insert(turns).values({ ...original, id: secondId }); + await service.recoverInterruptedTurn({ ...input, turnId: secondId }); + const request = { + ...input, + agentName: "builder", + model: "gpt-5.5", + usage: counters, + durationMs: 10, + status: "failed" as const, + }; + await costs.record(request); + await expect(costs.assertTurnAllowed(orgId, input.projectId, "gpt-5.5")).rejects.toMatchObject({ + code: "budget_usage_unconfirmed", + }); + await costs.record({ ...request, turnId: secondId }); + expect((await costs.budgetState(orgId, input.projectId)).state).toBe("ok"); + const rows = await db.select().from(turnUsage).where(eq(turnUsage.projectId, input.projectId)); + expect(rows).toHaveLength(2); + expect(rows.every((row) => row.priced)).toBe(true); + }); + it("settles competing final reports once without allowing a replay to change the charge", async () => { + const input = await seed(); + await service.recoverInterruptedTurn(input); + const request = { + ...input, + agentName: "builder", + model: "gpt-5.5", + durationMs: 10, + status: "failed" as const, + }; + const results = await Promise.all([ + costs.record({ ...request, usage: { ...counters, reportedCostCents: 11 } }), + costs.record({ ...request, usage: { ...counters, reportedCostCents: 17 } }), + ]); + const written = results.filter((row) => row !== null); + expect(written).toHaveLength(1); + const before = await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId)); + expect(before).toHaveLength(1); + expect([11, 17]).toContain(before[0]?.costCents); + expect((await costs.budgetState(orgId, input.projectId)).spentCents).toBe(before[0]?.costCents); + expect( + await costs.record({ ...request, usage: { ...counters, reportedCostCents: 0 } }), + ).toBeNull(); + expect(await db.select().from(turnUsage).where(eq(turnUsage.turnId, input.turnId))).toEqual( + before, + ); + }); it("counts a previous-month worker's recovered charge when it is actually settled", async () => { const previous = new Date(Date.UTC(now.getUTCFullYear(), now.getUTCMonth(), 0, 12)); const input = await seed("codex", previous);