Files
agent-workflow/src/workflow-state.test.ts
T

1156 lines
45 KiB
TypeScript

import { execFileSync } from "node:child_process";
import fs from "node:fs";
import os from "node:os";
import path from "node:path";
import { afterEach, beforeEach, describe, expect, it } from "vitest";
import { loadProjectWorkflowConfig, resolveWorkflowConfig } from "./workflow-config.js";
import {
WorkflowLockError,
acquireLock,
findRunDir,
lockFilePath,
readLock,
releaseLock,
repoHash,
runDir,
workflowsRoot,
} from "./workflow-state.js";
describe("Workflow state root and locks", () => {
let tmpDir: string;
let repoRoot: string;
let stateRoot: string;
let originalAgentDir: string | undefined;
beforeEach(() => {
tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), "agent-workflow-state-test-"));
repoRoot = path.join(tmpDir, "repo");
stateRoot = path.join(tmpDir, "workflows");
fs.mkdirSync(repoRoot, { recursive: true });
originalAgentDir = process.env.AGENT_WORKFLOW_DIR;
process.env.AGENT_WORKFLOW_DIR = stateRoot;
});
afterEach(() => {
if (originalAgentDir === undefined) delete process.env.AGENT_WORKFLOW_DIR;
else process.env.AGENT_WORKFLOW_DIR = originalAgentDir;
fs.rmSync(tmpDir, { recursive: true, force: true });
});
it("uses AGENT_WORKFLOW_DIR for all state", () => {
expect(workflowsRoot()).toBe(stateRoot);
expect(runDir(repoRoot, "run-1")).toBe(path.join(stateRoot, repoHash(repoRoot), "run-1"));
expect(lockFilePath(repoRoot)).toBe(path.join(stateRoot, repoHash(repoRoot), "active.lock"));
});
it("finds a run with and without a repository hint", () => {
const expected = runDir(repoRoot, "run-1");
fs.mkdirSync(expected, { recursive: true });
fs.writeFileSync(path.join(expected, "state.json"), "{}\n");
expect(findRunDir("run-1", repoRoot)).toBe(expected);
expect(findRunDir("run-1")).toBe(expected);
expect(findRunDir("missing", repoRoot)).toBeUndefined();
});
it("reuses the same run lock and rejects a different run", () => {
acquireLock(repoRoot, "run-1");
expect(readLock(repoRoot)).toBe("run-1");
expect(() => acquireLock(repoRoot, "run-1")).not.toThrow();
expect(() => acquireLock(repoRoot, "run-2")).toThrow(WorkflowLockError);
});
it("releases only the owning run lock", () => {
acquireLock(repoRoot, "run-1");
releaseLock(repoRoot, "run-2");
expect(readLock(repoRoot)).toBe("run-1");
releaseLock(repoRoot, "run-1");
expect(readLock(repoRoot)).toBeUndefined();
});
});
describe("Project workflow config", () => {
let tmpDir: string;
beforeEach(() => {
tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), "agent-workflow-config-test-"));
});
afterEach(() => {
fs.rmSync(tmpDir, { recursive: true, force: true });
});
it("loads agent-workflow.json", () => {
fs.writeFileSync(
path.join(tmpDir, "agent-workflow.json"),
JSON.stringify({ version: "1", defaultExecutor: "claude", maxCycles: 10 }),
);
expect(loadProjectWorkflowConfig(tmpDir)).toMatchObject({
defaultExecutor: "claude",
maxCycles: 10,
});
});
it("returns undefined when no config exists", () => {
expect(loadProjectWorkflowConfig(tmpDir)).toBeUndefined();
});
});
describe("Resolved workflow config", () => {
let tmpDir: string;
let repoRoot: string;
let originalHome: string | undefined;
beforeEach(() => {
tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), "agent-workflow-resolve-test-"));
repoRoot = path.join(tmpDir, "repo");
fs.mkdirSync(repoRoot, { recursive: true });
originalHome = process.env.HOME;
process.env.HOME = tmpDir;
});
afterEach(() => {
if (originalHome === undefined) delete process.env.HOME;
else process.env.HOME = originalHome;
fs.rmSync(tmpDir, { recursive: true, force: true });
});
function writeGlobalConfig(config: unknown): void {
const globalDir = path.join(tmpDir, ".agent-workflow");
fs.mkdirSync(globalDir, { recursive: true });
fs.writeFileSync(path.join(globalDir, "config.json"), JSON.stringify(config));
}
function writeProjectConfig(config: unknown): void {
fs.writeFileSync(path.join(repoRoot, "agent-workflow.json"), JSON.stringify(config));
}
it("merges global < project < CLI per field", () => {
writeGlobalConfig({
version: "1",
defaultExecutor: "agy",
defaultVcs: "none",
maxCycles: 10,
timeoutSeconds: 999,
executors: {
reasonix: { binary: "global-reasonix", model: "global-model" },
},
});
writeProjectConfig({
version: "1",
defaultExecutor: "reasonix",
maxCycles: 3,
executors: {
reasonix: { model: "project-model" },
},
});
const config = resolveWorkflowConfig(repoRoot, { maxCycles: 4 });
expect(config.executor).toBe("reasonix");
expect(config.maxCycles).toBe(4);
expect(config.timeoutSeconds).toBe(999);
expect(config.vcs).toBe("none");
expect(config.executorConfig).toMatchObject({
binary: "global-reasonix",
model: "project-model",
});
});
it("treats explicit project defaultVcs auto as an override that suppresses global defaultVcs", () => {
writeGlobalConfig({
version: "1",
defaultVcs: "git",
});
writeProjectConfig({
version: "1",
defaultVcs: "auto",
});
const config = resolveWorkflowConfig(repoRoot, {});
expect(config.vcs).toBeUndefined();
});
it("uses global defaults when project config is absent", () => {
writeGlobalConfig({
version: "1",
defaultExecutor: "claude",
executors: {
claude: { model: "claude-sonnet" },
},
});
const config = resolveWorkflowConfig(repoRoot, {});
expect(config.executor).toBe("claude");
expect(config.executorConfig.model).toBe("claude-sonnet");
});
it("rejects invalid model values", () => {
writeProjectConfig({
version: "1",
defaultExecutor: "reasonix",
executors: { reasonix: { model: "" } },
});
expect(() => resolveWorkflowConfig(repoRoot, {})).toThrow(/non-empty string/);
});
});
describe("Retry-execute recovery", () => {
let tmpDir: string;
let repoRoot: string;
let stateRoot: string;
let originalAgentDir: string | undefined;
let originalHome: string | undefined;
beforeEach(() => {
tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), "agent-workflow-retry-test-"));
repoRoot = path.join(tmpDir, "repo");
stateRoot = path.join(tmpDir, "workflows");
fs.mkdirSync(repoRoot, { recursive: true });
execFileSync("git", ["init"], { cwd: repoRoot, stdio: "ignore" });
execFileSync("git", ["config", "user.email", "test@example.com"], { cwd: repoRoot });
execFileSync("git", ["config", "user.name", "Test User"], { cwd: repoRoot });
fs.writeFileSync(path.join(repoRoot, "test.txt"), "test");
execFileSync("git", ["add", "."], { cwd: repoRoot });
execFileSync("git", ["commit", "-m", "init"], { cwd: repoRoot });
originalAgentDir = process.env.AGENT_WORKFLOW_DIR;
process.env.AGENT_WORKFLOW_DIR = stateRoot;
originalHome = process.env.HOME;
process.env.HOME = tmpDir;
});
afterEach(() => {
if (originalAgentDir === undefined) delete process.env.AGENT_WORKFLOW_DIR;
else process.env.AGENT_WORKFLOW_DIR = originalAgentDir;
if (originalHome === undefined) delete process.env.HOME;
else process.env.HOME = originalHome;
fs.rmSync(tmpDir, { recursive: true, force: true });
});
function writeFakeExecutor(sessionFilePath: string): string {
const scriptPath = path.join(path.dirname(sessionFilePath), "fake-reasonix");
const script = `#!/bin/bash
set -e
SESSION_FILE="${sessionFilePath}"
echo '{"type":"session_started"}' > "\${SESSION_FILE}"
echo '{"session_id":"fake-session-123","conversation_id":"fake-conv-456"}'
`;
fs.writeFileSync(scriptPath, script, { mode: 0o755 });
return scriptPath;
}
function plan() {
return {
version: "1" as const,
title: "Test",
planMarkdown: "Test plan",
scope: ["test.txt"],
acceptanceCriteria: ["Recovery succeeds"],
verificationCommands: [["git", "status", "--short"]],
};
}
it("recovers a stale executing state without changing its cycle or session", async () => {
const { retryExecuteWorkflow } = await import("./workflow-engine.js");
const { gitHead, initWorkflowState, readState, writeState } = await import("./workflow-state.js");
const runId = "stale-run-123";
const dir = runDir(repoRoot, runId);
fs.mkdirSync(dir, { recursive: true });
const sessionFile = path.join(dir, "session.jsonl");
fs.writeFileSync(sessionFile, '{"type":"session"}');
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: plan(),
executor: "reasonix",
maxCycles: 3,
});
state.status = "executing";
state.enginePid = 99_999_999;
state.currentCycle = 1;
state.cycles = [{
cycleIndex: 1,
startedAt: new Date().toISOString(),
executorLogFile: "cycle-1-exec.log",
executorAttemptLogs: ["cycle-1-exec.log"],
}];
state.sessionHandle = { kind: "reasonix", sessionFilePath: sessionFile };
state.executorConfig = { binary: writeFakeExecutor(sessionFile) };
fs.writeFileSync(path.join(dir, "cycle-1-exec.log"), "initial attempt failed");
acquireLock(repoRoot, runId);
writeState(dir, state);
await retryExecuteWorkflow({ runId });
const recovered = readState(dir)!;
const cycle = recovered.cycles[0]!;
expect(recovered.status).toBe("awaiting_review");
expect(recovered.currentCycle).toBe(1);
expect(recovered.sessionHandle).toEqual(state.sessionHandle);
expect(cycle.executorAttemptLogs).toHaveLength(2);
expect(cycle.executorAttemptLogs?.[0]).toBe("cycle-1-exec.log");
expect(cycle.executorLogFile).toBe(cycle.executorAttemptLogs?.[1]);
const retryLog = fs.readFileSync(path.join(dir, cycle.executorLogFile!), "utf8");
expect(retryLog).not.toMatch(/Permission denied|command not found|No such file/);
const events = fs.readFileSync(path.join(dir, "events.jsonl"), "utf8")
.trim().split("\n").map((line) => JSON.parse(line));
expect(events.find((event) => event.kind === "executor_retried")?.data.recoveryMode)
.toBe("stale_execution");
});
it("rejects retry while the recorded engine process is alive", async () => {
const { retryExecuteWorkflow } = await import("./workflow-engine.js");
const { gitHead, initWorkflowState, writeState } = await import("./workflow-state.js");
const runId = "live-run-456";
const dir = runDir(repoRoot, runId);
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: plan(),
executor: "reasonix",
maxCycles: 3,
});
state.status = "executing";
state.enginePid = process.pid;
state.cycles = [{ cycleIndex: 1, startedAt: new Date().toISOString() }];
writeState(dir, state);
await expect(retryExecuteWorkflow({ runId }))
.rejects.toThrow(/cannot retry execution from status 'executing'/);
});
it("recovers a failed initial cycle and preserves attempt logs", async () => {
const { retryExecuteWorkflow } = await import("./workflow-engine.js");
const { gitHead, initWorkflowState, readState, writeState } = await import("./workflow-state.js");
const runId = "initial-retry-789";
const dir = runDir(repoRoot, runId);
fs.mkdirSync(dir, { recursive: true });
const sessionFile = path.join(dir, "session.jsonl");
fs.writeFileSync(sessionFile, '{"type":"session"}');
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: plan(),
executor: "reasonix",
maxCycles: 3,
});
state.status = "blocked";
state.stopDescription = "Executor exited with failure: test error";
state.cycles = [{
cycleIndex: 1,
startedAt: new Date().toISOString(),
executorLogFile: "cycle-1-exec.log",
executorAttemptLogs: ["cycle-1-exec.log"],
}];
state.sessionHandle = { kind: "reasonix", sessionFilePath: sessionFile };
state.executorConfig = { binary: writeFakeExecutor(sessionFile) };
fs.writeFileSync(path.join(dir, "cycle-1-exec.log"), "initial attempt failed");
writeState(dir, state);
await retryExecuteWorkflow({ runId });
const recovered = readState(dir)!;
const cycle = recovered.cycles[0]!;
expect(recovered.status).toBe("awaiting_review");
expect(recovered.currentCycle).toBe(1);
expect(recovered.sessionHandle).toEqual(state.sessionHandle);
expect(cycle.executorAttemptLogs).toHaveLength(2);
expect(cycle.executorLogFile).toBe(cycle.executorAttemptLogs?.[1]);
const events = fs.readFileSync(path.join(dir, "events.jsonl"), "utf8")
.trim().split("\n").map((line) => JSON.parse(line));
expect(events.find((event) => event.kind === "executor_retried")?.data.recoveryMode)
.toBe("initial_continuation");
});
it("blocks before executor launch when reloaded model config is invalid, then recovers after correction", async () => {
const { retryExecuteWorkflow } = await import("./workflow-engine.js");
const { gitHead, initWorkflowState, readState, writeState } = await import("./workflow-state.js");
const runId = "invalid-model-retry";
const dir = runDir(repoRoot, runId);
fs.mkdirSync(dir, { recursive: true });
const sessionFile = path.join(dir, "session.jsonl");
fs.writeFileSync(sessionFile, '{"type":"session"}');
const configRoot = path.join(tmpDir, "config-root");
fs.mkdirSync(configRoot, { recursive: true });
fs.writeFileSync(
path.join(configRoot, "agent-workflow.json"),
JSON.stringify({ version: "1", executors: { reasonix: { model: "" } } }),
);
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: plan(),
executor: "reasonix",
maxCycles: 3,
});
state.status = "blocked";
state.stopDescription = "Executor exited with failure: test error";
state.configSourceRoot = configRoot;
state.cycles = [{
cycleIndex: 1,
startedAt: new Date().toISOString(),
executorLogFile: "cycle-1-exec.log",
executorAttemptLogs: ["cycle-1-exec.log"],
}];
state.sessionHandle = { kind: "reasonix", sessionFilePath: sessionFile };
state.executorConfig = { binary: writeFakeExecutor(sessionFile) };
fs.writeFileSync(path.join(dir, "cycle-1-exec.log"), "initial attempt failed");
writeState(dir, state);
await expect(retryExecuteWorkflow({ runId })).rejects.toThrow(/non-empty string/);
let recovered = readState(dir)!;
expect(recovered.status).toBe("blocked");
expect(recovered.stopDescription).toMatch(/Executor configuration is invalid:/);
// Fix the config and retry again.
fs.writeFileSync(
path.join(configRoot, "agent-workflow.json"),
JSON.stringify({ version: "1", executors: { reasonix: { model: "fixed-model" } } }),
);
await retryExecuteWorkflow({ runId });
recovered = readState(dir)!;
expect(recovered.status).toBe("awaiting_review");
expect(recovered.executorConfig?.model).toBe("fixed-model");
const events = fs.readFileSync(path.join(dir, "events.jsonl"), "utf8")
.trim().split("\n").map((line) => JSON.parse(line));
expect(events.find((event) => event.kind === "executor_started")?.data.model)
.toBe("fixed-model");
expect(events.find((event) => event.kind === "executor_retried")?.data.model)
.toBe("fixed-model");
});
it("reloads model for the run's frozen executor even when defaultExecutor changes", async () => {
const { retryExecuteWorkflow } = await import("./workflow-engine.js");
const { gitHead, initWorkflowState, readState, writeState } = await import("./workflow-state.js");
const runId = "frozen-executor-model";
const dir = runDir(repoRoot, runId);
fs.mkdirSync(dir, { recursive: true });
const sessionFile = path.join(dir, "session.jsonl");
fs.writeFileSync(sessionFile, '{"type":"session"}');
const configRoot = path.join(tmpDir, "frozen-config-root");
fs.mkdirSync(configRoot, { recursive: true });
fs.writeFileSync(
path.join(configRoot, "agent-workflow.json"),
JSON.stringify({
version: "1",
defaultExecutor: "claude",
executors: {
claude: { model: "claude-other" },
reasonix: { model: "frozen-reasonix-model" },
},
}),
);
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: plan(),
executor: "reasonix",
maxCycles: 3,
});
state.status = "blocked";
state.stopDescription = "Executor exited with failure: test error";
state.configSourceRoot = configRoot;
state.cycles = [{
cycleIndex: 1,
startedAt: new Date().toISOString(),
executorLogFile: "cycle-1-exec.log",
executorAttemptLogs: ["cycle-1-exec.log"],
}];
state.sessionHandle = { kind: "reasonix", sessionFilePath: sessionFile };
state.executorConfig = { binary: writeFakeExecutor(sessionFile) };
fs.writeFileSync(path.join(dir, "cycle-1-exec.log"), "initial attempt failed");
writeState(dir, state);
await retryExecuteWorkflow({ runId });
const recovered = readState(dir)!;
expect(recovered.status).toBe("awaiting_review");
expect(recovered.executor).toBe("reasonix");
expect(recovered.executorConfig?.model).toBe("frozen-reasonix-model");
});
it("recovers a missing Reasonix handle for a blocked fix cycle", async () => {
const { retryExecuteWorkflow } = await import("./workflow-engine.js");
const { gitHead, initWorkflowState, readState, writeState } = await import("./workflow-state.js");
const runId = "reasonix-missing-handle";
const dir = runDir(repoRoot, runId);
const projectKey = path.resolve(repoRoot).replaceAll(path.sep, "-");
const sessionsDir = path.join(process.env.HOME!, ".reasonix", "projects", projectKey, "sessions");
const sessionFile = path.join(sessionsDir, "matching-session.jsonl");
fs.mkdirSync(sessionsDir, { recursive: true });
fs.writeFileSync(sessionFile, [
JSON.stringify({ role: "system", content: `Current workspace: "${path.dirname(repoRoot)}"` }),
JSON.stringify({ role: "user", content: "# Implementation Task: Test" }),
].join("\n"));
fs.writeFileSync(`${sessionFile}.meta`, "{}");
const fakeReasonix = path.join(tmpDir, "fake-reasonix-resume");
fs.writeFileSync(fakeReasonix, "#!/bin/bash\nexit 0\n", { mode: 0o755 });
const startedAt = new Date(Date.now() - 5000).toISOString();
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: plan(),
executor: "reasonix",
maxCycles: 3,
});
state.status = "blocked";
state.stopDescription = "No session handle available for resume.";
state.currentCycle = 2;
state.cycles = [
{
cycleIndex: 1,
startedAt,
completedAt: new Date().toISOString(),
decisionFile: "cycle-1-decision.json",
decisionOutcome: "fix",
},
{ cycleIndex: 2, startedAt: new Date().toISOString() },
];
state.executorConfig = { binary: fakeReasonix };
fs.mkdirSync(dir, { recursive: true });
fs.writeFileSync(path.join(dir, "cycle-1-decision.json"), JSON.stringify({
version: "1",
outcome: "fix",
findingDecisions: [],
verificationEvidence: [],
decidedAt: new Date().toISOString(),
}));
writeState(dir, state);
await retryExecuteWorkflow({ runId });
const recovered = readState(dir)!;
expect(recovered.status).toBe("awaiting_review");
expect(recovered.sessionHandle).toEqual({ kind: "reasonix", sessionFilePath: sessionFile });
const events = fs.readFileSync(path.join(dir, "events.jsonl"), "utf8")
.trim().split("\n").map((line) => JSON.parse(line));
expect(events.find((event) => event.kind === "executor_session_rebound")?.data).toEqual({
sessionFilePath: sessionFile,
recovery: "reasonix_project_session",
});
});
it("rebinds a validated Claude session before retrying the failed cycle", async () => {
const { retryExecuteWorkflow } = await import("./workflow-engine.js");
const { gitHead, initWorkflowState, readState, writeState } = await import("./workflow-state.js");
const runId = "claude-rebind-123";
const sessionId = "22778b47-e4f3-4645-84cd-1c6a74a912d8";
const dir = runDir(repoRoot, runId);
const claudeProject = path.join(tmpDir, ".claude", "projects", "test-project");
fs.mkdirSync(claudeProject, { recursive: true });
fs.writeFileSync(
path.join(claudeProject, `${sessionId}.jsonl`),
`${JSON.stringify({ sessionId, cwd: repoRoot, message: { content: "# Implementation Task: Test" } })}\n`,
);
const fakeClaude = path.join(tmpDir, "fake-claude");
fs.writeFileSync(fakeClaude, `#!/bin/bash
printf '{"type":"system","subtype":"init","session_id":"${sessionId}"}\\n'
printf '{"type":"result","result":"done","session_id":"${sessionId}"}\\n'
`, { mode: 0o755 });
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: plan(),
executor: "claude",
maxCycles: 3,
});
state.status = "blocked";
state.stopDescription = "Executor failed to resume: missing session";
state.currentCycle = 1;
state.cycles = [{
cycleIndex: 1,
startedAt: new Date().toISOString(),
executorLogFile: "cycle-1-exec.log",
executorAttemptLogs: ["cycle-1-exec.log"],
}];
state.sessionHandle = { kind: "claude", sessionId: "97a606e8-cfa3-4996-83a8-ddb7ece4234b" };
state.executorConfig = { binary: fakeClaude };
fs.mkdirSync(dir, { recursive: true });
fs.writeFileSync(path.join(dir, "cycle-1-exec.log"), "initial attempt failed");
writeState(dir, state);
await retryExecuteWorkflow({ runId, sessionId });
const recovered = readState(dir)!;
expect(recovered.status).toBe("awaiting_review");
expect(recovered.sessionHandle).toEqual({ kind: "claude", sessionId });
const events = fs.readFileSync(path.join(dir, "events.jsonl"), "utf8")
.trim().split("\n").map((line) => JSON.parse(line));
expect(events.find((event) => event.kind === "executor_session_rebound")?.data).toEqual({
previousSessionId: "97a606e8-cfa3-4996-83a8-ddb7ece4234b",
sessionId,
});
});
it("rejects a Claude session from another repository without changing state", async () => {
const { retryExecuteWorkflow } = await import("./workflow-engine.js");
const { gitHead, initWorkflowState, readState, writeState } = await import("./workflow-state.js");
const runId = "claude-rebind-reject";
const sessionId = "22778b47-e4f3-4645-84cd-1c6a74a912d8";
const dir = runDir(repoRoot, runId);
const claudeProject = path.join(tmpDir, ".claude", "projects", "other-project");
fs.mkdirSync(claudeProject, { recursive: true });
fs.writeFileSync(
path.join(claudeProject, `${sessionId}.jsonl`),
`${JSON.stringify({ sessionId, cwd: path.join(tmpDir, "other"), message: { content: "Test" } })}\n`,
);
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: plan(),
executor: "claude",
maxCycles: 3,
});
state.status = "blocked";
state.stopDescription = "Executor failed to resume: missing session";
state.currentCycle = 1;
state.cycles = [{ cycleIndex: 1, startedAt: new Date().toISOString() }];
state.sessionHandle = { kind: "claude", sessionId: "97a606e8-cfa3-4996-83a8-ddb7ece4234b" };
writeState(dir, state);
await expect(retryExecuteWorkflow({ runId, sessionId }))
.rejects.toThrow(/does not belong to repository/);
expect(readState(dir)?.sessionHandle).toEqual(state.sessionHandle);
expect(fs.existsSync(path.join(dir, "events.jsonl"))).toBe(false);
});
});
// ---------------------------------------------------------------------------
// Progress sidecar
// ---------------------------------------------------------------------------
describe("Progress sidecar", () => {
let tmpDir: string;
beforeEach(() => {
tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), "agent-workflow-sidecar-test-"));
});
afterEach(() => {
fs.rmSync(tmpDir, { recursive: true, force: true });
});
it("writes and reads progress sidecar", async () => {
const { writeProgressSidecar, readProgressSidecar } =
await import("./workflow-state.js");
const sidecar = {
executor: "reasonix" as const,
cycleIndex: 1,
attemptIndex: 0,
phase: "executing",
startedAt: new Date().toISOString(),
lastActivityAt: new Date().toISOString(),
logFilePath: "cycle-1-exec.log",
logBytes: 1024,
activitySummary: "Running build",
updatedAt: new Date().toISOString(),
};
writeProgressSidecar(tmpDir, sidecar);
const read = readProgressSidecar(tmpDir);
expect(read).toBeDefined();
expect(read!.executor).toBe("reasonix");
expect(read!.phase).toBe("executing");
expect(read!.logBytes).toBe(1024);
expect(read!.activitySummary).toBe("Running build");
});
it("returns undefined when sidecar file does not exist", async () => {
const { readProgressSidecar } = await import("./workflow-state.js");
expect(readProgressSidecar(tmpDir)).toBeUndefined();
});
it("reads attempt log path for initial attempt and retry", async () => {
const { attemptLogPath } = await import("./workflow-state.js");
const dir = tmpDir;
expect(path.basename(attemptLogPath(dir, 1, 0))).toBe("cycle-1-exec.log");
expect(path.basename(attemptLogPath(dir, 2, 0))).toBe("cycle-2-exec.log");
expect(path.basename(attemptLogPath(dir, 1, 1))).toBe("cycle-1-exec-retry-1.log");
expect(path.basename(attemptLogPath(dir, 1, 3))).toBe("cycle-1-exec-retry-3.log");
});
it("nextAttemptIndex finds the next free index", async () => {
const { nextAttemptIndex, attemptLogPath } = await import("./workflow-state.js");
fs.writeFileSync(attemptLogPath(tmpDir, 1, 0), "initial");
expect(nextAttemptIndex(tmpDir, 1)).toBe(1);
fs.writeFileSync(attemptLogPath(tmpDir, 1, 1), "retry1");
expect(nextAttemptIndex(tmpDir, 1)).toBe(2);
fs.writeFileSync(attemptLogPath(tmpDir, 1, 2), "retry2");
expect(nextAttemptIndex(tmpDir, 1)).toBe(3);
// Non-existent cycle returns 0
expect(nextAttemptIndex(tmpDir, 9)).toBe(0);
});
it("sidecar preserves Chinese UTF-8 in activitySummary", async () => {
const { writeProgressSidecar, readProgressSidecar } =
await import("./workflow-state.js");
writeProgressSidecar(tmpDir, {
executor: "claude",
cycleIndex: 1,
attemptIndex: 0,
phase: "executing",
startedAt: new Date().toISOString(),
lastActivityAt: new Date().toISOString(),
logFilePath: "cycle-1-exec.log",
logBytes: 42,
activitySummary: "读取文件 src/index.ts",
updatedAt: new Date().toISOString(),
});
const read = readProgressSidecar(tmpDir);
expect(read).toBeDefined();
expect(read!.activitySummary).toBe("读取文件 src/index.ts");
});
});
// ---------------------------------------------------------------------------
// Status live fields and sidecar integration
// ---------------------------------------------------------------------------
describe("Status and sidecar live fields", () => {
let tmpDir: string;
let stateRoot: string;
let repoRoot: string;
let originalAgentDir: string | undefined;
beforeEach(async () => {
tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), "agent-workflow-status-test-"));
repoRoot = path.join(tmpDir, "repo");
stateRoot = path.join(tmpDir, "workflows");
fs.mkdirSync(repoRoot, { recursive: true });
const { execFileSync } = await import("node:child_process");
execFileSync("git", ["init"], { cwd: repoRoot, stdio: "ignore" });
execFileSync("git", ["config", "user.email", "test@example.com"], { cwd: repoRoot });
execFileSync("git", ["config", "user.name", "Test User"], { cwd: repoRoot });
fs.writeFileSync(path.join(repoRoot, "test.txt"), "test");
execFileSync("git", ["add", "."], { cwd: repoRoot });
execFileSync("git", ["commit", "-m", "init"], { cwd: repoRoot });
originalAgentDir = process.env.AGENT_WORKFLOW_DIR;
process.env.AGENT_WORKFLOW_DIR = stateRoot;
});
afterEach(() => {
if (originalAgentDir === undefined) delete process.env.AGENT_WORKFLOW_DIR;
else process.env.AGENT_WORKFLOW_DIR = originalAgentDir;
fs.rmSync(tmpDir, { recursive: true, force: true });
});
it("statusWorkflow includes live fields from sidecar when executing", async () => {
const { runDir, writeState, initWorkflowState, gitHead, writeProgressSidecar, readState }
= await import("./workflow-state.js");
const runId = "status-live-fields-001";
const dir = runDir(repoRoot, runId);
fs.mkdirSync(dir, { recursive: true });
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: { version: "1", title: "Test", planMarkdown: "Test", scope: ["src/"], acceptanceCriteria: ["Ok"], verificationCommands: [["echo"]] },
executor: "claude",
maxCycles: 3,
});
state.status = "executing";
state.currentCycle = 1;
state.cycles = [{ cycleIndex: 1, startedAt: new Date().toISOString() }];
writeState(dir, state);
// Write a sidecar as if an adapter was running
writeProgressSidecar(dir, {
executor: "claude",
cycleIndex: 1,
attemptIndex: 0,
pid: 12345,
phase: "executing",
startedAt: new Date(Date.now() - 120_000).toISOString(), // 2 min ago
lastActivityAt: new Date().toISOString(),
logFilePath: "cycle-1-exec.log",
logBytes: 4096,
activitySummary: "Reading src/index.ts",
updatedAt: new Date().toISOString(),
});
fs.writeFileSync(path.join(dir, "cycle-1-exec.log"), "log content");
const { statusWorkflow } = await import("./workflow-engine.js");
const chunks: string[] = [];
const origWrite = process.stdout.write.bind(process.stdout);
process.stdout.write = (chunk: unknown) => { chunks.push(String(chunk)); return true; };
try {
statusWorkflow({ runId, json: true });
} finally {
process.stdout.write = origWrite;
}
const output = JSON.parse(chunks.join(""));
expect(output.liveLogFile).toContain("cycle-1-exec.log");
expect(output.executionElapsed).toBeDefined();
expect(output.lastActivity).toBeDefined();
expect(output.logBytes).toBe(4096);
expect(output.executorPid).toBe(12345);
expect(output.activitySummary).toBe("Reading src/index.ts");
});
it("status live fields absent when no sidecar", async () => {
const { runDir, writeState, initWorkflowState, gitHead } = await import("./workflow-state.js");
const runId = "status-no-sidecar-002";
const dir = runDir(repoRoot, runId);
fs.mkdirSync(dir, { recursive: true });
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: { version: "1", title: "Test", planMarkdown: "Test", scope: ["src/"], acceptanceCriteria: ["Ok"], verificationCommands: [["echo"]] },
executor: "reasonix",
maxCycles: 3,
});
state.status = "completed";
state.currentCycle = 1;
state.cycles = [{ cycleIndex: 1, startedAt: new Date().toISOString() }];
writeState(dir, state);
const { statusWorkflow } = await import("./workflow-engine.js");
const chunks: string[] = [];
const origWrite = process.stdout.write.bind(process.stdout);
process.stdout.write = (chunk: unknown) => { chunks.push(String(chunk)); return true; };
try {
statusWorkflow({ runId, json: true });
} finally {
process.stdout.write = origWrite;
}
const output = JSON.parse(chunks.join(""));
expect(output.liveLogFile).toBeUndefined();
expect(output.executionElapsed).toBeUndefined();
expect(output.activitySummary).toBeUndefined();
});
});
// ---------------------------------------------------------------------------
// logsWorkflow follow mode
// ---------------------------------------------------------------------------
describe("logsWorkflow", () => {
let tmpDir: string;
let stateRoot: string;
let repoRoot: string;
let originalAgentDir: string | undefined;
beforeEach(async () => {
tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), "agent-workflow-logs-test-"));
repoRoot = path.join(tmpDir, "repo");
stateRoot = path.join(tmpDir, "workflows");
fs.mkdirSync(repoRoot, { recursive: true });
const { execFileSync } = await import("node:child_process");
execFileSync("git", ["init"], { cwd: repoRoot, stdio: "ignore" });
execFileSync("git", ["config", "user.email", "test@example.com"], { cwd: repoRoot });
execFileSync("git", ["config", "user.name", "Test User"], { cwd: repoRoot });
fs.writeFileSync(path.join(repoRoot, "test.txt"), "test");
execFileSync("git", ["add", "."], { cwd: repoRoot });
execFileSync("git", ["commit", "-m", "init"], { cwd: repoRoot });
originalAgentDir = process.env.AGENT_WORKFLOW_DIR;
process.env.AGENT_WORKFLOW_DIR = stateRoot;
});
afterEach(() => {
if (originalAgentDir === undefined) delete process.env.AGENT_WORKFLOW_DIR;
else process.env.AGENT_WORKFLOW_DIR = originalAgentDir;
fs.rmSync(tmpDir, { recursive: true, force: true });
});
it("tail mode prints Chinese UTF-8 lines without duplication", async () => {
const { runDir, writeState, initWorkflowState, gitHead, readState } = await import("./workflow-state.js");
const runId = "logs-tail-utf8-001";
const dir = runDir(repoRoot, runId);
fs.mkdirSync(dir, { recursive: true });
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: { version: "1", title: "Test", planMarkdown: "Test", scope: ["src/"], acceptanceCriteria: ["Ok"], verificationCommands: [["echo"]] },
executor: "claude",
maxCycles: 3,
});
state.status = "executing";
state.currentCycle = 1;
state.cycles = [{ cycleIndex: 1, startedAt: new Date().toISOString(), executorLogFile: "cycle-1-exec.log" }];
writeState(dir, state);
// Write Chinese UTF-8 content to the log file
const logContent = "line1\n读取文件 src/index.ts\n日本語\n🚀\nlast line";
fs.writeFileSync(path.join(dir, "cycle-1-exec.log"), logContent);
const { logsWorkflow } = await import("./workflow-engine.js");
const chunks: string[] = [];
const origWrite = process.stdout.write.bind(process.stdout);
process.stdout.write = (chunk: unknown) => { chunks.push(String(chunk)); return true; };
try {
await logsWorkflow({ runId, follow: false, tail: 100 });
} finally {
process.stdout.write = origWrite;
}
const output = chunks.join("");
expect(output).toContain("读取文件 src/index.ts");
expect(output).toContain("日本語");
expect(output).toContain("🚀");
// No duplication
expect(output.split("读取文件 src/index.ts").length - 1).toBe(1);
});
it("follow mode exits when state changes from executing to awaiting_review", async () => {
const { runDir, writeState, initWorkflowState, gitHead, readState } = await import("./workflow-state.js");
const runId = "logs-follow-exit-002";
const dir = runDir(repoRoot, runId);
fs.mkdirSync(dir, { recursive: true });
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: { version: "1", title: "Test", planMarkdown: "Test", scope: ["src/"], acceptanceCriteria: ["Ok"], verificationCommands: [["echo"]] },
executor: "claude",
maxCycles: 3,
});
state.status = "executing";
state.currentCycle = 1;
state.cycles = [{ cycleIndex: 1, startedAt: new Date().toISOString(), executorLogFile: "cycle-1-exec.log" }];
writeState(dir, state);
fs.writeFileSync(path.join(dir, "cycle-1-exec.log"), "initial log content\n");
const { logsWorkflow } = await import("./workflow-engine.js");
const chunks: string[] = [];
const origWrite = process.stdout.write.bind(process.stdout);
process.stdout.write = (chunk: unknown) => { chunks.push(String(chunk)); return true; };
// Start logsWorkflow in follow mode, change state after a short delay
const logsPromise = logsWorkflow({ runId, follow: true, tail: 5 });
setTimeout(() => {
const curState = readState(dir)!;
curState.status = "awaiting_review";
writeState(dir, curState);
}, 200);
try {
await logsPromise;
} finally {
process.stdout.write = origWrite;
}
const output = chunks.join("");
expect(output).toContain("initial log content");
});
it("rejects a second follower from the same Codex thread", async () => {
const { runDir, writeState, initWorkflowState, gitHead, readState } = await import("./workflow-state.js");
const runId = "logs-follow-dedupe-004";
const dir = runDir(repoRoot, runId);
fs.mkdirSync(dir, { recursive: true });
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: { version: "1", title: "Test", planMarkdown: "Test", scope: ["src/"], acceptanceCriteria: ["Ok"], verificationCommands: [["echo"]] },
executor: "claude",
maxCycles: 3,
});
state.status = "executing";
state.currentCycle = 1;
state.cycles = [{ cycleIndex: 1, startedAt: new Date().toISOString(), executorLogFile: "cycle-1-exec.log" }];
writeState(dir, state);
fs.writeFileSync(path.join(dir, "cycle-1-exec.log"), "initial log content\n");
const previousThread = process.env.CODEX_THREAD_ID;
process.env.CODEX_THREAD_ID = "thread-for-follower-dedupe";
const { logsWorkflow } = await import("./workflow-engine.js");
const chunks: string[] = [];
const origWrite = process.stdout.write.bind(process.stdout);
process.stdout.write = (chunk: unknown) => { chunks.push(String(chunk)); return true; };
const firstFollower = logsWorkflow({ runId, follow: true, tail: 5 });
try {
await expect(logsWorkflow({ runId, follow: true, tail: 5 }))
.rejects.toThrow(/already active for this Codex thread/);
const curState = readState(dir)!;
curState.status = "awaiting_review";
writeState(dir, curState);
await firstFollower;
expect(chunks.join("")).toContain("initial log content");
} finally {
process.stdout.write = origWrite;
if (previousThread === undefined) delete process.env.CODEX_THREAD_ID;
else process.env.CODEX_THREAD_ID = previousThread;
}
});
it("allows different Codex threads to hold followers for the same run", async () => {
const {
acquireLogsFollowLock,
logsFollowLockPath,
releaseLogsFollowLock,
runDir,
} = await import("./workflow-state.js");
const dir = runDir(repoRoot, "logs-follow-different-threads-005");
fs.mkdirSync(dir, { recursive: true });
const first = acquireLogsFollowLock(dir, "thread-a");
const second = acquireLogsFollowLock(dir, "thread-b");
try {
expect(first).toBe(logsFollowLockPath(dir, "thread-a"));
expect(second).toBe(logsFollowLockPath(dir, "thread-b"));
expect(first).not.toBe(second);
} finally {
releaseLogsFollowLock(first);
releaseLogsFollowLock(second);
}
expect(fs.existsSync(first)).toBe(false);
expect(fs.existsSync(second)).toBe(false);
});
it("reclaims a follower lock whose process is dead", async () => {
const {
acquireLogsFollowLock,
logsFollowLockPath,
releaseLogsFollowLock,
runDir,
} = await import("./workflow-state.js");
const dir = runDir(repoRoot, "logs-follow-stale-006");
fs.mkdirSync(dir, { recursive: true });
const lockPath = logsFollowLockPath(dir, "stale-thread");
fs.writeFileSync(lockPath, JSON.stringify({
pid: 999_999_999,
codexThreadId: "stale-thread",
startedAt: new Date(0).toISOString(),
}));
const acquired = acquireLogsFollowLock(dir, "stale-thread");
expect(acquired).toBe(lockPath);
expect(JSON.parse(fs.readFileSync(lockPath, "utf8")).pid).toBe(process.pid);
releaseLogsFollowLock(acquired);
expect(fs.existsSync(lockPath)).toBe(false);
});
it("keeps manual follow behavior unchanged without CODEX_THREAD_ID", async () => {
const { runDir, writeState, initWorkflowState, gitHead, readState } = await import("./workflow-state.js");
const runId = "logs-follow-manual-007";
const dir = runDir(repoRoot, runId);
fs.mkdirSync(dir, { recursive: true });
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: { version: "1", title: "Test", planMarkdown: "Test", scope: ["src/"], acceptanceCriteria: ["Ok"], verificationCommands: [["echo"]] },
executor: "claude",
maxCycles: 3,
});
state.status = "executing";
state.currentCycle = 1;
state.cycles = [{ cycleIndex: 1, startedAt: new Date().toISOString(), executorLogFile: "cycle-1-exec.log" }];
writeState(dir, state);
fs.writeFileSync(path.join(dir, "cycle-1-exec.log"), "manual viewer\n");
const previousThread = process.env.CODEX_THREAD_ID;
delete process.env.CODEX_THREAD_ID;
const { logsWorkflow } = await import("./workflow-engine.js");
const origWrite = process.stdout.write.bind(process.stdout);
process.stdout.write = () => true;
const first = logsWorkflow({ runId, follow: true, tail: 5 });
const second = logsWorkflow({ runId, follow: true, tail: 5 });
setTimeout(() => {
const curState = readState(dir)!;
curState.status = "awaiting_review";
writeState(dir, curState);
}, 100);
try {
await Promise.all([first, second]);
expect(fs.readdirSync(dir).some((name) => name.startsWith(".logs-follow-"))).toBe(false);
} finally {
process.stdout.write = origWrite;
if (previousThread === undefined) delete process.env.CODEX_THREAD_ID;
else process.env.CODEX_THREAD_ID = previousThread;
}
});
it("follow mode switches to retry log when sidecar changes and emits new content", async () => {
const { runDir, writeState, initWorkflowState, gitHead, writeProgressSidecar, readState } = await import("./workflow-state.js");
const runId = "logs-rotation-003";
const dir = runDir(repoRoot, runId);
fs.mkdirSync(dir, { recursive: true });
const state = initWorkflowState({
runId,
repoRoot,
baselineHead: gitHead(repoRoot),
plan: { version: "1", title: "Test", planMarkdown: "Test", scope: ["src/"], acceptanceCriteria: ["Ok"], verificationCommands: [["echo"]] },
executor: "claude",
maxCycles: 3,
});
state.status = "executing";
state.currentCycle = 1;
state.cycles = [{ cycleIndex: 1, startedAt: new Date().toISOString(), executorLogFile: "cycle-1-exec.log" }];
writeState(dir, state);
// Create initial log
fs.writeFileSync(path.join(dir, "cycle-1-exec.log"), "initial log\n");
// Write a sidecar pointing to the original log
writeProgressSidecar(dir, {
executor: "claude",
cycleIndex: 1,
attemptIndex: 0,
phase: "executing",
startedAt: new Date().toISOString(),
lastActivityAt: new Date().toISOString(),
logFilePath: "cycle-1-exec.log",
logBytes: 12,
activitySummary: "Running",
updatedAt: new Date().toISOString(),
});
const { logsWorkflow } = await import("./workflow-engine.js");
const chunks: string[] = [];
const origWrite = process.stdout.write.bind(process.stdout);
process.stdout.write = (chunk: unknown) => { chunks.push(String(chunk)); return true; };
const logsPromise = logsWorkflow({ runId, follow: true, tail: 5 });
setTimeout(() => {
// Create retry log and update cycle record
fs.writeFileSync(path.join(dir, "cycle-1-exec-retry-1.log"), "retry content\n你好世界\n");
const curState = readState(dir)!;
const cycle = curState.cycles[0]!;
cycle.executorLogFile = "cycle-1-exec-retry-1.log";
cycle.executorAttemptLogs = ["cycle-1-exec.log", "cycle-1-exec-retry-1.log"];
writeState(dir, curState);
// Update sidecar to point to retry log
writeProgressSidecar(dir, {
executor: "claude",
cycleIndex: 1,
attemptIndex: 1,
phase: "executing",
startedAt: new Date().toISOString(),
lastActivityAt: new Date().toISOString(),
logFilePath: "cycle-1-exec-retry-1.log",
logBytes: 30,
activitySummary: "Retrying",
updatedAt: new Date().toISOString(),
});
}, 200);
setTimeout(() => {
const curState = readState(dir)!;
curState.status = "awaiting_review";
writeState(dir, curState);
}, 600);
try {
await logsPromise;
} finally {
process.stdout.write = origWrite;
}
const output = chunks.join("");
expect(output).toContain("initial log");
expect(output).toContain("retry content");
expect(output).toContain("你好世界");
});
});