FN-8760: publish workflow lanes for new tasks
Publish durable workflow lane metadata when new tasks wake triage. - Defer task-created events until selected workflow lanes are durable. - Cache and emit resolved lifecycle lanes while preserving listener compatibility. - Prevent proposal-claim replays from issuing duplicate triage wakes and cover custom lanes. Files changed: .changeset/fn-8760-planning-wake.md | 7 ++++ .../postgres/task-proposal-claim.pg.test.ts | 48 +++++++++++++++++++++- .../__tests__/task-updated-lanes-payload.test.ts | 20 ++++++++- packages/core/src/store.ts | 12 ++++-- packages/core/src/task-store/task-creation.ts | 44 ++++++++++++++++---- packages/engine/src/__tests__/triage.test.ts | 34 ++++++++++++++- 6 files changed, 152 insertions(+), 13 deletions(-) Fusion-Task-Id: FN-8760 Fusion-Task-Lineage: 06275d4a-3fc0-405e-baa5-3cdc8d195695 Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-8760-planning-wake.md
Normal file
7
.changeset/fn-8760-planning-wake.md
Normal file
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@runfusion/fusion": patch
|
||||
---
|
||||
|
||||
summary: Schedule Planning Mode-created tasks promptly in selected workflow lanes.
|
||||
category: fix
|
||||
dev: Publish task creation lifecycle lanes after durable workflow selection; replayed proposal claims do not re-wake triage.
|
||||
@@ -7,6 +7,7 @@ already-materialized task rather than surfacing 23505 or creating another row.
|
||||
*/
|
||||
|
||||
import { afterAll, afterEach, beforeAll, beforeEach, expect, it } from "vitest";
|
||||
import type { WorkflowIr } from "../../workflows/workflow-ir-types.js";
|
||||
import {
|
||||
createSharedPgTaskStoreTestHarness,
|
||||
pgDescribe,
|
||||
@@ -15,6 +16,36 @@ import {
|
||||
|
||||
const pgTest = pgDescribe;
|
||||
|
||||
function planningWorkflowIr(): WorkflowIr {
|
||||
return {
|
||||
version: "v2",
|
||||
name: "planning-wake-lanes",
|
||||
columns: [
|
||||
{ id: "planning-inbox", name: "Planning inbox", traits: [{ trait: "intake" }] },
|
||||
{ id: "ready-to-plan", name: "Ready to plan", traits: [{ trait: "hold", config: { release: "capacity" } }] },
|
||||
{ id: "building", name: "Building", traits: [{ trait: "wip", config: { limitSetting: "maxConcurrent" } }] },
|
||||
{ id: "checking", name: "Checking", traits: [{ trait: "merge" }, { trait: "merge-blocker" }] },
|
||||
{ id: "shipped", name: "Shipped", traits: [{ trait: "complete" }] },
|
||||
{ id: "archive", name: "Archive", traits: [{ trait: "archived" }] },
|
||||
],
|
||||
nodes: [
|
||||
{ id: "start", kind: "start", column: "planning-inbox" },
|
||||
{ id: "plan", kind: "prompt", column: "ready-to-plan", config: { name: "Plan", prompt: "Specify." } },
|
||||
{ id: "build", kind: "prompt", column: "building", config: { name: "Build", prompt: "Implement." } },
|
||||
{ id: "review", kind: "prompt", column: "checking", config: { name: "Review", prompt: "Review." } },
|
||||
{ id: "merge", kind: "merge-attempt", column: "checking", config: { capability: "task-merge" } },
|
||||
{ id: "end", kind: "end", column: "shipped" },
|
||||
],
|
||||
edges: [
|
||||
{ from: "start", to: "plan" },
|
||||
{ from: "plan", to: "build", condition: "success" },
|
||||
{ from: "build", to: "review", condition: "success" },
|
||||
{ from: "review", to: "merge", condition: "success" },
|
||||
{ from: "merge", to: "end", condition: "success" },
|
||||
],
|
||||
} as WorkflowIr;
|
||||
}
|
||||
|
||||
pgTest("TaskStore proposal claim idempotency", () => {
|
||||
const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({
|
||||
prefix: "fusion_proposal_claim",
|
||||
@@ -25,20 +56,31 @@ pgTest("TaskStore proposal claim idempotency", () => {
|
||||
afterEach(h.afterEach);
|
||||
afterAll(h.afterAll);
|
||||
|
||||
it("returns the existing task when a reclaim races the original proposal insert", async () => {
|
||||
/*
|
||||
FNXC:PlanningModeScheduling 2026-08-03-09:44:
|
||||
This uses the real PostgreSQL TaskStore boundary rather than a route mock: a selected workflow
|
||||
must be durable before task:created exposes its renamed intake/hold lanes, and a proposal replay
|
||||
must not emit another creation wake that could schedule duplicate planning work.
|
||||
*/
|
||||
it("publishes selected workflow lanes once across proposal-claim replay", async () => {
|
||||
const stableProposalKey = "proposal-reclaim-race-stable-key";
|
||||
const store = h.store();
|
||||
const definition = await store.createWorkflowDefinition({ name: "Planning wake lanes", ir: planningWorkflowIr() });
|
||||
const events: Array<{ id: string; lanes?: { intake?: string; hold?: string } }> = [];
|
||||
store.on("task:created", (task, meta) => events.push({ id: task.id, lanes: meta?.lanes }));
|
||||
|
||||
const [originalCreate, reclaimedCreate] = await Promise.all([
|
||||
store.createTask({
|
||||
title: "Original proposal materialization",
|
||||
description: "Original creator resumes after its lease was released.",
|
||||
proposalClaimId: stableProposalKey,
|
||||
workflowId: definition.id,
|
||||
}),
|
||||
store.createTask({
|
||||
title: "Reclaimed proposal materialization",
|
||||
description: "Reclaimed creator uses the same stable proposal key.",
|
||||
proposalClaimId: stableProposalKey,
|
||||
workflowId: definition.id,
|
||||
}),
|
||||
]);
|
||||
|
||||
@@ -47,5 +89,9 @@ pgTest("TaskStore proposal claim idempotency", () => {
|
||||
const persisted = (await store.listTasks()).filter((task) => task.proposalClaimId === stableProposalKey);
|
||||
expect(persisted).toHaveLength(1);
|
||||
expect(persisted[0]?.id).toBe(originalCreate.id);
|
||||
expect(events).toEqual([{
|
||||
id: originalCreate.id,
|
||||
lanes: expect.objectContaining({ intake: "planning-inbox", hold: "ready-to-plan" }),
|
||||
}]);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -4,7 +4,7 @@ import type { Task } from "../types.js";
|
||||
|
||||
const task = { id: "FN-lanes", column: "building" } as Task;
|
||||
|
||||
describe("task:updated lane payload", () => {
|
||||
describe("task lifecycle lane payload", () => {
|
||||
it("decorates cache hits while keeping one-argument listeners and misses compatible", () => {
|
||||
const store = new TaskStore(process.cwd());
|
||||
const received: Array<{ lanes?: { wip?: string } } | undefined> = [];
|
||||
@@ -31,6 +31,24 @@ describe("task:updated lane payload", () => {
|
||||
expect(received).toEqual({ lanes: { wip: "building" } });
|
||||
});
|
||||
|
||||
/*
|
||||
FNXC:PlanningModeScheduling 2026-08-03-09:44:
|
||||
A created task needs the workflow lanes captured at the durable creation boundary; triage cannot
|
||||
synchronously resolve a custom selection after its wake handler receives the event.
|
||||
*/
|
||||
it("delivers resolved lanes with task:created without changing one-argument listeners", () => {
|
||||
const store = new TaskStore(process.cwd());
|
||||
const received: Array<{ lanes?: { intake?: string; hold?: string } } | undefined> = [];
|
||||
let oneArgumentCalls = 0;
|
||||
store.on("task:created", (_task, meta) => received.push(meta));
|
||||
store.on("task:created", () => { oneArgumentCalls += 1; });
|
||||
|
||||
store.emitTaskLifecycleEventSafely("task:created", [task, { lanes: { intake: "planning-inbox", hold: "ready-to-plan" } }]);
|
||||
|
||||
expect(received).toEqual([{ lanes: { intake: "planning-inbox", hold: "ready-to-plan" } }]);
|
||||
expect(oneArgumentCalls).toBe(1);
|
||||
});
|
||||
|
||||
it("preserves explicit metadata rather than replacing it from cache", () => {
|
||||
const store = new TaskStore(process.cwd());
|
||||
store.laneCache.set(task.id, { wip: "cached" });
|
||||
|
||||
@@ -144,7 +144,13 @@ import type { BranchGroupRow, PrEntityRow, TaskDocumentRow, ArtifactRow, TaskDoc
|
||||
/** Database row shape for the tasks table (all columns). */
|
||||
|
||||
export interface TaskStoreEvents {
|
||||
"task:created": [task: Task];
|
||||
/*
|
||||
FNXC:PlanningModeScheduling 2026-08-03-09:44:
|
||||
Task creation can select a custom workflow whose planning lanes do not use legacy names.
|
||||
Creation emits its resolved lanes after the selection is durable, so synchronous triage wake
|
||||
listeners observe the same authoritative lifecycle vocabulary as task:moved listeners.
|
||||
*/
|
||||
"task:created": [task: Task, meta?: { lanes?: TaskMoveLanes }];
|
||||
/*
|
||||
FNXC:WorkflowEvents 2026-07-31-21:00 (fleet — the emitter carries the lanes):
|
||||
`lanes` is the moving task's RESOLVED lifecycle columns, attached by the emitter.
|
||||
@@ -1046,7 +1052,7 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
||||
/**
|
||||
* FNXC:RuntimeTaskOrchestrationAsync 2026-06-24-13:25:
|
||||
*/
|
||||
public async _createTaskInternalBackend( input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; onProposalClaimConflict?: (task: Task) => void; }, ): Promise<Task> {
|
||||
public async _createTaskInternalBackend( input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; onProposalClaimConflict?: (task: Task) => void; deferTaskCreatedEvent?: boolean; onTaskInserted?: (task: Task) => void; }, ): Promise<Task> {
|
||||
return _createTaskInternalBackendImpl(this, input, title, resolvedWorkflowSteps, id, options);
|
||||
}
|
||||
|
||||
@@ -1062,7 +1068,7 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
||||
async createTaskWithReservedId( input: TaskCreateInput, options: { taskId: string; createdAt?: string; updatedAt?: string; prompt?: string; applyDefaultWorkflowSteps?: boolean; invokeTaskCreatedHook?: boolean; }, ): Promise<Task> {
|
||||
return createTaskWithReservedIdImpl(this, input, options);
|
||||
}
|
||||
public async _createTaskInternal( input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; onProposalClaimConflict?: (task: Task) => void; }, ): Promise<Task> {
|
||||
public async _createTaskInternal( input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; onProposalClaimConflict?: (task: Task) => void; deferTaskCreatedEvent?: boolean; onTaskInserted?: (task: Task) => void; }, ): Promise<Task> {
|
||||
/*
|
||||
FNXC:SqliteDualPathCleanup 2026-07-26-14:05:
|
||||
Task create is PostgreSQL-only (layer.transactionImmediate + insertTaskRowInTransaction). The former sync SQLite _createTaskInternalImpl arm is deleted; production always injects AsyncDataLayer.
|
||||
|
||||
@@ -24,8 +24,8 @@ import {getErrorMessage} from "../process/error-message.js";
|
||||
import {generateTaskLineageId} from "../tasks/task-lineage.js";
|
||||
import {archiveAsSameAgentDuplicate, findSameAgentDuplicates, flagSameAgentDuplicate, type SameAgentDuplicateCandidate} from "../duplicates/duplicate-intake.js";
|
||||
import {buildBootstrapPrompt} from "../mesh/mesh-task-replication.js";
|
||||
import {resolveWorkflowIrById} from "../workflows/workflow-ir-resolver.js";
|
||||
import {resolveTaskLifecycleColumns} from "../workflows/workflow-lifecycle-traits.js";
|
||||
import {resolveWorkflowIrById, resolveWorkflowIrForTask} from "../workflows/workflow-ir-resolver.js";
|
||||
import {resolveTaskLifecycleColumns, toTaskMoveLanes} from "../workflows/workflow-lifecycle-traits.js";
|
||||
import type {WorkflowIr} from "../workflows/workflow-ir-types.js";
|
||||
import {DEFAULT_WORKFLOW_ID} from "../workflows/builtin-workflows.js";
|
||||
import {columnsWithFlag} from "../workflows/workflow-lifecycle-traits.js";
|
||||
@@ -320,6 +320,7 @@ export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateI
|
||||
});
|
||||
|
||||
let task: Task;
|
||||
let insertedTask = false;
|
||||
try {
|
||||
await store.assertNoDependencyCycle(reservation.taskId, input.dependencies ?? [], "createTask");
|
||||
task = await store._createTaskInternalBackend(
|
||||
@@ -327,7 +328,13 @@ export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateI
|
||||
title,
|
||||
resolvedWorkflowSteps,
|
||||
reservation.taskId,
|
||||
{ invokeTaskCreatedHook: shouldInvokeTaskCreatedHook && !hasPendingSummarization, resolvedEntryColumn, onProposalClaimConflict: options?.onProposalClaimConflict },
|
||||
{
|
||||
deferTaskCreatedEvent: true,
|
||||
invokeTaskCreatedHook: shouldInvokeTaskCreatedHook && !hasPendingSummarization,
|
||||
onProposalClaimConflict: options?.onProposalClaimConflict,
|
||||
onTaskInserted: () => { insertedTask = true; },
|
||||
resolvedEntryColumn,
|
||||
},
|
||||
);
|
||||
await allocator.commitDistributedTaskIdReservation({
|
||||
reservationId: reservation.reservationId,
|
||||
@@ -354,6 +361,26 @@ export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateI
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
FNXC:PlanningModeScheduling 2026-08-03-09:44:
|
||||
A task:created listener runs synchronously and is the authoritative wake for triage admission.
|
||||
Planning Mode creates through a project-scoped TaskStore, so emitting before the selected
|
||||
workflow row exists leaves a custom intake lane indistinguishable from an unknown legacy lane.
|
||||
Persist selection first, then publish the resolved lanes through the shared store event; the
|
||||
listener only requests its normal poll, preserving pause, dependency, and capacity gates.
|
||||
|
||||
Proposal-claim replays return an existing row from the internal create path. They deliberately
|
||||
do not re-emit task:created, so idempotent Planning Mode retries cannot schedule duplicate work.
|
||||
*/
|
||||
if (insertedTask) {
|
||||
const lanes = toTaskMoveLanes(await resolveWorkflowIrForTask(store, task.id).catch(() => undefined));
|
||||
store.laneCache.set(task.id, lanes);
|
||||
store.emitTaskLifecycleEventSafely("task:created", [task, { lanes }]);
|
||||
if (shouldInvokeTaskCreatedHook && !hasPendingSummarization) {
|
||||
await store.invokeTaskCreatedHook(task);
|
||||
}
|
||||
}
|
||||
|
||||
// Deferred title summarization (same fire-and-forget pattern as SQLite path).
|
||||
if (hasPendingSummarization && shouldInvokeTaskCreatedHook) {
|
||||
const id = task.id;
|
||||
@@ -411,7 +438,7 @@ export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateI
|
||||
return task;
|
||||
}
|
||||
|
||||
export async function _createTaskInternalBackendImpl(store: TaskStore, input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; onProposalClaimConflict?: (task: Task) => void; },): Promise<Task> {
|
||||
export async function _createTaskInternalBackendImpl(store: TaskStore, input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; onProposalClaimConflict?: (task: Task) => void; deferTaskCreatedEvent?: boolean; onTaskInserted?: (task: Task) => void; },): Promise<Task> {
|
||||
const layer = store.asyncLayer!;
|
||||
const now = options?.createdAt ?? new Date().toISOString();
|
||||
const normalizedTitle = normalizeTitleForTaskId(title, id);
|
||||
@@ -707,10 +734,13 @@ export async function _createTaskInternalBackendImpl(store: TaskStore, input: Ta
|
||||
await store._maybeAutoArchiveSameAgentDuplicateBackend(task, input);
|
||||
}
|
||||
|
||||
store.emitTaskLifecycleEventSafely("task:created", [task]);
|
||||
if (options?.invokeTaskCreatedHook !== false) {
|
||||
await store.invokeTaskCreatedHook(task);
|
||||
if (!options?.deferTaskCreatedEvent) {
|
||||
store.emitTaskLifecycleEventSafely("task:created", [task]);
|
||||
if (options?.invokeTaskCreatedHook !== false) {
|
||||
await store.invokeTaskCreatedHook(task);
|
||||
}
|
||||
}
|
||||
options?.onTaskInserted?.(task);
|
||||
return task;
|
||||
}
|
||||
|
||||
|
||||
@@ -882,7 +882,7 @@ describe("canonical triage policy prompt", () => {
|
||||
|
||||
describe("FN-5893 invariant regression wording", () => {
|
||||
const corePromptSource = readFileSync(
|
||||
fileURLToPath(new URL("../../../core/src/agent-prompts.ts", import.meta.url)),
|
||||
fileURLToPath(new URL("../../../core/src/agents/agent-prompts.ts", import.meta.url)),
|
||||
"utf8",
|
||||
);
|
||||
|
||||
@@ -1759,6 +1759,38 @@ Planner rewrote mission without the raw request.
|
||||
// Should not throw
|
||||
});
|
||||
|
||||
/*
|
||||
FNXC:PlanningModeScheduling 2026-08-03-09:44:
|
||||
Planning Mode may create into a selected workflow whose intake and hold lanes are renamed.
|
||||
The creation event must carry that durable lane answer so the normal wake reaches triage without
|
||||
a route-specific scheduler call; paused work remains gated before the poll is requested.
|
||||
*/
|
||||
it("wakes normal planning admission for a created task in a custom workflow lane", () => {
|
||||
const listeners = new Map<string, (...args: unknown[]) => void>();
|
||||
const triageStore = createMockStore({
|
||||
on: vi.fn((event: string, listener: (...args: unknown[]) => void) => {
|
||||
listeners.set(event, listener);
|
||||
return triageStore;
|
||||
}),
|
||||
off: vi.fn(),
|
||||
});
|
||||
const triageProcessor = trackProcessor(new TriageProcessor(triageStore, rootDir));
|
||||
const requestImmediatePoll = vi.spyOn(triageProcessor, "requestImmediatePoll").mockReturnValue(true);
|
||||
|
||||
triageProcessor.start();
|
||||
listeners.get("task:created")?.(
|
||||
createTriageTask({ id: "FN-PLANNING-CREATED", column: "planning-inbox" }),
|
||||
{ lanes: { intake: "planning-inbox", hold: "ready-to-plan" } },
|
||||
);
|
||||
listeners.get("task:created")?.(
|
||||
createTriageTask({ id: "FN-PLANNING-PAUSED", column: "planning-inbox", paused: true }),
|
||||
{ lanes: { intake: "planning-inbox", hold: "ready-to-plan" } },
|
||||
);
|
||||
|
||||
expect(requestImmediatePoll).toHaveBeenCalledTimes(1);
|
||||
triageProcessor.stop();
|
||||
});
|
||||
|
||||
it("handles settings:updated event for globalPause", () => {
|
||||
const handler = vi.fn();
|
||||
(store.on as ReturnType<typeof vi.fn>).mockImplementation(
|
||||
|
||||
Reference in New Issue
Block a user