feat(FN-2677): persist step and single-session token usage stats
- Capture prompt/completion/total token usage from step-scoped sessions in the step session executor - Persist per-step token usage in executor run context so stats survive across task execution - Record single-session token usage totals alongside run context stats logging for consistent aggregation - Expand executor and step-session executor tests to validate token usage persistence and fixture behavior
This commit is contained in:
@@ -10939,161 +10939,276 @@ describe("StepSessionExecutor integration", () => {
|
||||
expect(onComplete).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("persists task tokenUsage on successful completion when agent has token totals", async () => {
|
||||
it("persists aggregated tokenUsage from step-session results", async () => {
|
||||
const { store } = createTokenUsageStepSessionStore();
|
||||
const totals = { input: 120, output: 45 };
|
||||
mockExecuteAll.mockImplementation(async () => {
|
||||
totals.input = 170;
|
||||
totals.output = 66;
|
||||
return [
|
||||
{ stepIndex: 0, success: true, retries: 0 },
|
||||
{ stepIndex: 1, success: true, retries: 0 },
|
||||
];
|
||||
});
|
||||
mockExecuteAll.mockResolvedValue([
|
||||
{
|
||||
stepIndex: 0,
|
||||
success: true,
|
||||
retries: 0,
|
||||
tokenUsage: { inputTokens: 22, outputTokens: 8, cachedTokens: 3, totalTokens: 33 },
|
||||
},
|
||||
{
|
||||
stepIndex: 1,
|
||||
success: true,
|
||||
retries: 0,
|
||||
tokenUsage: { inputTokens: 28, outputTokens: 13, cachedTokens: 1, totalTokens: 42 },
|
||||
},
|
||||
]);
|
||||
|
||||
const agentStore = {
|
||||
getAgent: vi.fn().mockImplementation(async () => ({
|
||||
id: "agent-001",
|
||||
totalInputTokens: totals.input,
|
||||
totalOutputTokens: totals.output,
|
||||
})),
|
||||
};
|
||||
|
||||
const executor = new TaskExecutor(store, "/tmp/test", { agentStore: agentStore as any });
|
||||
await executor.execute(createTaskWithSteps());
|
||||
|
||||
const tokenUsageUpdate = store.updateTask.mock.calls.find(([, updates]: [string, Record<string, unknown>]) => updates.tokenUsage);
|
||||
expect(tokenUsageUpdate).toBeDefined();
|
||||
const tokenUsage = tokenUsageUpdate![1].tokenUsage as Record<string, unknown>;
|
||||
expect(tokenUsage.inputTokens).toBe(50);
|
||||
expect(tokenUsage.outputTokens).toBe(21);
|
||||
expect(tokenUsage.cachedTokens).toBe(0);
|
||||
expect(tokenUsage.totalTokens).toBe(71);
|
||||
expect(typeof tokenUsage.firstUsedAt).toBe("string");
|
||||
expect(typeof tokenUsage.lastUsedAt).toBe("string");
|
||||
expect(Number.isNaN(new Date(tokenUsage.firstUsedAt as string).getTime())).toBe(false);
|
||||
expect(Number.isNaN(new Date(tokenUsage.lastUsedAt as string).getTime())).toBe(false);
|
||||
});
|
||||
|
||||
it("persists tokenUsage incrementally during step execution before in-review transition", async () => {
|
||||
const { store } = createTokenUsageStepSessionStore();
|
||||
const totals = { input: 100, output: 40 };
|
||||
|
||||
mockedStepSessionExecutor.mockImplementationOnce(((options: any) => ({
|
||||
executeAll: vi.fn(async () => {
|
||||
totals.input = 120;
|
||||
totals.output = 50;
|
||||
options.onStepComplete(0, { stepIndex: 0, success: true, retries: 0 });
|
||||
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||
|
||||
totals.input = 150;
|
||||
totals.output = 70;
|
||||
options.onStepComplete(1, { stepIndex: 1, success: true, retries: 0 });
|
||||
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||
|
||||
return [
|
||||
{ stepIndex: 0, success: true, retries: 0 },
|
||||
{ stepIndex: 1, success: true, retries: 0 },
|
||||
];
|
||||
}),
|
||||
terminateAllSessions: mockTerminateAllSessions,
|
||||
cleanup: mockCleanup,
|
||||
})) as any);
|
||||
|
||||
const agentStore = {
|
||||
getAgent: vi.fn().mockImplementation(async () => ({
|
||||
id: "agent-001",
|
||||
totalInputTokens: totals.input,
|
||||
totalOutputTokens: totals.output,
|
||||
})),
|
||||
};
|
||||
|
||||
const executor = new TaskExecutor(store, "/tmp/test", { agentStore: agentStore as any });
|
||||
const executor = new TaskExecutor(store, "/tmp/test", {});
|
||||
await executor.execute(createTaskWithSteps());
|
||||
|
||||
const tokenUsageUpdates = store.updateTask.mock.calls
|
||||
.filter(([, updates]: [string, Record<string, unknown>]) => updates.tokenUsage)
|
||||
.map(([, updates]: [string, Record<string, unknown>]) => updates.tokenUsage as Record<string, unknown>);
|
||||
|
||||
expect(tokenUsageUpdates.length).toBeGreaterThan(0);
|
||||
expect(tokenUsageUpdates[tokenUsageUpdates.length - 1]).toEqual(
|
||||
expect.objectContaining({
|
||||
inputTokens: 50,
|
||||
outputTokens: 21,
|
||||
cachedTokens: 4,
|
||||
totalTokens: 75,
|
||||
firstUsedAt: expect.any(String),
|
||||
lastUsedAt: expect.any(String),
|
||||
}),
|
||||
);
|
||||
});
|
||||
|
||||
it("persists tokenUsage incrementally during step execution before in-review transition", async () => {
|
||||
const { store } = createTokenUsageStepSessionStore();
|
||||
|
||||
mockedStepSessionExecutor.mockImplementationOnce(((options: any) => ({
|
||||
executeAll: vi.fn(async () => {
|
||||
options.onStepComplete(0, {
|
||||
stepIndex: 0,
|
||||
success: true,
|
||||
retries: 0,
|
||||
tokenUsage: { inputTokens: 20, outputTokens: 10, cachedTokens: 2, totalTokens: 32 },
|
||||
});
|
||||
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||
|
||||
options.onStepComplete(1, {
|
||||
stepIndex: 1,
|
||||
success: true,
|
||||
retries: 0,
|
||||
tokenUsage: { inputTokens: 30, outputTokens: 5, cachedTokens: 1, totalTokens: 36 },
|
||||
});
|
||||
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||
|
||||
return [
|
||||
{
|
||||
stepIndex: 0,
|
||||
success: true,
|
||||
retries: 0,
|
||||
tokenUsage: { inputTokens: 20, outputTokens: 10, cachedTokens: 2, totalTokens: 32 },
|
||||
},
|
||||
{
|
||||
stepIndex: 1,
|
||||
success: true,
|
||||
retries: 0,
|
||||
tokenUsage: { inputTokens: 30, outputTokens: 5, cachedTokens: 1, totalTokens: 36 },
|
||||
},
|
||||
];
|
||||
}),
|
||||
terminateAllSessions: mockTerminateAllSessions,
|
||||
cleanup: mockCleanup,
|
||||
})) as any);
|
||||
|
||||
const executor = new TaskExecutor(store, "/tmp/test", {});
|
||||
await executor.execute(createTaskWithSteps());
|
||||
|
||||
const tokenUsageUpdates = store.updateTask.mock.calls
|
||||
.filter(([, updates]: [string, Record<string, unknown>]) => updates.tokenUsage)
|
||||
.map(([, updates]: [string, Record<string, unknown>]) => updates.tokenUsage as Record<string, unknown>);
|
||||
|
||||
expect(tokenUsageUpdates.length).toBeGreaterThanOrEqual(2);
|
||||
expect(tokenUsageUpdates).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
inputTokens: 20,
|
||||
outputTokens: 10,
|
||||
totalTokens: 30,
|
||||
cachedTokens: 2,
|
||||
totalTokens: 32,
|
||||
}),
|
||||
expect.objectContaining({
|
||||
inputTokens: 50,
|
||||
outputTokens: 15,
|
||||
cachedTokens: 3,
|
||||
totalTokens: 68,
|
||||
}),
|
||||
]),
|
||||
);
|
||||
expect(store.moveTask).toHaveBeenCalledWith("FN-200", "in-review");
|
||||
});
|
||||
|
||||
it("persists task tokenUsage on failure paths so partial usage is visible", async () => {
|
||||
const { store } = createTokenUsageStepSessionStore();
|
||||
const totals = { input: 30, output: 10 };
|
||||
mockExecuteAll.mockImplementation(async () => {
|
||||
totals.input = 44;
|
||||
totals.output = 19;
|
||||
return [
|
||||
{ stepIndex: 0, success: true, retries: 0 },
|
||||
{ stepIndex: 1, success: false, error: "lint failed", retries: 1 },
|
||||
];
|
||||
});
|
||||
|
||||
const agentStore = {
|
||||
getAgent: vi.fn().mockImplementation(async () => ({
|
||||
id: "agent-001",
|
||||
totalInputTokens: totals.input,
|
||||
totalOutputTokens: totals.output,
|
||||
})),
|
||||
};
|
||||
|
||||
const executor = new TaskExecutor(store, "/tmp/test", { agentStore: agentStore as any });
|
||||
await executor.execute(createTaskWithSteps());
|
||||
|
||||
expect(store.updateTask).toHaveBeenCalledWith(
|
||||
"FN-200",
|
||||
expect.objectContaining({
|
||||
tokenUsage: expect.objectContaining({
|
||||
inputTokens: 14,
|
||||
outputTokens: 9,
|
||||
totalTokens: 23,
|
||||
}),
|
||||
}),
|
||||
);
|
||||
});
|
||||
|
||||
it("skips tokenUsage persistence when agentStore is not configured", async () => {
|
||||
it("persists task tokenUsage on step-session failure paths so partial usage is visible", async () => {
|
||||
const { store } = createTokenUsageStepSessionStore();
|
||||
mockExecuteAll.mockResolvedValue([
|
||||
{ stepIndex: 0, success: true, retries: 0 },
|
||||
{ stepIndex: 1, success: true, retries: 0 },
|
||||
{
|
||||
stepIndex: 0,
|
||||
success: true,
|
||||
retries: 0,
|
||||
tokenUsage: { inputTokens: 14, outputTokens: 6, cachedTokens: 2, totalTokens: 22 },
|
||||
},
|
||||
{
|
||||
stepIndex: 1,
|
||||
success: false,
|
||||
error: "lint failed",
|
||||
retries: 1,
|
||||
tokenUsage: { inputTokens: 7, outputTokens: 3, cachedTokens: 1, totalTokens: 11 },
|
||||
},
|
||||
]);
|
||||
|
||||
const executor = new TaskExecutor(store, "/tmp/test", {});
|
||||
await executor.execute(createTaskWithSteps());
|
||||
|
||||
const tokenUsageUpdate = store.updateTask.mock.calls.find(([, updates]: [string, Record<string, unknown>]) => updates.tokenUsage);
|
||||
expect(tokenUsageUpdate).toBeUndefined();
|
||||
const tokenUsageUpdates = store.updateTask.mock.calls
|
||||
.filter(([, updates]: [string, Record<string, unknown>]) => updates.tokenUsage)
|
||||
.map(([, updates]: [string, Record<string, unknown>]) => updates.tokenUsage as Record<string, unknown>);
|
||||
|
||||
expect(tokenUsageUpdates[tokenUsageUpdates.length - 1]).toEqual(
|
||||
expect.objectContaining({
|
||||
inputTokens: 21,
|
||||
outputTokens: 9,
|
||||
cachedTokens: 3,
|
||||
totalTokens: 33,
|
||||
}),
|
||||
);
|
||||
});
|
||||
|
||||
it("skips tokenUsage persistence when task has no assignedAgentId", async () => {
|
||||
const { store } = createTokenUsageStepSessionStore({ assignedAgentId: undefined });
|
||||
mockExecuteAll.mockResolvedValue([
|
||||
{ stepIndex: 0, success: true, retries: 0 },
|
||||
{ stepIndex: 1, success: true, retries: 0 },
|
||||
]);
|
||||
it("persists tokenUsage from session.getSessionStats in single-session mode", async () => {
|
||||
const store = createMockStore();
|
||||
const taskState = createTaskWithSteps({
|
||||
description: "# test\n## Steps\n### Step 0: Preflight\n- [ ] check",
|
||||
steps: [{ name: "Step 0", status: "done" }],
|
||||
currentStep: 0,
|
||||
});
|
||||
|
||||
const agentStore = {
|
||||
getAgent: vi.fn().mockResolvedValue({ id: "agent-001", totalInputTokens: 20, totalOutputTokens: 20 }),
|
||||
store.getSettings.mockResolvedValue({
|
||||
maxConcurrent: 2,
|
||||
maxWorktrees: 4,
|
||||
pollIntervalMs: 15000,
|
||||
groupOverlappingFiles: false,
|
||||
autoMerge: false,
|
||||
runStepsInNewSessions: false,
|
||||
});
|
||||
store.getTask.mockImplementation(async () => ({ ...taskState }));
|
||||
store.updateTask.mockImplementation(async (_taskId: string, updates: Record<string, unknown>) => {
|
||||
if (updates.tokenUsage !== undefined) {
|
||||
(taskState as Task).tokenUsage = updates.tokenUsage as Task["tokenUsage"];
|
||||
}
|
||||
if (updates.status !== undefined) {
|
||||
(taskState as Task).status = updates.status as Task["status"];
|
||||
}
|
||||
return {};
|
||||
});
|
||||
|
||||
const session = {
|
||||
prompt: vi.fn().mockResolvedValue(undefined),
|
||||
dispose: vi.fn(),
|
||||
subscribe: vi.fn(),
|
||||
on: vi.fn(),
|
||||
abortBash: vi.fn(),
|
||||
state: {},
|
||||
sessionManager: { getLeafId: vi.fn().mockReturnValue("leaf-1") },
|
||||
getSessionStats: vi.fn().mockReturnValue({
|
||||
tokens: {
|
||||
input: 31,
|
||||
output: 17,
|
||||
cacheRead: 5,
|
||||
cacheWrite: 2,
|
||||
total: 55,
|
||||
},
|
||||
}),
|
||||
};
|
||||
|
||||
const executor = new TaskExecutor(store, "/tmp/test", { agentStore: agentStore as any });
|
||||
await executor.execute(createTaskWithSteps({ assignedAgentId: undefined }));
|
||||
mockedCreateFnAgent.mockResolvedValue({ session } as any);
|
||||
|
||||
const tokenUsageUpdate = store.updateTask.mock.calls.find(([, updates]: [string, Record<string, unknown>]) => updates.tokenUsage);
|
||||
expect(tokenUsageUpdate).toBeUndefined();
|
||||
expect(agentStore.getAgent).not.toHaveBeenCalled();
|
||||
const executor = new TaskExecutor(store, "/tmp/test", {});
|
||||
await executor.execute(taskState);
|
||||
|
||||
expect(store.updateTask).toHaveBeenCalledWith(
|
||||
"FN-200",
|
||||
expect.objectContaining({
|
||||
tokenUsage: expect.objectContaining({
|
||||
inputTokens: 31,
|
||||
outputTokens: 17,
|
||||
cachedTokens: 7,
|
||||
totalTokens: 55,
|
||||
}),
|
||||
}),
|
||||
);
|
||||
expect(session.getSessionStats).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("persists single-session tokenUsage on failure so partial usage is visible", async () => {
|
||||
const store = createMockStore();
|
||||
const taskState = createTaskWithSteps({
|
||||
description: "# test\n## Steps\n### Step 0: Preflight\n- [ ] check",
|
||||
steps: [{ name: "Step 0", status: "pending" }],
|
||||
currentStep: 0,
|
||||
});
|
||||
|
||||
store.getSettings.mockResolvedValue({
|
||||
maxConcurrent: 2,
|
||||
maxWorktrees: 4,
|
||||
pollIntervalMs: 15000,
|
||||
groupOverlappingFiles: false,
|
||||
autoMerge: false,
|
||||
runStepsInNewSessions: false,
|
||||
});
|
||||
store.getTask.mockImplementation(async () => ({ ...taskState }));
|
||||
store.updateTask.mockImplementation(async (_taskId: string, updates: Record<string, unknown>) => {
|
||||
if (updates.tokenUsage !== undefined) {
|
||||
(taskState as Task).tokenUsage = updates.tokenUsage as Task["tokenUsage"];
|
||||
}
|
||||
if (updates.status !== undefined) {
|
||||
(taskState as Task).status = updates.status as Task["status"];
|
||||
}
|
||||
if (updates.error !== undefined) {
|
||||
(taskState as Task).error = updates.error as Task["error"];
|
||||
}
|
||||
return {};
|
||||
});
|
||||
|
||||
const session = {
|
||||
prompt: vi.fn().mockRejectedValue(new Error("session failed")),
|
||||
dispose: vi.fn(),
|
||||
subscribe: vi.fn(),
|
||||
on: vi.fn(),
|
||||
abortBash: vi.fn(),
|
||||
state: {},
|
||||
sessionManager: { getLeafId: vi.fn().mockReturnValue("leaf-1") },
|
||||
getSessionStats: vi.fn().mockReturnValue({
|
||||
tokens: {
|
||||
input: 12,
|
||||
output: 4,
|
||||
cacheRead: 1,
|
||||
cacheWrite: 0,
|
||||
total: 17,
|
||||
},
|
||||
}),
|
||||
};
|
||||
|
||||
mockedCreateFnAgent.mockResolvedValue({ session } as any);
|
||||
|
||||
const executor = new TaskExecutor(store, "/tmp/test", {});
|
||||
await executor.execute(taskState);
|
||||
|
||||
const tokenUsageUpdates = store.updateTask.mock.calls
|
||||
.filter(([, updates]: [string, Record<string, unknown>]) => updates.tokenUsage)
|
||||
.map(([, updates]: [string, Record<string, unknown>]) => updates.tokenUsage as Record<string, unknown>);
|
||||
|
||||
expect(tokenUsageUpdates[tokenUsageUpdates.length - 1]).toEqual(
|
||||
expect.objectContaining({
|
||||
inputTokens: 12,
|
||||
outputTokens: 4,
|
||||
cachedTokens: 1,
|
||||
totalTokens: 17,
|
||||
}),
|
||||
);
|
||||
expect(session.getSessionStats).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("moves task to in-review when step-session execution fails", async () => {
|
||||
|
||||
@@ -757,6 +757,128 @@ describe("StepSessionExecutor", () => {
|
||||
expect(onStepStart).toHaveBeenNthCalledWith(3, 2);
|
||||
});
|
||||
|
||||
it("includes token usage from session stats on successful step completion", async () => {
|
||||
const prompt = makeStepPrompt("FN-001", 1);
|
||||
const task = makeTaskDetail({
|
||||
prompt,
|
||||
steps: [{ name: "Step 0", status: "pending" }],
|
||||
});
|
||||
const settings = makeSettings({ maxParallelSteps: 1 });
|
||||
|
||||
const session = {
|
||||
...makeMockSession(),
|
||||
getSessionStats: vi.fn().mockReturnValue({
|
||||
tokens: {
|
||||
input: 25,
|
||||
output: 11,
|
||||
cacheRead: 4,
|
||||
cacheWrite: 2,
|
||||
total: 42,
|
||||
},
|
||||
}),
|
||||
};
|
||||
mockedCreateFnAgent.mockResolvedValue({ session } as any);
|
||||
|
||||
const executor = new StepSessionExecutor({
|
||||
taskDetail: task,
|
||||
worktreePath: "/project/.worktrees/main",
|
||||
rootDir: "/project",
|
||||
settings,
|
||||
});
|
||||
|
||||
const results = await executor.executeAll();
|
||||
|
||||
expect(results).toHaveLength(1);
|
||||
expect(results[0]).toMatchObject({
|
||||
stepIndex: 0,
|
||||
success: true,
|
||||
retries: 0,
|
||||
tokenUsage: {
|
||||
inputTokens: 25,
|
||||
outputTokens: 11,
|
||||
cachedTokens: 6,
|
||||
totalTokens: 42,
|
||||
},
|
||||
});
|
||||
expect(session.getSessionStats).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("includes token usage from session stats on failed step completion", async () => {
|
||||
const prompt = makeStepPrompt("FN-001", 1);
|
||||
const task = makeTaskDetail({
|
||||
prompt,
|
||||
steps: [{ name: "Step 0", status: "pending" }],
|
||||
});
|
||||
const settings = makeSettings({ maxParallelSteps: 1 });
|
||||
|
||||
const session = {
|
||||
...makeMockSession(() => Promise.reject(new Error("step failed"))),
|
||||
getSessionStats: vi.fn().mockReturnValue({
|
||||
tokens: {
|
||||
input: 19,
|
||||
output: 7,
|
||||
cacheRead: 3,
|
||||
cacheWrite: 0,
|
||||
total: 29,
|
||||
},
|
||||
}),
|
||||
};
|
||||
mockedCreateFnAgent.mockResolvedValue({ session } as any);
|
||||
|
||||
const executor = new StepSessionExecutor({
|
||||
taskDetail: task,
|
||||
worktreePath: "/project/.worktrees/main",
|
||||
rootDir: "/project",
|
||||
settings,
|
||||
});
|
||||
|
||||
const executePromise = executor.executeAll();
|
||||
await vi.runAllTimersAsync();
|
||||
const results = await executePromise;
|
||||
|
||||
expect(results).toHaveLength(1);
|
||||
expect(results[0]).toMatchObject({
|
||||
stepIndex: 0,
|
||||
success: false,
|
||||
error: "step failed",
|
||||
tokenUsage: {
|
||||
inputTokens: 19,
|
||||
outputTokens: 7,
|
||||
cachedTokens: 3,
|
||||
totalTokens: 29,
|
||||
},
|
||||
});
|
||||
expect(session.getSessionStats).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("keeps token usage undefined when session stats are unavailable", async () => {
|
||||
const prompt = makeStepPrompt("FN-001", 1);
|
||||
const task = makeTaskDetail({
|
||||
prompt,
|
||||
steps: [{ name: "Step 0", status: "pending" }],
|
||||
});
|
||||
const settings = makeSettings({ maxParallelSteps: 1 });
|
||||
|
||||
const session = {
|
||||
...makeMockSession(),
|
||||
getSessionStats: vi.fn().mockReturnValue(undefined),
|
||||
};
|
||||
mockedCreateFnAgent.mockResolvedValue({ session } as any);
|
||||
|
||||
const executor = new StepSessionExecutor({
|
||||
taskDetail: task,
|
||||
worktreePath: "/project/.worktrees/main",
|
||||
rootDir: "/project",
|
||||
settings,
|
||||
});
|
||||
|
||||
const results = await executor.executeAll();
|
||||
|
||||
expect(results).toHaveLength(1);
|
||||
expect(results[0].tokenUsage).toBeUndefined();
|
||||
expect(session.getSessionStats).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("step failure with retry: fails first 2 attempts, succeeds on 3rd", async () => {
|
||||
const prompt = makeStepPrompt("FN-001", 2);
|
||||
const task = makeTaskDetail({ prompt, steps: [
|
||||
|
||||
@@ -5,7 +5,7 @@ const execAsync = promisify(exec);
|
||||
import { isAbsolute, join, relative, resolve as resolvePath } from "node:path";
|
||||
import { existsSync } from "node:fs";
|
||||
import { readFile, writeFile } from "node:fs/promises";
|
||||
import type { TaskStore, Task, TaskDetail, StepStatus, Settings, WorkflowStep, MissionStore, Slice, AgentState, AgentCapability, RunMutationContext } from "@fusion/core";
|
||||
import type { TaskStore, Task, TaskDetail, TaskTokenUsage, StepStatus, Settings, WorkflowStep, MissionStore, Slice, AgentState, AgentCapability, RunMutationContext } from "@fusion/core";
|
||||
import { buildExecutionMemoryInstructions, getTaskMergeBlocker, resolveAgentPrompt, type RunCommandResult } from "@fusion/core";
|
||||
import { findWorktreeUser } from "./merger.js";
|
||||
import { generateWorktreeName, slugify } from "./worktree-names.js";
|
||||
@@ -489,8 +489,8 @@ export class TaskExecutor {
|
||||
private loopRecoveryState = new Map<string, { attempts: number; pending: boolean }>();
|
||||
/** Spawned child agent IDs per parent task ID. Used for lifecycle tracking. */
|
||||
private spawnedAgents = new Map<string, Set<string>>();
|
||||
/** Per-task baseline of agent cumulative token counters used for delta persistence. */
|
||||
private tokenUsageBaselines = new Map<string, { agentId: string; inputTokens: number; outputTokens: number }>();
|
||||
/** 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 async finalizeAlreadyReviewedTask(taskId: string): Promise<"merged" | "blocked" | "missing"> {
|
||||
const latestTask = await this.store.getTask(taskId);
|
||||
@@ -877,78 +877,107 @@ export class TaskExecutor {
|
||||
return getTaskCompletionBlockerForStore(this.store, task);
|
||||
}
|
||||
|
||||
private async initializeTokenUsageBaseline(taskId: string, task?: Task): Promise<void> {
|
||||
if (!this.options.agentStore) return;
|
||||
private accumulateTokenUsage(
|
||||
existing: TaskTokenUsage | undefined,
|
||||
delta: Pick<TaskTokenUsage, "inputTokens" | "outputTokens" | "cachedTokens" | "totalTokens"> | undefined,
|
||||
timestamp = new Date().toISOString(),
|
||||
): TaskTokenUsage | undefined {
|
||||
if (!delta) return existing;
|
||||
|
||||
const currentTask = task ?? await this.store.getTask(taskId);
|
||||
const assignedAgentId = currentTask.assignedAgentId?.trim();
|
||||
if (!assignedAgentId) return;
|
||||
const merged: TaskTokenUsage = {
|
||||
inputTokens: (existing?.inputTokens ?? 0) + delta.inputTokens,
|
||||
outputTokens: (existing?.outputTokens ?? 0) + delta.outputTokens,
|
||||
cachedTokens: (existing?.cachedTokens ?? 0) + delta.cachedTokens,
|
||||
totalTokens: (existing?.totalTokens ?? 0) + delta.totalTokens,
|
||||
firstUsedAt: existing?.firstUsedAt ?? timestamp,
|
||||
lastUsedAt: timestamp,
|
||||
};
|
||||
|
||||
const existing = this.tokenUsageBaselines.get(taskId);
|
||||
if (existing?.agentId === assignedAgentId) return;
|
||||
|
||||
const agent = await this.options.agentStore.getAgent(assignedAgentId);
|
||||
if (!agent) return;
|
||||
|
||||
this.tokenUsageBaselines.set(taskId, {
|
||||
agentId: assignedAgentId,
|
||||
inputTokens: agent.totalInputTokens ?? 0,
|
||||
outputTokens: agent.totalOutputTokens ?? 0,
|
||||
});
|
||||
return merged;
|
||||
}
|
||||
|
||||
private async persistTokenUsage(taskId: string): Promise<void> {
|
||||
if (!this.options.agentStore) return;
|
||||
private async extractSessionTokenUsage(
|
||||
session: AgentSession | undefined,
|
||||
): Promise<Pick<TaskTokenUsage, "inputTokens" | "outputTokens" | "cachedTokens" | "totalTokens"> | undefined> {
|
||||
if (!session) return undefined;
|
||||
|
||||
const task = await this.store.getTask(taskId);
|
||||
const assignedAgentId = task.assignedAgentId?.trim();
|
||||
if (!assignedAgentId) return;
|
||||
try {
|
||||
const statsResult = (session as AgentSession & {
|
||||
getSessionStats?: () =>
|
||||
| {
|
||||
tokens?: {
|
||||
input?: number;
|
||||
output?: number;
|
||||
cacheRead?: number;
|
||||
cacheWrite?: number;
|
||||
total?: number;
|
||||
};
|
||||
}
|
||||
| Promise<{
|
||||
tokens?: {
|
||||
input?: number;
|
||||
output?: number;
|
||||
cacheRead?: number;
|
||||
cacheWrite?: number;
|
||||
total?: number;
|
||||
};
|
||||
}>;
|
||||
}).getSessionStats?.();
|
||||
const stats = await Promise.resolve(statsResult);
|
||||
const tokens = stats?.tokens;
|
||||
if (!tokens) return undefined;
|
||||
|
||||
const agent = await this.options.agentStore.getAgent(assignedAgentId);
|
||||
if (!agent) return;
|
||||
const inputTokens = tokens.input ?? 0;
|
||||
const outputTokens = tokens.output ?? 0;
|
||||
const cacheReadTokens = tokens.cacheRead ?? 0;
|
||||
const cacheWriteTokens = tokens.cacheWrite ?? 0;
|
||||
const cachedTokens = cacheReadTokens + cacheWriteTokens;
|
||||
const totalTokens = tokens.total ?? (inputTokens + outputTokens + cachedTokens);
|
||||
|
||||
const currentInputTokens = agent.totalInputTokens ?? 0;
|
||||
const currentOutputTokens = agent.totalOutputTokens ?? 0;
|
||||
const baseline = this.tokenUsageBaselines.get(taskId);
|
||||
|
||||
if (!baseline || baseline.agentId !== assignedAgentId) {
|
||||
this.tokenUsageBaselines.set(taskId, {
|
||||
agentId: assignedAgentId,
|
||||
inputTokens: currentInputTokens,
|
||||
outputTokens: currentOutputTokens,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
const inputDelta = Math.max(0, currentInputTokens - baseline.inputTokens);
|
||||
const outputDelta = Math.max(0, currentOutputTokens - baseline.outputTokens);
|
||||
|
||||
this.tokenUsageBaselines.set(taskId, {
|
||||
agentId: assignedAgentId,
|
||||
inputTokens: currentInputTokens,
|
||||
outputTokens: currentOutputTokens,
|
||||
});
|
||||
|
||||
if (inputDelta === 0 && outputDelta === 0 && !task.tokenUsage) {
|
||||
return;
|
||||
}
|
||||
|
||||
const now = new Date().toISOString();
|
||||
const mergedInputTokens = (task.tokenUsage?.inputTokens ?? 0) + inputDelta;
|
||||
const mergedOutputTokens = (task.tokenUsage?.outputTokens ?? 0) + outputDelta;
|
||||
const cachedTokens = task.tokenUsage?.cachedTokens ?? 0;
|
||||
const totalTokens = mergedInputTokens + mergedOutputTokens + cachedTokens;
|
||||
|
||||
await this.store.updateTask(taskId, {
|
||||
tokenUsage: {
|
||||
inputTokens: mergedInputTokens,
|
||||
outputTokens: mergedOutputTokens,
|
||||
return {
|
||||
inputTokens,
|
||||
outputTokens,
|
||||
cachedTokens,
|
||||
totalTokens,
|
||||
firstUsedAt: task.tokenUsage?.firstUsedAt ?? now,
|
||||
lastUsedAt: now,
|
||||
},
|
||||
});
|
||||
};
|
||||
} catch (err: unknown) {
|
||||
const message = err instanceof Error ? err.message : String(err);
|
||||
executorLog.warn(`Failed to read session stats for token usage: ${message}`);
|
||||
return undefined;
|
||||
}
|
||||
}
|
||||
|
||||
private async persistTokenUsage(taskId: string, session?: AgentSession): Promise<void> {
|
||||
const activeSession = session ?? this.activeSessions.get(taskId)?.session;
|
||||
const currentUsage = await this.extractSessionTokenUsage(activeSession);
|
||||
if (!currentUsage) return;
|
||||
|
||||
const baseline = this.tokenUsageBaselines.get(taskId);
|
||||
this.tokenUsageBaselines.set(taskId, currentUsage);
|
||||
|
||||
const delta = baseline
|
||||
? {
|
||||
inputTokens: Math.max(0, currentUsage.inputTokens - baseline.inputTokens),
|
||||
outputTokens: Math.max(0, currentUsage.outputTokens - baseline.outputTokens),
|
||||
cachedTokens: Math.max(0, currentUsage.cachedTokens - baseline.cachedTokens),
|
||||
totalTokens: Math.max(0, currentUsage.totalTokens - baseline.totalTokens),
|
||||
}
|
||||
: currentUsage;
|
||||
|
||||
if (
|
||||
delta.inputTokens === 0
|
||||
&& delta.outputTokens === 0
|
||||
&& delta.cachedTokens === 0
|
||||
&& delta.totalTokens === 0
|
||||
) {
|
||||
return;
|
||||
}
|
||||
|
||||
const task = await this.store.getTask(taskId);
|
||||
const merged = this.accumulateTokenUsage(task.tokenUsage, delta);
|
||||
if (!merged) return;
|
||||
|
||||
await this.store.updateTask(taskId, { tokenUsage: merged });
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1543,7 +1572,6 @@ export class TaskExecutor {
|
||||
|
||||
const detail = await this.store.getTask(task.id);
|
||||
executorLog.log(`${task.id}: fetched task detail (${detail.steps.length} steps, prompt length=${detail.prompt?.length ?? 0})`);
|
||||
await this.initializeTokenUsageBaseline(task.id, detail);
|
||||
|
||||
// Initialize steps from PROMPT.md if empty
|
||||
if (detail.steps.length === 0) {
|
||||
@@ -1575,6 +1603,9 @@ export class TaskExecutor {
|
||||
: null;
|
||||
const stepSessionRuntimeHint = extractRuntimeHint(stepSessionAgent?.runtimeConfig);
|
||||
|
||||
let accumulatedStepTokenUsage = detail.tokenUsage;
|
||||
const tokenUsageRecordedSteps = new Set<number>();
|
||||
|
||||
const stepExecutor = new StepSessionExecutor({
|
||||
store: this.store,
|
||||
taskDetail: detail,
|
||||
@@ -1609,7 +1640,18 @@ export class TaskExecutor {
|
||||
} catch (err) {
|
||||
executorLog.warn(`${task.id}: failed to update step ${stepIndex} status: ${err}`);
|
||||
}
|
||||
this.persistTokenUsage(task.id).catch((err) => {
|
||||
|
||||
if (!result.tokenUsage) {
|
||||
return;
|
||||
}
|
||||
|
||||
accumulatedStepTokenUsage = this.accumulateTokenUsage(accumulatedStepTokenUsage, result.tokenUsage);
|
||||
tokenUsageRecordedSteps.add(stepIndex);
|
||||
if (!accumulatedStepTokenUsage) {
|
||||
return;
|
||||
}
|
||||
|
||||
this.store.updateTask(task.id, { tokenUsage: accumulatedStepTokenUsage }).catch((err) => {
|
||||
executorLog.warn(`${task.id}: failed to persist token usage on step ${stepIndex} complete: ${err}`);
|
||||
});
|
||||
},
|
||||
@@ -1637,6 +1679,17 @@ export class TaskExecutor {
|
||||
return;
|
||||
}
|
||||
|
||||
for (const result of results) {
|
||||
if (!result.tokenUsage || tokenUsageRecordedSteps.has(result.stepIndex)) {
|
||||
continue;
|
||||
}
|
||||
accumulatedStepTokenUsage = this.accumulateTokenUsage(accumulatedStepTokenUsage, result.tokenUsage);
|
||||
}
|
||||
|
||||
if (accumulatedStepTokenUsage) {
|
||||
await this.store.updateTask(task.id, { tokenUsage: accumulatedStepTokenUsage });
|
||||
}
|
||||
|
||||
const allSuccess = results.every(r => r.success);
|
||||
if (allSuccess) {
|
||||
const updatedTask = await this.store.getTask(task.id);
|
||||
@@ -1674,7 +1727,6 @@ export class TaskExecutor {
|
||||
// Reset retry counters on success
|
||||
await this.store.updateTask(task.id, { workflowStepRetries: undefined, taskDoneRetryCount: null });
|
||||
|
||||
await this.persistTokenUsage(task.id);
|
||||
await this.store.moveTask(task.id, "in-review");
|
||||
// Audit trail: record task move (FN-1404)
|
||||
await audit.database({ type: "task:move", target: task.id, metadata: { to: "in-review" } });
|
||||
@@ -1684,7 +1736,6 @@ export class TaskExecutor {
|
||||
const failedSteps = results.filter(r => !r.success);
|
||||
const errorSummary = failedSteps.map(r => `Step ${r.stepIndex}: ${r.error || "unknown error"}`).join("; ");
|
||||
await this.store.updateTask(task.id, { status: "failed", error: errorSummary });
|
||||
await this.persistTokenUsage(task.id);
|
||||
await this.store.moveTask(task.id, "in-review");
|
||||
executorLog.log(`✗ ${task.id} step-session failed → in-review: ${errorSummary}`);
|
||||
this.options.onError?.(task, new Error(errorSummary));
|
||||
@@ -1766,7 +1817,9 @@ export class TaskExecutor {
|
||||
recoveryRetryCount: null,
|
||||
nextRecoveryAt: null,
|
||||
});
|
||||
await this.persistTokenUsage(task.id);
|
||||
if (accumulatedStepTokenUsage) {
|
||||
await this.store.updateTask(task.id, { tokenUsage: accumulatedStepTokenUsage });
|
||||
}
|
||||
await this.store.moveTask(task.id, "in-review");
|
||||
executorLog.log(`✗ ${task.id} transient retries exhausted → in-review`);
|
||||
this.options.onError?.(task, err instanceof Error ? err : new Error(errorMessage));
|
||||
@@ -1774,7 +1827,9 @@ export class TaskExecutor {
|
||||
executorLog.error(`✗ ${task.id} step-session execution failed:`, errorDetail);
|
||||
await this.store.logEntry(task.id, `Step-session execution failed: ${errorMessage}`, errorStack ?? errorDetail, this.currentRunContext);
|
||||
await this.store.updateTask(task.id, { status: "failed", error: errorMessage });
|
||||
await this.persistTokenUsage(task.id);
|
||||
if (accumulatedStepTokenUsage) {
|
||||
await this.store.updateTask(task.id, { tokenUsage: accumulatedStepTokenUsage });
|
||||
}
|
||||
await this.store.moveTask(task.id, "in-review");
|
||||
executorLog.log(`✗ ${task.id} step-session execution failed → in-review`);
|
||||
this.options.onError?.(task, err instanceof Error ? err : new Error(errorMessage));
|
||||
@@ -2183,6 +2238,7 @@ export class TaskExecutor {
|
||||
|
||||
// Dispose old session and create a fresh one
|
||||
this.activeSessions.delete(task.id);
|
||||
this.tokenUsageBaselines.delete(task.id);
|
||||
session.dispose();
|
||||
|
||||
const { session: retrySession, sessionFile: retrySessionFile } = await createResolvedAgentSession({
|
||||
@@ -2312,6 +2368,10 @@ export class TaskExecutor {
|
||||
this.activeSessions.delete(task.id);
|
||||
stuckDetector?.untrackTask(task.id);
|
||||
await agentLogger.flush();
|
||||
await this.persistTokenUsage(task.id, session).catch((err: unknown) => {
|
||||
const msg = err instanceof Error ? err.message : String(err);
|
||||
executorLog.warn(`${task.id}: failed to persist final single-session token usage before dispose: ${msg}`);
|
||||
});
|
||||
session.dispose();
|
||||
// Terminate all spawned child agents when parent session ends
|
||||
await this.terminateAllChildren(task.id);
|
||||
|
||||
@@ -58,6 +58,13 @@ export interface StepResult {
|
||||
error?: string;
|
||||
/** Number of retry attempts made (0 = first attempt succeeded). */
|
||||
retries: number;
|
||||
/** Optional per-step token usage extracted from session stats. */
|
||||
tokenUsage?: {
|
||||
inputTokens: number;
|
||||
outputTokens: number;
|
||||
cachedTokens: number;
|
||||
totalTokens: number;
|
||||
};
|
||||
}
|
||||
|
||||
/** A group of step indices that can run in parallel. */
|
||||
@@ -793,6 +800,56 @@ export class StepSessionExecutor {
|
||||
stepExecLog.log(`Cleanup complete for task ${this.options.taskDetail.id}`);
|
||||
}
|
||||
|
||||
private async extractTokenUsageFromSession(session: AgentSession | null | undefined): Promise<StepResult["tokenUsage"] | undefined> {
|
||||
if (!session) return undefined;
|
||||
|
||||
try {
|
||||
const statsResult = (session as AgentSession & {
|
||||
getSessionStats?: () =>
|
||||
| {
|
||||
tokens?: {
|
||||
input?: number;
|
||||
output?: number;
|
||||
cacheRead?: number;
|
||||
cacheWrite?: number;
|
||||
total?: number;
|
||||
};
|
||||
}
|
||||
| Promise<{
|
||||
tokens?: {
|
||||
input?: number;
|
||||
output?: number;
|
||||
cacheRead?: number;
|
||||
cacheWrite?: number;
|
||||
total?: number;
|
||||
};
|
||||
}>;
|
||||
}).getSessionStats?.();
|
||||
|
||||
const stats = await Promise.resolve(statsResult);
|
||||
const tokens = stats?.tokens;
|
||||
if (!tokens) return undefined;
|
||||
|
||||
const inputTokens = tokens.input ?? 0;
|
||||
const outputTokens = tokens.output ?? 0;
|
||||
const cacheReadTokens = tokens.cacheRead ?? 0;
|
||||
const cacheWriteTokens = tokens.cacheWrite ?? 0;
|
||||
const cachedTokens = cacheReadTokens + cacheWriteTokens;
|
||||
const totalTokens = tokens.total ?? (inputTokens + outputTokens + cachedTokens);
|
||||
|
||||
return {
|
||||
inputTokens,
|
||||
outputTokens,
|
||||
cachedTokens,
|
||||
totalTokens,
|
||||
};
|
||||
} catch (err: unknown) {
|
||||
const message = err instanceof Error ? err.message : String(err);
|
||||
stepExecLog.warn(`Failed to read session stats for step token usage: ${message}`);
|
||||
return undefined;
|
||||
}
|
||||
}
|
||||
|
||||
// ── Internal: Step Execution ────────────────────────────────────────
|
||||
|
||||
/**
|
||||
@@ -973,7 +1030,12 @@ export class StepSessionExecutor {
|
||||
// the error is stored on session.state.error instead of being thrown.
|
||||
checkSessionError(session);
|
||||
|
||||
const result: StepResult = { stepIndex, success: true, retries };
|
||||
const result: StepResult = {
|
||||
stepIndex,
|
||||
success: true,
|
||||
retries,
|
||||
tokenUsage: await this.extractTokenUsageFromSession(session),
|
||||
};
|
||||
this.options.onStepComplete?.(stepIndex, result);
|
||||
return result;
|
||||
} catch (err: unknown) {
|
||||
@@ -1008,7 +1070,12 @@ export class StepSessionExecutor {
|
||||
`[step-exec] Reduced-prompt recovery succeeded for step ${stepIndex}`,
|
||||
"text",
|
||||
);
|
||||
const result: StepResult = { stepIndex, success: true, retries };
|
||||
const result: StepResult = {
|
||||
stepIndex,
|
||||
success: true,
|
||||
retries,
|
||||
tokenUsage: await this.extractTokenUsageFromSession(session),
|
||||
};
|
||||
this.options.onStepComplete?.(stepIndex, result);
|
||||
return result;
|
||||
} catch (reducedErr: unknown) {
|
||||
@@ -1034,6 +1101,7 @@ export class StepSessionExecutor {
|
||||
success: false,
|
||||
error: errorMessage,
|
||||
retries,
|
||||
tokenUsage: await this.extractTokenUsageFromSession(session),
|
||||
};
|
||||
this.options.onStepComplete?.(stepIndex, result);
|
||||
return result;
|
||||
|
||||
Reference in New Issue
Block a user