feat(FN-2855): route scheduled tasks using effective node resolution
- Add an effective node resolver with task override, project default, and local fallback precedence. - Wire scheduler dispatch to persist effectiveNodeId/effectiveNodeSource and log resolved node routing. - Add coverage for effective node resolution and scheduler node routing integration behavior. - Stabilize workspace test resolution by adding @fusion/core and @fusion/plugin-sdk aliases across Vitest configs.
This commit is contained in:
67
packages/engine/src/__tests__/effective-node.test.ts
Normal file
67
packages/engine/src/__tests__/effective-node.test.ts
Normal file
@@ -0,0 +1,67 @@
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { resolveEffectiveNode } from "../effective-node.js";
|
||||
|
||||
describe("resolveEffectiveNode", () => {
|
||||
it("uses task override when both task and project default are set", () => {
|
||||
expect(resolveEffectiveNode({ nodeId: "node-task" }, { defaultNodeId: "node-project" })).toEqual({
|
||||
nodeId: "node-task",
|
||||
source: "task-override",
|
||||
});
|
||||
});
|
||||
|
||||
it("uses project default when task override is not set", () => {
|
||||
expect(resolveEffectiveNode({ nodeId: undefined }, { defaultNodeId: "node-project" })).toEqual({
|
||||
nodeId: "node-project",
|
||||
source: "project-default",
|
||||
});
|
||||
});
|
||||
|
||||
it("uses local when neither task override nor project default is set", () => {
|
||||
expect(resolveEffectiveNode({ nodeId: undefined }, { defaultNodeId: undefined })).toEqual({
|
||||
nodeId: undefined,
|
||||
source: "local",
|
||||
});
|
||||
});
|
||||
|
||||
it("treats empty task override as unset and falls through to project default", () => {
|
||||
expect(resolveEffectiveNode({ nodeId: "" }, { defaultNodeId: "node-project" })).toEqual({
|
||||
nodeId: "node-project",
|
||||
source: "project-default",
|
||||
});
|
||||
});
|
||||
|
||||
it("treats null task nodeId as unset", () => {
|
||||
expect(resolveEffectiveNode({ nodeId: null as unknown as string }, { defaultNodeId: "node-project" })).toEqual({
|
||||
nodeId: "node-project",
|
||||
source: "project-default",
|
||||
});
|
||||
});
|
||||
|
||||
it("treats empty project default as unset and falls through to local", () => {
|
||||
expect(resolveEffectiveNode({ nodeId: undefined }, { defaultNodeId: "" })).toEqual({
|
||||
nodeId: undefined,
|
||||
source: "local",
|
||||
});
|
||||
});
|
||||
|
||||
it("treats null project default as unset", () => {
|
||||
expect(resolveEffectiveNode({ nodeId: undefined }, { defaultNodeId: null as unknown as string })).toEqual({
|
||||
nodeId: undefined,
|
||||
source: "local",
|
||||
});
|
||||
});
|
||||
|
||||
it("uses local when both task and project values are empty", () => {
|
||||
expect(resolveEffectiveNode({ nodeId: "" }, { defaultNodeId: "" })).toEqual({
|
||||
nodeId: undefined,
|
||||
source: "local",
|
||||
});
|
||||
});
|
||||
|
||||
it("uses task override when set even if project default is empty", () => {
|
||||
expect(resolveEffectiveNode({ nodeId: "node-task" }, { defaultNodeId: "" })).toEqual({
|
||||
nodeId: "node-task",
|
||||
source: "task-override",
|
||||
});
|
||||
});
|
||||
});
|
||||
130
packages/engine/src/__tests__/scheduler-node-routing.test.ts
Normal file
130
packages/engine/src/__tests__/scheduler-node-routing.test.ts
Normal file
@@ -0,0 +1,130 @@
|
||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import type { Task, TaskStore } from "@fusion/core";
|
||||
import { Scheduler } from "../scheduler.js";
|
||||
import { existsSync } from "node:fs";
|
||||
import { readFile } from "node:fs/promises";
|
||||
import { schedulerLog } from "../logger.js";
|
||||
|
||||
vi.mock("node:fs", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("node:fs")>();
|
||||
return {
|
||||
...actual,
|
||||
existsSync: vi.fn(),
|
||||
};
|
||||
});
|
||||
|
||||
vi.mock("node:fs/promises", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("node:fs/promises")>();
|
||||
return {
|
||||
...actual,
|
||||
readFile: vi.fn(),
|
||||
};
|
||||
});
|
||||
|
||||
vi.mock("../logger.js", () => ({
|
||||
schedulerLog: {
|
||||
log: vi.fn(),
|
||||
warn: vi.fn(),
|
||||
error: vi.fn(),
|
||||
},
|
||||
}));
|
||||
|
||||
function createMockTask(overrides: Partial<Task> = {}): Task {
|
||||
return {
|
||||
id: "FN-100",
|
||||
description: "Node routing task",
|
||||
column: "todo",
|
||||
dependencies: [],
|
||||
steps: [],
|
||||
currentStep: 0,
|
||||
log: [],
|
||||
createdAt: "2026-01-01T00:00:00Z",
|
||||
updatedAt: "2026-01-01T00:00:00Z",
|
||||
prompt: "",
|
||||
...overrides,
|
||||
} as Task;
|
||||
}
|
||||
|
||||
function createMockStore(task: Task, settings: Record<string, unknown> = {}): TaskStore {
|
||||
return {
|
||||
listTasks: vi.fn().mockResolvedValue([task]),
|
||||
getSettings: vi.fn().mockResolvedValue(settings),
|
||||
getTask: vi.fn().mockResolvedValue(task),
|
||||
updateTask: vi.fn().mockResolvedValue(undefined),
|
||||
moveTask: vi.fn().mockResolvedValue(undefined),
|
||||
parseFileScopeFromPrompt: vi.fn().mockResolvedValue([]),
|
||||
logEntry: vi.fn().mockResolvedValue(undefined),
|
||||
getRootDir: vi.fn().mockReturnValue("/tmp/test"),
|
||||
getTasksDir: vi.fn().mockReturnValue("/tmp/test/.fusion/tasks"),
|
||||
on: vi.fn(),
|
||||
off: vi.fn(),
|
||||
} as unknown as TaskStore;
|
||||
}
|
||||
|
||||
describe("Scheduler node routing", () => {
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks();
|
||||
vi.mocked(existsSync).mockReturnValue(true);
|
||||
vi.mocked(readFile).mockResolvedValue("# Task\nNode routing");
|
||||
});
|
||||
|
||||
it("stores task override as effective node", async () => {
|
||||
const task = createMockTask({ id: "FN-101", nodeId: "node-task" });
|
||||
const store = createMockStore(task, { maxConcurrent: 1, maxWorktrees: 1, defaultNodeId: "node-project" });
|
||||
const scheduler = new Scheduler(store);
|
||||
(scheduler as unknown as { running: boolean }).running = true;
|
||||
|
||||
await scheduler.schedule();
|
||||
|
||||
expect(store.updateTask).toHaveBeenCalledWith(task.id, expect.objectContaining({
|
||||
effectiveNodeId: "node-task",
|
||||
effectiveNodeSource: "task-override",
|
||||
}));
|
||||
expect(store.logEntry).toHaveBeenCalledWith(task.id, "Node routing resolved: node-task (source: task-override)");
|
||||
expect(schedulerLog.log).toHaveBeenCalledWith("Task FN-101 routed to node=node-task (source=task-override)");
|
||||
});
|
||||
|
||||
it("uses project default when task nodeId is unset", async () => {
|
||||
const task = createMockTask({ id: "FN-102", nodeId: undefined });
|
||||
const store = createMockStore(task, { maxConcurrent: 1, maxWorktrees: 1, defaultNodeId: "node-project" });
|
||||
const scheduler = new Scheduler(store);
|
||||
(scheduler as unknown as { running: boolean }).running = true;
|
||||
|
||||
await scheduler.schedule();
|
||||
|
||||
expect(store.updateTask).toHaveBeenCalledWith(task.id, expect.objectContaining({
|
||||
effectiveNodeId: "node-project",
|
||||
effectiveNodeSource: "project-default",
|
||||
}));
|
||||
expect(store.logEntry).toHaveBeenCalledWith(task.id, "Node routing resolved: node-project (source: project-default)");
|
||||
});
|
||||
|
||||
it("uses local when neither task nor project default are set", async () => {
|
||||
const task = createMockTask({ id: "FN-103", nodeId: undefined });
|
||||
const store = createMockStore(task, { maxConcurrent: 1, maxWorktrees: 1 });
|
||||
const scheduler = new Scheduler(store, { nodeHealthMonitor: undefined });
|
||||
(scheduler as unknown as { running: boolean }).running = true;
|
||||
|
||||
await scheduler.schedule();
|
||||
|
||||
expect(store.updateTask).toHaveBeenCalledWith(task.id, expect.objectContaining({
|
||||
effectiveNodeId: null,
|
||||
effectiveNodeSource: "local",
|
||||
}));
|
||||
expect(store.logEntry).toHaveBeenCalledWith(task.id, "Node routing resolved: local (source: local)");
|
||||
expect(schedulerLog.log).toHaveBeenCalledWith("Task FN-103 routed to node=local (source=local)");
|
||||
});
|
||||
|
||||
it("accepts nodeHealthMonitor option at construction", () => {
|
||||
const task = createMockTask();
|
||||
const store = createMockStore(task);
|
||||
|
||||
const scheduler = new Scheduler(store, {
|
||||
nodeHealthMonitor: {
|
||||
getNodeHealth: vi.fn(),
|
||||
} as unknown as import("../node-health-monitor.js").NodeHealthMonitor,
|
||||
});
|
||||
|
||||
expect(scheduler).toBeDefined();
|
||||
});
|
||||
});
|
||||
@@ -818,6 +818,8 @@ describe("Scheduler", () => {
|
||||
blockedBy: null,
|
||||
baseBranch: undefined,
|
||||
worktree: "/test/project/.worktrees/fn-010",
|
||||
effectiveNodeId: null,
|
||||
effectiveNodeSource: "local",
|
||||
});
|
||||
expect(moveTask).toHaveBeenCalledWith("FN-010", "in-progress");
|
||||
expect(updateTask.mock.invocationCallOrder[0]).toBeLessThan(moveTask.mock.invocationCallOrder[0]);
|
||||
@@ -854,12 +856,16 @@ describe("Scheduler", () => {
|
||||
blockedBy: null,
|
||||
baseBranch: undefined,
|
||||
worktree: "/test/project/.worktrees/amber-aspen",
|
||||
effectiveNodeId: null,
|
||||
effectiveNodeSource: "local",
|
||||
});
|
||||
expect(updateTask).toHaveBeenNthCalledWith(2, "FN-012", {
|
||||
status: null,
|
||||
blockedBy: null,
|
||||
baseBranch: undefined,
|
||||
worktree: "/test/project/.worktrees/amber-aspen-2",
|
||||
effectiveNodeId: null,
|
||||
effectiveNodeSource: "local",
|
||||
});
|
||||
|
||||
randomSpy.mockRestore();
|
||||
|
||||
27
packages/engine/src/effective-node.ts
Normal file
27
packages/engine/src/effective-node.ts
Normal file
@@ -0,0 +1,27 @@
|
||||
import type { ProjectSettings, Task } from "@fusion/core";
|
||||
|
||||
export type EffectiveNodeSource = "task-override" | "project-default" | "local";
|
||||
|
||||
export interface EffectiveNode {
|
||||
nodeId: string | undefined;
|
||||
source: EffectiveNodeSource;
|
||||
}
|
||||
|
||||
function isSetNodeId(nodeId: string | null | undefined): nodeId is string {
|
||||
return typeof nodeId === "string" && nodeId.trim().length > 0;
|
||||
}
|
||||
|
||||
export function resolveEffectiveNode(
|
||||
task: Pick<Task, "nodeId">,
|
||||
settings: Pick<ProjectSettings, "defaultNodeId">,
|
||||
): EffectiveNode {
|
||||
if (isSetNodeId(task.nodeId)) {
|
||||
return { nodeId: task.nodeId, source: "task-override" };
|
||||
}
|
||||
|
||||
if (isSetNodeId(settings.defaultNodeId)) {
|
||||
return { nodeId: settings.defaultNodeId, source: "project-default" };
|
||||
}
|
||||
|
||||
return { nodeId: undefined, source: "local" };
|
||||
}
|
||||
@@ -17,6 +17,7 @@ import { schedulerLog } from "./logger.js";
|
||||
import { type PrMonitor, type PrComment } from "./pr-monitor.js";
|
||||
import { reconcileMissionFeatureState } from "./mission-feature-sync.js";
|
||||
import { evaluateSpecStaleness, getPromptPath } from "./spec-staleness.js";
|
||||
import { resolveEffectiveNode } from "./effective-node.js";
|
||||
|
||||
/**
|
||||
* Check whether two sets of file scope paths overlap.
|
||||
@@ -129,6 +130,10 @@ export interface SchedulerOptions {
|
||||
) => void | Promise<void>;
|
||||
/** Optional MissionExecutionLoop for validation cycle handling */
|
||||
missionExecutionLoop?: import("./mission-execution-loop.js").MissionExecutionLoop;
|
||||
/** Optional NodeHealthMonitor for node health checks during dispatch.
|
||||
* Reserved for FN-2722-C (unavailable node policy enforcement).
|
||||
* Accepted here so the option can be wired at construction time. */
|
||||
nodeHealthMonitor?: import("./node-health-monitor.js").NodeHealthMonitor;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -748,6 +753,10 @@ export class Scheduler {
|
||||
continue;
|
||||
}
|
||||
|
||||
// Resolve effective node for routing
|
||||
const effectiveNode = resolveEffectiveNode(freshTask, settings);
|
||||
schedulerLog.log(`Task ${task.id} routed to node=${effectiveNode.nodeId ?? "local"} (source=${effectiveNode.source})`);
|
||||
|
||||
// Clear status, reserve worktree path, and then move to in-progress
|
||||
schedulerLog.log(`Starting ${task.id}: ${task.title || task.id} (deps satisfied)`);
|
||||
await this.store.updateTask(task.id, {
|
||||
@@ -755,8 +764,11 @@ export class Scheduler {
|
||||
blockedBy: null,
|
||||
baseBranch: baseBranch ?? undefined,
|
||||
worktree: plannedWorktree,
|
||||
effectiveNodeId: effectiveNode.nodeId ?? null,
|
||||
effectiveNodeSource: effectiveNode.source,
|
||||
});
|
||||
await this.store.moveTask(task.id, "in-progress");
|
||||
await this.store.logEntry(task.id, `Node routing resolved: ${effectiveNode.nodeId ?? "local"} (source: ${effectiveNode.source})`);
|
||||
this.options.onSchedule?.(task);
|
||||
started++;
|
||||
|
||||
|
||||
@@ -12,6 +12,8 @@ export default defineConfig({
|
||||
alias: {
|
||||
"@fusion/core": resolve(__dirname, "../core/src/index.ts"),
|
||||
"@fusion/test-utils": resolve(__dirname, "../core/src/__test-utils__/workspace.ts"),
|
||||
"@fusion/engine": resolve(__dirname, "./src/index.ts"),
|
||||
"@fusion/plugin-sdk": resolve(__dirname, "../plugin-sdk/src/index.ts"),
|
||||
},
|
||||
},
|
||||
test: {
|
||||
|
||||
Reference in New Issue
Block a user