feat(engine): gate executor and heartbeat on allowParallelExecution
When a permanent agent has allowParallelExecution=false, TaskExecutor.execute() defers if the agent has an active heartbeat run, and HeartbeatScheduler defers a heartbeat if the agent's bound task has an active executor session. Each side re-dispatches the other's deferred work on completion via resumeTaskForAgent and the in-process runtime's onRunCompleted hook. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
10
.changeset/agent-allow-parallel-execution.md
Normal file
10
.changeset/agent-allow-parallel-execution.md
Normal file
@@ -0,0 +1,10 @@
|
|||||||
|
---
|
||||||
|
"@fusion/core": minor
|
||||||
|
"@fusion/engine": minor
|
||||||
|
"@runfusion/fusion": minor
|
||||||
|
"@fusion/dashboard": minor
|
||||||
|
---
|
||||||
|
|
||||||
|
Per-agent setting `allowParallelExecution` (default true, permanent agents only): when disabled, an agent's heartbeat runs and task executor sessions serialize — a heartbeat will not start while the agent's bound task has an active executor session, and an executor session will not start while the agent has an active heartbeat run.
|
||||||
|
|
||||||
|
Field added to `AgentHeartbeatConfig` in `@fusion/core`. UI toggle surfaces in the agent's Heartbeat Settings tab alongside `runMissedHeartbeatOnStartup`. Engine gating: `TaskExecutor.execute()` defers when the assigned agent has an active heartbeat run; `HeartbeatMonitor` defers a heartbeat when the agent's bound task has an active executor session. After either side completes, the deferred work is re-dispatched.
|
||||||
@@ -14138,3 +14138,103 @@ describe("Executor verification gate (FN-3345)", () => {
|
|||||||
);
|
);
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
// allowParallelExecution gate
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
describe("allowParallelExecution heartbeat gate", () => {
|
||||||
|
const TASK_BASE: Omit<Task, "id"> = {
|
||||||
|
title: "Gated task",
|
||||||
|
description: "Test task",
|
||||||
|
column: "in-progress",
|
||||||
|
dependencies: [],
|
||||||
|
steps: [],
|
||||||
|
currentStep: 0,
|
||||||
|
log: [],
|
||||||
|
createdAt: new Date().toISOString(),
|
||||||
|
updatedAt: new Date().toISOString(),
|
||||||
|
};
|
||||||
|
|
||||||
|
function makeAgentStore(opts: {
|
||||||
|
ephemeral: boolean;
|
||||||
|
allowParallelExecution?: boolean;
|
||||||
|
hasActiveRun: boolean;
|
||||||
|
}) {
|
||||||
|
const agent = {
|
||||||
|
id: "agent-perm-1",
|
||||||
|
name: "Permanent Agent",
|
||||||
|
role: "executor",
|
||||||
|
state: "running",
|
||||||
|
metadata: opts.ephemeral ? { agentKind: "task-worker" } : {},
|
||||||
|
runtimeConfig: opts.allowParallelExecution !== undefined
|
||||||
|
? { allowParallelExecution: opts.allowParallelExecution }
|
||||||
|
: {},
|
||||||
|
};
|
||||||
|
return {
|
||||||
|
getAgent: vi.fn().mockResolvedValue(agent),
|
||||||
|
getActiveHeartbeatRun: vi.fn().mockResolvedValue(
|
||||||
|
opts.hasActiveRun ? { id: "run-1", status: "active" } : null,
|
||||||
|
),
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
beforeEach(() => {
|
||||||
|
vi.clearAllMocks();
|
||||||
|
});
|
||||||
|
|
||||||
|
it.each([
|
||||||
|
{
|
||||||
|
label: "permanent agent, allowParallelExecution=false, active heartbeat run → skipped",
|
||||||
|
ephemeral: false,
|
||||||
|
allowParallelExecution: false as boolean | undefined,
|
||||||
|
hasActiveRun: true,
|
||||||
|
expectExecute: false,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
label: "permanent agent, allowParallelExecution=true, active heartbeat run → proceeds",
|
||||||
|
ephemeral: false,
|
||||||
|
allowParallelExecution: true as boolean | undefined,
|
||||||
|
hasActiveRun: true,
|
||||||
|
expectExecute: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
label: "permanent agent, allowParallelExecution=false, no heartbeat run → proceeds",
|
||||||
|
ephemeral: false,
|
||||||
|
allowParallelExecution: false as boolean | undefined,
|
||||||
|
hasActiveRun: false,
|
||||||
|
expectExecute: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
label: "ephemeral agent, allowParallelExecution=false, active heartbeat run → proceeds (flag ignored)",
|
||||||
|
ephemeral: true,
|
||||||
|
allowParallelExecution: false as boolean | undefined,
|
||||||
|
hasActiveRun: true,
|
||||||
|
expectExecute: true,
|
||||||
|
},
|
||||||
|
])("$label", async ({ ephemeral, allowParallelExecution, hasActiveRun, expectExecute }) => {
|
||||||
|
const agentStore = makeAgentStore({ ephemeral, allowParallelExecution, hasActiveRun });
|
||||||
|
const store = createMockStore();
|
||||||
|
|
||||||
|
mockedCreateFnAgent.mockResolvedValue({
|
||||||
|
session: {
|
||||||
|
prompt: vi.fn().mockResolvedValue(undefined),
|
||||||
|
dispose: vi.fn(),
|
||||||
|
},
|
||||||
|
} as any);
|
||||||
|
|
||||||
|
const executor = new TaskExecutor(store, "/tmp/test", { agentStore: agentStore as any });
|
||||||
|
|
||||||
|
await executor.execute({
|
||||||
|
...TASK_BASE,
|
||||||
|
id: "FN-GATE-1",
|
||||||
|
assignedAgentId: "agent-perm-1",
|
||||||
|
});
|
||||||
|
|
||||||
|
if (expectExecute) {
|
||||||
|
expect(mockedCreateFnAgent).toHaveBeenCalled();
|
||||||
|
} else {
|
||||||
|
expect(mockedCreateFnAgent).not.toHaveBeenCalled();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|||||||
@@ -7,7 +7,6 @@ import { tmpdir } from "node:os";
|
|||||||
import {
|
import {
|
||||||
HeartbeatMonitor,
|
HeartbeatMonitor,
|
||||||
HeartbeatTriggerScheduler,
|
HeartbeatTriggerScheduler,
|
||||||
isBlockedStateDuplicate,
|
|
||||||
type AgentSession,
|
type AgentSession,
|
||||||
type HeartbeatExecutionOptions,
|
type HeartbeatExecutionOptions,
|
||||||
HEARTBEAT_SYSTEM_PROMPT,
|
HEARTBEAT_SYSTEM_PROMPT,
|
||||||
|
|||||||
@@ -7,7 +7,6 @@ import { tmpdir } from "node:os";
|
|||||||
import {
|
import {
|
||||||
HeartbeatMonitor,
|
HeartbeatMonitor,
|
||||||
HeartbeatTriggerScheduler,
|
HeartbeatTriggerScheduler,
|
||||||
isBlockedStateDuplicate,
|
|
||||||
type AgentSession,
|
type AgentSession,
|
||||||
type HeartbeatExecutionOptions,
|
type HeartbeatExecutionOptions,
|
||||||
HEARTBEAT_SYSTEM_PROMPT,
|
HEARTBEAT_SYSTEM_PROMPT,
|
||||||
@@ -74,31 +73,6 @@ describe("constructor", () => {
|
|||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
describe("isBlockedStateDuplicate", () => {
|
|
||||||
it("returns true when blockedBy and contextHash match", () => {
|
|
||||||
expect(
|
|
||||||
isBlockedStateDuplicate(
|
|
||||||
{ taskId: "FN-1", blockedBy: "FN-0", recordedAt: "2026-01-01T00:00:00.000Z", contextHash: "abc" },
|
|
||||||
{ taskId: "FN-1", blockedBy: "FN-0", recordedAt: "2026-01-02T00:00:00.000Z", contextHash: "abc" },
|
|
||||||
),
|
|
||||||
).toBe(true);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("returns false when blockedBy differs or contextHash differs", () => {
|
|
||||||
expect(
|
|
||||||
isBlockedStateDuplicate(
|
|
||||||
{ taskId: "FN-1", blockedBy: "FN-0", recordedAt: "2026-01-01T00:00:00.000Z", contextHash: "abc" },
|
|
||||||
{ taskId: "FN-1", blockedBy: "FN-2", recordedAt: "2026-01-02T00:00:00.000Z", contextHash: "abc" },
|
|
||||||
),
|
|
||||||
).toBe(false);
|
|
||||||
expect(
|
|
||||||
isBlockedStateDuplicate(
|
|
||||||
{ taskId: "FN-1", blockedBy: "FN-0", recordedAt: "2026-01-01T00:00:00.000Z", contextHash: "abc" },
|
|
||||||
{ taskId: "FN-1", blockedBy: "FN-0", recordedAt: "2026-01-02T00:00:00.000Z", contextHash: "xyz" },
|
|
||||||
),
|
|
||||||
).toBe(false);
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
describe("start", () => {
|
describe("start", () => {
|
||||||
it("initiates polling interval", () => {
|
it("initiates polling interval", () => {
|
||||||
|
|||||||
@@ -7,7 +7,6 @@ import { tmpdir } from "node:os";
|
|||||||
import {
|
import {
|
||||||
HeartbeatMonitor,
|
HeartbeatMonitor,
|
||||||
HeartbeatTriggerScheduler,
|
HeartbeatTriggerScheduler,
|
||||||
isBlockedStateDuplicate,
|
|
||||||
type AgentSession,
|
type AgentSession,
|
||||||
type HeartbeatExecutionOptions,
|
type HeartbeatExecutionOptions,
|
||||||
HEARTBEAT_SYSTEM_PROMPT,
|
HEARTBEAT_SYSTEM_PROMPT,
|
||||||
@@ -1187,5 +1186,106 @@ describe("HeartbeatTriggerScheduler", () => {
|
|||||||
);
|
);
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
describe("allowParallelExecution gate", () => {
|
||||||
|
function makeAgentWithConfig(overrides: Record<string, unknown> = {}) {
|
||||||
|
return {
|
||||||
|
id: "agent-par",
|
||||||
|
name: "Parallel Agent",
|
||||||
|
role: "executor",
|
||||||
|
state: "active",
|
||||||
|
taskId: "FN-TASK-1",
|
||||||
|
createdAt: "2026-01-01T00:00:00.000Z",
|
||||||
|
updatedAt: "2026-01-01T00:00:00.000Z",
|
||||||
|
metadata: {},
|
||||||
|
runtimeConfig: overrides,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
it("timer tick skips when allowParallelExecution=false and task is executing", async () => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
const isTaskExecuting = vi.fn().mockReturnValue(true);
|
||||||
|
|
||||||
|
const parallelStore = {
|
||||||
|
getAgent: vi.fn().mockResolvedValue(makeAgentWithConfig({ allowParallelExecution: false })),
|
||||||
|
getActiveHeartbeatRun: vi.fn().mockResolvedValue(null),
|
||||||
|
getBudgetStatus: vi.fn().mockResolvedValue(createBudgetStatus()),
|
||||||
|
on: vi.fn(),
|
||||||
|
off: vi.fn(),
|
||||||
|
} as unknown as AgentStore;
|
||||||
|
|
||||||
|
scheduler = new HeartbeatTriggerScheduler(parallelStore, callback, undefined, { isTaskExecuting });
|
||||||
|
scheduler.start();
|
||||||
|
scheduler.registerAgent("agent-par", { heartbeatIntervalMs: 1000 });
|
||||||
|
|
||||||
|
await vi.advanceTimersByTimeAsync(1100);
|
||||||
|
|
||||||
|
expect(callback).not.toHaveBeenCalled();
|
||||||
|
expect(isTaskExecuting).toHaveBeenCalledWith("FN-TASK-1");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("timer tick fires when allowParallelExecution=false and task is NOT executing", async () => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
const isTaskExecuting = vi.fn().mockReturnValue(false);
|
||||||
|
|
||||||
|
const parallelStore = {
|
||||||
|
getAgent: vi.fn().mockResolvedValue(makeAgentWithConfig({ allowParallelExecution: false })),
|
||||||
|
getActiveHeartbeatRun: vi.fn().mockResolvedValue(null),
|
||||||
|
getBudgetStatus: vi.fn().mockResolvedValue(createBudgetStatus()),
|
||||||
|
on: vi.fn(),
|
||||||
|
off: vi.fn(),
|
||||||
|
} as unknown as AgentStore;
|
||||||
|
|
||||||
|
scheduler = new HeartbeatTriggerScheduler(parallelStore, callback, undefined, { isTaskExecuting });
|
||||||
|
scheduler.start();
|
||||||
|
scheduler.registerAgent("agent-par", { heartbeatIntervalMs: 1000 });
|
||||||
|
|
||||||
|
await vi.advanceTimersByTimeAsync(1100);
|
||||||
|
|
||||||
|
expect(callback).toHaveBeenCalledOnce();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("timer tick fires when allowParallelExecution=true even while task is executing", async () => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
const isTaskExecuting = vi.fn().mockReturnValue(true);
|
||||||
|
|
||||||
|
const parallelStore = {
|
||||||
|
getAgent: vi.fn().mockResolvedValue(makeAgentWithConfig({ allowParallelExecution: true })),
|
||||||
|
getActiveHeartbeatRun: vi.fn().mockResolvedValue(null),
|
||||||
|
getBudgetStatus: vi.fn().mockResolvedValue(createBudgetStatus()),
|
||||||
|
on: vi.fn(),
|
||||||
|
off: vi.fn(),
|
||||||
|
} as unknown as AgentStore;
|
||||||
|
|
||||||
|
scheduler = new HeartbeatTriggerScheduler(parallelStore, callback, undefined, { isTaskExecuting });
|
||||||
|
scheduler.start();
|
||||||
|
scheduler.registerAgent("agent-par", { heartbeatIntervalMs: 1000 });
|
||||||
|
|
||||||
|
await vi.advanceTimersByTimeAsync(1100);
|
||||||
|
|
||||||
|
expect(callback).toHaveBeenCalledOnce();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("timer tick fires when allowParallelExecution is unset (default) even while task is executing", async () => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
const isTaskExecuting = vi.fn().mockReturnValue(true);
|
||||||
|
|
||||||
|
const parallelStore = {
|
||||||
|
getAgent: vi.fn().mockResolvedValue(makeAgentWithConfig({})),
|
||||||
|
getActiveHeartbeatRun: vi.fn().mockResolvedValue(null),
|
||||||
|
getBudgetStatus: vi.fn().mockResolvedValue(createBudgetStatus()),
|
||||||
|
on: vi.fn(),
|
||||||
|
off: vi.fn(),
|
||||||
|
} as unknown as AgentStore;
|
||||||
|
|
||||||
|
scheduler = new HeartbeatTriggerScheduler(parallelStore, callback, undefined, { isTaskExecuting });
|
||||||
|
scheduler.start();
|
||||||
|
scheduler.registerAgent("agent-par", { heartbeatIntervalMs: 1000 });
|
||||||
|
|
||||||
|
await vi.advanceTimersByTimeAsync(1100);
|
||||||
|
|
||||||
|
expect(callback).toHaveBeenCalledOnce();
|
||||||
|
});
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
@@ -7,7 +7,6 @@ import { tmpdir } from "node:os";
|
|||||||
import {
|
import {
|
||||||
HeartbeatMonitor,
|
HeartbeatMonitor,
|
||||||
HeartbeatTriggerScheduler,
|
HeartbeatTriggerScheduler,
|
||||||
isBlockedStateDuplicate,
|
|
||||||
type AgentSession,
|
type AgentSession,
|
||||||
type HeartbeatExecutionOptions,
|
type HeartbeatExecutionOptions,
|
||||||
HEARTBEAT_SYSTEM_PROMPT,
|
HEARTBEAT_SYSTEM_PROMPT,
|
||||||
|
|||||||
@@ -7,7 +7,6 @@ import { tmpdir } from "node:os";
|
|||||||
import {
|
import {
|
||||||
HeartbeatMonitor,
|
HeartbeatMonitor,
|
||||||
HeartbeatTriggerScheduler,
|
HeartbeatTriggerScheduler,
|
||||||
isBlockedStateDuplicate,
|
|
||||||
type AgentSession,
|
type AgentSession,
|
||||||
type HeartbeatExecutionOptions,
|
type HeartbeatExecutionOptions,
|
||||||
HEARTBEAT_SYSTEM_PROMPT,
|
HEARTBEAT_SYSTEM_PROMPT,
|
||||||
|
|||||||
@@ -17,7 +17,7 @@
|
|||||||
* - onTerminated: Called when an unresponsive agent is terminated
|
* - onTerminated: Called when an unresponsive agent is terminated
|
||||||
*/
|
*/
|
||||||
|
|
||||||
import type { AgentStore, AgentHeartbeatRun, HeartbeatInvocationSource, AgentHeartbeatConfig, AgentBudgetStatus, Message, MessageStore, TaskStore, TaskDetail, AgentRole, Agent, InboxTask, BlockedStateSnapshot, RunMutationContext, Settings, AgentConfigRevision, ReflectionStore } from "@fusion/core";
|
import type { AgentStore, AgentHeartbeatRun, HeartbeatInvocationSource, AgentHeartbeatConfig, AgentBudgetStatus, Message, MessageStore, TaskStore, TaskDetail, AgentRole, Agent, InboxTask, RunMutationContext, Settings, AgentConfigRevision, ReflectionStore } from "@fusion/core";
|
||||||
import { buildExecutionMemoryInstructions, isEphemeralAgent, hasAgentIdentity } from "@fusion/core";
|
import { buildExecutionMemoryInstructions, isEphemeralAgent, hasAgentIdentity } from "@fusion/core";
|
||||||
import type { ToolDefinition } from "@mariozechner/pi-coding-agent";
|
import type { ToolDefinition } from "@mariozechner/pi-coding-agent";
|
||||||
import { Type, type Static } from "@mariozechner/pi-ai";
|
import { Type, type Static } from "@mariozechner/pi-ai";
|
||||||
@@ -189,11 +189,6 @@ function taskRelevanceScore(agent: Agent, task: TaskDetail): number {
|
|||||||
return score;
|
return score;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Compare blocked-state snapshots to decide whether blocked messaging is duplicate noise. */
|
|
||||||
export function isBlockedStateDuplicate(current: BlockedStateSnapshot, previous: BlockedStateSnapshot): boolean {
|
|
||||||
return current.blockedBy === previous.blockedBy && current.contextHash === previous.contextHash;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* System prompt for heartbeat agent sessions.
|
* System prompt for heartbeat agent sessions.
|
||||||
* Instructs the agent to perform a single-pass check on its assigned task
|
* Instructs the agent to perform a single-pass check on its assigned task
|
||||||
|
|||||||
@@ -5,10 +5,11 @@ const execAsync = promisify(exec);
|
|||||||
import { isAbsolute, join, relative, resolve as resolvePath } from "node:path";
|
import { isAbsolute, join, relative, resolve as resolvePath } from "node:path";
|
||||||
import { existsSync } from "node:fs";
|
import { existsSync } from "node:fs";
|
||||||
import { readFile, writeFile } from "node:fs/promises";
|
import { readFile, writeFile } from "node:fs/promises";
|
||||||
import type { TaskStore, Task, TaskDetail, TaskTokenUsage, 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, AgentHeartbeatConfig } from "@fusion/core";
|
||||||
import {
|
import {
|
||||||
buildExecutionMemoryInstructions,
|
buildExecutionMemoryInstructions,
|
||||||
getTaskMergeBlocker,
|
getTaskMergeBlocker,
|
||||||
|
isEphemeralAgent,
|
||||||
resolveAgentPrompt,
|
resolveAgentPrompt,
|
||||||
resolveProjectDefaultModel,
|
resolveProjectDefaultModel,
|
||||||
type RunCommandResult,
|
type RunCommandResult,
|
||||||
@@ -1839,6 +1840,50 @@ export class TaskExecutor {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Returns true when execute() should be deferred because the agent bound to
|
||||||
|
* this task has an active heartbeat run and allowParallelExecution=false.
|
||||||
|
*
|
||||||
|
* Only applies to permanent (non-ephemeral) agents. Always returns false
|
||||||
|
* when agentStore is unavailable or the agent cannot be resolved.
|
||||||
|
*/
|
||||||
|
private async shouldDeferForHeartbeat(agentId: string): Promise<boolean> {
|
||||||
|
if (!this.options.agentStore) return false;
|
||||||
|
const agent = await this.options.agentStore.getAgent(agentId).catch(() => null);
|
||||||
|
if (!agent) return false;
|
||||||
|
if (isEphemeralAgent(agent)) return false;
|
||||||
|
const rc = (agent.runtimeConfig ?? {}) as AgentHeartbeatConfig;
|
||||||
|
if (rc.allowParallelExecution !== false) return false;
|
||||||
|
const activeRun = await this.options.agentStore.getActiveHeartbeatRun(agentId).catch(() => null);
|
||||||
|
return activeRun !== null;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Re-dispatch execute() for any unstarted in-progress task belonging to the
|
||||||
|
* given agent. Called after a heartbeat run completes to unblock tasks that
|
||||||
|
* were deferred by the allowParallelExecution=false gate.
|
||||||
|
*/
|
||||||
|
async resumeTaskForAgent(agentId: string): Promise<void> {
|
||||||
|
const settings = await this.store.getSettings();
|
||||||
|
if (settings.globalPause || settings.enginePaused) return;
|
||||||
|
const tasks = await this.store.listTasks({ slim: true, column: "in-progress" });
|
||||||
|
for (const task of tasks) {
|
||||||
|
if (
|
||||||
|
task.assignedAgentId === agentId
|
||||||
|
&& !task.paused
|
||||||
|
&& !this.executing.has(task.id)
|
||||||
|
&& !this.activeSessions.has(task.id)
|
||||||
|
&& !this.activeStepExecutors.has(task.id)
|
||||||
|
&& !this.activeWorkflowStepSessions.has(task.id)
|
||||||
|
) {
|
||||||
|
executorLog.log(`${task.id}: re-dispatching execute() after heartbeat completion for agent ${agentId}`);
|
||||||
|
this.execute(task).catch((err) =>
|
||||||
|
executorLog.error(`Failed to resume ${task.id} after heartbeat completion:`, err),
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Resume orphaned in-progress tasks (e.g., after crash/restart).
|
* Resume orphaned in-progress tasks (e.g., after crash/restart).
|
||||||
* Call once after engine startup.
|
* Call once after engine startup.
|
||||||
@@ -1987,6 +2032,13 @@ export class TaskExecutor {
|
|||||||
async execute(task: Task): Promise<void> {
|
async execute(task: Task): Promise<void> {
|
||||||
executorLog.log(`execute() called for ${task.id} (already executing=${this.executing.has(task.id)})`);
|
executorLog.log(`execute() called for ${task.id} (already executing=${this.executing.has(task.id)})`);
|
||||||
if (this.executing.has(task.id)) return;
|
if (this.executing.has(task.id)) return;
|
||||||
|
|
||||||
|
const assignedAgentId = task.assignedAgentId;
|
||||||
|
if (assignedAgentId && await this.shouldDeferForHeartbeat(assignedAgentId)) {
|
||||||
|
executorLog.log(`${task.id}: skipping execute — agent ${assignedAgentId} has active heartbeat run (allowParallelExecution=false)`);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
this.executing.add(task.id);
|
this.executing.add(task.id);
|
||||||
|
|
||||||
executorLog.log(`Starting ${task.id}: ${task.title || task.description.slice(0, 60)}`);
|
executorLog.log(`Starting ${task.id}: ${task.title || task.description.slice(0, 60)}`);
|
||||||
|
|||||||
@@ -448,6 +448,13 @@ export class InProcessRuntime
|
|||||||
onTerminated: (agentId, reason) => {
|
onTerminated: (agentId, reason) => {
|
||||||
runtimeLog.warn(`Agent ${agentId} terminated (unresponsive): ${reason}`);
|
runtimeLog.warn(`Agent ${agentId} terminated (unresponsive): ${reason}`);
|
||||||
},
|
},
|
||||||
|
onRunCompleted: (agentId) => {
|
||||||
|
if (this.executor) {
|
||||||
|
void this.executor.resumeTaskForAgent(agentId).catch((err) => {
|
||||||
|
runtimeLog.warn(`resumeTaskForAgent failed for ${agentId}: ${err instanceof Error ? err.message : String(err)}`);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
},
|
||||||
});
|
});
|
||||||
this.heartbeatMonitor.start();
|
this.heartbeatMonitor.start();
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user