feat(FN-1661): merge fusion/fn-1661 (auto-resolved)
- feat(FN-1661): complete Step 4 — document task-worker health fix - feat(FN-1661): complete Step 3 — verify task-worker agent health fix
This commit is contained in:
@@ -2905,6 +2905,26 @@ describe("HeartbeatTriggerScheduler", () => {
|
||||
expect(callback).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("skips trigger when agent heartbeat is disabled", async () => {
|
||||
const agent: import("@fusion/core").Agent = {
|
||||
id: "agent-test",
|
||||
name: "executor-FN-1661",
|
||||
role: "executor",
|
||||
state: "active",
|
||||
taskId: "FN-1661",
|
||||
metadata: {},
|
||||
runtimeConfig: { enabled: false },
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
};
|
||||
eventStore.emit("agent:assigned", agent, "FN-1661");
|
||||
|
||||
await new Promise((resolve) => setTimeout(resolve, 10));
|
||||
|
||||
expect(callback).not.toHaveBeenCalled();
|
||||
expect(eventStore.getActiveHeartbeatRun).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("skips trigger when agent has active run", async () => {
|
||||
(eventStore.getActiveHeartbeatRun as ReturnType<typeof vi.fn>).mockResolvedValue({
|
||||
id: "run-active",
|
||||
|
||||
@@ -20,15 +20,13 @@
|
||||
import type { AgentStore, AgentHeartbeatRun, HeartbeatInvocationSource, AgentHeartbeatConfig, AgentBudgetStatus, Message, MessageStore, TaskStore, TaskDetail, AgentRole, Agent, InboxTask, BlockedStateSnapshot, RunMutationContext } from "@fusion/core";
|
||||
import type { ToolDefinition } from "@mariozechner/pi-coding-agent";
|
||||
import { Type, type Static } from "@mariozechner/pi-ai";
|
||||
import { createTaskCreateTool, createTaskLogTool, createTaskLogToolWithContext, taskCreateParams } from "./agent-tools.js";
|
||||
import { createTaskCreateTool, createTaskLogToolWithContext, taskCreateParams } from "./agent-tools.js";
|
||||
import { AgentLogger } from "./agent-logger.js";
|
||||
import { heartbeatLog } from "./logger.js";
|
||||
import { createRunAuditor, type EngineRunContext } from "./run-audit.js";
|
||||
|
||||
// Lazy import for pi — avoids pulling the pi SDK into the module graph
|
||||
// when heartbeat execution isn't needed.
|
||||
type CreateKbAgentFn = (options: import("./pi.js").AgentOptions) => Promise<import("./pi.js").AgentResult>;
|
||||
type PromptWithFallbackFn = (session: import("@mariozechner/pi-coding-agent").AgentSession, prompt: string) => Promise<void>;
|
||||
|
||||
/** Resolved per-agent heartbeat config after validation and fallback */
|
||||
interface ResolvedHeartbeatConfig {
|
||||
@@ -1081,8 +1079,8 @@ export class HeartbeatMonitor {
|
||||
const baseCreateTool = createTaskCreateTool(taskStore);
|
||||
const trackedCreateTool: ToolDefinition = {
|
||||
...baseCreateTool,
|
||||
execute: async (id: string, params: Static<typeof taskCreateParams>, _signal?: unknown, _onUpdate?: unknown, _ctx?: unknown) => {
|
||||
const result = await baseCreateTool.execute(id, params, undefined as any, undefined as any, undefined as any);
|
||||
execute: async (id: string, params: Static<typeof taskCreateParams>, signal, onUpdate, ctx) => {
|
||||
const result = await baseCreateTool.execute(id, params, signal, onUpdate, ctx);
|
||||
|
||||
// Extract created task ID from the response text ("Created FN-XXX: ...")
|
||||
const firstContent = result.content[0];
|
||||
@@ -1404,6 +1402,11 @@ export class HeartbeatTriggerScheduler {
|
||||
if (!this.running) return;
|
||||
|
||||
try {
|
||||
if (agent.runtimeConfig?.enabled === false) {
|
||||
heartbeatLog.log(`Assignment trigger skipped for ${agent.id} (heartbeat disabled)`);
|
||||
return;
|
||||
}
|
||||
|
||||
// Guard: skip if agent already has an active run
|
||||
const activeRun = await this.store.getActiveHeartbeatRun(agent.id);
|
||||
if (activeRun) {
|
||||
|
||||
@@ -1,9 +1,8 @@
|
||||
import { describe, it, expect, vi, beforeEach, afterEach } from "vitest";
|
||||
import { EventEmitter } from "node:events";
|
||||
import { mkdtempSync, rmSync } from "node:fs";
|
||||
import { join } from "node:path";
|
||||
import { randomUUID } from "node:crypto";
|
||||
import type { Task, TaskStore, CentralCore } from "@fusion/core";
|
||||
import type { Task, TaskStore, CentralCore, AgentStore, Agent } from "@fusion/core";
|
||||
import { InProcessRuntime } from "./in-process-runtime.js";
|
||||
import type { ProjectRuntimeConfig } from "../project-runtime.js";
|
||||
|
||||
@@ -135,6 +134,21 @@ vi.mock("../executor.js", async () => {
|
||||
};
|
||||
});
|
||||
|
||||
type RuntimeInternals = {
|
||||
agentStore?: AgentStore;
|
||||
stuckTaskDetector?: unknown;
|
||||
};
|
||||
|
||||
function getRuntimeInternals(runtime: InProcessRuntime): RuntimeInternals {
|
||||
return runtime as unknown as RuntimeInternals;
|
||||
}
|
||||
|
||||
function getAgentStore(runtime: InProcessRuntime): AgentStore {
|
||||
const store = getRuntimeInternals(runtime).agentStore;
|
||||
expect(store).toBeDefined();
|
||||
return store!;
|
||||
}
|
||||
|
||||
describe("InProcessRuntime", () => {
|
||||
let runtime: InProcessRuntime;
|
||||
let mockCentralCore: CentralCore;
|
||||
@@ -222,7 +236,7 @@ describe("InProcessRuntime", () => {
|
||||
stuckTaskDetector: expect.any(Object),
|
||||
}),
|
||||
);
|
||||
expect((runtime as any).stuckTaskDetector).toBeDefined();
|
||||
expect(getRuntimeInternals(runtime).stuckTaskDetector).toBeDefined();
|
||||
});
|
||||
|
||||
it("should transition to 'stopped' after stop", async () => {
|
||||
@@ -407,8 +421,7 @@ describe("InProcessRuntime", () => {
|
||||
await runtime.start();
|
||||
|
||||
// Create an agent with heartbeat config
|
||||
const store = (runtime as any).agentStore;
|
||||
expect(store).toBeDefined();
|
||||
const store = getAgentStore(runtime);
|
||||
|
||||
const createdAgent = await store.createAgent({
|
||||
name: "Configured Agent",
|
||||
@@ -434,12 +447,13 @@ describe("InProcessRuntime", () => {
|
||||
|
||||
const monitor = runtime.getHeartbeatMonitor();
|
||||
expect(monitor).toBeDefined();
|
||||
const heartbeatMonitor = monitor!;
|
||||
const executeResult = { id: "run-test" } as Awaited<ReturnType<typeof heartbeatMonitor.executeHeartbeat>>;
|
||||
const executeSpy = vi
|
||||
.spyOn(monitor!, "executeHeartbeat")
|
||||
.mockResolvedValue({ id: "run-test" } as any);
|
||||
.spyOn(heartbeatMonitor, "executeHeartbeat")
|
||||
.mockResolvedValue(executeResult);
|
||||
|
||||
const store = (runtime as any).agentStore;
|
||||
expect(store).toBeDefined();
|
||||
const store = getAgentStore(runtime);
|
||||
|
||||
const agent = await store.createAgent({
|
||||
name: "Assignable",
|
||||
@@ -462,6 +476,72 @@ describe("InProcessRuntime", () => {
|
||||
);
|
||||
});
|
||||
}, 30000);
|
||||
|
||||
it("creates runtime task-worker agents with disabled heartbeat metadata and running state", async () => {
|
||||
await runtime.start();
|
||||
|
||||
const store = getAgentStore(runtime);
|
||||
|
||||
const assignTaskSpy = vi.spyOn(store, "assignTask");
|
||||
const updateStateSpy = vi.spyOn(store, "updateAgentState");
|
||||
const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as {
|
||||
onStart?: (task: Task, worktreePath: string) => void;
|
||||
};
|
||||
expect(executorOptions.onStart).toBeTypeOf("function");
|
||||
|
||||
executorOptions.onStart?.({ id: "FN-1661" } as Task, join(testDir, "worktree-FN-1661"));
|
||||
|
||||
await vi.waitFor(async () => {
|
||||
const agents = await store.listAgents();
|
||||
expect(agents).toHaveLength(1);
|
||||
expect(agents[0]).toMatchObject({
|
||||
name: "executor-FN-1661",
|
||||
role: "executor",
|
||||
state: "running",
|
||||
taskId: "FN-1661",
|
||||
metadata: {
|
||||
agentKind: "task-worker",
|
||||
taskWorker: true,
|
||||
managedBy: "task-executor",
|
||||
},
|
||||
runtimeConfig: {
|
||||
enabled: false,
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
expect(assignTaskSpy).toHaveBeenCalledWith(expect.any(String), "FN-1661");
|
||||
expect(updateStateSpy).toHaveBeenNthCalledWith(1, expect.any(String), "active");
|
||||
expect(updateStateSpy).toHaveBeenNthCalledWith(2, expect.any(String), "running");
|
||||
expect(assignTaskSpy.mock.invocationCallOrder[0]).toBeLessThan(updateStateSpy.mock.invocationCallOrder[0]);
|
||||
}, 30000);
|
||||
|
||||
it("does not wake executeHeartbeat for runtime task-worker assignment events", async () => {
|
||||
await runtime.start();
|
||||
|
||||
const monitor = runtime.getHeartbeatMonitor();
|
||||
expect(monitor).toBeDefined();
|
||||
const heartbeatMonitor = monitor!;
|
||||
const executeResult = { id: "run-task-worker" } as Awaited<ReturnType<typeof heartbeatMonitor.executeHeartbeat>>;
|
||||
const executeSpy = vi
|
||||
.spyOn(heartbeatMonitor, "executeHeartbeat")
|
||||
.mockResolvedValue(executeResult);
|
||||
|
||||
const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as {
|
||||
onStart?: (task: Task, worktreePath: string) => void;
|
||||
};
|
||||
executorOptions.onStart?.({ id: "FN-2001" } as Task, join(testDir, "worktree-FN-2001"));
|
||||
|
||||
const store = getAgentStore(runtime);
|
||||
|
||||
await vi.waitFor(async () => {
|
||||
const agents = await store.listAgents();
|
||||
expect(agents.some((agent: Agent) => agent.name === "executor-FN-2001")).toBe(true);
|
||||
});
|
||||
|
||||
await new Promise((resolve) => setTimeout(resolve, 25));
|
||||
expect(executeSpy).not.toHaveBeenCalled();
|
||||
}, 30000);
|
||||
});
|
||||
|
||||
describe("configuration", () => {
|
||||
|
||||
@@ -193,7 +193,7 @@ export class InProcessRuntime
|
||||
missionStore,
|
||||
missionAutopilot: missionAutopilot
|
||||
? {
|
||||
notifyValidationComplete: async (featureId: string, _status: "passed" | "failed" | "blocked" | "error") => {
|
||||
notifyValidationComplete: async (featureId: string) => {
|
||||
// Pass the feature's linked taskId to handleTaskCompletion, not the featureId
|
||||
const feature = missionStore.getFeature(featureId);
|
||||
if (feature?.taskId) {
|
||||
@@ -254,15 +254,26 @@ export class InProcessRuntime
|
||||
onStart: (task, worktreePath) => {
|
||||
this.recordActivity();
|
||||
runtimeLog.log(`Started executing task ${task.id} in ${worktreePath}`);
|
||||
// Create agent in AgentStore for lifecycle tracking
|
||||
// Create a runtime-managed task worker agent for lifecycle tracking.
|
||||
// These workers are not heartbeat-managed dashboard agents, so mark them
|
||||
// explicitly and disable heartbeat triggers/timers.
|
||||
if (this.agentStore) {
|
||||
this.agentStore.createAgent({
|
||||
name: `executor-${task.id}`,
|
||||
role: "executor",
|
||||
metadata: {
|
||||
agentKind: "task-worker",
|
||||
taskWorker: true,
|
||||
managedBy: "task-executor",
|
||||
},
|
||||
runtimeConfig: {
|
||||
enabled: false,
|
||||
},
|
||||
}).then(async (agent: { id: string }) => {
|
||||
this.taskAgentMap.set(task.id, agent.id);
|
||||
await this.agentStore!.assignTask(agent.id, task.id);
|
||||
await this.agentStore!.updateAgentState(agent.id, "active");
|
||||
await this.agentStore!.updateAgentState(agent.id, "running");
|
||||
}).catch((err: unknown) => {
|
||||
runtimeLog.warn(`Failed to create agent for task ${task.id}:`, err);
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user