feat(FN-4398): complete Step 3 — wire retry-burned logging and caps

Fusion-Task-Id: FN-4398
Fusion-Task-Lineage: 8b898b2e-3468-4fa7-8546-b8f6ba52cf46
This commit is contained in:
Fusion
2026-05-14 14:55:31 -07:00
committed by gsxdsm
parent 831e4969e9
commit 9230c80fb0
5 changed files with 280 additions and 4 deletions

View File

@@ -0,0 +1,96 @@
import { describe, expect, it, vi, beforeEach, afterEach } from "vitest";
import { RetryStormError, type TaskDetail } from "@fusion/core";
import { recordRetry } from "../retry-burned-logger.js";
function makeTask(overrides: Partial<TaskDetail> = {}): TaskDetail {
return {
id: "FN-1",
lineageId: "lineage-1",
description: "desc",
column: "in-progress",
dependencies: [],
steps: [],
currentStep: 0,
log: [],
prompt: "prompt",
createdAt: "2026-01-01T00:00:00.000Z",
updatedAt: "2026-01-01T00:00:00.000Z",
...overrides,
};
}
describe("recordRetry", () => {
const baseSettings = {
maxBranchConflictRecoveries: 5,
maxReviewerContextRetries: 2,
maxReviewerFallbackRetries: 2,
maxTotalRetriesBeforeFail: 25,
};
let task: TaskDetail;
let store: { updateTask: ReturnType<typeof vi.fn>; getTask: ReturnType<typeof vi.fn> };
const consoleSpy = vi.spyOn(console, "error").mockImplementation(() => undefined);
beforeEach(() => {
task = makeTask();
store = {
updateTask: vi.fn(async (_id: string, patch: Record<string, number>) => {
Object.assign(task, patch);
}),
getTask: vi.fn(async () => task),
};
});
afterEach(() => {
consoleSpy.mockClear();
});
it("increments new-category counters and logs payload", async () => {
await recordRetry({
store: store as never,
settings: baseSettings,
task,
category: "reviewerContext",
role: "reviewer",
attempt: 1,
});
expect(task.reviewerContextRetryCount).toBe(1);
expect(store.updateTask).toHaveBeenCalled();
expect(consoleSpy).toHaveBeenCalledWith(
expect.stringContaining("[retry-burned] retry-burned"),
expect.objectContaining({ taskId: "FN-1", category: "reviewerContext", attempt: 1, total: 1 }),
);
});
it("does not rewrite persisted counters when skipIncrement=true", async () => {
task.stuckKillCount = 3;
await recordRetry({
store: store as never,
settings: { ...baseSettings, maxTotalRetriesBeforeFail: 3 },
task,
category: "stuckKill",
role: "self-healing",
skipIncrement: true,
});
expect(task.stuckKillCount).toBe(3);
expect(store.updateTask).not.toHaveBeenCalled();
expect(consoleSpy).toHaveBeenCalled();
});
it("throws RetryStormError when master cap is exceeded", async () => {
task.stuckKillCount = 26;
await expect(
recordRetry({
store: store as never,
settings: { ...baseSettings, maxTotalRetriesBeforeFail: 25 },
task,
category: "stuckKill",
role: "self-healing",
skipIncrement: true,
}),
).rejects.toBeInstanceOf(RetryStormError);
});
});

View File

@@ -314,8 +314,11 @@ describe("reviewStep — context-limit retry", () => {
},
} as any);
const task = { id: "FN-4082", column: "in-progress", description: "d", dependencies: [], steps: [], currentStep: 0, log: [], prompt: "# prompt", createdAt: "2026-01-01T00:00:00.000Z", updatedAt: "2026-01-01T00:00:00.000Z", reviewerContextRetryCount: 0 };
const store = {
getSettings: vi.fn().mockResolvedValue({}),
getSettings: vi.fn().mockResolvedValue({ maxReviewerContextRetries: 2, maxTotalRetriesBeforeFail: 25 }),
getTask: vi.fn().mockImplementation(async () => task),
updateTask: vi.fn().mockImplementation(async (_id: string, patch: Record<string, unknown>) => Object.assign(task, patch)),
logEntry: vi.fn().mockResolvedValue(undefined),
appendAgentLog: vi.fn().mockResolvedValue(undefined),
};
@@ -359,6 +362,7 @@ describe("reviewStep — context-limit retry", () => {
"FN-4082",
"code review hit context limit — retrying with compacted request",
);
expect(task.reviewerContextRetryCount).toBe(1);
});
it("returns UNAVAILABLE when both attempts hit the context limit", async () => {
@@ -401,8 +405,11 @@ describe("reviewStep — fallback retry for terminal unavailable", () => {
.mockResolvedValueOnce(createMockSession("No parseable verdict here."))
.mockResolvedValueOnce(createMockSession("### Verdict: APPROVE\n### Summary\nRecovered on fallback."));
const task = { id: "FN-4092", column: "in-progress", description: "d", dependencies: [], steps: [], currentStep: 0, log: [], prompt: "# prompt", createdAt: "2026-01-01T00:00:00.000Z", updatedAt: "2026-01-01T00:00:00.000Z", reviewerFallbackRetryCount: 0 };
const store = {
getSettings: vi.fn().mockResolvedValue({}),
getSettings: vi.fn().mockResolvedValue({ maxReviewerFallbackRetries: 2, maxTotalRetriesBeforeFail: 25 }),
getTask: vi.fn().mockImplementation(async () => task),
updateTask: vi.fn().mockImplementation(async (_id: string, patch: Record<string, unknown>) => Object.assign(task, patch)),
logEntry: vi.fn().mockResolvedValue(undefined),
appendAgentLog: vi.fn().mockResolvedValue(undefined),
};
@@ -423,6 +430,7 @@ describe("reviewStep — fallback retry for terminal unavailable", () => {
"FN-4092",
expect.stringContaining("review retry with fallback model after UNAVAILABLE verdict"),
);
expect(task.reviewerFallbackRetryCount).toBe(1);
});
it("retries once after non-context reviewer error", async () => {

View File

@@ -6,6 +6,7 @@ import { delimiter, isAbsolute, join, relative, resolve as resolvePath } from "n
import { existsSync, realpathSync } from "node:fs";
import { readFile, writeFile } from "node:fs/promises";
import type { TaskStore, Task, TaskDetail, TaskTokenUsage, StepStatus, Settings, WorkflowStep, MissionStore, Slice, AgentState, AgentCapability, RunMutationContext, AgentHeartbeatConfig, Agent, AgentMemoryInclusionMode } from "@fusion/core";
import { RetryStormError, serializeRetryStormError } from "@fusion/core";
import {
ApprovalRequestStore,
buildExecutionMemoryInstructions,
@@ -99,6 +100,7 @@ import {
import { createFusionAuthStorage, getModelRegistryModelsPath } from "./auth-storage.js";
import { createRunVerificationTool } from "./run-verification-tool.js";
import { createFallbackModelObserver } from "./fallback-model-observer.js";
import { recordRetry } from "./retry-burned-logger.js";
import type { AgentActionGateContext } from "./agent-action-gate.js";
// Re-export for backward compatibility (tests import from executor.ts)
@@ -4285,6 +4287,16 @@ export class TaskExecutor {
outcome = await this.handleBranchConflict(task, err);
if (outcome !== "retry") break;
await this.store.logEntry(task.id, `[recovery] ${task.id} branch-conflict auto-retry requested (${attempt}/${this.MAX_AUTO_RECOVERY_ATTEMPTS})`, undefined, this.currentRunContext);
const taskForRetry = await this.store.getTask(task.id);
await recordRetry({
store: this.store,
settings: await this.store.getSettings(),
task: taskForRetry,
category: "branchConflict",
role: "executor",
agentId: this.agentId,
attempt,
});
}
if (outcome === "retry") {
await this.store.updateTask(task.id, {
@@ -4350,9 +4362,12 @@ export class TaskExecutor {
this.options.onError?.(task, err instanceof Error ? err : new Error(errorMessage));
return;
}
const terminalError = err instanceof RetryStormError
? JSON.stringify(serializeRetryStormError(err))
: errorMessage;
executorLog.error(`${task.id} execution failed:`, errorDetail);
await this.store.logEntry(task.id, `Execution failed: ${errorMessage}`, errorStack ?? errorDetail, this.currentRunContext);
await this.store.updateTask(task.id, { status: "failed", error: errorMessage });
await this.store.logEntry(task.id, `Execution failed: ${terminalError}`, errorStack ?? errorDetail, this.currentRunContext);
await this.store.updateTask(task.id, { status: "failed", error: terminalError });
await this.persistTokenUsage(task.id);
await this.store.moveTask(task.id, "in-review");
executorLog.log(`${task.id} execution failed → in-review`);

View File

@@ -0,0 +1,103 @@
import {
computeRetrySummary,
RetryStormError,
type Settings,
type TaskDetail,
type TaskStore,
} from "@fusion/core";
import { createLogger } from "./logger.js";
const retryBurnedLog = createLogger("retry-burned");
type RetryCategory =
| "branchConflict"
| "reviewerContext"
| "reviewerFallback"
| "stuckKill"
| "recovery"
| "taskDone"
| "workflowStep"
| "verification"
| "postReviewFix"
| "mergeConflict";
const CATEGORY_COLUMN: Record<RetryCategory, keyof TaskDetail> = {
branchConflict: "branchConflictRecoveryCount",
reviewerContext: "reviewerContextRetryCount",
reviewerFallback: "reviewerFallbackRetryCount",
stuckKill: "stuckKillCount",
recovery: "recoveryRetryCount",
taskDone: "taskDoneRetryCount",
workflowStep: "workflowStepRetries",
verification: "verificationFailureCount",
postReviewFix: "postReviewFixCount",
mergeConflict: "mergeConflictBounceCount",
};
const CATEGORY_CAP = (category: RetryCategory, settings: Settings): number | undefined => {
switch (category) {
case "branchConflict":
return settings.maxBranchConflictRecoveries;
case "reviewerContext":
return settings.maxReviewerContextRetries;
case "reviewerFallback":
return settings.maxReviewerFallbackRetries;
default:
return undefined;
}
};
export async function recordRetry(options: {
store: Pick<TaskStore, "updateTask" | "getTask">;
settings: Settings;
task: TaskDetail;
category: RetryCategory;
role: string;
agentId?: string;
attempt?: number;
skipIncrement?: boolean;
}): Promise<TaskDetail> {
const { store, settings, task, category, role, agentId, attempt, skipIncrement } = options;
const column = CATEGORY_COLUMN[category];
if (!skipIncrement) {
const current = (task[column] as number | undefined) ?? 0;
await store.updateTask(task.id, { [column]: current + 1 });
}
const refreshed = await store.getTask(task.id);
const breakdown = computeRetrySummary(refreshed);
const categoryCount = (refreshed[column] as number | undefined) ?? 0;
const categoryCap = CATEGORY_CAP(category, settings);
const totalCap = settings.maxTotalRetriesBeforeFail;
retryBurnedLog.log("retry-burned", {
taskId: task.id,
agentId,
role,
category,
attempt,
total: breakdown.total,
breakdown,
});
if (typeof categoryCap === "number" && categoryCount > categoryCap) {
throw new RetryStormError({
category,
total: breakdown.total,
cap: categoryCap,
breakdown,
});
}
if (typeof totalCap === "number" && breakdown.total > totalCap) {
throw new RetryStormError({
category,
total: breakdown.total,
cap: totalCap,
breakdown,
});
}
return refreshed;
}

View File

@@ -10,6 +10,7 @@
import type { TaskStore, TaskComment, AgentPromptsConfig, Settings } from "@fusion/core";
import { buildReviewerMemoryInstructions, resolveAgentPrompt, resolvePersistAgentThinkingLog, resolveAgentMemoryInclusionMode } from "@fusion/core";
import { recordRetry } from "./retry-burned-logger.js";
import { describeModel, promptWithFallback } from "./pi.js";
import { isContextLimitError } from "./context-limit-detector.js";
import { createResolvedAgentSession, extractRuntimeHint } from "./agent-session-helpers.js";
@@ -596,6 +597,15 @@ export async function reviewStep(
reviewerLog.warn(`${taskId}: ${retryLogMessage}`);
if (options.store && options.taskId) {
await options.store.logEntry(options.taskId, retryLogMessage).catch(() => undefined);
const taskForRetry = await options.store.getTask(options.taskId);
await recordRetry({
store: options.store,
settings: liveSettings ?? options.settings ?? {},
task: taskForRetry,
category: "reviewerContext",
role: "reviewer",
agentId: options.agentId,
});
}
reviewText = "";
@@ -654,6 +664,17 @@ export async function reviewStep(
} catch (err) {
if (hasConfiguredFallback) {
await logFallbackRetry("reviewer error", `${validatorFallbackProvider}/${validatorFallbackModelId}`);
if (options.store && options.taskId) {
const taskForRetry = await options.store.getTask(options.taskId);
await recordRetry({
store: options.store,
settings: liveSettings ?? options.settings ?? {},
task: taskForRetry,
category: "reviewerFallback",
role: "reviewer",
agentId: options.agentId,
});
}
try {
return await runAttempt(request, {
forceProvider: validatorFallbackProvider,
@@ -665,6 +686,17 @@ export async function reviewStep(
}
await logFallbackRetry("reviewer error", "same-model strict prompt");
if (options.store && options.taskId) {
const taskForRetry = await options.store.getTask(options.taskId);
await recordRetry({
store: options.store,
settings: liveSettings ?? options.settings ?? {},
task: taskForRetry,
category: "reviewerFallback",
role: "reviewer",
agentId: options.agentId,
});
}
try {
return await runAttempt(fallbackReviewRequest);
} catch {
@@ -678,6 +710,17 @@ export async function reviewStep(
if (hasConfiguredFallback) {
await logFallbackRetry("UNAVAILABLE verdict", `${validatorFallbackProvider}/${validatorFallbackModelId}`);
if (options.store && options.taskId) {
const taskForRetry = await options.store.getTask(options.taskId);
await recordRetry({
store: options.store,
settings: liveSettings ?? options.settings ?? {},
task: taskForRetry,
category: "reviewerFallback",
role: "reviewer",
agentId: options.agentId,
});
}
return runAttempt(request, {
forceProvider: validatorFallbackProvider,
forceModelId: validatorFallbackModelId,
@@ -685,6 +728,17 @@ export async function reviewStep(
}
await logFallbackRetry("UNAVAILABLE verdict", "same-model strict prompt");
if (options.store && options.taskId) {
const taskForRetry = await options.store.getTask(options.taskId);
await recordRetry({
store: options.store,
settings: liveSettings ?? options.settings ?? {},
task: taskForRetry,
category: "reviewerFallback",
role: "reviewer",
agentId: options.agentId,
});
}
return runAttempt(fallbackReviewRequest);
}