feat(FN-4389): complete Step 8 — run verification and fix regressions
Fusion-Task-Id: FN-4389 Fusion-Task-Lineage: 366f9d57-c476-448f-96c0-6e753fc8b090
This commit is contained in:
@@ -2145,7 +2145,8 @@ describe("StepSessionExecutor integration", () => {
|
||||
tokenUsage: expect.objectContaining({
|
||||
inputTokens: 31,
|
||||
outputTokens: 17,
|
||||
cachedTokens: 7,
|
||||
cachedTokens: 5,
|
||||
cacheWriteTokens: 2,
|
||||
totalTokens: 55,
|
||||
}),
|
||||
}),
|
||||
|
||||
@@ -530,7 +530,7 @@ describe("Budget Governance", () => {
|
||||
|
||||
await monitor.completeRun("agent-001", "run-budget-001", {
|
||||
status: "completed",
|
||||
usageJson: { inputTokens: 0, outputTokens: 100, cachedTokens: 0 },
|
||||
usageJson: { inputTokens: 0, outputTokens: 100, cachedTokens: 0, cacheWriteTokens: 0 },
|
||||
});
|
||||
|
||||
expect(store.updateAgentState).toHaveBeenCalledWith("agent-001", "paused");
|
||||
@@ -553,7 +553,7 @@ describe("Budget Governance", () => {
|
||||
|
||||
await monitor.completeRun("agent-001", "run-budget-001", {
|
||||
status: "completed",
|
||||
usageJson: { inputTokens: 10, outputTokens: 50, cachedTokens: 0 },
|
||||
usageJson: { inputTokens: 10, outputTokens: 50, cachedTokens: 0, cacheWriteTokens: 0 },
|
||||
});
|
||||
|
||||
expect(store.updateAgentState).toHaveBeenCalledWith("agent-001", "active");
|
||||
@@ -568,7 +568,7 @@ describe("Budget Governance", () => {
|
||||
|
||||
await monitor.completeRun("agent-001", "run-budget-001", {
|
||||
status: "failed",
|
||||
usageJson: { inputTokens: 10, outputTokens: 50, cachedTokens: 0 },
|
||||
usageJson: { inputTokens: 10, outputTokens: 50, cachedTokens: 0, cacheWriteTokens: 0 },
|
||||
stderrExcerpt: "failure",
|
||||
});
|
||||
|
||||
@@ -585,7 +585,7 @@ describe("Budget Governance", () => {
|
||||
|
||||
await monitor.completeRun("agent-001", "run-budget-001", {
|
||||
status: "terminated",
|
||||
usageJson: { inputTokens: 10, outputTokens: 50, cachedTokens: 0 },
|
||||
usageJson: { inputTokens: 10, outputTokens: 50, cachedTokens: 0, cacheWriteTokens: 0 },
|
||||
});
|
||||
|
||||
expect(store.getBudgetStatus).not.toHaveBeenCalled();
|
||||
|
||||
@@ -38,9 +38,10 @@ describe("accumulateSessionTokenUsage", () => {
|
||||
expect(store.updateTask).toHaveBeenCalledTimes(1);
|
||||
const call = store.updateTask.mock.calls[0]![1] as { tokenUsage: Task["tokenUsage"] };
|
||||
expect(call.tokenUsage).toMatchObject({
|
||||
inputTokens: 102, // input + cacheWrite
|
||||
inputTokens: 100,
|
||||
outputTokens: 30,
|
||||
cachedTokens: 5,
|
||||
cacheWriteTokens: 2,
|
||||
totalTokens: 137,
|
||||
});
|
||||
expect(typeof call.tokenUsage!.firstUsedAt).toBe("string");
|
||||
@@ -102,6 +103,7 @@ describe("accumulateSessionTokenUsage", () => {
|
||||
inputTokens: 50,
|
||||
outputTokens: 20,
|
||||
cachedTokens: 0,
|
||||
cacheWriteTokens: 0,
|
||||
totalTokens: 70,
|
||||
firstUsedAt: "2024-01-01T00:00:00.000Z",
|
||||
lastUsedAt: "2024-01-01T00:00:00.000Z",
|
||||
|
||||
@@ -843,7 +843,8 @@ describe("StepSessionExecutor", () => {
|
||||
tokenUsage: {
|
||||
inputTokens: 25,
|
||||
outputTokens: 11,
|
||||
cachedTokens: 6,
|
||||
cachedTokens: 4,
|
||||
cacheWriteTokens: 2,
|
||||
totalTokens: 42,
|
||||
},
|
||||
});
|
||||
|
||||
@@ -1095,7 +1095,7 @@ export class HeartbeatMonitor {
|
||||
status: "completed" | "failed" | "terminated";
|
||||
exitCode?: number;
|
||||
sessionIdAfter?: string;
|
||||
usageJson?: { inputTokens: number; outputTokens: number; cachedTokens: number };
|
||||
usageJson?: { inputTokens: number; outputTokens: number; cachedTokens: number; cacheWriteTokens: number };
|
||||
resultJson?: Record<string, unknown>;
|
||||
stdoutExcerpt?: string;
|
||||
stderrExcerpt?: string;
|
||||
@@ -2446,15 +2446,17 @@ export class HeartbeatMonitor {
|
||||
let usageInput = 0;
|
||||
let usageOutput = Math.ceil(outputLength / 4);
|
||||
let usageCached = 0;
|
||||
let usageCacheWrite = 0;
|
||||
try {
|
||||
const sessionStats = (session as unknown as {
|
||||
getSessionStats?: () => { tokens?: { input?: number; output?: number; cacheRead?: number; cacheWrite?: number } };
|
||||
}).getSessionStats?.();
|
||||
const tokens = sessionStats?.tokens;
|
||||
if (tokens) {
|
||||
usageInput = (tokens.input ?? 0) + (tokens.cacheWrite ?? 0);
|
||||
usageInput = tokens.input ?? 0;
|
||||
usageOutput = tokens.output ?? usageOutput;
|
||||
usageCached = tokens.cacheRead ?? 0;
|
||||
usageCacheWrite = tokens.cacheWrite ?? 0;
|
||||
}
|
||||
} catch (statsErr) {
|
||||
heartbeatLog.warn(`Agent ${agentId} session stats read failed: ${statsErr instanceof Error ? statsErr.message : String(statsErr)} — using estimated tokens`);
|
||||
@@ -2496,7 +2498,7 @@ export class HeartbeatMonitor {
|
||||
|
||||
await this.completeRun(agentId, run.id, {
|
||||
status: "completed",
|
||||
usageJson: { inputTokens: usageInput, outputTokens: usageOutput, cachedTokens: usageCached },
|
||||
usageJson: { inputTokens: usageInput, outputTokens: usageOutput, cachedTokens: usageCached, cacheWriteTokens: usageCacheWrite },
|
||||
resultJson: completionResultJson,
|
||||
stdoutExcerpt: stdoutExcerpt || undefined,
|
||||
});
|
||||
@@ -2509,7 +2511,7 @@ export class HeartbeatMonitor {
|
||||
}
|
||||
}
|
||||
|
||||
heartbeatLog.log(`Heartbeat completed for ${agentId} (${toolCallCount} tool calls, ${usageInput} input + ${usageOutput} output + ${usageCached} cached tokens)`);
|
||||
heartbeatLog.log(`Heartbeat completed for ${agentId} (${toolCallCount} tool calls, ${usageInput} input + ${usageOutput} output + ${usageCached} cache-read + ${usageCacheWrite} cache-write tokens)`);
|
||||
} catch (err) {
|
||||
const errorDetail = formatError(err).detail;
|
||||
heartbeatLog.error(`Heartbeat execution failed for ${agentId}: ${errorDetail}`);
|
||||
|
||||
@@ -752,7 +752,7 @@ export class TaskExecutor {
|
||||
/** Spawned child agent IDs per parent task ID. Used for lifecycle tracking. */
|
||||
private spawnedAgents = new Map<string, Set<string>>();
|
||||
/** Per-task baseline of session stats used for delta persistence across repeated updates. */
|
||||
private tokenUsageBaselines = new Map<string, { inputTokens: number; outputTokens: number; cachedTokens: number; totalTokens: number }>();
|
||||
private tokenUsageBaselines = new Map<string, { inputTokens: number; outputTokens: number; cachedTokens: number; cacheWriteTokens: number; totalTokens: number }>();
|
||||
/** In-memory branch conflict error counters per task for tripwire protection. */
|
||||
private branchConflictErrorCount = new Map<string, number>();
|
||||
/** One-shot watchdogs for completed tasks that should have transitioned to in-review. */
|
||||
@@ -1883,7 +1883,7 @@ export class TaskExecutor {
|
||||
|
||||
private accumulateTokenUsage(
|
||||
existing: TaskTokenUsage | undefined,
|
||||
delta: Pick<TaskTokenUsage, "inputTokens" | "outputTokens" | "cachedTokens" | "totalTokens"> | undefined,
|
||||
delta: Pick<TaskTokenUsage, "inputTokens" | "outputTokens" | "cachedTokens" | "cacheWriteTokens" | "totalTokens"> | undefined,
|
||||
timestamp = new Date().toISOString(),
|
||||
): TaskTokenUsage | undefined {
|
||||
if (!delta) return existing;
|
||||
@@ -1892,6 +1892,7 @@ export class TaskExecutor {
|
||||
inputTokens: (existing?.inputTokens ?? 0) + delta.inputTokens,
|
||||
outputTokens: (existing?.outputTokens ?? 0) + delta.outputTokens,
|
||||
cachedTokens: (existing?.cachedTokens ?? 0) + delta.cachedTokens,
|
||||
cacheWriteTokens: (existing?.cacheWriteTokens ?? 0) + delta.cacheWriteTokens,
|
||||
totalTokens: (existing?.totalTokens ?? 0) + delta.totalTokens,
|
||||
firstUsedAt: existing?.firstUsedAt ?? timestamp,
|
||||
lastUsedAt: timestamp,
|
||||
@@ -1902,7 +1903,7 @@ export class TaskExecutor {
|
||||
|
||||
private async extractSessionTokenUsage(
|
||||
session: AgentSession | undefined,
|
||||
): Promise<Pick<TaskTokenUsage, "inputTokens" | "outputTokens" | "cachedTokens" | "totalTokens"> | undefined> {
|
||||
): Promise<Pick<TaskTokenUsage, "inputTokens" | "outputTokens" | "cachedTokens" | "cacheWriteTokens" | "totalTokens"> | undefined> {
|
||||
if (!session) return undefined;
|
||||
|
||||
try {
|
||||
@@ -1933,15 +1934,15 @@ export class TaskExecutor {
|
||||
|
||||
const inputTokens = tokens.input ?? 0;
|
||||
const outputTokens = tokens.output ?? 0;
|
||||
const cacheReadTokens = tokens.cacheRead ?? 0;
|
||||
const cachedTokens = tokens.cacheRead ?? 0;
|
||||
const cacheWriteTokens = tokens.cacheWrite ?? 0;
|
||||
const cachedTokens = cacheReadTokens + cacheWriteTokens;
|
||||
const totalTokens = tokens.total ?? (inputTokens + outputTokens + cachedTokens);
|
||||
const totalTokens = tokens.total ?? (inputTokens + outputTokens + cachedTokens + cacheWriteTokens);
|
||||
|
||||
return {
|
||||
inputTokens,
|
||||
outputTokens,
|
||||
cachedTokens,
|
||||
cacheWriteTokens,
|
||||
totalTokens,
|
||||
};
|
||||
} catch (err: unknown) {
|
||||
@@ -1964,6 +1965,7 @@ export class TaskExecutor {
|
||||
inputTokens: Math.max(0, currentUsage.inputTokens - baseline.inputTokens),
|
||||
outputTokens: Math.max(0, currentUsage.outputTokens - baseline.outputTokens),
|
||||
cachedTokens: Math.max(0, currentUsage.cachedTokens - baseline.cachedTokens),
|
||||
cacheWriteTokens: Math.max(0, currentUsage.cacheWriteTokens - baseline.cacheWriteTokens),
|
||||
totalTokens: Math.max(0, currentUsage.totalTokens - baseline.totalTokens),
|
||||
}
|
||||
: currentUsage;
|
||||
@@ -1972,6 +1974,7 @@ export class TaskExecutor {
|
||||
delta.inputTokens === 0
|
||||
&& delta.outputTokens === 0
|
||||
&& delta.cachedTokens === 0
|
||||
&& delta.cacheWriteTokens === 0
|
||||
&& delta.totalTokens === 0
|
||||
) {
|
||||
return;
|
||||
|
||||
@@ -8,6 +8,7 @@ interface SessionBaseline {
|
||||
input: number;
|
||||
output: number;
|
||||
cached: number;
|
||||
cacheWrite: number;
|
||||
}
|
||||
|
||||
// Per-session cumulative-token baselines so repeated calls only persist deltas.
|
||||
@@ -50,37 +51,40 @@ export async function accumulateSessionTokenUsage(
|
||||
const tokens = stats?.tokens;
|
||||
if (!tokens) return;
|
||||
|
||||
// Treat cache-write tokens as input (they're billed as input on first write
|
||||
// and read back at a discount on subsequent turns).
|
||||
const currentInput = (tokens.input ?? 0) + (tokens.cacheWrite ?? 0);
|
||||
const currentInput = tokens.input ?? 0;
|
||||
const currentOutput = tokens.output ?? 0;
|
||||
const currentCached = tokens.cacheRead ?? 0;
|
||||
const currentCacheWrite = tokens.cacheWrite ?? 0;
|
||||
|
||||
const baseline = sessionBaselines.get(session) ?? { input: 0, output: 0, cached: 0 };
|
||||
const baseline = sessionBaselines.get(session) ?? { input: 0, output: 0, cached: 0, cacheWrite: 0 };
|
||||
const inputDelta = Math.max(0, currentInput - baseline.input);
|
||||
const outputDelta = Math.max(0, currentOutput - baseline.output);
|
||||
const cachedDelta = Math.max(0, currentCached - baseline.cached);
|
||||
const cacheWriteDelta = Math.max(0, currentCacheWrite - baseline.cacheWrite);
|
||||
|
||||
sessionBaselines.set(session, {
|
||||
input: currentInput,
|
||||
output: currentOutput,
|
||||
cached: currentCached,
|
||||
cacheWrite: currentCacheWrite,
|
||||
});
|
||||
|
||||
if (inputDelta === 0 && outputDelta === 0 && cachedDelta === 0) return;
|
||||
if (inputDelta === 0 && outputDelta === 0 && cachedDelta === 0 && cacheWriteDelta === 0) return;
|
||||
|
||||
const task = await store.getTask(taskId);
|
||||
const now = new Date().toISOString();
|
||||
const newInput = (task.tokenUsage?.inputTokens ?? 0) + inputDelta;
|
||||
const newOutput = (task.tokenUsage?.outputTokens ?? 0) + outputDelta;
|
||||
const newCached = (task.tokenUsage?.cachedTokens ?? 0) + cachedDelta;
|
||||
const newCacheWrite = (task.tokenUsage?.cacheWriteTokens ?? 0) + cacheWriteDelta;
|
||||
|
||||
await store.updateTask(taskId, {
|
||||
tokenUsage: {
|
||||
inputTokens: newInput,
|
||||
outputTokens: newOutput,
|
||||
cachedTokens: newCached,
|
||||
totalTokens: newInput + newOutput + newCached,
|
||||
cacheWriteTokens: newCacheWrite,
|
||||
totalTokens: newInput + newOutput + newCached + newCacheWrite,
|
||||
firstUsedAt: task.tokenUsage?.firstUsedAt ?? now,
|
||||
lastUsedAt: now,
|
||||
},
|
||||
@@ -95,10 +99,8 @@ export async function accumulateSessionTokenUsage(
|
||||
* Compute the cache hit ratio: `cachedTokens / (inputTokens + cachedTokens)`.
|
||||
* Returns a number in [0, 1], or 0 when both arguments are 0.
|
||||
*
|
||||
* Compatible with stored `task.tokenUsage` fields: pass `inputTokens` (which
|
||||
* includes cache-write tokens per `accumulateSessionTokenUsage`) and
|
||||
* `cachedTokens` (cache-read tokens). Note this differs slightly from the
|
||||
* Anthropic console metric, which excludes cache-write from the denominator.
|
||||
* Compatible with canonical stored `task.tokenUsage` fields: pass raw
|
||||
* `inputTokens` and cache-read `cachedTokens`.
|
||||
*/
|
||||
export function computeCacheHitRatio(
|
||||
inputTokens: number,
|
||||
|
||||
@@ -68,6 +68,7 @@ export interface StepResult {
|
||||
inputTokens: number;
|
||||
outputTokens: number;
|
||||
cachedTokens: number;
|
||||
cacheWriteTokens: number;
|
||||
totalTokens: number;
|
||||
};
|
||||
}
|
||||
@@ -811,15 +812,15 @@ export class StepSessionExecutor {
|
||||
|
||||
const inputTokens = tokens.input ?? 0;
|
||||
const outputTokens = tokens.output ?? 0;
|
||||
const cacheReadTokens = tokens.cacheRead ?? 0;
|
||||
const cachedTokens = tokens.cacheRead ?? 0;
|
||||
const cacheWriteTokens = tokens.cacheWrite ?? 0;
|
||||
const cachedTokens = cacheReadTokens + cacheWriteTokens;
|
||||
const totalTokens = tokens.total ?? (inputTokens + outputTokens + cachedTokens);
|
||||
const totalTokens = tokens.total ?? (inputTokens + outputTokens + cachedTokens + cacheWriteTokens);
|
||||
|
||||
return {
|
||||
inputTokens,
|
||||
outputTokens,
|
||||
cachedTokens,
|
||||
cacheWriteTokens,
|
||||
totalTokens,
|
||||
};
|
||||
} catch (err: unknown) {
|
||||
|
||||
Reference in New Issue
Block a user