feat(HAI-016): complete Step 3 — integrate semaphore into TaskExecutor

This commit is contained in:
Dustin Byrne
2026-03-25 21:19:07 -04:00
parent 640c288832
commit d487ebee51
2 changed files with 201 additions and 24 deletions

View File

@@ -0,0 +1,167 @@
import { describe, it, expect, vi, beforeEach } from "vitest";
import { AgentSemaphore } from "./concurrency.js";
// Mock external dependencies
vi.mock("./pi.js", () => ({
createHaiAgent: vi.fn(),
}));
vi.mock("./reviewer.js", () => ({
reviewStep: vi.fn(),
}));
// Mock node modules used by executor
vi.mock("node:child_process", () => ({
execSync: vi.fn(),
}));
vi.mock("node:fs", () => ({
existsSync: vi.fn().mockReturnValue(true),
}));
import { TaskExecutor } from "./executor.js";
import { createHaiAgent } from "./pi.js";
const mockedCreateHaiAgent = vi.mocked(createHaiAgent);
function createMockStore() {
const listeners = new Map<string, Function[]>();
return {
on: vi.fn((event: string, fn: Function) => {
const existing = listeners.get(event) || [];
existing.push(fn);
listeners.set(event, existing);
}),
emit: vi.fn(),
listTasks: vi.fn().mockResolvedValue([]),
getTask: vi.fn().mockResolvedValue({
id: "HAI-001",
title: "Test",
description: "Test task",
column: "in-progress",
dependencies: [],
steps: [],
currentStep: 0,
log: [],
prompt: "# test\n## Steps\n### Step 0: Preflight\n- [ ] check",
createdAt: new Date().toISOString(),
updatedAt: new Date().toISOString(),
}),
updateTask: vi.fn().mockResolvedValue({}),
moveTask: vi.fn().mockResolvedValue({}),
logEntry: vi.fn().mockResolvedValue(undefined),
parseStepsFromPrompt: vi.fn().mockResolvedValue([]),
updateStep: vi.fn().mockResolvedValue({}),
} as any;
}
describe("TaskExecutor with semaphore", () => {
beforeEach(() => {
vi.clearAllMocks();
});
it("acquires semaphore before creating agent and releases after", async () => {
const sem = new AgentSemaphore(2);
const store = createMockStore();
const acquireSpy = vi.spyOn(sem, "acquire");
const releaseSpy = vi.spyOn(sem, "release");
mockedCreateHaiAgent.mockResolvedValue({
session: {
prompt: vi.fn().mockResolvedValue(undefined),
dispose: vi.fn(),
},
} as any);
const executor = new TaskExecutor(store, "/tmp/test", { semaphore: sem });
await executor.execute({
id: "HAI-001",
title: "Test",
description: "Test",
column: "in-progress",
dependencies: [],
steps: [],
currentStep: 0,
log: [],
createdAt: new Date().toISOString(),
updatedAt: new Date().toISOString(),
});
expect(acquireSpy).toHaveBeenCalledOnce();
expect(releaseSpy).toHaveBeenCalledOnce();
expect(sem.activeCount).toBe(0);
});
it("releases semaphore on agent error", async () => {
const sem = new AgentSemaphore(1);
const store = createMockStore();
mockedCreateHaiAgent.mockRejectedValue(new Error("agent failed"));
const onError = vi.fn();
const executor = new TaskExecutor(store, "/tmp/test", {
semaphore: sem,
onError,
});
await executor.execute({
id: "HAI-001",
title: "Test",
description: "Test",
column: "in-progress",
dependencies: [],
steps: [],
currentStep: 0,
log: [],
createdAt: new Date().toISOString(),
updatedAt: new Date().toISOString(),
});
expect(sem.activeCount).toBe(0);
expect(onError).toHaveBeenCalled();
});
it("concurrent executions respect semaphore limit", async () => {
const sem = new AgentSemaphore(1);
const store = createMockStore();
let concurrent = 0;
let maxConcurrent = 0;
mockedCreateHaiAgent.mockImplementation(async () => {
concurrent++;
maxConcurrent = Math.max(maxConcurrent, concurrent);
return {
session: {
prompt: vi.fn().mockImplementation(async () => {
await new Promise((r) => setTimeout(r, 10));
concurrent--;
}),
dispose: vi.fn(),
},
} as any;
});
const executor = new TaskExecutor(store, "/tmp/test", { semaphore: sem });
const task = (id: string) => ({
id,
title: "Test",
description: "Test",
column: "in-progress" as const,
dependencies: [],
steps: [],
currentStep: 0,
log: [],
createdAt: new Date().toISOString(),
updatedAt: new Date().toISOString(),
});
await Promise.all([
executor.execute(task("HAI-001")),
executor.execute(task("HAI-002")),
executor.execute(task("HAI-003")),
]);
expect(maxConcurrent).toBe(1);
expect(sem.activeCount).toBe(0);
});
});

View File

@@ -6,6 +6,7 @@ import { Type } from "@mariozechner/pi-ai";
import { createHaiAgent } from "./pi.js";
import { reviewStep } from "./reviewer.js";
import type { ToolDefinition } from "@mariozechner/pi-coding-agent";
import type { AgentSemaphore } from "./concurrency.js";
const STEP_STATUSES: StepStatus[] = ["pending", "in-progress", "done", "skipped"];
@@ -79,6 +80,7 @@ echo "done" > .DONE
\`\`\``;
export interface TaskExecutorOptions {
semaphore?: AgentSemaphore;
onStart?: (task: Task, worktreePath: string) => void;
onComplete?: (task: Task) => void;
onError?: (task: Task, error: Error) => void;
@@ -152,33 +154,41 @@ export class TaskExecutor {
this.createReviewStepTool(task.id, worktreePath, detail.prompt),
];
const { session } = await createHaiAgent({
cwd: worktreePath,
systemPrompt: EXECUTOR_SYSTEM_PROMPT,
tools: "coding",
customTools,
onText: (delta) => this.options.onAgentText?.(task.id, delta),
onToolStart: (name) => this.options.onAgentTool?.(task.id, name),
});
const agentWork = async () => {
const { session } = await createHaiAgent({
cwd: worktreePath,
systemPrompt: EXECUTOR_SYSTEM_PROMPT,
tools: "coding",
customTools,
onText: (delta) => this.options.onAgentText?.(task.id, delta),
onToolStart: (name) => this.options.onAgentTool?.(task.id, name),
});
try {
const agentPrompt = buildExecutionPrompt(detail);
await session.prompt(agentPrompt);
try {
const agentPrompt = buildExecutionPrompt(detail);
await session.prompt(agentPrompt);
const doneCwd = join(worktreePath, ".DONE");
if (existsSync(doneCwd)) {
await this.store.logEntry(task.id, "Execution complete — .DONE created");
await this.store.moveTask(task.id, "in-review");
console.log(`[executor] ✓ ${task.id} completed → in-review`);
this.options.onComplete?.(task);
} else {
await this.store.logEntry(task.id, "Agent finished without .DONE — moved to in-review for inspection");
await this.store.moveTask(task.id, "in-review");
console.log(`[executor] ⚠ ${task.id} agent finished without .DONE → in-review`);
this.options.onComplete?.(task);
const doneCwd = join(worktreePath, ".DONE");
if (existsSync(doneCwd)) {
await this.store.logEntry(task.id, "Execution complete — .DONE created");
await this.store.moveTask(task.id, "in-review");
console.log(`[executor] ✓ ${task.id} completed → in-review`);
this.options.onComplete?.(task);
} else {
await this.store.logEntry(task.id, "Agent finished without .DONE — moved to in-review for inspection");
await this.store.moveTask(task.id, "in-review");
console.log(`[executor] ⚠ ${task.id} agent finished without .DONE → in-review`);
this.options.onComplete?.(task);
}
} finally {
session.dispose();
}
} finally {
session.dispose();
};
if (this.options.semaphore) {
await this.options.semaphore.run(agentWork);
} else {
await agentWork();
}
} catch (err: any) {
console.error(`[executor] ✗ ${task.id} execution failed:`, err.message);