From ad91795dace9fe0ec97a410992a15a93aacf9011 Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Sun, 9 Aug 2026 14:12:49 -0700 Subject: [PATCH] FN-8843: assign eligible executors to new tasks Assign durable executor ownership during shared task intake. - Resolve eligible executor owners for API, dashboard, and reserved-ID task creation. - Preserve explicit assignments, reject invalid owners, and audit unresolved ownerless intake. - Cover routing and PostgreSQL intake behavior with regression tests and document the contract. Files changed: .changeset/fn-8843-intake-agent-assignment.md | 7 + docs/architecture.md | 6 + docs/workflow-steps.md | 12 +- .../postgres/create-task-reserved-id.pg.test.ts | 187 +++++++++++++++ .../__tests__/task-intake-owner-resolver.test.ts | 141 ++++++++++++ packages/core/src/store.ts | 29 ++- packages/core/src/task-store/task-creation.ts | 251 +++++++++++++++++---- .../core/src/tasks/task-intake-owner-resolver.ts | 239 ++++++++++++++++++++ .../dashboard/src/__tests__/routes-tasks.test.ts | 47 ++++ .../__tests__/task-create-intake-owner.pg.test.ts | 195 ++++++++++++++++ .../src/routes/register-task-workflow-routes.ts | 12 +- .../src/__tests__/workflow-agent-routing.test.ts | 19 +- .../engine/src/agents/workflow-agent-router.ts | 37 ++- 13 files changed, 1120 insertions(+), 62 deletions(-) Fusion-Task-Id: FN-8843 Fusion-Task-Lineage: b94728b4-d9e8-4a26-8fbc-7096eb797839 Co-authored-by: Fusion (runfusion.ai) --- .changeset/fn-8843-intake-agent-assignment.md | 7 + docs/architecture.md | 6 + docs/workflow-steps.md | 12 +- .../create-task-reserved-id.pg.test.ts | 187 +++++++++++++ .../task-intake-owner-resolver.test.ts | 141 ++++++++++ packages/core/src/store.ts | 29 +- packages/core/src/task-store/task-creation.ts | 251 ++++++++++++++---- .../src/tasks/task-intake-owner-resolver.ts | 239 +++++++++++++++++ .../src/__tests__/routes-tasks.test.ts | 47 ++++ .../task-create-intake-owner.pg.test.ts | 195 ++++++++++++++ .../routes/register-task-workflow-routes.ts | 12 +- .../__tests__/workflow-agent-routing.test.ts | 19 +- .../src/agents/workflow-agent-router.ts | 37 ++- 13 files changed, 1120 insertions(+), 62 deletions(-) create mode 100644 .changeset/fn-8843-intake-agent-assignment.md create mode 100644 packages/core/src/__tests__/task-intake-owner-resolver.test.ts create mode 100644 packages/core/src/tasks/task-intake-owner-resolver.ts create mode 100644 packages/dashboard/src/routes/__tests__/task-create-intake-owner.pg.test.ts diff --git a/.changeset/fn-8843-intake-agent-assignment.md b/.changeset/fn-8843-intake-agent-assignment.md new file mode 100644 index 0000000000..f9a3b060b3 --- /dev/null +++ b/.changeset/fn-8843-intake-agent-assignment.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Assign eligible executor owners to newly created tasks automatically. +category: fix +dev: Resolves role-safe owners before the shared insert boundary; public payloads cannot forge exemption. diff --git a/docs/architecture.md b/docs/architecture.md index b1ad0ea1a9..5a4838de70 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -2379,3 +2379,9 @@ Workflow session capacity is acquired separately from heartbeat capacity and rel A scheduler pass batches missing task workflow selections once and shares a strictly pass-scoped selection cache with hold-release and reservation resolution. The cache is never retained because selections are mutable. Only one sweep per project identity may run: concurrent calls are skipped, not joined, since their clocks, budgets, and slot reservations are caller-owned. Each sweep has a 10-second default budget. It does not claim to cancel database calls: after the deadline it starts no further sweep-owned await, while an in-flight read or started release can overrun and is reported as `budgetOverrunMs`. Dependency evaluation returns `{ satisfied, truncated }`; truncation is `sweep-budget-exhausted` only for an already-classified dependency hold, never `deps-unsatisfied`, and does not reset its held clock. Cards not reached are represented only by `unevaluatedCount`, never fabricated hold reasons. A truncated sweep never releases on partial evidence and never prunes an unreached card's held clock. + +## Durable intake executor ownership + +FN-8843 resolves task ownership at `_createTaskInternalBackendImpl`, the one pre-insert boundary shared by normal and reserved-ID task creation. A store-owned, project-scoped `AgentStore` supplies the durable-agent snapshot, so each intake does not open its own backend. It resolves an effective workflow independently of optional-step materialization; `workflowId: null` means no workflow steps, not no owner, and uses the eligible executor pool directly. An owner is a non-ephemeral durable executor with runtime enabled, a non-paused/non-error state, and a permitted implementation assignment policy. Explicit eligible ownership wins, followed by a reachable execute-node column binding and the pool ordered by active durable session, creation time, then ID. + +The resolver deliberately separates four outcomes: `selected` writes the owner on insert; internal options-only `exempt` writes deliberate null for terminal/historical/fixture creators; `rejected` fails before row/reservation/event publication; and `unowned` succeeds only for a zero-eligible-executor project, with null owner plus a `task:intake-owner-unresolved` run-audit event. Named terminal, historical, and fixture-only factories issue the opaque in-process exemption capability; its module-private symbol token makes it non-serializable, so public API, CLI, tool, automation, and remote-node payloads cannot forge it. Per-stage workflow work items still own planning/review principals; stable task ownership only participates in executor routing. diff --git a/docs/workflow-steps.md b/docs/workflow-steps.md index 89ca4f35b9..24192398b4 100644 --- a/docs/workflow-steps.md +++ b/docs/workflow-steps.md @@ -306,6 +306,10 @@ A v2 column can optionally name a **permanent agent** from the agent registry, s **Missing-agent fallback.** A missing or deleted agent at resolution time logs and falls back to normal resolution — a live session is never aborted because its column agent was deleted mid-flight. +### Durable intake executor owner + +New tasks receive a stable `assignedAgentId` before insertion. Creation resolves an explicit eligible executor first, then the first reachable execute-node column binding, then the eligible durable executor pool ordered by active session, creation time, and ID. Planning and review principals remain separately fenced by workflow work items; a planner or reviewer binding never becomes the durable implementation owner. `workflowId: null` disables graph steps but still resolves an owner from the executor pool. The only successful null-owner case is a project with no eligible durable executor, which emits `task:intake-owner-unresolved`; public request and tool payloads cannot opt out of ownership resolution. + **Flag requirements.** Column agents act only when **both** `experimentalFeatures.workflowColumns` and `experimentalFeatures.workflowGraphExecutor` are on; with either off the binding is inert (config is still stored and round-trips — only execution is gated), and the editor surfaces that the picker is disabled with a tooltip naming both flags. **Write-time validation.** Saving a workflow validates agent references: an unknown `agentId` is rejected with a typed 4xx naming the column. Binding an agent whose permission policy is broader than the project default requires an explicit policy-escalation confirmation (`confirmPolicyEscalation`) at save time, so override cannot silently re-key action gates to a more-privileged agent. @@ -979,4 +983,10 @@ Findings persist through both ordinary-node and optional-group result writers in Only session-launching nodes acquire a workflow principal: planning prompts use `triage`; implementation, remediation, script/CLI-agent and completion prompts use `executor`; `step-review` and review prompts use `reviewer`; merge prompts use `merger`. Control and lifecycle nodes, including start/end, hold, split/join, containers, parse-steps, notifications, asks, gates, and `review-handoff`, have no principal. -A review node may persist `reviewerAgentId` in its IR. It is exact-node scoped and may name any permanent agent. Resolution is reviewer override (review only), task owner, column binding, then a matching role pool. Claimed `workflow_work_items` persist principal, role, authority kind, and node-instance fence. Named unavailable principals hold closed; pool exhaustion is reported separately. +A review node may persist `reviewerAgentId` in its IR. It is exact-node scoped and may name any permanent agent. Resolution is reviewer override (review only), column binding, then a matching role pool. Execute nodes additionally prefer the durable task owner when it remains an eligible executor. Claimed `workflow_work_items` persist principal, role, authority kind, and node-instance fence. Named unavailable principals hold closed; pool exhaustion is reported separately. + +## Intake executor ownership + +Task creation resolves ownership once at the shared pre-insert boundary used by ordinary and reserved-ID creates. A durable owner must be a non-ephemeral, runtime-enabled executor not paused or errored and permitted by implementation assignment policy. A valid explicit owner wins; otherwise the first reachable execute-node column binding is used, then the deterministic executor pool. Planning and review principals remain work-item-scoped and never rewrite this owner. + +`workflowId: null` disables workflow-step materialization only: it still resolves from the executor pool. The only internal exemption is an options-bag reason for terminal, historical, or fixture creation; no HTTP/tool/CLI payload can set it. Resolution outcomes are distinct: `selected` persists an owner; internal `exempt` deliberately persists null; `rejected` fails before insertion; and `unowned` succeeds only when no eligible executor exists, emitting `task:intake-owner-unresolved` for ordinary later assignment. diff --git a/packages/core/src/__tests__/postgres/create-task-reserved-id.pg.test.ts b/packages/core/src/__tests__/postgres/create-task-reserved-id.pg.test.ts index 96670aac01..de72ba0573 100644 --- a/packages/core/src/__tests__/postgres/create-task-reserved-id.pg.test.ts +++ b/packages/core/src/__tests__/postgres/create-task-reserved-id.pg.test.ts @@ -14,6 +14,9 @@ import { createTaskStoreForTest, type PgTestHarness, } from "../../__test-utils__/pg-test-harness.js"; +import { AgentStore } from "../../agents/agent-store.js"; +import { createFixtureIntakeOwnershipExemption } from "../../tasks/task-intake-owner-resolver.js"; +import { BUILTIN_CODING_WORKFLOW_IR } from "../../workflows/builtin-coding-workflow-ir.js"; pgDescribe("createTaskWithReservedId backend mode (PostgreSQL)", () => { let harness: PgTestHarness | null = null; @@ -30,6 +33,190 @@ pgDescribe("createTaskWithReservedId backend mode (PostgreSQL)", () => { } } + /* + FNXC:IntakeOwnership 2026-08-09-16:28: + Both public creation gateways must carry the same effective workflow into the + shared insert boundary. This PostgreSQL regression protects the production + path where workflow-step materialization is intentionally disabled. + */ + it("assigns the durable executor through ordinary and no-workflow reserved creates", async () => { + const h = await makeHarness(); + const agents = new AgentStore({ rootDir: h.rootDir, asyncLayer: h.layer, projectId: h.layer.projectId }); + try { + const executor = await agents.createAgent({ name: "Intake executor", role: "executor" }); + const ordinary = await h.store.createTask({ description: "ordinary intake owner" }); + const reserved = await h.store.createTaskWithReservedId( + { description: "reserved no-materialization intake owner" }, + { taskId: "FN-OWNER-RESERVED", applyDefaultWorkflowSteps: false }, + ); + const noWorkflow = await h.store.createTaskWithReservedId( + { description: "no-workflow still has intake owner", workflowId: null }, + { taskId: "FN-OWNER-NO-WORKFLOW" }, + ); + // A public caller can type-cast an options object, but cannot forge the + // module-private symbol capability that is the only exemption authority. + const forgedExemption = await h.store.createTask( + { description: "forged exemption remains owned" }, + { ownershipExemption: { reason: "fixture" } } as never, + ); + + for (const task of [ordinary, reserved, noWorkflow, forgedExemption]) { + expect(task.assignedAgentId).toBe(executor.id); + expect((await h.store.getTask(task.id))?.assignedAgentId).toBe(executor.id); + } + } finally { + agents.close(); + await teardown(); + } + }); + + it("fails rejected owner outcomes before either gateway inserts a row", async () => { + const h = await makeHarness(); + const agents = new AgentStore({ rootDir: h.rootDir, asyncLayer: h.layer, projectId: h.layer.projectId }); + try { + const executor = await agents.createAgent({ name: "Fallback executor", role: "executor" }); + const unavailableBindingWorkflow = await h.store.createWorkflowDefinition({ + name: "Unavailable executor binding", + ir: { + ...BUILTIN_CODING_WORKFLOW_IR, + columns: BUILTIN_CODING_WORKFLOW_IR.columns.map((column) => column.id === "in-progress" + ? { ...column, agent: { agentId: "missing-executor", mode: "override" as const } } + : column), + }, + }); + + await expect(h.store.createTask({ + description: "explicit triage agent is invalid", + assignedAgentId: "missing-triage-agent", + })).rejects.toMatchObject({ code: "task-intake-owner-resolution", reason: "explicit-assignee-ineligible" }); + await expect(h.store.createTaskWithReservedId( + { description: "named binding cannot fall through", workflowId: unavailableBindingWorkflow.id }, + { taskId: "FN-REJECT-NAMED-BINDING" }, + )).rejects.toMatchObject({ code: "task-intake-owner-resolution", reason: "named-execute-binding-unavailable" }); + await expect(h.store.createTaskWithReservedId( + { description: "unknown workflow cannot become unowned", workflowId: "WF-NOT-RESOLVABLE" }, + { taskId: "FN-REJECT-WORKFLOW" }, + )).rejects.toMatchObject({ code: "task-intake-owner-resolution", reason: "workflow-unresolvable" }); + + expect(executor.id).toBeTruthy(); + await expect(h.store.getTask("FN-REJECT-NAMED-BINDING")).rejects.toThrow("not found"); + await expect(h.store.getTask("FN-REJECT-WORKFLOW")).rejects.toThrow("not found"); + expect((await h.store.listTasks()).map((task) => task.description)).not.toContain("explicit triage agent is invalid"); + } finally { + agents.close(); + await teardown(); + } + }); + + it("permits only the named internal fixture capability to create a deliberate unowned row", async () => { + const h = await makeHarness(); + try { + const task = await h.store.createTaskWithReservedId( + { description: "fixture-only historical row" }, + { + taskId: "FN-FIXTURE-EXEMPT", + applyDefaultWorkflowSteps: false, + ownershipExemption: createFixtureIntakeOwnershipExemption(), + }, + ); + + expect(task.assignedAgentId).toBeUndefined(); + expect((await h.store.getTask(task.id))?.assignedAgentId).toBeUndefined(); + const audit = await h.store.getRunAuditEventsAsync({ taskId: task.id }); + expect(audit.filter((event) => event.mutationType === "task:intake-owner-unresolved")) + .toHaveLength(0); + } finally { + await teardown(); + } + }); + + it("returns a normal-gateway proposal replay before resolving its new owner", async () => { + const h = await makeHarness(); + const agents = new AgentStore({ rootDir: h.rootDir, asyncLayer: h.layer, projectId: h.layer.projectId }); + try { + const executor = await agents.createAgent({ name: "Proposal executor", role: "executor" }); + const first = await h.store.createTask({ + description: "proposal claim owner is stable", + proposalClaimId: "proposal:intake-owner-replay", + }); + expect(first.assignedAgentId).toBe(executor.id); + + /* + FNXC:IntakeOwnership 2026-08-09-18:57: + A retry deliberately carries an invalid new assignee. A replay is not a + new intake and must return the canonical task instead of rerunning the + resolver and rejecting that changed payload. + */ + const replay = await h.store.createTask({ + description: "proposal retry changed payload", + proposalClaimId: "proposal:intake-owner-replay", + assignedAgentId: "deleted-or-invalid-agent", + }); + expect(replay.id).toBe(first.id); + expect(replay.assignedAgentId).toBe(executor.id); + expect((await h.store.listTasks()).filter((task) => task.proposalClaimId === "proposal:intake-owner-replay")) + .toHaveLength(1); + } finally { + agents.close(); + await teardown(); + } + }); + + it("replays concurrent reserved-ID proposals without reaching the removed SQLite backend", async () => { + const h = await makeHarness(); + const agents = new AgentStore({ rootDir: h.rootDir, asyncLayer: h.layer, projectId: h.layer.projectId }); + try { + const executor = await agents.createAgent({ name: "Reserved proposal executor", role: "executor" }); + const createdIds: string[] = []; + h.store.on("task:created", (task) => { createdIds.push(task.id); }); + const [first, replay] = await Promise.all([ + h.store.createTaskWithReservedId( + { description: "reserved proposal owner", proposalClaimId: "proposal:reserved-intake-owner-replay" }, + { taskId: "FN-RESERVED-PROPOSAL-ONE" }, + ), + h.store.createTaskWithReservedId( + { description: "reserved proposal retry", proposalClaimId: "proposal:reserved-intake-owner-replay" }, + { taskId: "FN-RESERVED-PROPOSAL-TWO" }, + ), + ]); + + expect(first.assignedAgentId).toBe(executor.id); + expect(replay.id).toBe(first.id); + expect(replay.assignedAgentId).toBe(executor.id); + expect((await h.store.listTasks()).filter((task) => task.proposalClaimId === "proposal:reserved-intake-owner-replay")) + .toHaveLength(1); + expect(createdIds).toEqual([first.id]); + } finally { + agents.close(); + await teardown(); + } + }); + + it("makes concurrent proposal retries a single insertion and event", async () => { + const h = await makeHarness(); + const agents = new AgentStore({ rootDir: h.rootDir, asyncLayer: h.layer, projectId: h.layer.projectId }); + try { + const executor = await agents.createAgent({ name: "Concurrent proposal executor", role: "executor" }); + const createdIds: string[] = []; + h.store.on("task:created", (task) => { createdIds.push(task.id); }); + + const [first, second] = await Promise.all([ + h.store.createTask({ description: "concurrent proposal owner", proposalClaimId: "proposal:intake-owner-concurrent" }), + h.store.createTask({ description: "concurrent proposal owner", proposalClaimId: "proposal:intake-owner-concurrent" }), + ]); + + expect(first.id).toBe(second.id); + expect(first.assignedAgentId).toBe(executor.id); + expect(second.assignedAgentId).toBe(executor.id); + expect((await h.store.listTasks()).filter((task) => task.proposalClaimId === "proposal:intake-owner-concurrent")) + .toHaveLength(1); + expect(createdIds).toEqual([first.id]); + } finally { + agents.close(); + await teardown(); + } + }); + it("createTaskWithReservedId persists a task with the reserved id in backend mode", async () => { const h = await makeHarness(); try { diff --git a/packages/core/src/__tests__/task-intake-owner-resolver.test.ts b/packages/core/src/__tests__/task-intake-owner-resolver.test.ts new file mode 100644 index 0000000000..264733a76d --- /dev/null +++ b/packages/core/src/__tests__/task-intake-owner-resolver.test.ts @@ -0,0 +1,141 @@ +import { describe, expect, it } from "vitest"; +import { + createFixtureIntakeOwnershipExemption, + getInternalIntakeOwnershipExemptionReason, + resolveTaskIntakeOwner, +} from "../tasks/task-intake-owner-resolver.js"; +import type { Agent } from "../types.js"; +import type { WorkflowIr } from "../workflows/workflow-ir-types.js"; + +const agent = (id: string, overrides: Partial = {}): Agent => ({ + id, name: id, roles: ["executor"], role: "executor", state: "idle", + createdAt: "2026-01-01T00:00:00.000Z", updatedAt: "2026-01-01T00:00:00.000Z", metadata: {}, ...overrides, +}); + +const workflow: WorkflowIr = { + version: "v2", + nodes: [ + { id: "start", kind: "start", column: "plan" }, + { id: "plan", kind: "prompt", column: "plan", config: { seam: "planning" } }, + { id: "execute", kind: "prompt", column: "work", config: { seam: "execute" } }, + ], + edges: [{ from: "start", to: "plan" }, { from: "plan", to: "execute" }], + columns: [ + { id: "plan", name: "Plan", traits: [], agent: { agentId: "triage", mode: "override" } }, + { id: "work", name: "Work", traits: [], agent: { agentId: "bound", mode: "defer" } }, + ], +}; + +describe("resolveTaskIntakeOwner", () => { + it("uses an eligible explicit executor before workflow and pool candidates", () => { + expect(resolveTaskIntakeOwner({ workflow, explicitAssigneeId: "explicit", agents: [agent("bound"), agent("explicit")] })).toEqual({ status: "selected", agentId: "explicit", source: "explicit" }); + }); + + it("rejects an ineligible explicit owner and unavailable named execute binding", () => { + expect(resolveTaskIntakeOwner({ workflow, explicitAssigneeId: "triage", agents: [agent("triage", { roles: ["triage"], role: "triage" })] })).toEqual({ status: "rejected", reason: "explicit-assignee-ineligible" }); + expect(resolveTaskIntakeOwner({ workflow, agents: [agent("pool")] })).toEqual({ status: "rejected", reason: "named-execute-binding-unavailable" }); + }); + + it("uses only the first reachable executor binding, never intake or unreachable bindings", () => { + const graph: WorkflowIr = { + ...workflow, + nodes: [ + ...workflow.nodes, + { id: "unreachable-execute", kind: "prompt", column: "unreachable", config: { seam: "execute" } }, + ], + columns: [...workflow.columns, { id: "unreachable", name: "Unreachable", traits: [], agent: { agentId: "wrong", mode: "override" } }], + }; + expect(resolveTaskIntakeOwner({ workflow: graph, agents: [agent("bound"), agent("wrong"), agent("triage", { roles: ["triage"], role: "triage" })] })) + .toEqual({ status: "selected", agentId: "bound", source: "execute-binding" }); + }); + + it("resolves a reachable foreach execute template through its inherited binding", () => { + const graph: WorkflowIr = { + version: "v2", + nodes: [ + { id: "start", kind: "start", column: "plan" }, + { id: "steps", kind: "foreach", column: "work", config: { source: "task-steps", template: { nodes: [{ id: "run", kind: "prompt", config: { seam: "execute" } }], edges: [] } } }, + ], + edges: [{ from: "start", to: "steps" }], + columns: [{ id: "plan", name: "Plan", traits: [] }, { id: "work", name: "Work", traits: [], agent: { agentId: "bound", mode: "override" } }], + }; + expect(resolveTaskIntakeOwner({ workflow: graph, agents: [agent("bound")] })) + .toEqual({ status: "selected", agentId: "bound", source: "execute-binding" }); + }); + + it("honors defer bindings by falling back to the eligible pool when task model settings exist", () => { + expect(resolveTaskIntakeOwner({ workflow, ownModelProvider: "openai", ownModelId: "gpt", agents: [agent("bound"), agent("pool", { createdAt: "2025-01-01T00:00:00.000Z" })] })) + .toEqual({ status: "selected", agentId: "pool", source: "executor-pool" }); + }); + + it("ignores execute bindings inside disabled optional groups", () => { + const optionalExecute: WorkflowIr = { + version: "v2", + nodes: [ + { id: "start", kind: "start", column: "plan" }, + { id: "optional-execute", kind: "optional-group", column: "work", config: { + defaultOn: false, + template: { nodes: [{ id: "run", kind: "prompt", config: { seam: "execute" } }], edges: [] }, + } }, + ], + edges: [{ from: "start", to: "optional-execute" }], + columns: [{ id: "plan", name: "Plan", traits: [] }, { id: "work", name: "Work", traits: [], agent: { agentId: "missing-bound", mode: "override" } }], + }; + expect(resolveTaskIntakeOwner({ workflow: optionalExecute, agents: [agent("pool")] })) + .toEqual({ status: "selected", agentId: "pool", source: "executor-pool" }); + expect(resolveTaskIntakeOwner({ workflow: optionalExecute, enabledWorkflowSteps: ["optional-execute"], agents: [agent("pool")] })) + .toEqual({ status: "rejected", reason: "named-execute-binding-unavailable" }); + }); + + it("pool-resolves no-workflow creates and reports only a genuinely empty pool as unowned", () => { + expect(resolveTaskIntakeOwner({ workflow: "no-workflow-context", agents: [agent("executor")] })).toEqual({ status: "selected", agentId: "executor", source: "executor-pool" }); + expect(resolveTaskIntakeOwner({ workflow: "no-workflow-context", agents: [] })).toEqual({ status: "unowned", reason: "no-eligible-executor" }); + }); + + it("keeps automatic selection role-safe and orders its pool by load, creation time, then id", () => { + const explicitOnly = agent("explicit-only", { runtimeConfig: { assignmentPolicy: "explicit-only" } }); + const paused = agent("paused", { state: "paused" }); + const policyDenied = agent("denied", { runtimeConfig: { assignmentPolicy: "none" } }); + const olderBusy = agent("older-busy", { createdAt: "2025-01-01T00:00:00.000Z" }); + const newerIdle = agent("newer-idle", { createdAt: "2026-01-01T00:00:00.000Z" }); + expect(resolveTaskIntakeOwner({ + workflow: "no-workflow-context", + agents: [explicitOnly, paused, policyDenied, olderBusy, newerIdle], + activeSessions: new Map([["older-busy", 1]]), + })).toEqual({ status: "selected", agentId: "newer-idle", source: "executor-pool" }); + expect(resolveTaskIntakeOwner({ workflow: "no-workflow-context", explicitAssigneeId: "explicit-only", agents: [explicitOnly] })) + .toEqual({ status: "selected", agentId: "explicit-only", source: "explicit" }); + }); + + it("finds a reachable executor nested inside a container template", () => { + const graph: WorkflowIr = { + version: "v2", + nodes: [ + { id: "start", kind: "start", column: "plan" }, + { id: "container", kind: "foreach", column: "work", config: { source: "task-steps", template: { + nodes: [{ id: "nested", kind: "parse-steps", config: { template: { nodes: [{ id: "execute", kind: "prompt", column: "nested-work", config: { seam: "execute" } }], edges: [] } } }], + edges: [], + } } }, + ], + edges: [{ from: "start", to: "container" }], + columns: [ + { id: "plan", name: "Plan", traits: [] }, + { id: "work", name: "Work", traits: [], agent: { agentId: "outer", mode: "override" } }, + { id: "nested-work", name: "Nested work", traits: [], agent: { agentId: "bound", mode: "override" } }, + ], + }; + expect(resolveTaskIntakeOwner({ workflow: graph, agents: [agent("outer"), agent("bound")] })) + .toEqual({ status: "selected", agentId: "bound", source: "execute-binding" }); + }); + + it("keeps rejection, exemption, and unowned outcomes distinct", () => { + expect(resolveTaskIntakeOwner({ workflow: "unresolvable", agents: [] })).toEqual({ status: "rejected", reason: "workflow-unresolvable" }); + expect(resolveTaskIntakeOwner({ workflow: "no-workflow-context", agents: [], ownershipExemptionReason: "fixture" })).toEqual({ status: "exempt", reason: "fixture" }); + }); + + it("accepts the named fixture capability but rejects payload-shaped exemptions", () => { + expect(getInternalIntakeOwnershipExemptionReason(createFixtureIntakeOwnershipExemption())).toBe("fixture"); + expect(getInternalIntakeOwnershipExemptionReason({ reason: "fixture", ownershipExemption: true })).toBeUndefined(); + expect(getInternalIntakeOwnershipExemptionReason("fixture")).toBeUndefined(); + }); +}); diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index 110f40483d..04d3703b13 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -85,6 +85,8 @@ import { type UsageEventInput } from "./tasks/usage-events.js"; import { assertNotLinkedWorktreeOfExistingProject, assertProjectRootDir } from "./central/project-root-guard.js"; import { type DistributedTaskIdAllocator } from "./tasks/distributed-task-id.js"; import { type TaskIdIntegrityReport } from "./tasks/task-id-integrity.js"; +import { AgentStore } from "./agents/agent-store.js"; +import type { IntakeOwnershipExemption } from "./tasks/task-intake-owner-resolver.js"; // file. These are pure behavior-invariant moves — the extracted symbols are // byte-identical to their pre-extraction form. store.ts remains the facade and @@ -503,6 +505,23 @@ export class TaskStore extends EventEmitter { */ public asyncDistributedTaskIdAllocator: DistributedTaskIdAllocator | null = null; + /* + FNXC:IntakeOwnership 2026-08-09-18:04: + Every task create must read the same project-scoped durable-agent snapshot without + constructing a disposable AgentStore per row. Reuse this store-owned seam so normal + Dashboard, API, CLI, and reserved-ID intake share agent-backend lifecycle ownership. + */ + private intakeOwnerAgentStore: AgentStore | null = null; + + public getIntakeOwnerAgentStore(): AgentStore { + if (!this.asyncLayer) throw new Error("Task intake ownership requires an async agent backend"); + return this.intakeOwnerAgentStore ??= new AgentStore({ + rootDir: this.rootDir, + asyncLayer: this.asyncLayer, + projectId: this.asyncLayer.projectId, + }); + } + public agentLogBuffer: Array<{ taskId: string; timestamp: string; @@ -1090,7 +1109,7 @@ export class TaskStore extends EventEmitter { /** * 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; deferTaskCreatedEvent?: boolean; onTaskInserted?: (task: Task) => void; }, ): Promise { + public async _createTaskInternalBackend( input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; resolvedWorkflowIdForOwnership?: string; onProposalClaimConflict?: (task: Task) => void; deferTaskCreatedEvent?: boolean; onTaskInserted?: (task: Task) => void; ownershipExemption?: IntakeOwnershipExemption; }, ): Promise { return _createTaskInternalBackendImpl(this, input, title, resolvedWorkflowSteps, id, options); } @@ -1100,13 +1119,13 @@ export class TaskStore extends EventEmitter { public async _maybeAutoArchiveSameAgentDuplicateBackend( task: Task, input: TaskCreateInput, ): Promise { return _maybeAutoArchiveSameAgentDuplicateBackendImpl(this, task, input); } - async createTask( input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; } ): Promise { + async createTask( input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; ownershipExemption?: IntakeOwnershipExemption; } ): Promise { return createTaskImpl(this, input, options); } - async createTaskWithReservedId( input: TaskCreateInput, options: { taskId: string; createdAt?: string; updatedAt?: string; prompt?: string; applyDefaultWorkflowSteps?: boolean; invokeTaskCreatedHook?: boolean; }, ): Promise { + async createTaskWithReservedId( input: TaskCreateInput, options: { taskId: string; createdAt?: string; updatedAt?: string; prompt?: string; applyDefaultWorkflowSteps?: boolean; invokeTaskCreatedHook?: boolean; ownershipExemption?: IntakeOwnershipExemption; }, ): Promise { 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; deferTaskCreatedEvent?: boolean; onTaskInserted?: (task: Task) => void; }, ): Promise { + public async _createTaskInternal( input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; resolvedWorkflowIdForOwnership?: string; onProposalClaimConflict?: (task: Task) => void; deferTaskCreatedEvent?: boolean; onTaskInserted?: (task: Task) => void; ownershipExemption?: IntakeOwnershipExemption; }, ): Promise { /* 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. @@ -3084,6 +3103,8 @@ Issue #2149 requires read-only type filtering to occur in the file-store before return clearTaskWorkflowSelectionImpl(this, taskId); } async close(): Promise { + this.intakeOwnerAgentStore?.close(); + this.intakeOwnerAgentStore = null; return closeImpl(this); } get fts5Available(): boolean { diff --git a/packages/core/src/task-store/task-creation.ts b/packages/core/src/task-store/task-creation.ts index 54fe2251df..bc310ceb29 100644 --- a/packages/core/src/task-store/task-creation.ts +++ b/packages/core/src/task-store/task-creation.ts @@ -27,7 +27,7 @@ import {buildBootstrapPrompt} from "../mesh/mesh-task-replication.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 {DEFAULT_WORKFLOW_ID, getBuiltinWorkflow, isBuiltinWorkflowId} from "../workflows/builtin-workflows.js"; import {columnsWithFlag} from "../workflows/workflow-lifecycle-traits.js"; import {validateFileScopeInPromptContent} from "../task-store/file-scope.js"; import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js"; @@ -35,9 +35,15 @@ import {withTaskBranchContextInSourceMetadata} from "../task-store/branch-contex import {resolveCreateDeclaredSymbols} from "../tasks/task-symbol-resolution.js"; import {softDeleteTaskRow as softDeleteTaskRowAsync, insertTaskRowInTransaction, isTaskIdConflictError} from "../task-store/async/async-persistence.js"; import {recordRunAuditEvent as recordRunAuditEventAsync} from "../task-store/async/async-audit.js"; +import {recordRunAuditEventWithinTransaction} from "../postgres/data-layer.js"; import type {DbTransaction} from "../postgres/data-layer.js"; import { resolveTaskPrefix } from "./task-prefix.js"; import {assertValidProviderInstanceId} from "../provider-instance.js"; +import { + getInternalIntakeOwnershipExemptionReason, + resolveTaskIntakeOwner, + type IntakeOwnershipExemption, +} from "../tasks/task-intake-owner-resolver.js"; type CreateTaskWithAfterInsert = TaskCreateInput & { /** Internal transaction hook; never persisted in task source metadata. */ @@ -49,24 +55,6 @@ type CreateTaskWithAfterInsert = TaskCreateInput & { skipSameAgentDuplicateIntake?: boolean; }; -function ensureSqliteProposalClaimUniqueness(store: TaskStore): void { - /* - FNXC:EphemeralAgentTaskCreation 2026-07-30-19:10: - The legacy SQLite store remains a supported MessageStore/task-materialization - backend. Its durable partial unique index is the same at-most-once anchor as - PostgreSQL: release/reclaim reuses one stable key, so concurrent creators can - only insert one task and the loser returns that persisted task. - */ - const columns = store.db.prepare("PRAGMA table_info(tasks)").all() as Array<{ name: string }>; - if (!columns.some((column) => column.name === "proposalClaimId")) { - store.db.exec("ALTER TABLE tasks ADD COLUMN proposalClaimId TEXT"); - } - store.db.exec( - "CREATE UNIQUE INDEX IF NOT EXISTS uq_tasks_proposal_claim_id ON tasks(proposalClaimId) WHERE proposalClaimId IS NOT NULL", - ); -} - - /* FNXC:MergedPlanningColumn 2026-07-28-12:55 (U11 precondition): The intake column used to be resolved ONLY as a by-product of materializing workflow steps, so a @@ -150,7 +138,7 @@ async function resolveDefaultWorkflowIntakeColumn(store: TaskStore): Promise Promise; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; },): Promise { + onSummarize?: (description: string) => Promise; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; ownershipExemption?: IntakeOwnershipExemption; },): Promise { /* FNXC:CredentialInstanceSelection 2026-08-01-05:43: validate task authoring input before persistence; ids are stored but runtime credential resolution remains unchanged. */ for (const key of ["credentialInstanceId", "validatorCredentialInstanceId", "planningCredentialInstanceId", "mergerCredentialInstanceId"] as const) { const value = (input as unknown as Record)[key]; @@ -163,6 +151,21 @@ export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateI throw new Error("Description is required and cannot be empty"); } + /* + FNXC:IntakeOwnership 2026-08-09-18:57: + A proposal-claim replay is an idempotent read of its canonical row, not a + new intake. Return before workflow or owner resolution so a later invalid + assignee, unavailable agent backend, or changed executor pool cannot turn a + successful prior proposal into a failed retry or emit a second owner signal. + */ + if (input.proposalClaimId) { + const existing = (await store.listTasks()).find((task) => task.proposalClaimId === input.proposalClaimId); + if (existing) { + options?.onProposalClaimConflict?.(existing); + return existing; + } + } + const selfDefeatingDep = detectSelfDefeatingDependency(input.title, input.dependencies ?? []); if (selfDefeatingDep) { throw new SelfDefeatingDependencyError( @@ -245,18 +248,25 @@ export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateI // Explicit "No workflow": skip default materialization entirely. resolvedWorkflowSteps = undefined; } else { - // Compile + materialize up front so unknown/fragment ids throw BEFORE - // the task row is created (no orphaned steps, no half-created task). - const selected = await store.materializeExplicitWorkflowSteps(explicitWorkflowId); - const explicitStepIds = input.enabledWorkflowSteps !== undefined - ? (resolvedWorkflowSteps ?? []) - : undefined; - resolvedWorkflowSteps = explicitStepIds ?? selected.stepIds; - resolvedEntryColumn = selected.entryColumnId; - pendingWorkflowSelection = { - workflowId: selected.workflowId, - stepIds: explicitStepIds ?? selected.stepIds, - }; + try { + const selected = await store.materializeExplicitWorkflowSteps(explicitWorkflowId); + const explicitStepIds = input.enabledWorkflowSteps !== undefined + ? (resolvedWorkflowSteps ?? []) + : undefined; + resolvedWorkflowSteps = explicitStepIds ?? selected.stepIds; + resolvedEntryColumn = selected.entryColumnId; + pendingWorkflowSelection = { + workflowId: selected.workflowId, + stepIds: explicitStepIds ?? selected.stepIds, + }; + } catch { + /* + FNXC:IntakeOwnership 2026-08-09-19:51: + Explicit workflow materialization is an optimization, not a second ownership authority. + Let the one shared pre-insert boundary compile this id and return its typed + `workflow-unresolvable` rejection, preserving a zero-row create failure for both gateways. + */ + } } } else if (input.enabledWorkflowSteps === undefined) { try { @@ -334,8 +344,26 @@ export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateI onProposalClaimConflict: options?.onProposalClaimConflict, onTaskInserted: () => { insertedTask = true; }, resolvedEntryColumn, + resolvedWorkflowIdForOwnership: pendingWorkflowSelection?.workflowId, + ownershipExemption: options?.ownershipExemption, }, ); + if (!insertedTask) { + /* + FNXC:IntakeOwnership 2026-08-09-20:04: + A concurrent proposal-claim winner returns its canonical task through the + shared boundary. Abort this attempt's unused distributed ID and discard + its unreferenced materialized steps; only the winning insertion may + commit an ID, persist a workflow selection, or publish task:created. + */ + await allocator.abortDistributedTaskIdReservation({ + reservationId: reservation.reservationId, + nodeId, + reason: "failed-create", + }); + await store.cleanupOrphanedMaterializedSteps(pendingWorkflowSelection?.stepIds); + return task; + } await allocator.commitDistributedTaskIdReservation({ reservationId: reservation.reservationId, nodeId, @@ -346,6 +374,11 @@ export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateI nodeId, reason: "failed-create", }).catch(() => undefined); + // FNXC:IntakeOwnership 2026-08-09-17:10: Owner resolution can reject after + // workflow materialization but before insertion. Remove those unreferenced + // step rows just as the reserved-ID gateway does, so rejected creates leave + // neither a task nor orphaned workflow state. + await store.cleanupOrphanedMaterializedSteps(pendingWorkflowSelection?.stepIds); throw err; } @@ -438,9 +471,89 @@ 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; deferTaskCreatedEvent?: boolean; onTaskInserted?: (task: Task) => void; },): Promise { +export class TaskIntakeOwnerResolutionError extends Error { + readonly code = "task-intake-owner-resolution" as const; + constructor(readonly reason: "explicit-assignee-ineligible" | "named-execute-binding-unavailable" | "workflow-unresolvable" | "agent-backend-unavailable") { + super(`Task intake owner resolution failed: ${reason}`); + this.name = "TaskIntakeOwnerResolutionError"; + } +} + +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; resolvedWorkflowIdForOwnership?: string; onProposalClaimConflict?: (task: Task) => void; deferTaskCreatedEvent?: boolean; onTaskInserted?: (task: Task) => void; ownershipExemption?: IntakeOwnershipExemption; },): Promise { const layer = store.asyncLayer!; const now = options?.createdAt ?? new Date().toISOString(); + /* + FNXC:IntakeOwnership 2026-08-09-08:49: + This is the sole shared pre-insert boundary for ordinary and reserved-ID creates. `workflowId: null` + suppresses workflow routing only; it still selects an executor from the pool. Exemption is options-only. + + FNXC:IntakeOwnership 2026-08-09-18:45: + Options may carry only an opaque in-process exemption capability, never a public string or body field. + The runtime symbol check below ignores forged payload-shaped values so normal creators cannot silently + bypass durable executor selection. + */ + let workflow: WorkflowIr | "no-workflow-context" | "unresolvable"; + /* + FNXC:IntakeOwnership 2026-08-09-16:28: + Materialization can canonicalize the selected workflow before the row exists. + Forward that exact selection so ownership binding and the later selection row + cannot diverge; public `workflowId: null` remains the pool-only path. + */ + let effectiveWorkflowIdForOwnership: string | undefined; + if (input.workflowId === null) { + workflow = "no-workflow-context"; + } else { + try { + const workflowId = options?.resolvedWorkflowIdForOwnership + ?? input.workflowId + ?? (await store.getDefaultWorkflowId()) + ?? DEFAULT_WORKFLOW_ID; + effectiveWorkflowIdForOwnership = workflowId; + /* + FNXC:IntakeOwnership 2026-08-09-19:51: + General workflow readers intentionally fall back to builtin:coding for a missing definition. + Intake may not turn a client-named missing workflow into an owned default-workflow task, so + prove the requested definition exists before compiling it at this universal insert boundary. + */ + const exists = isBuiltinWorkflowId(workflowId) + ? getBuiltinWorkflow(workflowId) !== undefined + : (await store.getWorkflowDefinition(workflowId)) !== undefined; + workflow = exists ? await resolveWorkflowIrById(store, workflowId) : "unresolvable"; + } catch { + workflow = "unresolvable"; + } + } + let agents; + try { + /* + FNXC:IntakeOwnership 2026-08-09-18:04: + The TaskStore owns this project-scoped reader. Reusing it prevents each intake row + from opening a fresh durable-agent backend while preserving a visible rejected + outcome if that backend cannot be constructed or read. + */ + agents = await store.getIntakeOwnerAgentStore().listAgents(); + } catch { + agents = undefined; + } + // `taskId` is the durable active-session link. It orders equally eligible pool + // agents without treating a temporary session load as an eligibility failure. + const activeSessions = new Map( + (agents ?? []).flatMap((agent) => agent.taskId ? [[agent.id, 1] as const] : []), + ); + const ownership = resolveTaskIntakeOwner({ + workflow, + // Runtime callers predating the public type may pass null to mean "no + // explicit owner". Normalize it to omission so it follows pool/binding + // resolution rather than becoming an ineligible named assignee. + explicitAssigneeId: input.assignedAgentId ?? undefined, + agents, + enabledWorkflowSteps: resolvedWorkflowSteps, + activeSessions, + ownershipExemptionReason: getInternalIntakeOwnershipExemptionReason(options?.ownershipExemption), + ownModelProvider: input.modelProvider, + ownModelId: input.modelId, + }); + if (ownership.status === "rejected") throw new TaskIntakeOwnerResolutionError(ownership.reason); const normalizedTitle = normalizeTitleForTaskId(title, id); /* FNXC:MergedPlanningColumn 2026-07-29-14:30 (U11 post-merge audit): @@ -502,7 +615,7 @@ export async function _createTaskInternalBackendImpl(store: TaskStore, input: Ta noCommitsExpected: input.noCommitsExpected === true ? true : undefined, enabledWorkflowSteps: resolvedWorkflowSteps, modelPresetId: input.modelPresetId, - assignedAgentId: input.assignedAgentId, + assignedAgentId: ownership.status === "selected" ? ownership.agentId : undefined, assigneeUserId: input.assigneeUserId, scopeOverride: input.scopeOverride === true ? true : undefined, scopeOverrideReason: input.scopeOverrideReason, @@ -699,6 +812,19 @@ export async function _createTaskInternalBackendImpl(store: TaskStore, input: Ta ownsStagingDirectory = false; ownsPromotedTaskDirectory = true; await (input as CreateTaskWithAfterInsert).afterTaskInsert?.(tx, task); + if (ownership.status === "unowned") { + /* + FNXC:IntakeOwnership 2026-08-09-09:30: + The successful no-executor outcome is observable evidence, not a best-effort + post-insert side effect. Write its audit row in this transaction so an audit + failure rolls back the null-owner task and its staged artifacts together. + */ + await recordRunAuditEventWithinTransaction(tx, { + taskId: task.id, agentId: "system", runId: `store:intake-owner:${task.id}`, + domain: "database", mutationType: "task:intake-owner-unresolved", target: task.id, + metadata: { reason: ownership.reason, workflowId: effectiveWorkflowIdForOwnership, source: "intake" }, + }); + } }); } catch (error) { await cleanupPreparedTaskFiles(); @@ -746,7 +872,7 @@ export async function _createTaskInternalBackendImpl(store: TaskStore, input: Ta return task; } -export async function createTaskImpl(store: TaskStore, input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; }): Promise { +export async function createTaskImpl(store: TaskStore, input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; ownershipExemption?: IntakeOwnershipExemption; }): Promise { // U8/R6: apply the reviewLevel creation-time preset (maps level -> enabledWorkflowSteps; explicit wins). input = applyReviewLevelPreset(input); // FNXC:RuntimeTaskOrchestrationAsync 2026-06-24-13:10: @@ -760,7 +886,7 @@ export async function createTaskImpl(store: TaskStore, input: TaskCreateInput, o return store.createTaskBackend(input, options); } -export async function createTaskWithReservedIdImpl(store: TaskStore, input: TaskCreateInput, options: { taskId: string; createdAt?: string; updatedAt?: string; prompt?: string; applyDefaultWorkflowSteps?: boolean; invokeTaskCreatedHook?: boolean; },): Promise { +export async function createTaskWithReservedIdImpl(store: TaskStore, input: TaskCreateInput, options: { taskId: string; createdAt?: string; updatedAt?: string; prompt?: string; applyDefaultWorkflowSteps?: boolean; invokeTaskCreatedHook?: boolean; ownershipExemption?: IntakeOwnershipExemption; },): Promise { // U8/R6: apply the reviewLevel creation-time preset (maps level -> enabledWorkflowSteps; explicit wins). input = applyReviewLevelPreset(input); if (!input.description?.trim()) { @@ -777,7 +903,12 @@ export async function createTaskWithReservedIdImpl(store: TaskStore, input: Task } if (input.proposalClaimId) { - ensureSqliteProposalClaimUniqueness(store); + /* + FNXC:IntakeOwnership 2026-08-09-20:18: + Reserved-ID creation is PostgreSQL-only. Proposal-claim uniqueness belongs + to its project-scoped partial index, so this replay read must not touch the + removed SQLite backend before the shared insert boundary handles a race. + */ const existing = (await store.listTasks()).find((task) => task.proposalClaimId === input.proposalClaimId); if (existing) return existing; } @@ -822,18 +953,24 @@ export async function createTaskWithReservedIdImpl(store: TaskStore, input: Task // Explicit "No workflow": skip default materialization entirely. resolvedWorkflowSteps = undefined; } else { - // Compile + materialize up front so unknown/fragment ids throw BEFORE - // the task row is created (no orphaned steps, no half-created task). - const selected = await store.materializeExplicitWorkflowSteps(explicitWorkflowId); - const explicitStepIds = input.enabledWorkflowSteps !== undefined - ? (resolvedWorkflowSteps ?? []) - : undefined; - resolvedWorkflowSteps = explicitStepIds ?? selected.stepIds; - resolvedEntryColumn = selected.entryColumnId; - pendingWorkflowSelection = { - workflowId: selected.workflowId, - stepIds: explicitStepIds ?? selected.stepIds, - }; + try { + const selected = await store.materializeExplicitWorkflowSteps(explicitWorkflowId); + const explicitStepIds = input.enabledWorkflowSteps !== undefined + ? (resolvedWorkflowSteps ?? []) + : undefined; + resolvedWorkflowSteps = explicitStepIds ?? selected.stepIds; + resolvedEntryColumn = selected.entryColumnId; + pendingWorkflowSelection = { + workflowId: selected.workflowId, + stepIds: explicitStepIds ?? selected.stepIds, + }; + } catch { + /* + FNXC:IntakeOwnership 2026-08-09-19:51: + Reserved-ID creation must reach the same typed owner-resolution rejection as ordinary + creation when an explicit workflow cannot compile; materialization may not bypass it. + */ + } } } else if (input.enabledWorkflowSteps === undefined && options.applyDefaultWorkflowSteps !== false) { // Mirror createTask: a configured project default workflow takes @@ -883,6 +1020,7 @@ export async function createTaskWithReservedIdImpl(store: TaskStore, input: Task } let createdTask: Task; + let proposalReplay = false; try { createdTask = await store._createTaskInternal(input, title, resolvedWorkflowSteps, id, { createdAt: options.createdAt, @@ -890,6 +1028,9 @@ export async function createTaskWithReservedIdImpl(store: TaskStore, input: Task promptOverride: options.prompt, invokeTaskCreatedHook: options.invokeTaskCreatedHook, resolvedEntryColumn, + resolvedWorkflowIdForOwnership: pendingWorkflowSelection?.workflowId, + ownershipExemption: options.ownershipExemption, + onProposalClaimConflict: () => { proposalReplay = true; }, }); } catch (err) { // The task row was never created, so any default-workflow steps we @@ -902,6 +1043,18 @@ export async function createTaskWithReservedIdImpl(store: TaskStore, input: Task throw err; } + if (proposalReplay) { + /* + FNXC:IntakeOwnership 2026-08-09-20:04: + The reserved-ID gateway can race after its preflight replay lookup. The + shared boundary reports the winner through this callback, so never attach + this loser's materialized selection to the canonical task or leave its + generated step rows behind. + */ + await store.cleanupOrphanedMaterializedSteps(pendingWorkflowSelection?.stepIds); + return createdTask; + } + // Record the inherited workflow selection now that the task row exists. if (pendingWorkflowSelection) { try { diff --git a/packages/core/src/tasks/task-intake-owner-resolver.ts b/packages/core/src/tasks/task-intake-owner-resolver.ts new file mode 100644 index 0000000000..a32c2c9862 --- /dev/null +++ b/packages/core/src/tasks/task-intake-owner-resolver.ts @@ -0,0 +1,239 @@ +import type { Agent } from "../types.js"; +import type { WorkflowIr, WorkflowIrNode } from "../workflows/workflow-ir-types.js"; +import { classifyWorkflowAgentNode } from "../workflows/workflow-ir-types.js"; +import { instanceNodeId, resolveColumnAgentBinding, resolveEffectiveAgent } from "../agents/column-agent-resolver.js"; +import { isEphemeralAgent } from "../types.js"; +import { canAgentReceiveImplementationTasks, isAgentAutoAssignable } from "../agents/agent-role-policy.js"; + +export type TaskIntakeOwnerResolution = + | { status: "selected"; agentId: string; source: "explicit" | "execute-binding" | "executor-pool" } + | { status: "exempt"; reason: "terminal" | "historical" | "fixture" } + | { status: "rejected"; reason: "explicit-assignee-ineligible" | "named-execute-binding-unavailable" | "workflow-unresolvable" | "agent-backend-unavailable" } + | { status: "unowned"; reason: "no-eligible-executor" }; + +type IntakeOwnershipExemptionReason = "terminal" | "historical" | "fixture"; + +declare const intakeOwnershipExemptionBrand: unique symbol; + +/** + * Opaque capability accepted only by TaskStore's internal creation plumbing. + * JSON and public tool/API inputs cannot carry its symbol brand. + */ +export type IntakeOwnershipExemption = { + readonly reason: IntakeOwnershipExemptionReason; + readonly [intakeOwnershipExemptionBrand]: true; +}; + +const intakeOwnershipExemptionToken = Symbol("fusion.task-intake-ownership-exemption"); + +type RuntimeIntakeOwnershipExemption = IntakeOwnershipExemption & { + readonly [intakeOwnershipExemptionToken]: true; +}; + +/* +FNXC:IntakeOwnership 2026-08-09-19:49: +The exemption has named in-process factories for terminal, historical, and fixture-only creation so +its deliberate-null outcome is reachable without becoming a generic public opt-out. The symbol token +keeps JSON bodies, tool arguments, and forwarded payloads unable to forge any of these capabilities. +*/ +function createIntakeOwnershipExemption(reason: IntakeOwnershipExemptionReason): IntakeOwnershipExemption { + return { + reason, + [intakeOwnershipExemptionToken]: true, + } as unknown as IntakeOwnershipExemption; +} + +/** Creates the closed capability for a terminal task that must never be dispatched. */ +export function createTerminalIntakeOwnershipExemption(): IntakeOwnershipExemption { + return createIntakeOwnershipExemption("terminal"); +} + +/** Creates the closed capability for an imported historical record. */ +export function createHistoricalIntakeOwnershipExemption(): IntakeOwnershipExemption { + return createIntakeOwnershipExemption("historical"); +} + +/** Creates the closed capability for test fixtures that are not executable intake. */ +export function createFixtureIntakeOwnershipExemption(): IntakeOwnershipExemption { + return createIntakeOwnershipExemption("fixture"); +} + +/** Converts a closed internal capability into the resolver's simple pure-input reason. */ +export function getInternalIntakeOwnershipExemptionReason(value: unknown): IntakeOwnershipExemptionReason | undefined { + if (!value || typeof value !== "object") return undefined; + const candidate = value as Partial; + return candidate[intakeOwnershipExemptionToken] === true + && (candidate.reason === "terminal" || candidate.reason === "historical" || candidate.reason === "fixture") + ? candidate.reason + : undefined; +} + +export interface ResolveTaskIntakeOwnerInput { + workflow: WorkflowIr | "no-workflow-context" | "unresolvable"; + explicitAssigneeId?: string; + agents?: readonly Agent[]; + /** Internal plumbing supplies this only after validating an opaque exemption capability. */ + ownershipExemptionReason?: IntakeOwnershipExemptionReason; + ownModelProvider?: string; + ownModelId?: string; + /** Per-task optional-group enablement. Omitted means no optional groups run. */ + enabledWorkflowSteps?: readonly string[]; + /** Current workflow-session load snapshot; it only orders pool candidates. */ + activeSessions?: ReadonlyMap; +} + +/* +FNXC:IntakeOwnership 2026-08-09-08:49: +Intake ownership is a durable executor identity, not the mutable workflow-stage principal. `unowned` +means a genuinely empty eligible executor pool and inserts a visible null owner; `rejected` means the +create must fail before insertion. Neither outcome is called `held`, which is reserved for routing waits. +*/ +export function resolveTaskIntakeOwner(input: ResolveTaskIntakeOwnerInput): TaskIntakeOwnerResolution { + if (input.ownershipExemptionReason) return { status: "exempt", reason: input.ownershipExemptionReason }; + if (!input.agents) return { status: "rejected", reason: "agent-backend-unavailable" }; + if (input.workflow === "unresolvable") return { status: "rejected", reason: "workflow-unresolvable" }; + + const eligible = (agent: Agent | undefined, automatic: boolean): agent is Agent => Boolean( + agent + && !isEphemeralAgent(agent) + && agent.roles.includes("executor") + && agent.runtimeConfig?.enabled !== false + && agent.state !== "paused" + && agent.state !== "error" + && canAgentReceiveImplementationTasks(agent) + // `explicit-only` is valid only for a caller's deliberate assignee. A binding + // or pool choice is automatic routing and must retain that policy boundary. + && (!automatic || isAgentAutoAssignable(agent)), + ); + const byId = new Map(input.agents.map((agent) => [agent.id, agent])); + if (input.explicitAssigneeId !== undefined) { + const explicit = byId.get(input.explicitAssigneeId); + return eligible(explicit, false) + ? { status: "selected", agentId: explicit.id, source: "explicit" } + : { status: "rejected", reason: "explicit-assignee-ineligible" }; + } + + if (input.workflow !== "no-workflow-context") { + const execute = firstReachableExecuteNode(input.workflow, new Set(input.enabledWorkflowSteps)); + if (execute) { + const effective = resolveEffectiveAgent({ + binding: execute.binding, + ownModelProvider: input.ownModelProvider, + ownModelId: input.ownModelId, + }); + if (effective.source === "column-agent") { + const bound = byId.get(effective.agentId); + return eligible(bound, true) + ? { status: "selected", agentId: bound.id, source: "execute-binding" } + : { status: "rejected", reason: "named-execute-binding-unavailable" }; + } + } + } + + const activeSessions = input.activeSessions ?? new Map(); + const pool = input.agents.filter((agent) => eligible(agent, true)).sort((a, b) => + (activeSessions.get(a.id) ?? 0) - (activeSessions.get(b.id) ?? 0) + || a.createdAt.localeCompare(b.createdAt) + || a.id.localeCompare(b.id), + ); + return pool[0] + ? { status: "selected", agentId: pool[0].id, source: "executor-pool" } + : { status: "unowned", reason: "no-eligible-executor" }; +} + +type ReachableExecuteNode = { node: WorkflowIrNode; binding: ReturnType }; + +/* +FNXC:IntakeOwnership 2026-08-09-09:20: +The durable owner follows graph reachability, not authoring-array order: an unreachable executor +binding must never steal a new task from the first execute node its workflow can actually dispatch. +Container templates are searched through their local entry walk; foreach bindings use the canonical +instance-node identity so template inheritance preserves the same column-agent policy as runtime. +*/ +function firstReachableExecuteNode(ir: WorkflowIr, enabledWorkflowSteps: ReadonlySet): ReachableExecuteNode | undefined { + const topLevel = reachableNodes(ir.nodes, ir.edges); + for (const node of topLevel) { + const found = findExecuteInNode(ir, node, enabledWorkflowSteps); + if (found) return found; + } + return undefined; +} + +function findExecuteInNode( + ir: WorkflowIr, + node: WorkflowIrNode, + enabledWorkflowSteps: ReadonlySet, +): ReachableExecuteNode | undefined { + if (classifyWorkflowAgentNode(node) === "executor") { + return { node, binding: resolveColumnAgentBinding(ir, node.id) }; + } + const template = (node.config as { template?: { nodes?: WorkflowIrNode[]; edges?: WorkflowIr["edges"] } } | undefined)?.template; + if (!template?.nodes?.length) return undefined; + // FNXC:IntakeOwnership 2026-08-09-16:38: Optional-group templates are not + // reachable when their group is disabled for this task. Never let a disabled + // execute binding select or reject a durable owner for a path that cannot run. + if (node.kind === "optional-group" && !enabledWorkflowSteps.has(node.id)) return undefined; + + for (const inner of reachableNodes(template.nodes, template.edges ?? [])) { + if (classifyWorkflowAgentNode(inner) === "executor") { + // FNXC:IntakeOwnership 2026-08-09-18:15: A nested foreach has no + // top-level instance identity for resolveColumnAgentBinding. Retain its + // template-column fallback so an execute node's own column wins over its outer container. + const binding = node.kind === "foreach" + ? resolveColumnAgentBinding(ir, instanceNodeId(node.id, 0, inner.id)) + ?? bindingForTemplateNode(ir, node, inner) + : bindingForTemplateNode(ir, node, inner); + return { node: inner, binding }; + } + // Templates may nest foreach/containers. Continue the reachability walk rather + // than treating a non-executor wrapper as a terminal graph fragment. + const nested = findExecuteInNode(ir, inner, enabledWorkflowSteps); + if (nested) { + return { + ...nested, + binding: nested.binding ?? bindingForTemplateNode(ir, node, inner), + }; + } + } + return undefined; +} + +function bindingForTemplateNode( + ir: WorkflowIr, + container: WorkflowIrNode, + node: WorkflowIrNode, +): ReturnType { + const columnId = node.column ?? container.column; + return ir.version === "v2" && columnId !== undefined + ? ir.columns.find((column) => column.id === columnId)?.agent + : undefined; +} + +function reachableNodes(nodes: readonly WorkflowIrNode[], edges: readonly WorkflowIr["edges"][number][]): WorkflowIrNode[] { + const byId = new Map(nodes.map((node) => [node.id, node])); + const outgoing = new Map(); + for (const edge of edges) { + if (edge.kind === "rework" || !byId.has(edge.from) || !byId.has(edge.to)) continue; + const list = outgoing.get(edge.from) ?? []; + list.push(edge.to); + outgoing.set(edge.from, list); + } + const start = nodes.find((node) => node.kind === "start"); + // Invalid/legacy fragments have no start; retaining declaration order is a safe + // compatibility fallback while valid workflow IR always takes the graph walk. + if (!start) return [...nodes]; + + const ordered: WorkflowIrNode[] = []; + const seen = new Set(); + const queue = [start.id]; + while (queue.length > 0) { + const id = queue.shift()!; + if (seen.has(id)) continue; + seen.add(id); + const node = byId.get(id); + if (!node) continue; + ordered.push(node); + queue.push(...(outgoing.get(id) ?? [])); + } + return ordered; +} diff --git a/packages/dashboard/src/__tests__/routes-tasks.test.ts b/packages/dashboard/src/__tests__/routes-tasks.test.ts index e25a5b9175..b3c1c5296a 100644 --- a/packages/dashboard/src/__tests__/routes-tasks.test.ts +++ b/packages/dashboard/src/__tests__/routes-tasks.test.ts @@ -756,6 +756,53 @@ describe("POST /tasks", () => { ); }); + it("forwards no-workflow and explicit owner inputs but drops public exemption-shaped fields", async () => { + (store.createTask as ReturnType).mockResolvedValue({ + ...FAKE_TASK_DETAIL, + assignedAgentId: "executor-1", + }); + + const res = await REQUEST( + buildApp(), + "POST", + "/api/tasks", + JSON.stringify({ + description: "No workflow still needs an executor owner", + workflowId: null, + agentId: "executor-1", + ownershipExemption: true, + }), + { "Content-Type": "application/json" }, + ); + + expect(res.status).toBe(201); + const createInput = (store.createTask as ReturnType).mock.calls[0][0]; + expect(createInput).toMatchObject({ + description: "No workflow still needs an executor owner", + workflowId: null, + assignedAgentId: "executor-1", + }); + expect(createInput).not.toHaveProperty("ownershipExemption"); + }); + + it("returns a typed client failure when the universal owner resolver rejects an explicit owner", async () => { + (store.createTask as ReturnType).mockRejectedValue( + new Error("Task intake owner resolution failed: explicit-assignee-ineligible"), + ); + + const res = await REQUEST( + buildApp(), + "POST", + "/api/tasks", + JSON.stringify({ description: "Reject a triage-only owner", assignedAgentId: "triage-agent" }), + { "Content-Type": "application/json" }, + ); + + expect(res.status).toBe(400); + expect(res.body).toMatchObject({ error: "Task intake owner resolution failed: explicit-assignee-ineligible" }); + expect(store.createTask).toHaveBeenCalledTimes(1); + }); + it("does not synchronously create tracking issues in POST /tasks route", async () => { const createIssueSpy = vi.spyOn(GitHubClient.prototype, "createIssue").mockResolvedValue({ owner: "task", diff --git a/packages/dashboard/src/routes/__tests__/task-create-intake-owner.pg.test.ts b/packages/dashboard/src/routes/__tests__/task-create-intake-owner.pg.test.ts new file mode 100644 index 0000000000..8720578b93 --- /dev/null +++ b/packages/dashboard/src/routes/__tests__/task-create-intake-owner.pg.test.ts @@ -0,0 +1,195 @@ +// @vitest-environment node + +import { afterEach, beforeEach, expect, it, vi } from "vitest"; +import express from "express"; + +const { mockSession, mockCreateResolvedAgentSession } = vi.hoisted(() => ({ + mockSession: () => ({ + model: { provider: "mock", id: "intake-owner" }, + state: {}, + prompt: vi.fn().mockResolvedValue(undefined), + dispose: vi.fn(), + abort: vi.fn().mockResolvedValue(undefined), + subscribe: vi.fn(() => () => undefined), + sessionManager: { getLeafId: vi.fn().mockReturnValue(null) }, + }), + mockCreateResolvedAgentSession: vi.fn(), +})); + +vi.mock("../../../../engine/src/agents/agent-session-helpers.js", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + createResolvedAgentSession: mockCreateResolvedAgentSession, + }; +}); + +vi.mock("../../../../engine/src/mcp/mcp-resolution.js", async (importOriginal) => { + const actual = await importOriginal(); + return { ...actual, resolveMcpServersForStore: vi.fn(async () => ({ servers: [], errors: [] })) }; +}); + +vi.mock("../../../../engine/src/pi.js", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + describeModel: vi.fn(() => "mock/intake-owner"), + promptWithFallback: vi.fn(async (session: { prompt: (prompt: string) => Promise }, prompt: string) => await session.prompt(prompt)), + }; +}); + +vi.mock("../../../../engine/src/worktree/worktree-acquisition.js", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + acquireTaskWorktree: vi.fn(async () => ({ + worktreePath: "/tmp/fusion-intake-owner-heartbeat", + branch: "fusion/intake-owner", + source: "existing", + hydrated: false, + isResume: true, + })), + }; +}); + +import { AgentStore, type TaskStore } from "@fusion/core"; +import { BUILTIN_CODING_WORKFLOW_IR } from "../../../../core/src/workflows/builtin-coding-workflow-ir.js"; +import { TriageProcessor } from "../../../../engine/src/triage.js"; +import { HeartbeatMonitor } from "../../../../engine/src/agent-heartbeat.js"; +import { createTaskStoreForTest, pgDescribe, type PgTestHarness } from "../../../../core/src/__test-utils__/pg-test-harness.js"; +import { createApiRoutes } from "../../routes.js"; +import { request as REQUEST } from "../../test-request.js"; + +/* +FNXC:IntakeOwnership 2026-08-09-17:11: +Dashboard intake must persist an executor owner before its response makes the new +card visible. This production-route PostgreSQL fixture protects the reported +unowned-card symptom and keeps workflowId:null on the same pool-resolution path. +*/ +pgDescribe("POST /tasks intake owner", () => { + let harness: PgTestHarness; + let store: TaskStore; + let agents: AgentStore; + let app: express.Express; + + beforeEach(async () => { + mockCreateResolvedAgentSession.mockImplementation(async () => ({ + session: mockSession(), + settleFallbackDispatch: async () => undefined, + runtimeId: "mock", + })); + harness = await createTaskStoreForTest({ + prefix: "fusion_dashboard_intake_owner", + projectId: "dashboard-intake-owner", + }); + store = harness.store; + agents = new AgentStore({ rootDir: harness.rootDir, asyncLayer: harness.layer, projectId: "dashboard-intake-owner", taskStore: store }); + app = express(); + app.use(express.json()); + app.use("/api", createApiRoutes(store)); + }); + + afterEach(async () => { + agents.close(); + await harness.teardown(); + }); + + const post = (body: unknown) => REQUEST(app, "POST", "/api/tasks", JSON.stringify(body), { + "content-type": "application/json", + }); + + it("keeps Dashboard intake ownership stable through planning routing and the owner heartbeat claim", async () => { + const planner = await agents.createAgent({ name: "Dashboard planner", role: "triage" }); + const executor = await agents.createAgent({ name: "Dashboard executor", role: "executor" }); + + const response = await post({ description: "Dashboard task gets an implementation owner" }); + + expect(response.status).toBe(201); + const created = response.body as { id: string; assignedAgentId?: string }; + expect(created.assignedAgentId).toBe(executor.id); + expect(created.assignedAgentId).not.toBe(planner.id); + const persisted = await store.getTask(created.id); + expect(persisted?.assignedAgentId).toBe(executor.id); + + /* + FNXC:IntakeOwnership 2026-08-09-20:54: + The reported Dashboard/API failure must be proved through live lifecycle + entry points, not routing helpers. Planning receives the separately fenced + triage principal, while the executor's real heartbeat discovers the durable + owner inbox; neither stage may rewrite the persisted executor owner. + */ + const triage = new TriageProcessor(store, harness.rootDir, { agentStore: agents }); + await triage.specifyTask(persisted!); + + const afterPlanning = await store.getTask(created.id); + expect(afterPlanning?.assignedAgentId).toBe(executor.id); + expect(mockCreateResolvedAgentSession).toHaveBeenCalledWith(expect.objectContaining({ + sessionPurpose: "triage", + taskId: created.id, + actionGateContext: expect.objectContaining({ agentId: planner.id }), + })); + + const heartbeat = new HeartbeatMonitor({ store: agents, taskStore: store, rootDir: harness.rootDir }); + await heartbeat.executeHeartbeat({ agentId: executor.id, source: "on_demand" }); + + expect(mockCreateResolvedAgentSession).toHaveBeenCalledWith(expect.objectContaining({ + sessionPurpose: "heartbeat", + actionGateContext: expect.objectContaining({ agentId: executor.id }), + })); + expect((await agents.getAgent(executor.id))?.taskId).toBe(created.id); + expect((await store.getTask(created.id))?.assignedAgentId).toBe(executor.id); + }); + + it("rejects invalid explicit and named execute owners without creating an unowned card", async () => { + const executor = await agents.createAgent({ name: "Route fallback executor", role: "executor" }); + const unavailableBinding = await store.createWorkflowDefinition({ + name: "Route unavailable execute binding", + ir: { + ...BUILTIN_CODING_WORKFLOW_IR, + columns: BUILTIN_CODING_WORKFLOW_IR.columns.map((column) => column.id === "in-progress" + ? { ...column, agent: { agentId: "route-missing-executor", mode: "override" as const } } + : column), + }, + }); + + const invalidExplicit = await post({ + description: "Route invalid explicit owner", + assignedAgentId: "route-missing-explicit-owner", + }); + const invalidBinding = await post({ + description: "Route unavailable named execute owner", + workflowId: unavailableBinding.id, + }); + + expect(executor.id).toBeTruthy(); + expect(invalidExplicit.status).toBeGreaterThanOrEqual(400); + expect(invalidBinding.status).toBeGreaterThanOrEqual(400); + expect((await store.listTasks()).map((task) => task.description)).not.toContain("Route invalid explicit owner"); + expect((await store.listTasks()).map((task) => task.description)).not.toContain("Route unavailable named execute owner"); + }); + + it("keeps no-workflow creates owned and makes a genuine empty pool observable", async () => { + const executor = await agents.createAgent({ name: "No workflow executor", role: "executor" }); + const ownedResponse = await post({ + description: "No workflow still needs an executor", + workflowId: null, + ownershipExemption: true, + }); + + expect(ownedResponse.status).toBe(201); + const owned = ownedResponse.body as { id: string; assignedAgentId?: string }; + expect(owned.assignedAgentId).toBe(executor.id); + expect((await store.getTask(owned.id))?.assignedAgentId).toBe(executor.id); + + await agents.deleteAgent(executor.id); + const unownedResponse = await post({ description: "Fresh project has no executor" }); + + expect(unownedResponse.status).toBe(201); + const unowned = unownedResponse.body as { id: string; assignedAgentId?: string }; + expect(unowned.assignedAgentId).toBeUndefined(); + expect((await store.getTask(unowned.id))?.assignedAgentId).toBeUndefined(); + const audit = await store.getRunAuditEventsAsync({ taskId: unowned.id }); + expect(audit.filter((event) => event.mutationType === "task:intake-owner-unresolved")) + .toHaveLength(1); + }); +}); diff --git a/packages/dashboard/src/routes/register-task-workflow-routes.ts b/packages/dashboard/src/routes/register-task-workflow-routes.ts index 4e957cfa43..7ed2bf9f4d 100644 --- a/packages/dashboard/src/routes/register-task-workflow-routes.ts +++ b/packages/dashboard/src/routes/register-task-workflow-routes.ts @@ -1630,6 +1630,8 @@ export function registerTaskWorkflowRoutes(ctx: ApiRoutesContext, deps: TaskWork breakIntoSubtasks, enabledWorkflowSteps, workflowId, + agentId, + assignedAgentId, modelPresetId, modelProvider, modelId, @@ -1743,6 +1745,12 @@ export function registerTaskWorkflowRoutes(ctx: ApiRoutesContext, deps: TaskWork throw badRequest("workflowId must be a string or null"); } + // The public aliases normalize once; owner eligibility is enforced by the shared store boundary. + const requestedOwnerId = assignedAgentId ?? agentId; + if (requestedOwnerId !== undefined && (typeof requestedOwnerId !== "string" || requestedOwnerId.trim() === "")) { + throw badRequest("agentId must be a non-empty string"); + } + // Check for summarize flag in request const summarize = req.body.summarize === true; @@ -2076,6 +2084,7 @@ export function registerTaskWorkflowRoutes(ctx: ApiRoutesContext, deps: TaskWork // U6/R3: forward only when the client set it (string | null). Leaving it // absent preserves the project-default inheritance behavior. ...(workflowId !== undefined ? { workflowId: workflowId as string | null } : {}), + ...(typeof requestedOwnerId === "string" ? { assignedAgentId: requestedOwnerId.trim() } : {}), modelPresetId: validateOptionalModelField(modelPresetId, "modelPresetId"), modelProvider: executorModel.provider ?? undefined, modelId: executorModel.modelId ?? undefined, @@ -2241,7 +2250,8 @@ export function registerTaskWorkflowRoutes(ctx: ApiRoutesContext, deps: TaskWork message.includes("must be a string") || message.includes("must be an array of strings") || /^Workflow '.*' not found$/.test(message) - || /is a fragment and cannot be selected/.test(message); + || /is a fragment and cannot be selected/.test(message) + || message.startsWith("Task intake owner resolution failed:"); const status = isClientError ? 400 : 500; throw new ApiError(status, message); } diff --git a/packages/engine/src/__tests__/workflow-agent-routing.test.ts b/packages/engine/src/__tests__/workflow-agent-routing.test.ts index 6260b4ca32..29a7766922 100644 --- a/packages/engine/src/__tests__/workflow-agent-routing.test.ts +++ b/packages/engine/src/__tests__/workflow-agent-routing.test.ts @@ -7,13 +7,26 @@ const agent = (id: string, roles: string[], createdAt = "2026-01-01T00:00:00.000 const ir: any = { version: "v2", name: "test", columns: [{ id: "todo", name: "Todo", traits: [] }], nodes: [] }; describe("routeWorkflowPrincipal", () => { - it("uses exact review override and returns to task owner for execution", () => { - const owner = agent("owner", ["custom"]); - const reviewer = agent("reviewer", ["custom"]); + it("uses exact review override and returns to an executor owner for execution", () => { + const owner = agent("owner", ["executor"]); + const reviewer = agent("reviewer", ["reviewer"]); expect(routeWorkflowPrincipal({ task: { assignedAgentId: "owner" }, ir, node: { id: "r", kind: "prompt", reviewerAgentId: "reviewer", config: { workflowRole: "reviewer" } }, agents: [owner, reviewer] })).toMatchObject({ status: "routed", route: { agent: reviewer, authority: "review-node-override" } }); expect(routeWorkflowPrincipal({ task: { assignedAgentId: "owner" }, ir, node: { id: "e", kind: "prompt", config: { seam: "execute" } }, agents: [owner, reviewer] })).toMatchObject({ status: "routed", route: { agent: owner, authority: "task-assignee" } }); }); + it("holds rather than falling back when a named owner is not eligible to execute", () => { + const pool = agent("pool", ["executor"]); + const node = { id: "e", kind: "prompt", config: { seam: "execute" } }; + for (const owner of [ + agent("triage-owner", ["triage"]), + { ...agent("disabled-owner", ["executor"]), runtimeConfig: { enabled: false } }, + { ...agent("policy-denied-owner", ["executor"]), runtimeConfig: { assignmentPolicy: "none" } }, + ]) { + expect(routeWorkflowPrincipal({ task: { assignedAgentId: owner.id }, ir, node, agents: [owner, pool] })) + .toEqual({ status: "held", role: "executor", reason: "named-principal-unavailable" }); + } + }); + it("holds rather than falling back when a named principal is unavailable", () => { const paused = { ...agent("owner", ["executor"]), state: "paused" }; const pool = agent("pool", ["executor"]); diff --git a/packages/engine/src/agents/workflow-agent-router.ts b/packages/engine/src/agents/workflow-agent-router.ts index f24148e328..a3ec1ca2ab 100644 --- a/packages/engine/src/agents/workflow-agent-router.ts +++ b/packages/engine/src/agents/workflow-agent-router.ts @@ -1,4 +1,5 @@ import { + canAgentReceiveImplementationTasks, classifyWorkflowAgentNode, isEphemeralAgent, resolveColumnAgentBinding, @@ -94,7 +95,9 @@ export function validateFencedWorkflowPrincipal(input: { if (classifiedRole !== input.role) { return { status: "held", role: input.role, reason: "named-principal-unavailable" }; } - if (input.authority === "task-assignee" && input.task.assignedAgentId !== input.principalAgentId) { + if (input.authority === "task-assignee" && ( + input.role !== "executor" || input.task.assignedAgentId !== input.principalAgentId + )) { return { status: "held", role: input.role, reason: "named-principal-unavailable" }; } if (input.authority === "review-node-override" && ( @@ -115,6 +118,8 @@ export function validateFencedWorkflowPrincipal(input: { } const agent = input.agents.find((candidate) => candidate.id === input.principalAgentId); return available(agent, input.activeSessions ?? new Map()) + && agent.roles.includes(input.role) + && (input.role !== "executor" || canAgentReceiveImplementationTasks(agent)) ? { status: "routed", route: { agent, role: input.role, authority: input.authority } } : { status: "held", role: input.role, reason: "named-principal-unavailable" }; } @@ -128,7 +133,13 @@ function available(agent: Agent | undefined, activeSessions: ReadonlyMap { if (!id) return undefined; const agent = byId.get(id); + /* + FNXC:IntakeOwnership 2026-08-09-09:55: + A durable task owner is only authority for an executor node when it still has + the executor role. This revalidation protects legacy rows and concurrent + reassignment from dispatching implementation work through a stale non-executor owner. + */ + // FNXC:IntakeOwnership 2026-08-09-18:15: A persisted executor owner is + // revalidated at dispatch. Runtime disablement and assignment policy changes + // revoke implementation authority instead of letting a stale owner run work. return available(agent, activeSessions) + && agent.roles.includes(role) + && (role !== "executor" || canAgentReceiveImplementationTasks(agent)) ? { status: "routed", route: { agent, role, authority } } : { status: "held", role, reason: "named-principal-unavailable" }; }; @@ -168,8 +190,15 @@ export function routeWorkflowPrincipal(input: { const overridden = named(input.node.reviewerAgentId, "review-node-override"); if (overridden) return overridden; } - const owner = named(input.task.assignedAgentId, "task-assignee"); - if (owner) return owner; + /* + FNXC:IntakeOwnership 2026-08-09-08:49: + Durable task ownership names an executor only. Planning and review must retain their separately + fenced stage principals; treating the owner as authority there would route planners to executors. + */ + if (role === "executor") { + const owner = named(input.task.assignedAgentId, "task-assignee"); + if (owner) return owner; + } const column = resolveColumnAgentBinding(input.ir, input.node.id); const bound = named(column?.agentId, "column-binding"); if (bound) return bound;