From 93437ff09bd345192ab50fd6ba844c14bd6e54e4 Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Sun, 19 Jul 2026 22:07:46 -0700 Subject: [PATCH] FN-8306: add mission symbol lock scheduler admission Schedule approved mission work with durable symbol-level concurrency control. - Gate mission execution on active lineage and required plan approval - Acquire, renew, and release symbol locks across scheduler and workflow execution - Honor custom workflow WIP columns when maintaining symbol-lock leases - Add admission, contention, renewal, and custom-WIP regression coverage Files changed: .../fn-8306-mission-symbol-scheduler-admission.md | 7 ++ docs/architecture.md | 1 + packages/core/src/task-store/moves.ts | 23 ++++ .../src/__tests__/mission-symbol-admission.test.ts | 54 +++++++++ .../__tests__/scheduler-workflow-cutover.test.ts | 49 ++++++++ .../src/__tests__/workflow-work-processor.test.ts | 39 ++++++ .../src/__tests__/workflow-work-scheduler.test.ts | 44 +++++++ packages/engine/src/mission-symbol-admission.ts | 94 +++++++++++++++ packages/engine/src/scheduler.ts | 133 ++++++++++++++++++++- packages/engine/src/workflow-work-processor.ts | 22 +++- packages/engine/src/workflow-work-scheduler.ts | 62 +++++++++- 11 files changed, 520 insertions(+), 8 deletions(-) Fusion-Task-Id: FN-8306 Fusion-Task-Lineage: ec6a9928-d8c0-4cf7-99b5-47e1fc3a2dcc Co-authored-by: Fusion (runfusion.ai) --- ...8306-mission-symbol-scheduler-admission.md | 7 + docs/architecture.md | 1 + packages/core/src/task-store/moves.ts | 23 +++ .../mission-symbol-admission.test.ts | 54 +++++++ .../scheduler-workflow-cutover.test.ts | 49 +++++++ .../__tests__/workflow-work-processor.test.ts | 39 +++++ .../__tests__/workflow-work-scheduler.test.ts | 44 ++++++ .../engine/src/mission-symbol-admission.ts | 94 +++++++++++++ packages/engine/src/scheduler.ts | 133 +++++++++++++++++- .../engine/src/workflow-work-processor.ts | 22 ++- .../engine/src/workflow-work-scheduler.ts | 62 +++++++- 11 files changed, 520 insertions(+), 8 deletions(-) create mode 100644 .changeset/fn-8306-mission-symbol-scheduler-admission.md create mode 100644 packages/engine/src/__tests__/mission-symbol-admission.test.ts create mode 100644 packages/engine/src/__tests__/workflow-work-processor.test.ts create mode 100644 packages/engine/src/__tests__/workflow-work-scheduler.test.ts create mode 100644 packages/engine/src/mission-symbol-admission.ts diff --git a/.changeset/fn-8306-mission-symbol-scheduler-admission.md b/.changeset/fn-8306-mission-symbol-scheduler-admission.md new file mode 100644 index 0000000000..52ac09e2f5 --- /dev/null +++ b/.changeset/fn-8306-mission-symbol-scheduler-admission.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": minor +--- + +summary: Schedule approved mission work with symbol-level concurrency control. +category: feature +dev: Enforces mission lineage admission and releases durable symbol locks on lifecycle exits. diff --git a/docs/architecture.md b/docs/architecture.md index 062fa1df49..365dff5e1a 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -620,6 +620,7 @@ See [Memory Plugin Contract](./memory-plugin-contract.md) for the full plan. - `blockedBy` invariant (FN-3924/FN-4091): the field is only durable when it references a current unresolved explicit dependency (or, for dependency-free tasks, an active overlap blocker). Completion gating now validates `blockedBy` through live task resolution: missing blockers and blockers already in `done`/`archived` are treated as stale, while only still-active blockers continue to prevent `fn_task_done`. If no current blocker remains, scheduler/event reconciliation clears `blockedBy` to `null` and re-evaluates from live task state. - Dependency-cycle invariant (FN-5256): task dependency graphs are acyclic at write time (`DependencyCycleError` in `TaskStore` for `createTask`, `createTaskWithReservedId`, `updateTask`, and `applyReplicatedTaskCreate`) with `task:dependency-cycle-rejected` audit evidence. Self-healing batch 2 adds `reconcileDependencyCycles`, which emits `task:dependency-cycle-detected`, auto-repairs only bounded umbrella-back-edge loops via `task:auto-reconciled-dependency-cycle`, and leaves ambiguous cycles untouched with `task:dependency-cycle-unrepaired` for operator inspection. - Dependency-blocking lease invariant (FN-6292): an `in-progress` task with unmet scheduling dependencies must not contribute an active file-scope lease in scheduler lease maps. This prevents a holder from queueing its own dependency behind its lease and creating a circular wait. + - **Mission symbol admission (FN-8306):** autonomous mission implementation evaluates `evaluateMissionLineageApproval` (active Mission/Milestone/Slice, triaged or in-progress Feature, and required plan fingerprint). Approved work with durable declared symbols atomically acquires project-scoped locks after ordinary capacity, checkout, and dependency gates; same symbols queue with holder diagnostics while disjoint symbols can run in parallel. Mission-linked work lacking approved lineage is `lineage-blocked` and consumes neither symbol nor work lease. Non-mission work and approved mission work without resolvable symbols retain coarse file-scope serialization. Active scheduler heartbeats and long-running workflow processors renew their short crash-recoverable leases before expiry. Central `moveTask` exits from implementation release declared-symbol locks for review, cancellation/requeue, and terminal transitions; self-healing only expires terminal/missing/expired owners as a backstop. #### BlockedBy stamping invariants - Scheduler writes overlap-based `blockedBy` only when overlap gating is active and there is a live overlapping active scope; otherwise overlap logic does not stamp blockers. diff --git a/packages/core/src/task-store/moves.ts b/packages/core/src/task-store/moves.ts index 0dae1c38c8..8c95314bbf 100644 --- a/packages/core/src/task-store/moves.ts +++ b/packages/core/src/task-store/moves.ts @@ -35,6 +35,7 @@ import {getTaskMergeBlocker} from "../task-merge.js"; import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js"; import {readTaskRow as readTaskRowAsync, readTaskRowInTransaction, upsertTaskRowInTransaction} from "../task-store/async-persistence.js"; import {disposeTaskBeforeMove} from "../task-move-disposer.js"; +import {resolveTaskSymbolsForTask} from "../task-symbol-resolution.js"; /* FNXC:PostgresCutover 2026-07-05-19:50: @@ -1135,6 +1136,28 @@ export async function moveTaskInternalImpl(store: TaskStore, id: string, toColum } await store.writeTaskJsonFile(dir, task); + + /* + FNXC:MissionSymbolAdmission 2026-07-19-22:04: + FN-8306 makes the central lifecycle transition the primary symbol-lock + release authority. A durable lock belongs only to active implementation; + handoff to review, cancellation/requeue, and terminal moves therefore release + the task's declared symbols here. Active implementation is the workflow's + countsTowardWip trait, not the legacy `in-progress` id: custom WIP columns + must retain and renew locks between internal WIP transitions, then release + them only when leaving WIP. PostgreSQL owns durable locks; SQLite has no + symbol-lock table and remains on coarse scheduling behavior. + */ + const lifecycleWorkflowIr = workflowIr ?? await resolveTaskWorkflowIrForMove(store, id); + const fromIsImplementation = resolveTransitionColumnFacts(lifecycleWorkflowIr, fromColumn).flags.countsTowardWip === true; + const toIsImplementation = resolveTransitionColumnFacts(lifecycleWorkflowIr, toColumn).flags.countsTowardWip === true; + if (store.backendMode && fromIsImplementation && !toIsImplementation) { + const symbols = resolveTaskSymbolsForTask(task); + if (symbols.resolvable) { + await store.releaseSymbolLocks(symbols.symbols, id); + } + } + if (fromColumn === "in-review" && toColumn === "todo" && moveSource === "user") { const handoffAccepted = await store.getCompletionHandoffAcceptedMarker(id); const mergeRequest = await store.getMergeRequestRecordAsync(id); diff --git a/packages/engine/src/__tests__/mission-symbol-admission.test.ts b/packages/engine/src/__tests__/mission-symbol-admission.test.ts new file mode 100644 index 0000000000..5994ab0fda --- /dev/null +++ b/packages/engine/src/__tests__/mission-symbol-admission.test.ts @@ -0,0 +1,54 @@ +import { describe, expect, it } from "vitest"; +import type { Mission, Milestone, MissionFeature, Slice, Task } from "@fusion/core"; +import { decideMissionSymbolAdmission } from "../mission-symbol-admission.js"; + +const mission: Mission = { id: "M-1", title: "Mission", status: "active", interviewState: "completed", createdAt: "2026-01-01", updatedAt: "2026-01-01" }; +const milestone: Milestone = { id: "MS-1", missionId: mission.id, title: "Milestone", status: "active", orderIndex: 0, interviewState: "completed", dependencies: [], createdAt: "2026-01-01", updatedAt: "2026-01-01" }; +const slice: Slice = { id: "SL-1", milestoneId: milestone.id, title: "Slice", status: "active", orderIndex: 0, planState: "planned", createdAt: "2026-01-01", updatedAt: "2026-01-01" }; +const feature: MissionFeature = { id: "F-1", sliceId: slice.id, taskId: "FN-1", title: "Feature", status: "triaged", createdAt: "2026-01-01", updatedAt: "2026-01-01" }; + +function task(overrides: Partial = {}): Task { + return { id: "FN-1", title: "Task", description: "", column: "todo", priority: "normal", createdAt: "2026-01-01", updatedAt: "2026-01-01", steps: [], dependencies: [], missionId: mission.id, sliceId: slice.id, ...overrides } as Task; +} + +function store(overrides: Partial<{ mission: Mission | undefined; milestone: Milestone | undefined; slice: Slice | undefined; feature: MissionFeature | undefined }> = {}) { + const values = { mission, milestone, slice, feature, ...overrides }; + return { + getFeatureByTaskId: async () => values.feature, + getSlice: async () => values.slice, + getMilestone: async () => values.milestone, + getMission: async () => values.mission, + } as any; +} + +describe("decideMissionSymbolAdmission", () => { + it("uses symbol locking for approved mission lineage with normalized declarations", async () => { + await expect(decideMissionSymbolAdmission(task({ declaredSymbols: ["pkg/a.ts#A", " pkg/b.ts # B "] }), store())).resolves.toMatchObject({ + kind: "symbol-lock", symbols: ["pkg/a.ts#a", "pkg/b.ts#b"], reason: "approved", + }); + }); + + it("uses coarse fallback for non-mission and approved empty-symbol work", async () => { + await expect(decideMissionSymbolAdmission(task({ missionId: undefined, sliceId: undefined }), store({ feature: undefined }))).resolves.toEqual({ kind: "coarse-fallback", reason: "non-mission" }); + await expect(decideMissionSymbolAdmission(task({ declaredSymbols: [] }), store())).resolves.toEqual({ kind: "coarse-fallback", reason: "symbols-unresolvable" }); + }); + + it.each([ + ["missing-mission", { mission: undefined }], + ["missing-milestone", { milestone: undefined }], + // A missing slice also prevents its parent lookup; the canonical predicate + // intentionally reports missing-milestone first in that impossible graph. + ["missing-milestone", { slice: undefined }], + ["missing-feature", { feature: undefined }], + ["mission-not-active", { mission: { ...mission, status: "blocked" } }], + ["milestone-not-active", { milestone: { ...milestone, status: "planning" } }], + ["slice-not-active", { slice: { ...slice, status: "pending" } }], + ["feature-not-implementable", { feature: { ...feature, status: "defined" } }], + ])("blocks mission work for %s", async (reason, values) => { + await expect(decideMissionSymbolAdmission(task({ declaredSymbols: ["pkg/a.ts#A"] }), store(values as any))).resolves.toEqual({ kind: "lineage-blocked", reason }); + }); + + it("honors required plan fingerprints through the canonical predicate", async () => { + await expect(decideMissionSymbolAdmission(task({ declaredSymbols: ["pkg/a.ts#A"] }), store(), { planApprovalRequired: true })).resolves.toEqual({ kind: "lineage-blocked", reason: "plan-not-approved" }); + }); +}); diff --git a/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts b/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts index 76ccf297a8..778185b4bf 100644 --- a/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts +++ b/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts @@ -64,6 +64,7 @@ function storeWith( on: vi.fn(), off: vi.fn(), recordRunAuditEvent: vi.fn(async () => undefined), + renewSymbolLocks: vi.fn(async () => ({ renewed: [], lost: [] })), getMissionStore: vi.fn(() => ({ listMissions: () => [], listGoalIdsForMission: () => [], @@ -112,6 +113,54 @@ describe("Scheduler workflow cutover", () => { expect(store.updateSettings).toHaveBeenCalledTimes(2); }); + it("renews active mission symbol locks before their short admission lease expires", async () => { + const active = task({ + id: "FN-symbol-owner", + column: "in-progress", + missionId: "M-1", + sliceId: "SL-1", + declaredSymbols: ["pkg/a.ts#A"], + }); + const store = storeWith([active]); + const scheduler = new Scheduler(store); + (scheduler as unknown as { running: boolean }).running = true; + + await scheduler.schedule(); + + expect(store.renewSymbolLocks).toHaveBeenCalledWith(["pkg/a.ts#A"], "FN-symbol-owner", 10 * 60_000); + }); + + it("renews locks in a custom workflow WIP column", async () => { + const active = task({ + id: "FN-custom-symbol-owner", + column: "implementing", + missionId: "M-1", + sliceId: "SL-1", + declaredSymbols: ["pkg/a.ts#A"], + }); + const workflow: WorkflowIr = { + version: "v2", + name: "custom-wip", + columns: [ + { id: "todo", name: "Todo", traits: [{ trait: "hold" }] }, + { id: "implementing", name: "Implementing", traits: [{ trait: "wip", config: { limit: 1 } }] }, + { id: "done", name: "Done", traits: [{ trait: "complete" }] }, + ], + nodes: [], + edges: [], + } as WorkflowIr; + const store = storeWith([active], {}, { + selections: { "FN-custom-symbol-owner": "custom:wip" }, + definitions: { "custom:wip": workflow }, + }); + const scheduler = new Scheduler(store); + (scheduler as unknown as { running: boolean }).running = true; + + await scheduler.schedule(); + + expect(store.renewSymbolLocks).toHaveBeenCalledWith(["pkg/a.ts#A"], "FN-custom-symbol-owner", 10 * 60_000); + }); + it("uses the workflow sweep for todo pickup even when stale workflowColumns=false is persisted", async () => { const ready = task({ id: "FN-100" }); const store = storeWith([ready]); diff --git a/packages/engine/src/__tests__/workflow-work-processor.test.ts b/packages/engine/src/__tests__/workflow-work-processor.test.ts new file mode 100644 index 0000000000..1a2da6fdbe --- /dev/null +++ b/packages/engine/src/__tests__/workflow-work-processor.test.ts @@ -0,0 +1,39 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import { processDueWorkflowWorkItem } from "../workflow-work-processor.js"; + +const item = { id: "WW-renew", taskId: "FN-renew", runId: "run-renew", nodeId: "execute", kind: "execute" } as any; + +afterEach(() => vi.useRealTimers()); + +describe("processDueWorkflowWorkItem symbol lock renewal", () => { + it("renews a claimed mission symbol before its short admission lease can expire", async () => { + vi.useFakeTimers(); + let finish!: () => void; + const runWorkItem = vi.fn(() => new Promise((resolve) => { finish = () => resolve({ disposition: "completed", outcome: "success", visitedNodeIds: [], context: {} }); })); + const renewSymbolLocks = vi.fn(async () => ({ renewed: ["pkg/a.ts#a"], lost: [] })); + const store = { + listDueWorkflowWorkItems: () => [item], + acquireWorkflowWorkItemLease: () => item, + getTask: async () => ({ id: "FN-renew", missionId: "M-1", sliceId: "SL-1", declaredSymbols: ["pkg/a.ts#A"] }), + getMissionStore: () => ({ + getFeatureByTaskId: async () => ({ id: "F-1", sliceId: "SL-1", status: "triaged" }), + getSlice: async () => ({ id: "SL-1", milestoneId: "MS-1", status: "active" }), + getMilestone: async () => ({ id: "MS-1", missionId: "M-1", status: "active" }), + getMission: async () => ({ id: "M-1", status: "active" }), + }), + acquireSymbolLocks: async () => ({ acquired: true, conflicts: [] as [] }), + renewSymbolLocks, + }; + + const processing = processDueWorkflowWorkItem(store as any, { runWorkItem } as any, undefined, { + leaseOwner: "worker", leaseDurationMs: 1_000, + }); + await vi.advanceTimersByTimeAsync(200_001); + expect(renewSymbolLocks).toHaveBeenCalledWith(["pkg/a.ts#a"], "FN-renew", 10 * 60_000); + + finish(); + await processing; + await vi.advanceTimersByTimeAsync(10 * 60_000); + expect(renewSymbolLocks).toHaveBeenCalledOnce(); + }); +}); diff --git a/packages/engine/src/__tests__/workflow-work-scheduler.test.ts b/packages/engine/src/__tests__/workflow-work-scheduler.test.ts new file mode 100644 index 0000000000..59cd861a2d --- /dev/null +++ b/packages/engine/src/__tests__/workflow-work-scheduler.test.ts @@ -0,0 +1,44 @@ +import { describe, expect, it, vi } from "vitest"; +import { claimDueWorkflowWorkItem } from "../workflow-work-scheduler.js"; + +const item = { id: "WW-1", taskId: "FN-1", runId: "run-1", nodeId: "execute", kind: "execute" } as any; + +describe("claimDueWorkflowWorkItem", () => { + it("awaits normal coarse-fallback workflow lease acquisition", async () => { + const acquireWorkflowWorkItemLease = vi.fn(() => item); + const result = await claimDueWorkflowWorkItem({ listDueWorkflowWorkItems: () => [item], acquireWorkflowWorkItemLease }, { leaseOwner: "worker", leaseDurationMs: 1000 }); + expect(result).toMatchObject({ taskId: "FN-1", workItem: item }); + expect(acquireWorkflowWorkItemLease).toHaveBeenCalledOnce(); + }); + + it("does not consume a work lease when mission lineage is unapproved", async () => { + const acquireWorkflowWorkItemLease = vi.fn(() => item); + const logEntry = vi.fn(async () => undefined); + const result = await claimDueWorkflowWorkItem({ + listDueWorkflowWorkItems: () => [item], acquireWorkflowWorkItemLease, logEntry, + getTask: async () => ({ id: "FN-1", missionId: "M-1", sliceId: "SL-1", declaredSymbols: ["pkg/a.ts#A"] } as any), + getMissionStore: () => ({ getFeatureByTaskId: async () => undefined, getSlice: async () => undefined, getMilestone: async () => undefined, getMission: async () => undefined } as any), + acquireSymbolLocks: vi.fn(), + }, { leaseOwner: "worker", leaseDurationMs: 1000 }); + expect(result).toBeNull(); + expect(acquireWorkflowWorkItemLease).not.toHaveBeenCalled(); + expect(logEntry).toHaveBeenCalledWith("FN-1", expect.stringContaining("mission lineage blocked")); + }); + + it("releases an acquired symbol lock when the workflow lease races", async () => { + const releaseSymbolLocks = vi.fn(async () => undefined); + const result = await claimDueWorkflowWorkItem({ + listDueWorkflowWorkItems: () => [item], acquireWorkflowWorkItemLease: () => null, releaseSymbolLocks, + getTask: async () => ({ id: "FN-1", missionId: "M-1", sliceId: "SL-1", declaredSymbols: ["pkg/a.ts#A"] } as any), + getMissionStore: () => ({ + getFeatureByTaskId: async () => ({ id: "F-1", sliceId: "SL-1", status: "triaged" }), + getSlice: async () => ({ id: "SL-1", milestoneId: "MS-1", status: "active" }), + getMilestone: async () => ({ id: "MS-1", missionId: "M-1", status: "active" }), + getMission: async () => ({ id: "M-1", status: "active" }), + } as any), + acquireSymbolLocks: async () => ({ acquired: true, conflicts: [] }), + }, { leaseOwner: "worker", leaseDurationMs: 1000 }); + expect(result).toBeNull(); + expect(releaseSymbolLocks).toHaveBeenCalledWith(["pkg/a.ts#a"], "FN-1"); + }); +}); diff --git a/packages/engine/src/mission-symbol-admission.ts b/packages/engine/src/mission-symbol-admission.ts new file mode 100644 index 0000000000..2ccfe2bc09 --- /dev/null +++ b/packages/engine/src/mission-symbol-admission.ts @@ -0,0 +1,94 @@ +import { + evaluateMissionLineageApproval, + resolveTaskSymbolsForTask, + type AsyncMissionStore, + type MissionFeature, + type MissionLineageApprovalReason, + type MissionStore, + type Task, +} from "@fusion/core"; + +/** The only scheduler admission modes for implementation work. */ +export type MissionSymbolAdmissionDecision = + | { + kind: "symbol-lock"; + symbols: string[]; + feature: MissionFeature; + reason: "approved"; + } + | { + kind: "lineage-blocked"; + reason: Exclude; + } + | { + kind: "coarse-fallback"; + reason: "non-mission" | "symbols-unresolvable"; + }; + +export interface MissionSymbolAdmissionOptions { + /** Plan approval is policy-owned; the canonical predicate owns fingerprint validation. */ + planApprovalRequired?: boolean; +} + +type MissionReader = Pick< + MissionStore | AsyncMissionStore, + "getMission" | "getMilestone" | "getSlice" | "getFeatureByTaskId" +>; + +/** + * Resolve the canonical feature link rather than title matching: scheduler + * admission must fail closed for mission work whose hierarchy cannot be proven. + */ +async function resolveFeature(store: MissionReader, task: Task): Promise { + return await store.getFeatureByTaskId(task.id); +} + +/** + * FNXC:MissionSymbolAdmission 2026-07-31-12:00: + * FN-8306 makes autonomous implementation a three-way contract using + * evaluateMissionLineageApproval: approved mission lineage with durable symbols + * uses a symbol lock; mission-linked work that cannot prove active + * Mission→Milestone→Slice→Feature lineage is blocked; only non-mission work or + * approved work with no resolvable symbols retains coarse file-scope admission. + * This function is deliberately side-effect free so no blocked decision can + * consume a work or symbol lease. + */ +export async function decideMissionSymbolAdmission( + task: Task, + missionStore: MissionReader | undefined, + options: MissionSymbolAdmissionOptions = {}, +): Promise { + const declaredMissionLink = Boolean(task.missionId || task.sliceId); + if (!missionStore) { + return declaredMissionLink + ? { kind: "lineage-blocked", reason: "missing-feature" } + : { kind: "coarse-fallback", reason: "non-mission" }; + } + + const feature = await resolveFeature(missionStore, task); + const missionLinked = declaredMissionLink || Boolean(feature); + if (!missionLinked) return { kind: "coarse-fallback", reason: "non-mission" }; + // Resolve every stated lineage edge independently so diagnostics distinguish + // a missing feature from a missing parent rather than collapsing to the first + // child lookup that happened to be unavailable. + const sliceId = feature?.sliceId ?? task.sliceId; + const slice = sliceId ? await missionStore.getSlice(sliceId) : undefined; + const milestone = slice ? await missionStore.getMilestone(slice.milestoneId) : undefined; + const missionId = milestone?.missionId ?? task.missionId; + const mission = missionId ? await missionStore.getMission(missionId) : undefined; + const approval = evaluateMissionLineageApproval({ + task, + feature, + slice, + milestone, + mission, + planApprovalRequired: options.planApprovalRequired === true, + }); + if (!approval.approved) return { kind: "lineage-blocked", reason: approval.reason }; + + const symbols = resolveTaskSymbolsForTask(task); + if (!symbols.resolvable || symbols.symbols.length === 0) { + return { kind: "coarse-fallback", reason: "symbols-unresolvable" }; + } + return { kind: "symbol-lock", symbols: symbols.symbols, feature: feature!, reason: "approved" }; +} diff --git a/packages/engine/src/scheduler.ts b/packages/engine/src/scheduler.ts index 73c1ccdd80..b351572353 100644 --- a/packages/engine/src/scheduler.ts +++ b/packages/engine/src/scheduler.ts @@ -47,6 +47,9 @@ import type { WorkflowIr, WorkflowIrV2 } from "@fusion/core"; import { runHoldReleaseSweep, isUnplannedForExecution, type SlotReservation } from "./hold-release.js"; import { moveTaskToReplanColumn } from "./replan-target.js"; import { evaluateParkedAgentTaskLink } from "./task-agent-sync.js"; +import { decideMissionSymbolAdmission } from "./mission-symbol-admission.js"; + +const SYMBOL_LOCK_LEASE_MS = 10 * 60_000; /* FNXC:WorkflowScheduling 2026-07-15-12:55: @@ -1326,6 +1329,8 @@ export class Scheduler { // Refresh the poll interval if the persisted setting has changed this.refreshPollInterval(settings.pollIntervalMs); + await this.renewActiveMissionSymbolLocks(tasks); + // Global pause (hard stop): halt all scheduling activity if (settings.globalPause) { if (!this.wasGlobalPaused) { @@ -1761,8 +1766,24 @@ export class Scheduler { // If staleness evaluation was skipped (missing/unreadable file), continue to // existing scheduler logic which handles filesystem validation separately. - // Check file scope overlap when enabled - if (settings.groupOverlappingFiles) { + // FNXC:MissionSymbolAdmission 2026-08-01-00:00: Apply the same three-way + // admission before the early overlap/priority pass and the final release + // reservation. Otherwise this pass would serialize approved disjoint symbols + // (or misdiagnose unapproved mission work as an overlap) before lock admission. + const earlyMissionAdmission = await decideMissionSymbolAdmission( + task, + this.options.missionStore, + { planApprovalRequired: settings.planApprovalMode === "require-all" }, + ); + if (earlyMissionAdmission.kind === "lineage-blocked") { + await this.store.updateTask(task.id, { status: "queued", blockedBy: null, overlapBlockedBy: null }); + await this.logDispatchQueuedReason(task.id, `queued — mission lineage blocked: ${earlyMissionAdmission.reason}`); + continue; + } + + // Check file scope overlap when enabled. Approved, resolvable mission work + // bypasses this coarse pass and is mutually excluded by its durable symbols. + if (settings.groupOverlappingFiles && earlyMissionAdmission.kind === "coarse-fallback") { const taskScope = await getFilteredFileScope(task.id); const coordinationOnlyTask = isCoordinationOnlyTask(task, taskScope); if (taskScope.length > 0 && !coordinationOnlyTask) { @@ -2262,6 +2283,60 @@ export class Scheduler { * worktree allocation into the reservation-first ordering (KTD-10). Failures * are isolated so a sweep error never breaks the scheduling pass. */ + /** + * FNXC:MissionSymbolAdmission 2026-07-19-22:04: + * Admission leases are intentionally short so crash recovery can reclaim them, + * but every active mission owner renews on the scheduler heartbeat. Renewal runs + * before pause handling because a paused scheduler still permits existing work to + * finish and must not let same-symbol work enter after its original lease expires. + * Active implementation follows the workflow `countsTowardWip` trait, so renamed + * and multi-WIP workflows keep their declared symbols exclusively held. + */ + private async renewActiveMissionSymbolLocks(tasks: Task[]): Promise { + for (const task of tasks) { + let isImplementationColumn = task.column === "in-progress"; + try { + const ir = await resolveWorkflowIrForTask(this.store, task.id); + const column = (ir as WorkflowIrV2).columns?.find((candidate) => candidate.id === task.column); + if (column) isImplementationColumn = resolveColumnFlags(column).countsTowardWip === true; + } catch { + // Preserve the legacy in-progress fallback if workflow resolution is unavailable. + } + if (!isImplementationColumn) continue; + // TaskStore normalizes durable declarations on write; renewal deliberately + // reads that field rather than the prompt so an explicit declaration clear + // can never recreate a lock. + const symbols = [...new Set((task.declaredSymbols ?? []).filter((symbol): symbol is string => typeof symbol === "string" && symbol.trim().length > 0))]; + if (symbols.length === 0) continue; + + // A direct mission/slice link is sufficient for renewal. Feature-only links + // are resolved through the store when available, without re-evaluating + // approval: a running owner retains its admission lease until transition release. + let missionLinked = Boolean(task.missionId || task.sliceId); + if (!missionLinked && this.options.missionStore) { + try { + missionLinked = Boolean(await this.options.missionStore.getFeatureByTaskId(task.id)); + } catch (error) { + schedulerLog.warn(`Symbol-lock renewal lineage lookup failed for ${task.id}:`, error); + continue; + } + } + if (!missionLinked) continue; + + try { + const result = await this.store.renewSymbolLocks(symbols, task.id, SYMBOL_LOCK_LEASE_MS); + if (result.lost.length > 0) { + schedulerLog.warn(`Symbol-lock renewal lost ownership for ${task.id}: ${result.lost.join(", ")}`); + await this.store.logEntry(task.id, `symbol-lock renewal lost: ${result.lost.join(", ")}`); + } + } catch (error) { + // Do not abort capacity admission for unrelated tasks; the durable expiry + // and self-healing reconciliation remain the crash-safe backstop. + schedulerLog.warn(`Symbol-lock renewal failed for ${task.id}:`, error); + } + } + } + private async runHoldReleaseSweepPass(tasks: Task[], settings: Settings): Promise { try { const maxWorktrees = settings.maxWorktrees ?? this.options.maxWorktrees ?? 4; @@ -2766,7 +2841,25 @@ export class Scheduler { return null; } - if (settings.groupOverlappingFiles) { + /* + FNXC:MissionSymbolAdmission 2026-07-31-12:00: + FN-8306 blocks mission-linked work before any file-scope or work-slot + reservation unless the canonical lineage predicate is satisfied. A + symbol-lock candidate bypasses coarse overlap only; every capacity, + checkout, dependency, and semaphore gate below still applies. + */ + const missionAdmission = await decideMissionSymbolAdmission( + freshTask, + this.options.missionStore, + { planApprovalRequired: latestSettings.planApprovalMode === "require-all" }, + ); + if (missionAdmission.kind === "lineage-blocked") { + await this.store.updateTask(task.id, { status: "queued", blockedBy: null, overlapBlockedBy: null }); + await this.logDispatchQueuedReason(task.id, `queued — mission lineage blocked: ${missionAdmission.reason}`); + return null; + } + + if (settings.groupOverlappingFiles && missionAdmission.kind === "coarse-fallback") { const taskScope = await getFilteredFileScope(task.id); if (taskScope.length > 0 && !isCoordinationOnlyTask(task, taskScope)) { const overlappingTaskId = Array.from(activeScopes.entries()) @@ -2857,6 +2950,37 @@ export class Scheduler { registerPreHeldExecutorSlot(task.id); } + let acquiredSymbols: string[] | undefined; + if (missionAdmission.kind === "symbol-lock") { + /* + FNXC:MissionSymbolAdmission 2026-07-31-12:00: + Acquire after all capacity gates and immediately before hold release; + the reservation release path below returns this lock if moveTask + rejects, while a successful move transfers ownership to the task. + */ + const lockResult = await this.store.acquireSymbolLocks( + missionAdmission.symbols, + { ownerTaskId: task.id, missionId: freshTask.missionId, featureId: missionAdmission.feature.id, agentId: "scheduler" }, + SYMBOL_LOCK_LEASE_MS, + ); + if (!lockResult.acquired) { + dropPreHeldExecutorSlot(task.id, sem); + if (reservedScope) { + activeScopes.delete(task.id); + activeScopeColumns.delete(task.id); + } + const conflict = lockResult.conflicts[0]; + await this.store.updateTask(task.id, { status: "queued", blockedBy: null, overlapBlockedBy: null }); + await this.logDispatchQueuedReason( + task.id, + `queued — symbol contention: symbol=${conflict?.symbolKey ?? "unknown"} holder=${conflict?.ownerTaskId ?? "unknown"}`, + ); + return null; + } + acquiredSymbols = missionAdmission.symbols; + await this.store.logEntry(task.id, `symbol-lock admission acquired: ${missionAdmission.symbols.join(", ")}`); + } + dispatchPrepByTaskId.set(task.id, { baseBranch: this.resolveBaseBranch(freshTask, tasks), dispatchStormCount: nextDispatchStormCount, @@ -2881,6 +3005,9 @@ export class Scheduler { reservedConcurrentSlots = Math.max(0, reservedConcurrentSlots - 1); dispatchPrepByTaskId.delete(task.id); dropPreHeldExecutorSlot(task.id, sem); + if (acquiredSymbols) { + void this.store.releaseSymbolLocks(acquiredSymbols, task.id); + } }, }; }, diff --git a/packages/engine/src/workflow-work-processor.ts b/packages/engine/src/workflow-work-processor.ts index 76b8b0734b..733efa1426 100644 --- a/packages/engine/src/workflow-work-processor.ts +++ b/packages/engine/src/workflow-work-processor.ts @@ -30,7 +30,8 @@ export async function processDueWorkflowWorkItem( settings: (Pick & Partial) | undefined, opts: WorkflowWorkProcessorOptions, ): Promise { - const dispatch = claimDueWorkflowWorkItem(store, { + /* FNXC:MissionSymbolAdmission 2026-07-31-12:00: await the async symbol-lock admission before runtime may consume the workflow work lease. */ + const dispatch = await claimDueWorkflowWorkItem(store, { now: opts.now, leaseOwner: opts.leaseOwner, leaseDurationMs: opts.leaseDurationMs, @@ -39,6 +40,23 @@ export async function processDueWorkflowWorkItem( if (!dispatch) return { claimed: false }; let runtimeResult: WorkflowTaskRuntimeResult; + /* + FNXC:MissionSymbolAdmission 2026-08-01-01:00: + Workflow execution can outlive the ten-minute crash-recoverable lease. Renew + only locks acquired by this claim while its runtime is live; transition release + remains authoritative once the work reaches review, requeue, or terminal state. + */ + const renewInterval = dispatch.symbolLocks && store.renewSymbolLocks + ? setInterval(() => { + void store.renewSymbolLocks!(dispatch.symbolLocks!, dispatch.taskId, 10 * 60_000) + .then(async (result) => { + if (result.lost.length > 0) { + await store.logEntry?.(dispatch.taskId, `workflow symbol-lock renewal lost: ${result.lost.join(", ")}`); + } + }) + .catch(() => undefined); + }, (10 * 60_000) / 3) + : undefined; try { runtimeResult = await runtime.runWorkItem(dispatch.workItem, settings); } catch (err) { @@ -60,6 +78,8 @@ export async function processDueWorkflowWorkItem( context: {}, reason, }; + } finally { + if (renewInterval) clearInterval(renewInterval); } return { claimed: true, diff --git a/packages/engine/src/workflow-work-scheduler.ts b/packages/engine/src/workflow-work-scheduler.ts index 8ae7206a86..1bf25f7ddb 100644 --- a/packages/engine/src/workflow-work-scheduler.ts +++ b/packages/engine/src/workflow-work-scheduler.ts @@ -1,4 +1,7 @@ -import type { WorkflowWorkItem, WorkflowWorkItemDueFilter, WorkflowWorkItemKind } from "@fusion/core"; +import type { AsyncMissionStore, MissionStore, Task, WorkflowWorkItem, WorkflowWorkItemDueFilter, WorkflowWorkItemKind } from "@fusion/core"; +import { decideMissionSymbolAdmission } from "./mission-symbol-admission.js"; + +const WORKFLOW_SYMBOL_LOCK_LEASE_MS = 10 * 60_000; export interface WorkflowWorkSchedulerStore { listDueWorkflowWorkItems(filter?: WorkflowWorkItemDueFilter): WorkflowWorkItem[]; @@ -7,6 +10,18 @@ export interface WorkflowWorkSchedulerStore { leaseOwner: string, opts: { leaseDurationMs: number; now?: string }, ): WorkflowWorkItem | null; + /** TaskStore supplies these optional scheduler-admission capabilities. */ + getTask?(id: string): Promise; + getMissionStore?(): MissionStore | AsyncMissionStore; + getSettings?(): Promise<{ planApprovalMode?: "workflow" | "auto-approve-all" | "require-all" }>; + acquireSymbolLocks?( + symbols: readonly string[], + owner: { ownerTaskId: string; missionId?: string; featureId?: string; agentId?: string }, + leaseMs: number, + ): Promise<{ acquired: true; conflicts: [] } | { acquired: false; conflicts: Array<{ symbolKey: string; ownerTaskId: string }> }>; + releaseSymbolLocks?(symbols: readonly string[], ownerTaskId: string): Promise; + renewSymbolLocks?(symbols: readonly string[], ownerTaskId: string, leaseMs: number): Promise<{ renewed: string[]; lost: string[] }>; + logEntry?(taskId: string, message: string): Promise; } export interface WorkflowWorkDispatch { @@ -14,6 +29,8 @@ export interface WorkflowWorkDispatch { runId: string; taskId: string; nodeId: string; + /** Symbols held by this claim; the processor renews them while runtime work is live. */ + symbolLocks?: string[]; } export interface ClaimWorkflowWorkOptions { @@ -24,10 +41,17 @@ export interface ClaimWorkflowWorkOptions { kinds?: WorkflowWorkItemKind[]; } -export function claimDueWorkflowWorkItem( +/** + * FNXC:MissionSymbolAdmission 2026-07-31-12:00: + * Workflow work claiming is async because durable symbol acquisition must occur + * before its work lease is consumed. Unapproved mission work and contention are + * skipped without a lease, while deployments without the TaskStore admission + * seam retain their existing coarse workflow-lease behavior. + */ +export async function claimDueWorkflowWorkItem( store: WorkflowWorkSchedulerStore, opts: ClaimWorkflowWorkOptions, -): WorkflowWorkDispatch | null { +): Promise { const due = store.listDueWorkflowWorkItems({ now: opts.now, limit: opts.limit ?? 25, @@ -35,16 +59,46 @@ export function claimDueWorkflowWorkItem( }); for (const candidate of due) { + const task = store.getTask ? await store.getTask(candidate.taskId) : undefined; + let lockedSymbols: string[] | undefined; + if (task && store.getMissionStore && store.acquireSymbolLocks) { + const settings = await store.getSettings?.(); + const admission = await decideMissionSymbolAdmission(task, store.getMissionStore(), { + planApprovalRequired: settings?.planApprovalMode === "require-all", + }); + if (admission.kind === "lineage-blocked") { + await store.logEntry?.(task.id, `workflow work not claimed — mission lineage blocked: ${admission.reason}`); + continue; + } + if (admission.kind === "symbol-lock") { + const result = await store.acquireSymbolLocks( + admission.symbols, + { ownerTaskId: task.id, missionId: task.missionId, featureId: admission.feature.id, agentId: opts.leaseOwner }, + WORKFLOW_SYMBOL_LOCK_LEASE_MS, + ); + if (!result.acquired) { + const conflict = result.conflicts[0]; + await store.logEntry?.(task.id, `workflow work not claimed — symbol contention: symbol=${conflict?.symbolKey ?? "unknown"} holder=${conflict?.ownerTaskId ?? "unknown"}`); + continue; + } + lockedSymbols = admission.symbols; + } + } + const workItem = store.acquireWorkflowWorkItemLease(candidate.id, opts.leaseOwner, { now: opts.now, leaseDurationMs: opts.leaseDurationMs, }); - if (!workItem) continue; + if (!workItem) { + if (lockedSymbols) await store.releaseSymbolLocks?.(lockedSymbols, candidate.taskId); + continue; + } return { workItem, runId: workItem.runId, taskId: workItem.taskId, nodeId: workItem.nodeId, + symbolLocks: lockedSymbols, }; }