fix(FN-2384): schedule triage and todo tasks by priority-aware FIFO

- Sort eligible todo tasks in Scheduler by priority first, then createdAt/id for stable FIFO ordering within each tier
- Sort eligible triage tasks using the same priority-aware ordering while preserving existing pause/status/recovery gating
- Add regression tests for scheduler and triage ordering, including blocked/paused/recovery-gated edge cases
- Update architecture docs to describe priority-first task dispatch behavior
This commit is contained in:
Fusion
2026-04-24 04:08:17 -07:00
committed by gsxdsm
parent 6453fc771e
commit d74295f7dc
5 changed files with 219 additions and 6 deletions

View File

@@ -557,6 +557,102 @@ describe("Scheduler", () => {
});
});
describe("priority-aware todo dispatch", () => {
it("schedules eligible todo tasks by priority desc then createdAt asc", async () => {
vi.mocked(existsSync).mockReturnValue(true);
vi.mocked(readFile).mockResolvedValue("# Task\nDo something");
const tasks = [
createMockTask({ id: "FN-010", column: "todo", priority: "normal", createdAt: "2026-01-01T00:02:00.000Z" }),
createMockTask({ id: "FN-011", column: "todo", priority: "urgent", createdAt: "2026-01-01T00:10:00.000Z" }),
createMockTask({ id: "FN-012", column: "todo", priority: "high", createdAt: "2026-01-01T00:03:00.000Z" }),
createMockTask({ id: "FN-013", column: "todo", priority: "high", createdAt: "2026-01-01T00:01:00.000Z" }),
];
const store = createMockStore({
listTasks: vi.fn().mockResolvedValue(tasks),
getSettings: vi.fn().mockResolvedValue({
maxConcurrent: 10,
maxWorktrees: 10,
groupOverlappingFiles: false,
}),
updateTask: vi.fn().mockResolvedValue(undefined),
moveTask: vi.fn().mockResolvedValue(undefined),
});
const scheduler = new Scheduler(store);
(scheduler as any).running = true;
await scheduler.schedule();
expect((store.moveTask as ReturnType<typeof vi.fn>).mock.calls.map((call: unknown[]) => call[0])).toEqual([
"FN-011",
"FN-013",
"FN-012",
"FN-010",
]);
});
it("keeps blocked high-priority todo tasks unscheduled while scheduling ready lower-priority work", async () => {
vi.mocked(existsSync).mockReturnValue(true);
vi.mocked(readFile).mockResolvedValue("# Task\nDo something");
const future = new Date(Date.now() + 60_000).toISOString();
const tasks = [
createMockTask({ id: "FN-001", column: "in-progress" }),
createMockTask({ id: "FN-100", column: "todo", priority: "urgent", dependencies: ["FN-900"] }),
createMockTask({ id: "FN-900", column: "todo", priority: "low" }),
createMockTask({ id: "FN-101", column: "todo", priority: "urgent", paused: true }),
createMockTask({ id: "FN-102", column: "todo", priority: "urgent", nextRecoveryAt: future }),
createMockTask({ id: "FN-103", column: "todo", priority: "urgent" }),
createMockTask({ id: "FN-104", column: "todo", priority: "normal" }),
];
const parseScopeMock = vi.fn(async (taskId: string): Promise<string[]> => {
if (taskId === "FN-001" || taskId === "FN-103") {
return ["packages/engine/src/scheduler.ts"];
}
if (taskId === "FN-900") {
return ["packages/core/src/store.ts"];
}
if (taskId === "FN-104") {
return ["packages/engine/src/triage.ts"];
}
return ["packages/engine/src/logger.ts"];
});
const updateTask = vi.fn().mockResolvedValue(undefined);
const moveTask = vi.fn().mockResolvedValue(undefined);
const store = createMockStore({
listTasks: vi.fn().mockResolvedValue(tasks),
getSettings: vi.fn().mockResolvedValue({
maxConcurrent: 10,
maxWorktrees: 10,
groupOverlappingFiles: true,
}),
parseFileScopeFromPrompt: parseScopeMock,
updateTask,
moveTask,
});
const scheduler = new Scheduler(store);
(scheduler as any).running = true;
await scheduler.schedule();
// Dependency-blocked urgent task should be queued, not started.
expect(updateTask).toHaveBeenCalledWith("FN-100", { status: "queued" });
// Overlap-blocked urgent task should be queued with blocker id.
expect(updateTask).toHaveBeenCalledWith("FN-103", { status: "queued", blockedBy: "FN-001" });
// Paused and recovery-gated urgent tasks never enter scheduling.
expect(moveTask).not.toHaveBeenCalledWith("FN-101", "in-progress");
expect(moveTask).not.toHaveBeenCalledWith("FN-102", "in-progress");
// Lower-priority ready task still runs.
expect(moveTask).toHaveBeenCalledWith("FN-104", "in-progress");
// Overlap-blocked urgent task must not run.
expect(moveTask).not.toHaveBeenCalledWith("FN-103", "in-progress");
});
});
describe("worktree reservation", () => {
it("assigns a planned worktree path before moving a task to in-progress", async () => {
vi.mocked(existsSync).mockReturnValue(true);

View File

@@ -1,4 +1,13 @@
import { getCurrentRepo, resolveDependencyOrder, type TaskStore, type Task, type MissionStore, type MissionFeature, type PrInfo } from "@fusion/core";
import {
getCurrentRepo,
resolveDependencyOrder,
sortTasksByPriorityThenAgeAndId,
type TaskStore,
type Task,
type MissionStore,
type MissionFeature,
type PrInfo,
} from "@fusion/core";
import { existsSync } from "node:fs";
import { readFile } from "node:fs/promises";
import { join } from "node:path";
@@ -563,6 +572,8 @@ export class Scheduler {
if (todo.length === 0) return;
todo = sortTasksByPriorityThenAgeAndId(todo);
/**
* Pre-compute file scopes for all currently active tasks (in-progress
* AND in-review with unmerged worktrees) so that todo tasks are never

View File

@@ -103,6 +103,21 @@ const mockTaskDetail: TaskDetail = {
attachments: [],
};
function createTriageTask(overrides: Partial<Task> = {}): Task {
return {
id: "FN-001",
description: "Triage task",
column: "triage",
dependencies: [],
steps: [],
currentStep: 0,
log: [],
createdAt: "2026-01-01T00:00:00.000Z",
updatedAt: "2026-01-01T00:00:00.000Z",
...overrides,
};
}
describe("buildSpecificationPrompt", () => {
const baseTask: TaskDetail = {
...mockTaskDetail,
@@ -710,7 +725,93 @@ describe("TriageProcessor", () => {
expect(store.on).toHaveBeenCalledWith("settings:updated", expect.any(Function));
});
it("re-reads settings when fn_review_spec runs so reviewer uses the latest validator model", async () => {
describe("poll ordering", () => {
it("dispatches eligible triage tasks by priority desc then createdAt asc", async () => {
const tasks: Task[] = [
createTriageTask({
id: "FN-100",
priority: "normal",
createdAt: "2026-01-01T00:01:00.000Z",
}),
createTriageTask({
id: "FN-101",
priority: "urgent",
createdAt: "2026-01-01T00:10:00.000Z",
}),
createTriageTask({
id: "FN-102",
priority: "high",
createdAt: "2026-01-01T00:03:00.000Z",
}),
createTriageTask({
id: "FN-103",
priority: "high",
createdAt: "2026-01-01T00:02:00.000Z",
}),
];
const triageStore = createMockStore({
listTasks: vi.fn().mockResolvedValue(tasks),
getSettings: vi.fn().mockResolvedValue({
maxConcurrent: 10,
maxTriageConcurrent: 10,
pollIntervalMs: 10_000,
groupOverlappingFiles: false,
autoMerge: true,
}),
});
const triageProcessor = new TriageProcessor(triageStore, rootDir);
const specifySpy = vi
.spyOn(triageProcessor, "specifyTask")
.mockResolvedValue(undefined);
(triageProcessor as any).running = true;
await (triageProcessor as any).poll();
expect(specifySpy).toHaveBeenCalledTimes(4);
expect(specifySpy.mock.calls.map(([task]) => task.id)).toEqual([
"FN-101",
"FN-103",
"FN-102",
"FN-100",
]);
});
it("excludes paused, awaiting-approval, failed, stuck-killed, and recovery-gated tasks from ordered candidates", async () => {
const future = new Date(Date.now() + 60_000).toISOString();
const tasks: Task[] = [
createTriageTask({ id: "FN-200", priority: "urgent" }),
createTriageTask({ id: "FN-201", priority: "urgent", paused: true }),
createTriageTask({ id: "FN-202", priority: "urgent", status: "awaiting-approval" }),
createTriageTask({ id: "FN-203", priority: "urgent", status: "failed" }),
createTriageTask({ id: "FN-204", priority: "urgent", status: "stuck-killed" }),
createTriageTask({ id: "FN-205", priority: "urgent", nextRecoveryAt: future }),
];
const triageStore = createMockStore({
listTasks: vi.fn().mockResolvedValue(tasks),
getSettings: vi.fn().mockResolvedValue({
maxConcurrent: 10,
maxTriageConcurrent: 10,
pollIntervalMs: 10_000,
groupOverlappingFiles: false,
autoMerge: true,
}),
});
const triageProcessor = new TriageProcessor(triageStore, rootDir);
const specifySpy = vi
.spyOn(triageProcessor, "specifyTask")
.mockResolvedValue(undefined);
(triageProcessor as any).running = true;
await (triageProcessor as any).poll();
expect(specifySpy).toHaveBeenCalledTimes(1);
expect(specifySpy).toHaveBeenCalledWith(expect.objectContaining({ id: "FN-200" }));
});
});
it("re-reads settings when review_spec runs so reviewer uses the latest validator model", async () => {
const taskId = "FN-001";
const testRootDir = await createTriageFixtureRoot("fusion-triage-review-spec-");
try {

View File

@@ -6,7 +6,11 @@ import type {
TaskAttachment,
Settings,
} from "@fusion/core";
import { buildTriageMemoryInstructions, resolveAgentPrompt } from "@fusion/core";
import {
buildTriageMemoryInstructions,
resolveAgentPrompt,
sortTasksByPriorityThenAgeAndId,
} from "@fusion/core";
import type { ImageContent } from "@mariozechner/pi-ai";
import { Type, type Static } from "@mariozechner/pi-ai";
import type {
@@ -577,7 +581,7 @@ export class TriageProcessor {
// Fetch all tasks (not just triage) to count active agents across columns.
const allTasks = await this.store.listTasks({ slim: true, includeArchived: false });
const now = Date.now();
const triageTasks = allTasks.filter(
const eligibleTriageTasks = allTasks.filter(
(t) => t.column === "triage" && !this.processing.has(t.id) && !t.paused
// Skip tasks awaiting manual plan approval — they should not be auto-discovered
&& t.status !== "awaiting-approval"
@@ -587,6 +591,7 @@ export class TriageProcessor {
// Skip tasks with a recovery backoff that hasn't elapsed yet
&& !(t.nextRecoveryAt && new Date(t.nextRecoveryAt).getTime() > now),
);
const triageTasks = sortTasksByPriorityThenAgeAndId(eligibleTriageTasks);
// Respect both per-project maxTriageConcurrent and the global semaphore.
// Only specifying tasks count against the triage limit; execution is governed by maxConcurrent.