feat(FN-5264): abort in-flight executor/merge/triage work on soft-delete (F
Implements in-flight abort for soft-deleted tasks across all three execution lanes: executor, merger, and triage now check for `deletedAt` before proceeding and emit `task:soft-delete-in-flight-abort` audits rather than continuing work on a deleted task. The 829-line addition is dominated by integra Fusion-Task-Id: FN-5264 Fusion-Task-Lineage: 4ee8e63b-abf0-43e2-9130-39a05434d8f9
This commit is contained in:
committed by
gsxdsm
parent
93b11c6c0c
commit
ab0a253004
137
packages/engine/src/__tests__/executor-soft-delete-abort.test.ts
Normal file
137
packages/engine/src/__tests__/executor-soft-delete-abort.test.ts
Normal file
@@ -0,0 +1,137 @@
|
||||
import "./executor-test-helpers.js";
|
||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import type { Task } from "@fusion/core";
|
||||
import { TaskExecutor } from "../executor.js";
|
||||
import { executorLog } from "../logger.js";
|
||||
import { resetExecutorMocks } from "./executor-test-helpers.js";
|
||||
|
||||
type Listener = (...args: any[]) => void;
|
||||
|
||||
function createEventedStore() {
|
||||
const listeners = new Map<string, Set<Listener>>();
|
||||
return {
|
||||
store: {
|
||||
on: vi.fn((event: string, listener: Listener) => {
|
||||
const set = listeners.get(event) ?? new Set<Listener>();
|
||||
set.add(listener);
|
||||
listeners.set(event, set);
|
||||
}),
|
||||
off: vi.fn((event: string, listener: Listener) => {
|
||||
listeners.get(event)?.delete(listener);
|
||||
}),
|
||||
getSettings: vi.fn().mockResolvedValue({ globalPause: false, enginePaused: false }),
|
||||
listTasks: vi.fn().mockResolvedValue([]),
|
||||
} as any,
|
||||
emit(event: string, ...args: any[]) {
|
||||
for (const listener of listeners.get(event) ?? []) {
|
||||
listener(...args);
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function makeTask(id: string): Task {
|
||||
return {
|
||||
id,
|
||||
title: id,
|
||||
description: "desc",
|
||||
status: "open",
|
||||
column: "in-progress",
|
||||
createdAt: "2026-01-01T00:00:00.000Z",
|
||||
updatedAt: "2026-01-01T00:00:00.000Z",
|
||||
dependencies: [],
|
||||
comments: [],
|
||||
steps: [],
|
||||
currentStep: 0,
|
||||
log: [],
|
||||
} as unknown as Task;
|
||||
}
|
||||
|
||||
describe("TaskExecutor soft-delete aborts", () => {
|
||||
beforeEach(() => {
|
||||
resetExecutorMocks();
|
||||
vi.clearAllMocks();
|
||||
});
|
||||
|
||||
it("aborts and disposes an active agent session on task:deleted", async () => {
|
||||
const { store, emit } = createEventedStore();
|
||||
const stuckTaskDetector = { untrackTask: vi.fn() };
|
||||
const executor = new TaskExecutor(store, "/tmp/test", { stuckTaskDetector } as any);
|
||||
const abort = vi.fn().mockResolvedValue(undefined);
|
||||
const dispose = vi.fn();
|
||||
|
||||
(executor as any).activeSessions.set("FN-TEST-1", {
|
||||
session: { abort, dispose },
|
||||
seenSteeringIds: new Set<string>(),
|
||||
});
|
||||
|
||||
emit("task:deleted", makeTask("FN-TEST-1"));
|
||||
await Promise.resolve();
|
||||
|
||||
expect(abort).toHaveBeenCalledTimes(1);
|
||||
expect(dispose).toHaveBeenCalledTimes(1);
|
||||
expect((executor as any).activeSessions.has("FN-TEST-1")).toBe(false);
|
||||
expect((executor as any).pausedAborted.has("FN-TEST-1")).toBe(true);
|
||||
expect((executor as any).userCanceledTaskIds.has("FN-TEST-1")).toBe(true);
|
||||
expect(stuckTaskDetector.untrackTask).toHaveBeenCalledWith("FN-TEST-1");
|
||||
});
|
||||
|
||||
it("aborts and removes an active step-session executor on task:deleted", async () => {
|
||||
const { store, emit } = createEventedStore();
|
||||
const executor = new TaskExecutor(store, "/tmp/test");
|
||||
const abortAllSessionBash = vi.fn();
|
||||
const terminateAllSessions = vi.fn().mockResolvedValue(undefined);
|
||||
|
||||
(executor as any).activeStepExecutors.set("FN-TEST-2", {
|
||||
abortAllSessionBash,
|
||||
terminateAllSessions,
|
||||
});
|
||||
|
||||
emit("task:deleted", makeTask("FN-TEST-2"));
|
||||
await Promise.resolve();
|
||||
|
||||
expect(abortAllSessionBash).toHaveBeenCalledTimes(1);
|
||||
expect(terminateAllSessions).toHaveBeenCalledTimes(1);
|
||||
expect((executor as any).activeStepExecutors.has("FN-TEST-2")).toBe(false);
|
||||
});
|
||||
|
||||
it("aborts and disposes an active workflow session on task:deleted", async () => {
|
||||
const { store, emit } = createEventedStore();
|
||||
const executor = new TaskExecutor(store, "/tmp/test");
|
||||
const abort = vi.fn().mockResolvedValue(undefined);
|
||||
const dispose = vi.fn();
|
||||
|
||||
(executor as any).activeWorkflowStepSessions.set("FN-TEST-3", { abort, dispose });
|
||||
|
||||
emit("task:deleted", makeTask("FN-TEST-3"));
|
||||
await Promise.resolve();
|
||||
|
||||
expect(abort).toHaveBeenCalledTimes(1);
|
||||
expect(dispose).toHaveBeenCalledTimes(1);
|
||||
expect((executor as any).activeWorkflowStepSessions.has("FN-TEST-3")).toBe(false);
|
||||
});
|
||||
|
||||
it("disposes reviewer subagents on task:deleted", () => {
|
||||
const { store, emit } = createEventedStore();
|
||||
const executor = new TaskExecutor(store, "/tmp/test");
|
||||
const dispose = vi.fn();
|
||||
|
||||
(executor as any).registerSubagentSession("FN-TEST-4", { dispose });
|
||||
|
||||
emit("task:deleted", makeTask("FN-TEST-4"));
|
||||
|
||||
expect(dispose).toHaveBeenCalledTimes(1);
|
||||
expect((executor as any).activeSubagentSessions.has("FN-TEST-4")).toBe(false);
|
||||
});
|
||||
|
||||
it("is a silent no-op when the deleted task has no active surfaces", () => {
|
||||
const { store, emit } = createEventedStore();
|
||||
const executor = new TaskExecutor(store, "/tmp/test");
|
||||
const errorSpy = vi.spyOn(executorLog, "error");
|
||||
|
||||
expect(() => emit("task:deleted", makeTask("FN-TEST-5"))).not.toThrow();
|
||||
expect(errorSpy).not.toHaveBeenCalled();
|
||||
expect((executor as any).pausedAborted.has("FN-TEST-5")).toBe(true);
|
||||
expect((executor as any).userCanceledTaskIds.has("FN-TEST-5")).toBe(true);
|
||||
});
|
||||
});
|
||||
@@ -190,6 +190,13 @@ vi.mock("node:child_process", async () => {
|
||||
}
|
||||
});
|
||||
|
||||
const execFileFn: any = vi.fn((_file: string, _args: string[] | undefined, opts: any, cb: any) => {
|
||||
const callback = typeof opts === "function" ? opts : cb;
|
||||
if (typeof callback === "function") {
|
||||
callback(null, { stdout: "", stderr: "" });
|
||||
}
|
||||
});
|
||||
|
||||
execFn[promisify.custom] = (cmd: string, opts?: any) =>
|
||||
new Promise((resolve, reject) => {
|
||||
execFn(cmd, opts, (err: any, stdout: string, stderr: string) => {
|
||||
@@ -202,7 +209,10 @@ vi.mock("node:child_process", async () => {
|
||||
}
|
||||
});
|
||||
});
|
||||
return { execSync: execSyncFn, exec: execFn, spawn: spawnFn };
|
||||
execFileFn[promisify.custom] = (_file: string, _args?: string[], _opts?: any) =>
|
||||
Promise.resolve({ stdout: "", stderr: "" });
|
||||
|
||||
return { execSync: execSyncFn, exec: execFn, execFile: execFileFn, spawn: spawnFn };
|
||||
});
|
||||
vi.mock("node:fs", () => ({
|
||||
existsSync: vi.fn().mockReturnValue(true),
|
||||
|
||||
@@ -0,0 +1,187 @@
|
||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { ProjectEngine } from "../project-engine.js";
|
||||
import { runtimeLog } from "../logger.js";
|
||||
|
||||
const mocks = vi.hoisted(() => ({
|
||||
runtimeStart: vi.fn(async () => undefined),
|
||||
runtimeStop: vi.fn(async () => undefined),
|
||||
runtimeResumeAfterUnpause: vi.fn(async () => undefined),
|
||||
runtimeConfigurePrMonitoring: vi.fn(),
|
||||
currentStore: null as Record<string, unknown> | null,
|
||||
aiMergeTask: vi.fn(),
|
||||
execFile: vi.fn(),
|
||||
}));
|
||||
|
||||
vi.mock("@fusion/core", async (importOriginal) => {
|
||||
const { createEngineCoreMock } = await import("../test/mockCore.js");
|
||||
return createEngineCoreMock(() => importOriginal<typeof import("@fusion/core")>(), {});
|
||||
});
|
||||
|
||||
vi.mock("../merger.js", () => ({ aiMergeTask: mocks.aiMergeTask, sweepStaleAutostashes: vi.fn(async () => undefined) }));
|
||||
vi.mock("node:child_process", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("node:child_process")>();
|
||||
return { ...actual, execFile: mocks.execFile };
|
||||
});
|
||||
vi.mock("../pr-monitor.js", () => ({ PrMonitor: vi.fn().mockImplementation(() => ({ onNewComments: vi.fn() })) }));
|
||||
vi.mock("../pr-comment-handler.js", () => ({ PrCommentHandler: vi.fn().mockImplementation(() => ({ handleNewComments: vi.fn() })) }));
|
||||
vi.mock("../auth-storage.js", () => ({
|
||||
createFusionAuthStorage: vi.fn(() => ({ reload: vi.fn(), getOAuthProviders: vi.fn(() => []), get: vi.fn(() => undefined) })),
|
||||
}));
|
||||
vi.mock("../notifier.js", () => ({ NtfyNotifier: vi.fn().mockImplementation(() => ({ start: vi.fn(), stop: vi.fn() })) }));
|
||||
vi.mock("../notification/index.js", () => ({
|
||||
NotificationService: vi.fn().mockImplementation(() => ({ start: vi.fn(), stop: vi.fn() })),
|
||||
OAuthExpiryMonitor: vi.fn().mockImplementation(() => ({ start: vi.fn(), stop: vi.fn() })),
|
||||
}));
|
||||
vi.mock("../cron-runner.js", () => ({
|
||||
CronRunner: vi.fn().mockImplementation(() => ({ start: vi.fn(), stop: vi.fn() })),
|
||||
createAiPromptExecutor: vi.fn(async () => vi.fn()),
|
||||
}));
|
||||
vi.mock("../runtimes/in-process-runtime.js", () => ({
|
||||
InProcessRuntime: vi.fn().mockImplementation(() => ({
|
||||
start: mocks.runtimeStart,
|
||||
stop: mocks.runtimeStop,
|
||||
resumeAfterUnpause: mocks.runtimeResumeAfterUnpause,
|
||||
getTaskStore: () => mocks.currentStore,
|
||||
getAgentStore: vi.fn(),
|
||||
getMessageStore: vi.fn(),
|
||||
getRoutineStore: vi.fn(),
|
||||
getRoutineRunner: vi.fn(),
|
||||
getHeartbeatMonitor: vi.fn(),
|
||||
getTriggerScheduler: vi.fn(),
|
||||
configurePrMonitoring: mocks.runtimeConfigurePrMonitoring,
|
||||
setActiveMergeTaskIdProvider: vi.fn(),
|
||||
setMergeEnqueuer: vi.fn(),
|
||||
setMergeActiveClearer: vi.fn(),
|
||||
})),
|
||||
}));
|
||||
|
||||
type Listener = (...args: any[]) => void | Promise<void>;
|
||||
|
||||
function createMockStore() {
|
||||
const listeners = new Map<string, Set<Listener>>();
|
||||
return {
|
||||
store: {
|
||||
getSettings: vi.fn(async () => ({ autoMerge: true, globalPause: false, enginePaused: false })),
|
||||
listTasks: vi.fn(async () => []),
|
||||
getTask: vi.fn(async (taskId: string) => ({ id: taskId, column: "in-review", paused: false, mergeRetries: 0, status: null })),
|
||||
updateTask: vi.fn(async () => undefined),
|
||||
moveTask: vi.fn(async () => undefined),
|
||||
logEntry: vi.fn(async () => undefined),
|
||||
addTaskComment: vi.fn(async () => undefined),
|
||||
emit: vi.fn(),
|
||||
getActiveMergingTask: vi.fn(() => null),
|
||||
on: vi.fn((event: string, listener: Listener) => {
|
||||
const set = listeners.get(event) ?? new Set<Listener>();
|
||||
set.add(listener);
|
||||
listeners.set(event, set);
|
||||
}),
|
||||
off: vi.fn((event: string, listener: Listener) => {
|
||||
listeners.get(event)?.delete(listener);
|
||||
}),
|
||||
},
|
||||
emit(event: string, ...args: any[]) {
|
||||
for (const listener of listeners.get(event) ?? []) {
|
||||
listener(...args);
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function createEngine() {
|
||||
return new ProjectEngine(
|
||||
{
|
||||
projectId: "proj_test",
|
||||
workingDirectory: "/tmp/proj_test",
|
||||
isolationMode: "in-process",
|
||||
maxConcurrent: 2,
|
||||
maxWorktrees: 2,
|
||||
},
|
||||
{} as never,
|
||||
{ skipNotifier: true },
|
||||
);
|
||||
}
|
||||
|
||||
describe("ProjectEngine soft-delete merge interruption", () => {
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks();
|
||||
});
|
||||
|
||||
it("aborts and disposes the active merge when the task is soft-deleted", async () => {
|
||||
const mockStore = createMockStore();
|
||||
mocks.currentStore = mockStore.store;
|
||||
const logSpy = vi.spyOn(runtimeLog, "log").mockImplementation(() => {});
|
||||
const engine = createEngine();
|
||||
const privateEngine = engine as any;
|
||||
const abort = vi.fn();
|
||||
const dispose = vi.fn();
|
||||
|
||||
await engine.start();
|
||||
privateEngine.activeMergeTaskId = "FN-TEST-1";
|
||||
privateEngine.activeMergeSession = { dispose };
|
||||
privateEngine.mergeAbortController = { abort };
|
||||
privateEngine.mergeActive.add("FN-TEST-1");
|
||||
|
||||
mockStore.emit("task:deleted", { id: "FN-TEST-1" });
|
||||
|
||||
expect(abort).toHaveBeenCalledTimes(1);
|
||||
expect(dispose).toHaveBeenCalledTimes(1);
|
||||
expect(privateEngine.activeMergeTaskId).toBeNull();
|
||||
expect(privateEngine.activeMergeSession).toBeNull();
|
||||
expect(privateEngine.mergeAbortController).toBeNull();
|
||||
expect(privateEngine.mergeActive.has("FN-TEST-1")).toBe(false);
|
||||
expect(logSpy).toHaveBeenCalledWith("Soft-deleted task interrupting active merge: FN-TEST-1");
|
||||
|
||||
await engine.stop();
|
||||
});
|
||||
|
||||
it("removes queued soft-deleted tasks without touching another active merge", async () => {
|
||||
const mockStore = createMockStore();
|
||||
mocks.currentStore = mockStore.store;
|
||||
const engine = createEngine();
|
||||
const privateEngine = engine as any;
|
||||
|
||||
await engine.start();
|
||||
privateEngine.mergeQueue = ["FN-TEST-2", "FN-OTHER"];
|
||||
privateEngine.mergeActive.add("FN-TEST-2");
|
||||
privateEngine.activeMergeTaskId = "FN-ACTIVE";
|
||||
|
||||
mockStore.emit("task:deleted", { id: "FN-TEST-2" });
|
||||
|
||||
expect(privateEngine.mergeQueue).toEqual(["FN-OTHER"]);
|
||||
expect(privateEngine.mergeActive.has("FN-TEST-2")).toBe(false);
|
||||
expect(privateEngine.activeMergeTaskId).toBe("FN-ACTIVE");
|
||||
|
||||
await engine.stop();
|
||||
});
|
||||
|
||||
it("removes soft-deleted tasks from pausedReviewTaskIds", async () => {
|
||||
const mockStore = createMockStore();
|
||||
mocks.currentStore = mockStore.store;
|
||||
const engine = createEngine();
|
||||
const privateEngine = engine as any;
|
||||
|
||||
await engine.start();
|
||||
privateEngine.pausedReviewTaskIds.add("FN-TEST-3");
|
||||
|
||||
mockStore.emit("task:deleted", { id: "FN-TEST-3" });
|
||||
|
||||
expect(privateEngine.pausedReviewTaskIds.has("FN-TEST-3")).toBe(false);
|
||||
|
||||
await engine.stop();
|
||||
});
|
||||
|
||||
it("detaches the task:deleted handler on stop", async () => {
|
||||
const mockStore = createMockStore();
|
||||
mocks.currentStore = mockStore.store;
|
||||
const engine = createEngine();
|
||||
const privateEngine = engine as any;
|
||||
|
||||
await engine.start();
|
||||
await engine.stop();
|
||||
privateEngine.mergeQueue = ["FN-TEST-4"];
|
||||
|
||||
mockStore.emit("task:deleted", { id: "FN-TEST-4" });
|
||||
|
||||
expect(privateEngine.mergeQueue).toEqual(["FN-TEST-4"]);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,250 @@
|
||||
import "../executor-test-helpers.js";
|
||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { AutoClaimSnapshotManager } from "../../auto-claim-snapshot.js";
|
||||
import { TaskExecutor } from "../../executor.js";
|
||||
import { ProjectEngine } from "../../project-engine.js";
|
||||
import { Scheduler } from "../../scheduler.js";
|
||||
import { resetExecutorMocks } from "../executor-test-helpers.js";
|
||||
|
||||
const projectEngineMocks = vi.hoisted(() => ({
|
||||
runtimeStart: vi.fn(async () => undefined),
|
||||
runtimeStop: vi.fn(async () => undefined),
|
||||
runtimeResumeAfterUnpause: vi.fn(async () => undefined),
|
||||
runtimeConfigurePrMonitoring: vi.fn(),
|
||||
currentStore: null as Record<string, unknown> | null,
|
||||
aiMergeTask: vi.fn(),
|
||||
execFile: vi.fn(),
|
||||
}));
|
||||
|
||||
vi.mock("@fusion/core", async (importOriginal) => {
|
||||
const { createEngineCoreMock } = await import("../../test/mockCore.js");
|
||||
return createEngineCoreMock(() => importOriginal<typeof import("@fusion/core")>(), {});
|
||||
});
|
||||
vi.mock("../../merger.js", () => ({ aiMergeTask: projectEngineMocks.aiMergeTask, sweepStaleAutostashes: vi.fn(async () => undefined) }));
|
||||
vi.mock("node:child_process", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("node:child_process")>();
|
||||
return { ...actual, execFile: projectEngineMocks.execFile };
|
||||
});
|
||||
vi.mock("../../pr-monitor.js", () => ({ PrMonitor: vi.fn().mockImplementation(() => ({ onNewComments: vi.fn() })) }));
|
||||
vi.mock("../../pr-comment-handler.js", () => ({ PrCommentHandler: vi.fn().mockImplementation(() => ({ handleNewComments: vi.fn() })) }));
|
||||
vi.mock("../../auth-storage.js", () => ({
|
||||
createFusionAuthStorage: vi.fn(() => ({ reload: vi.fn(), getOAuthProviders: vi.fn(() => []), get: vi.fn(() => undefined) })),
|
||||
}));
|
||||
vi.mock("../../notifier.js", () => ({ NtfyNotifier: vi.fn().mockImplementation(() => ({ start: vi.fn(), stop: vi.fn() })) }));
|
||||
vi.mock("../../notification/index.js", () => ({
|
||||
NotificationService: vi.fn().mockImplementation(() => ({ start: vi.fn(), stop: vi.fn() })),
|
||||
OAuthExpiryMonitor: vi.fn().mockImplementation(() => ({ start: vi.fn(), stop: vi.fn() })),
|
||||
}));
|
||||
vi.mock("../../cron-runner.js", () => ({
|
||||
CronRunner: vi.fn().mockImplementation(() => ({ start: vi.fn(), stop: vi.fn() })),
|
||||
createAiPromptExecutor: vi.fn(async () => vi.fn()),
|
||||
}));
|
||||
vi.mock("../../runtimes/in-process-runtime.js", () => ({
|
||||
InProcessRuntime: vi.fn().mockImplementation(() => ({
|
||||
start: projectEngineMocks.runtimeStart,
|
||||
stop: projectEngineMocks.runtimeStop,
|
||||
resumeAfterUnpause: projectEngineMocks.runtimeResumeAfterUnpause,
|
||||
getTaskStore: () => projectEngineMocks.currentStore,
|
||||
getAgentStore: vi.fn(),
|
||||
getMessageStore: vi.fn(),
|
||||
getRoutineStore: vi.fn(),
|
||||
getRoutineRunner: vi.fn(),
|
||||
getHeartbeatMonitor: vi.fn(),
|
||||
getTriggerScheduler: vi.fn(),
|
||||
configurePrMonitoring: projectEngineMocks.runtimeConfigurePrMonitoring,
|
||||
setActiveMergeTaskIdProvider: vi.fn(),
|
||||
setMergeEnqueuer: vi.fn(),
|
||||
setMergeActiveClearer: vi.fn(),
|
||||
})),
|
||||
}));
|
||||
|
||||
type Listener = (...args: any[]) => void | Promise<void>;
|
||||
|
||||
type TestTask = {
|
||||
id: string;
|
||||
title: string;
|
||||
description: string;
|
||||
status: string;
|
||||
column: string;
|
||||
createdAt: string;
|
||||
updatedAt: string;
|
||||
dependencies: string[];
|
||||
comments: unknown[];
|
||||
steps: unknown[];
|
||||
currentStep: number;
|
||||
log: unknown[];
|
||||
deletedAt?: string | null;
|
||||
paused?: boolean;
|
||||
};
|
||||
|
||||
function createEventedStore(initialTasks: TestTask[] = []) {
|
||||
const listeners = new Map<string, Set<Listener>>();
|
||||
let sequence = 1;
|
||||
const tasks = initialTasks.map((task) => ({ ...task }));
|
||||
const nextTimestamp = () => new Date(1_716_000_000_000 + sequence++).toISOString();
|
||||
|
||||
const store = {
|
||||
on: vi.fn((event: string, listener: Listener) => {
|
||||
const set = listeners.get(event) ?? new Set<Listener>();
|
||||
set.add(listener);
|
||||
listeners.set(event, set);
|
||||
}),
|
||||
off: vi.fn((event: string, listener: Listener) => {
|
||||
listeners.get(event)?.delete(listener);
|
||||
}),
|
||||
getSettings: vi.fn().mockResolvedValue({ globalPause: false, enginePaused: false, autoMerge: true, maxConcurrent: 2, maxWorktrees: 2 }),
|
||||
getRootDir: vi.fn().mockReturnValue("/test/project"),
|
||||
listTasks: vi.fn(async (options?: { column?: string }) =>
|
||||
tasks
|
||||
.filter((task) => !task.deletedAt)
|
||||
.filter((task) => (options?.column ? task.column === options.column : true))
|
||||
.map((task) => ({ ...task })),
|
||||
),
|
||||
getTask: vi.fn(async (taskId: string) => {
|
||||
const task = tasks.find((entry) => entry.id === taskId);
|
||||
if (!task || task.deletedAt) throw new Error(`Task ${taskId} not found`);
|
||||
return { ...task };
|
||||
}),
|
||||
updateTask: vi.fn(async () => undefined),
|
||||
moveTask: vi.fn(async (id: string, column: string) => {
|
||||
const task = tasks.find((entry) => entry.id === id);
|
||||
if (!task) throw new Error(`Task ${id} not found`);
|
||||
const from = task.column;
|
||||
task.column = column;
|
||||
task.updatedAt = nextTimestamp();
|
||||
emit("task:moved", { task: { ...task }, from, to: column, source: "user" });
|
||||
return { ...task };
|
||||
}),
|
||||
deleteTask: vi.fn(async (id: string) => {
|
||||
const task = tasks.find((entry) => entry.id === id);
|
||||
if (!task) throw new Error(`Task ${id} not found`);
|
||||
if (!task.deletedAt) {
|
||||
task.deletedAt = nextTimestamp();
|
||||
task.updatedAt = task.deletedAt;
|
||||
emit("task:deleted", { ...task });
|
||||
}
|
||||
return { ...task };
|
||||
}),
|
||||
addTaskComment: vi.fn(async () => undefined),
|
||||
logEntry: vi.fn(async () => undefined),
|
||||
getActiveMergingTask: vi.fn(() => null),
|
||||
} as any;
|
||||
|
||||
const emit = (event: string, ...args: any[]) => {
|
||||
for (const listener of listeners.get(event) ?? []) {
|
||||
listener(...args);
|
||||
}
|
||||
};
|
||||
|
||||
store.emit = emit;
|
||||
|
||||
return { store, emit, tasks };
|
||||
}
|
||||
|
||||
function createProjectEngine() {
|
||||
return new ProjectEngine(
|
||||
{
|
||||
projectId: "proj_test",
|
||||
workingDirectory: "/tmp/proj_test",
|
||||
isolationMode: "in-process",
|
||||
maxConcurrent: 2,
|
||||
maxWorktrees: 2,
|
||||
},
|
||||
{} as never,
|
||||
{ skipNotifier: true },
|
||||
);
|
||||
}
|
||||
|
||||
function makeTask(id: string, column: string = "in-progress"): TestTask {
|
||||
return {
|
||||
id,
|
||||
title: id,
|
||||
description: id,
|
||||
status: "open",
|
||||
column,
|
||||
createdAt: "2026-01-01T00:00:00.000Z",
|
||||
updatedAt: "2026-01-01T00:00:00.000Z",
|
||||
dependencies: [],
|
||||
comments: [],
|
||||
steps: [],
|
||||
currentStep: 0,
|
||||
log: [],
|
||||
deletedAt: null,
|
||||
};
|
||||
}
|
||||
|
||||
describe("reliability interactions: soft-delete in-flight abort", () => {
|
||||
beforeEach(() => {
|
||||
resetExecutorMocks();
|
||||
vi.clearAllMocks();
|
||||
});
|
||||
|
||||
it("aborts executor work while invalidating the auto-claim snapshot and leaving deleted tasks undispatchable", async () => {
|
||||
const { store } = createEventedStore([makeTask("FN-EXEC", "in-progress")]);
|
||||
const snapshotManager = new AutoClaimSnapshotManager({ taskStore: store as any });
|
||||
const invalidateSpy = vi.spyOn(snapshotManager, "invalidate");
|
||||
new Scheduler(store as any, { snapshotManager } as any);
|
||||
const executor = new TaskExecutor(store as any, "/tmp/test");
|
||||
const abort = vi.fn().mockResolvedValue(undefined);
|
||||
const dispose = vi.fn();
|
||||
|
||||
(executor as any).activeSessions.set("FN-EXEC", {
|
||||
session: { abort, dispose },
|
||||
seenSteeringIds: new Set<string>(),
|
||||
});
|
||||
|
||||
await store.deleteTask("FN-EXEC");
|
||||
await Promise.resolve();
|
||||
|
||||
expect(abort).toHaveBeenCalledTimes(1);
|
||||
expect(dispose).toHaveBeenCalledTimes(1);
|
||||
expect((executor as any).activeSessions.has("FN-EXEC")).toBe(false);
|
||||
expect(invalidateSpy).toHaveBeenCalledWith("task:deleted");
|
||||
expect((await snapshotManager.getSnapshot()).tasks.map((task) => task.id)).not.toContain("FN-EXEC");
|
||||
});
|
||||
|
||||
it("aborts an active merge without double-abort errors when executor has nothing active", async () => {
|
||||
const { store } = createEventedStore([makeTask("FN-MERGE", "in-review")]);
|
||||
projectEngineMocks.currentStore = store;
|
||||
const engine = createProjectEngine();
|
||||
const executor = new TaskExecutor(store as any, "/tmp/test");
|
||||
const abort = vi.fn();
|
||||
const dispose = vi.fn();
|
||||
|
||||
await engine.start();
|
||||
(engine as any).activeMergeTaskId = "FN-MERGE";
|
||||
(engine as any).activeMergeSession = { dispose };
|
||||
(engine as any).mergeAbortController = { abort };
|
||||
(engine as any).mergeActive.add("FN-MERGE");
|
||||
|
||||
store.emit("task:deleted", { ...makeTask("FN-MERGE", "in-review"), deletedAt: "2026-01-02T00:00:00.000Z" });
|
||||
await Promise.resolve();
|
||||
|
||||
expect(abort).toHaveBeenCalledTimes(1);
|
||||
expect(dispose).toHaveBeenCalledTimes(1);
|
||||
expect((engine as any).activeMergeTaskId).toBeNull();
|
||||
expect((executor as any).activeSessions.size).toBe(0);
|
||||
|
||||
await engine.stop();
|
||||
});
|
||||
|
||||
it("treats task:moved followed by task:deleted as idempotent cleanup", async () => {
|
||||
const { store } = createEventedStore([makeTask("FN-RACE", "in-progress")]);
|
||||
const executor = new TaskExecutor(store as any, "/tmp/test");
|
||||
const abort = vi.fn().mockResolvedValue(undefined);
|
||||
const dispose = vi.fn();
|
||||
|
||||
(executor as any).activeSessions.set("FN-RACE", {
|
||||
session: { abort, dispose },
|
||||
seenSteeringIds: new Set<string>(),
|
||||
});
|
||||
|
||||
await store.moveTask("FN-RACE", "todo");
|
||||
await store.deleteTask("FN-RACE");
|
||||
await Promise.resolve();
|
||||
|
||||
expect(abort).toHaveBeenCalledTimes(1);
|
||||
expect(dispose).toHaveBeenCalledTimes(1);
|
||||
expect((executor as any).activeSessions.has("FN-RACE")).toBe(false);
|
||||
});
|
||||
});
|
||||
102
packages/engine/src/__tests__/triage-soft-delete-abort.test.ts
Normal file
102
packages/engine/src/__tests__/triage-soft-delete-abort.test.ts
Normal file
@@ -0,0 +1,102 @@
|
||||
import "./executor-test-helpers.js";
|
||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { TriageProcessor } from "../triage.js";
|
||||
import { resetExecutorMocks } from "./executor-test-helpers.js";
|
||||
|
||||
type Listener = (...args: any[]) => void;
|
||||
|
||||
function createEventedStore() {
|
||||
const listeners = new Map<string, Set<Listener>>();
|
||||
return {
|
||||
store: {
|
||||
getSettings: vi.fn().mockResolvedValue({ pollIntervalMs: 60_000, maxConcurrent: 1, maxWorktrees: 1, autoMerge: true }),
|
||||
listTasks: vi.fn().mockResolvedValue([]),
|
||||
updateTask: vi.fn().mockResolvedValue(undefined),
|
||||
on: vi.fn((event: string, listener: Listener) => {
|
||||
const set = listeners.get(event) ?? new Set<Listener>();
|
||||
set.add(listener);
|
||||
listeners.set(event, set);
|
||||
}),
|
||||
off: vi.fn((event: string, listener: Listener) => {
|
||||
listeners.get(event)?.delete(listener);
|
||||
}),
|
||||
} as any,
|
||||
emit(event: string, ...args: any[]) {
|
||||
for (const listener of listeners.get(event) ?? []) {
|
||||
listener(...args);
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
describe("TriageProcessor soft-delete aborts", () => {
|
||||
beforeEach(() => {
|
||||
resetExecutorMocks();
|
||||
vi.clearAllMocks();
|
||||
});
|
||||
|
||||
it("aborts and disposes an active specify session on task:deleted", async () => {
|
||||
const { store, emit } = createEventedStore();
|
||||
const stuckTaskDetector = { untrackTask: vi.fn() };
|
||||
const processor = new TriageProcessor(store, "/tmp/root", { stuckTaskDetector } as any);
|
||||
const abort = vi.fn().mockResolvedValue(undefined);
|
||||
const dispose = vi.fn();
|
||||
|
||||
processor.start();
|
||||
(processor as any).activeSessions.set("FN-TEST-1", { abort, dispose });
|
||||
|
||||
emit("task:deleted", { id: "FN-TEST-1" });
|
||||
await Promise.resolve();
|
||||
|
||||
expect(abort).toHaveBeenCalledTimes(1);
|
||||
expect(dispose).toHaveBeenCalledTimes(1);
|
||||
expect((processor as any).activeSessions.has("FN-TEST-1")).toBe(false);
|
||||
expect((processor as any).pauseAborted.has("FN-TEST-1")).toBe(true);
|
||||
expect(stuckTaskDetector.untrackTask).toHaveBeenCalledWith("FN-TEST-1");
|
||||
|
||||
processor.stop();
|
||||
});
|
||||
|
||||
it("disposes reviewer subagents for a soft-deleted task", () => {
|
||||
const { store, emit } = createEventedStore();
|
||||
const processor = new TriageProcessor(store, "/tmp/root");
|
||||
const dispose = vi.fn();
|
||||
|
||||
processor.start();
|
||||
(processor as any).registerSubagentSession("FN-TEST-2", { dispose });
|
||||
|
||||
emit("task:deleted", { id: "FN-TEST-2" });
|
||||
|
||||
expect(dispose).toHaveBeenCalledTimes(1);
|
||||
expect((processor as any).activeSubagentSessions.has("FN-TEST-2")).toBe(false);
|
||||
|
||||
processor.stop();
|
||||
});
|
||||
|
||||
it("is a no-op for unknown soft-deleted ids", () => {
|
||||
const { store, emit } = createEventedStore();
|
||||
const processor = new TriageProcessor(store, "/tmp/root");
|
||||
|
||||
processor.start();
|
||||
expect(() => emit("task:deleted", { id: "FN-UNKNOWN" })).not.toThrow();
|
||||
processor.stop();
|
||||
});
|
||||
|
||||
it("detaches the task:deleted listener on stop", () => {
|
||||
const { store, emit } = createEventedStore();
|
||||
const processor = new TriageProcessor(store, "/tmp/root");
|
||||
const abort = vi.fn().mockResolvedValue(undefined);
|
||||
const dispose = vi.fn();
|
||||
|
||||
processor.start();
|
||||
(processor as any).activeSessions.set("FN-TEST-3", { abort, dispose });
|
||||
processor.stop();
|
||||
const abortCallsAfterStop = abort.mock.calls.length;
|
||||
const disposeCallsAfterStop = dispose.mock.calls.length;
|
||||
|
||||
emit("task:deleted", { id: "FN-TEST-3" });
|
||||
|
||||
expect(abort).toHaveBeenCalledTimes(abortCallsAfterStop);
|
||||
expect(dispose).toHaveBeenCalledTimes(disposeCallsAfterStop);
|
||||
});
|
||||
});
|
||||
@@ -1437,6 +1437,74 @@ export class TaskExecutor {
|
||||
this.activeSubagentSessions.delete(taskId);
|
||||
}
|
||||
|
||||
private abortInFlightTaskWork(taskId: string, reason: string, options: { userCanceled?: boolean } = {}): void {
|
||||
let hadActiveSurface = false;
|
||||
|
||||
if (options.userCanceled) {
|
||||
this.userCanceledTaskIds.add(taskId);
|
||||
}
|
||||
this.pausedAborted.add(taskId);
|
||||
this.options.stuckTaskDetector?.untrackTask(taskId);
|
||||
this.clearWorkflowRerunWatchdog(taskId);
|
||||
this.clearCompletedTaskWatchdog(taskId);
|
||||
|
||||
if (this.activeSessions.has(taskId)) {
|
||||
hadActiveSurface = true;
|
||||
const { session } = this.activeSessions.get(taskId)!;
|
||||
const sessionWithAbort = session as AgentSession & { abort?: () => Promise<void> };
|
||||
if (typeof sessionWithAbort.abort === "function") {
|
||||
void sessionWithAbort.abort().catch((err) => {
|
||||
executorLog.warn(`Failed to abort agent session for ${taskId}: ${err}`);
|
||||
});
|
||||
}
|
||||
session.dispose();
|
||||
this.deleteActiveSession(taskId);
|
||||
}
|
||||
|
||||
if (this.activeStepExecutors.has(taskId)) {
|
||||
hadActiveSurface = true;
|
||||
const stepExecutor = this.activeStepExecutors.get(taskId)!;
|
||||
const stepExecutorWithAbort = stepExecutor as StepSessionExecutor & { abortAllSessionBash?: () => void };
|
||||
if (typeof stepExecutorWithAbort.abortAllSessionBash === "function") {
|
||||
try {
|
||||
stepExecutorWithAbort.abortAllSessionBash();
|
||||
} catch (err) {
|
||||
executorLog.warn(`Failed to abort step-session bash for ${taskId}: ${err}`);
|
||||
}
|
||||
}
|
||||
stepExecutor.terminateAllSessions().catch((err) =>
|
||||
executorLog.error(`Failed to terminate step sessions for ${taskId}:`, err),
|
||||
);
|
||||
this.deleteActiveStepExecutor(taskId);
|
||||
}
|
||||
|
||||
if (this.activeWorkflowStepSessions.has(taskId)) {
|
||||
hadActiveSurface = true;
|
||||
const workflowSession = this.activeWorkflowStepSessions.get(taskId)!;
|
||||
const sessionWithAbort = workflowSession as AgentSession & { abort?: () => Promise<void> };
|
||||
if (typeof sessionWithAbort.abort === "function") {
|
||||
void sessionWithAbort.abort().catch((err) => {
|
||||
executorLog.warn(`Failed to abort workflow step session for ${taskId}: ${err}`);
|
||||
});
|
||||
}
|
||||
workflowSession.dispose();
|
||||
this.deleteActiveWorkflowStepSession(taskId);
|
||||
}
|
||||
|
||||
if (this.activeSubagentSessions.has(taskId)) {
|
||||
hadActiveSurface = true;
|
||||
this.disposeSubagentsForTask(taskId, reason);
|
||||
}
|
||||
|
||||
this.loopRecoveryState.delete(taskId);
|
||||
this.spawnedAgents.delete(taskId);
|
||||
this.stuckAborted.delete(taskId);
|
||||
|
||||
if (hadActiveSurface) {
|
||||
executorLog.log(`${taskId}: aborting in-flight work — ${reason}`);
|
||||
}
|
||||
}
|
||||
|
||||
abortAllSessionBash(): void {
|
||||
for (const [taskId, { session }] of this.activeSessions) {
|
||||
try {
|
||||
@@ -1497,73 +1565,16 @@ export class TaskExecutor {
|
||||
executorLog.error(`Failed to start ${task.id}:`, err),
|
||||
);
|
||||
} else if (from === "in-progress") {
|
||||
if (source === "user" && to === "todo") {
|
||||
this.userCanceledTaskIds.add(task.id);
|
||||
}
|
||||
this.clearCompletedTaskWatchdog(task.id);
|
||||
// Task moved away from in-progress — terminate any active sessions
|
||||
if (this.activeSessions.has(task.id)) {
|
||||
executorLog.log(`${task.id} moved from in-progress to ${to} — terminating agent session`);
|
||||
this.pausedAborted.add(task.id);
|
||||
this.options.stuckTaskDetector?.untrackTask(task.id);
|
||||
const { session } = this.activeSessions.get(task.id)!;
|
||||
const sessionWithAbort = session as AgentSession & { abort?: () => Promise<void> };
|
||||
if (typeof sessionWithAbort.abort === "function") {
|
||||
void sessionWithAbort.abort().catch((err) => {
|
||||
executorLog.warn(`Failed to abort agent session for ${task.id}: ${err}`);
|
||||
});
|
||||
}
|
||||
session.dispose();
|
||||
this.deleteActiveSession(task.id);
|
||||
}
|
||||
if (this.activeStepExecutors.has(task.id)) {
|
||||
executorLog.log(`${task.id} moved from in-progress to ${to} — terminating step sessions`);
|
||||
this.pausedAborted.add(task.id);
|
||||
this.options.stuckTaskDetector?.untrackTask(task.id);
|
||||
const stepExecutor = this.activeStepExecutors.get(task.id)!;
|
||||
const stepExecutorWithAbort = stepExecutor as StepSessionExecutor & { abortAllSessionBash?: () => void };
|
||||
if (typeof stepExecutorWithAbort.abortAllSessionBash === "function") {
|
||||
try {
|
||||
stepExecutorWithAbort.abortAllSessionBash();
|
||||
} catch (err) {
|
||||
executorLog.warn(`Failed to abort step-session bash for ${task.id}: ${err}`);
|
||||
}
|
||||
}
|
||||
stepExecutor.terminateAllSessions().catch((err) =>
|
||||
executorLog.error(`Failed to terminate step sessions for ${task.id}:`, err),
|
||||
);
|
||||
this.deleteActiveStepExecutor(task.id);
|
||||
}
|
||||
if (this.activeWorkflowStepSessions.has(task.id)) {
|
||||
executorLog.log(`${task.id} moved from in-progress to ${to} — terminating workflow step session`);
|
||||
this.pausedAborted.add(task.id);
|
||||
this.options.stuckTaskDetector?.untrackTask(task.id);
|
||||
const workflowSession = this.activeWorkflowStepSessions.get(task.id)!;
|
||||
const sessionWithAbort = workflowSession as AgentSession & { abort?: () => Promise<void> };
|
||||
if (typeof sessionWithAbort.abort === "function") {
|
||||
void sessionWithAbort.abort().catch((err) => {
|
||||
executorLog.warn(`Failed to abort workflow step session for ${task.id}: ${err}`);
|
||||
});
|
||||
}
|
||||
workflowSession.dispose();
|
||||
this.deleteActiveWorkflowStepSession(task.id);
|
||||
}
|
||||
// Reviewer subagents run in their own sessions outside `activeSessions`
|
||||
// and `activeStepExecutors`, so the loops above don't reach them.
|
||||
// Without this, a reviewer keeps running (and emitting verdicts that
|
||||
// trigger step transitions) even after the parent task was moved out
|
||||
// of in-progress.
|
||||
this.disposeSubagentsForTask(task.id, `parent moved from in-progress to ${to}`);
|
||||
// Clean up all in-memory state for this task so nothing leaks across runs.
|
||||
// This prevents zombie state from persisting when a task moves away from
|
||||
// in-progress while execute() is still unwinding, or when the scheduler
|
||||
// moves a task out before the executor's task:moved fires.
|
||||
this.loopRecoveryState.delete(task.id);
|
||||
this.spawnedAgents.delete(task.id);
|
||||
this.stuckAborted.delete(task.id);
|
||||
this.abortInFlightTaskWork(task.id, `parent moved from in-progress to ${to}`, {
|
||||
userCanceled: source === "user" && to === "todo",
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
store.on("task:deleted", (task) => {
|
||||
this.abortInFlightTaskWork(task.id, "task soft-deleted", { userCanceled: true });
|
||||
});
|
||||
|
||||
// When a task is paused while executing, terminate the agent session.
|
||||
// When steering comments are added during execution, inject them into the running session.
|
||||
//
|
||||
|
||||
@@ -236,6 +236,7 @@ export class ProjectEngine {
|
||||
private settingsHandlers: Array<(...args: any[]) => void> = [];
|
||||
private taskMovedHandler?: (...args: any[]) => void;
|
||||
private taskUpdatedHandler?: (...args: any[]) => void;
|
||||
private taskDeletedHandler?: (...args: any[]) => void;
|
||||
private autostashOrphansHandler?: (...args: any[]) => void;
|
||||
|
||||
constructor(
|
||||
@@ -566,6 +567,9 @@ export class ProjectEngine {
|
||||
if (this.taskUpdatedHandler) {
|
||||
store.off("task:updated", this.taskUpdatedHandler);
|
||||
}
|
||||
if (this.taskDeletedHandler) {
|
||||
store.off("task:deleted", this.taskDeletedHandler);
|
||||
}
|
||||
if (this.autostashOrphansHandler) {
|
||||
store.off("merger:autostashOrphans", this.autostashOrphansHandler as any);
|
||||
}
|
||||
@@ -2247,7 +2251,39 @@ export class ProjectEngine {
|
||||
}
|
||||
};
|
||||
|
||||
this.taskDeletedHandler = (task: Task) => {
|
||||
this.pausedReviewTaskIds.delete(task.id);
|
||||
|
||||
const queueLengthBefore = this.mergeQueue.length;
|
||||
this.mergeQueue = this.mergeQueue.filter((queuedTaskId) => queuedTaskId !== task.id);
|
||||
const removedFromQueue = this.mergeQueue.length !== queueLengthBefore;
|
||||
|
||||
if (removedFromQueue) {
|
||||
if (this.activeMergeTaskId !== task.id) {
|
||||
this.mergeActive.delete(task.id);
|
||||
}
|
||||
runtimeLog.log(`Soft-deleted task removed from merge queue: ${task.id}`);
|
||||
}
|
||||
|
||||
if (this.activeMergeTaskId !== task.id) {
|
||||
return;
|
||||
}
|
||||
|
||||
runtimeLog.log(`Soft-deleted task interrupting active merge: ${task.id}`);
|
||||
this.mergeAbortController?.abort();
|
||||
this.mergeAbortController = null;
|
||||
|
||||
if (this.activeMergeSession) {
|
||||
this.activeMergeSession.dispose();
|
||||
this.activeMergeSession = null;
|
||||
}
|
||||
|
||||
this.mergeActive.delete(task.id);
|
||||
this.activeMergeTaskId = null;
|
||||
};
|
||||
|
||||
store.on("task:updated", this.taskUpdatedHandler);
|
||||
store.on("task:deleted", this.taskDeletedHandler);
|
||||
}
|
||||
|
||||
private async startupMergeSweep(store: TaskStore): Promise<void> {
|
||||
|
||||
@@ -580,6 +580,7 @@ export class TriageProcessor {
|
||||
private pauseAborted = new Set<string>();
|
||||
/** Tasks killed by the stuck task detector (to avoid reporting as errors). */
|
||||
private stuckAborted = new Set<string>();
|
||||
private taskDeletedHandler?: (task: Task) => void;
|
||||
|
||||
/**
|
||||
* @param store — Task store instance (also used to listen for `settings:updated` events)
|
||||
@@ -629,11 +630,37 @@ export class TriageProcessor {
|
||||
this.poll();
|
||||
}
|
||||
});
|
||||
|
||||
this.taskDeletedHandler = (task: Task) => {
|
||||
if (this.activeSubagentSessions.has(task.id)) {
|
||||
this.disposeSubagentsForTask(task.id, "task soft-deleted");
|
||||
}
|
||||
if (this.activeSessions.has(task.id)) {
|
||||
const session = this.activeSessions.get(task.id)!;
|
||||
planLog.log(`task soft-deleted — terminating triage session for ${task.id}`);
|
||||
this.pauseAborted.add(task.id);
|
||||
this.options.stuckTaskDetector?.untrackTask(task.id);
|
||||
const sessionWithAbort = session as {
|
||||
abort?: () => Promise<void>;
|
||||
dispose: () => void;
|
||||
};
|
||||
if (typeof sessionWithAbort.abort === "function") {
|
||||
void sessionWithAbort.abort().catch((err) => {
|
||||
planLog.warn(`Failed to abort triage session for ${task.id}: ${err}`);
|
||||
});
|
||||
}
|
||||
session.dispose();
|
||||
this.activeSessions.delete(task.id);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
start(): void {
|
||||
if (this.running) return;
|
||||
this.running = true;
|
||||
if (this.taskDeletedHandler && typeof this.store.on === "function") {
|
||||
this.store.on("task:deleted", this.taskDeletedHandler);
|
||||
}
|
||||
|
||||
// Clear stale "planning" statuses left by a prior crash/restart.
|
||||
// No triage agent is actually running at startup, so any task still
|
||||
@@ -672,6 +699,9 @@ export class TriageProcessor {
|
||||
this.pollInterval = null;
|
||||
this.activePollMs = null;
|
||||
}
|
||||
if (this.taskDeletedHandler && typeof this.store.off === "function") {
|
||||
this.store.off("task:deleted", this.taskDeletedHandler);
|
||||
}
|
||||
// Tear down any in-flight specify sessions and reviewer subagents so they
|
||||
// don't keep streaming LLM tokens / tool calls past engine shutdown.
|
||||
this.abortAndDisposeActiveSessions("engine stop");
|
||||
|
||||
Reference in New Issue
Block a user