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) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-8843-intake-agent-assignment.md
Normal file
7
.changeset/fn-8843-intake-agent-assignment.md
Normal file
@@ -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.
|
||||
@@ -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.
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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 {
|
||||
|
||||
141
packages/core/src/__tests__/task-intake-owner-resolver.test.ts
Normal file
141
packages/core/src/__tests__/task-intake-owner-resolver.test.ts
Normal file
@@ -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> = {}): 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();
|
||||
});
|
||||
});
|
||||
@@ -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<TaskStoreEvents> {
|
||||
*/
|
||||
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<TaskStoreEvents> {
|
||||
/**
|
||||
* FNXC:RuntimeTaskOrchestrationAsync 2026-06-24-13:25:
|
||||
*/
|
||||
public async _createTaskInternalBackend( input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; onProposalClaimConflict?: (task: Task) => void; deferTaskCreatedEvent?: boolean; onTaskInserted?: (task: Task) => void; }, ): Promise<Task> {
|
||||
public async _createTaskInternalBackend( input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; resolvedWorkflowIdForOwnership?: string; onProposalClaimConflict?: (task: Task) => void; deferTaskCreatedEvent?: boolean; onTaskInserted?: (task: Task) => void; ownershipExemption?: IntakeOwnershipExemption; }, ): Promise<Task> {
|
||||
return _createTaskInternalBackendImpl(this, input, title, resolvedWorkflowSteps, id, options);
|
||||
}
|
||||
|
||||
@@ -1100,13 +1119,13 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
||||
public async _maybeAutoArchiveSameAgentDuplicateBackend( task: Task, input: TaskCreateInput, ): Promise<void> {
|
||||
return _maybeAutoArchiveSameAgentDuplicateBackendImpl(this, task, input);
|
||||
}
|
||||
async createTask( input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise<string | null>; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; } ): Promise<Task> {
|
||||
async createTask( input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise<string | null>; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; ownershipExemption?: IntakeOwnershipExemption; } ): Promise<Task> {
|
||||
return createTaskImpl(this, input, options);
|
||||
}
|
||||
async createTaskWithReservedId( input: TaskCreateInput, options: { taskId: string; createdAt?: string; updatedAt?: string; prompt?: string; applyDefaultWorkflowSteps?: boolean; invokeTaskCreatedHook?: boolean; }, ): Promise<Task> {
|
||||
async createTaskWithReservedId( input: TaskCreateInput, options: { taskId: string; createdAt?: string; updatedAt?: string; prompt?: string; applyDefaultWorkflowSteps?: boolean; invokeTaskCreatedHook?: boolean; ownershipExemption?: IntakeOwnershipExemption; }, ): Promise<Task> {
|
||||
return createTaskWithReservedIdImpl(this, input, options);
|
||||
}
|
||||
public async _createTaskInternal( input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; onProposalClaimConflict?: (task: Task) => void; deferTaskCreatedEvent?: boolean; onTaskInserted?: (task: Task) => void; }, ): Promise<Task> {
|
||||
public async _createTaskInternal( input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; resolvedWorkflowIdForOwnership?: string; onProposalClaimConflict?: (task: Task) => void; deferTaskCreatedEvent?: boolean; onTaskInserted?: (task: Task) => void; ownershipExemption?: IntakeOwnershipExemption; }, ): Promise<Task> {
|
||||
/*
|
||||
FNXC:SqliteDualPathCleanup 2026-07-26-14:05:
|
||||
Task create is PostgreSQL-only (layer.transactionImmediate + insertTaskRowInTransaction). The former sync SQLite _createTaskInternalImpl arm is deleted; production always injects AsyncDataLayer.
|
||||
@@ -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<void> {
|
||||
this.intakeOwnerAgentStore?.close();
|
||||
this.intakeOwnerAgentStore = null;
|
||||
return closeImpl(this);
|
||||
}
|
||||
get fts5Available(): boolean {
|
||||
|
||||
@@ -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<str
|
||||
}
|
||||
|
||||
export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateInput, options?: {
|
||||
onSummarize?: (description: string) => Promise<string | null>; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; },): Promise<Task> {
|
||||
onSummarize?: (description: string) => Promise<string | null>; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; ownershipExemption?: IntakeOwnershipExemption; },): Promise<Task> {
|
||||
/* 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<string, unknown>)[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<Task> {
|
||||
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<Task> {
|
||||
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<string | null>; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; }): Promise<Task> {
|
||||
export async function createTaskImpl(store: TaskStore, input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise<string | null>; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; ownershipExemption?: IntakeOwnershipExemption; }): Promise<Task> {
|
||||
// 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<Task> {
|
||||
export async function createTaskWithReservedIdImpl(store: TaskStore, input: TaskCreateInput, options: { taskId: string; createdAt?: string; updatedAt?: string; prompt?: string; applyDefaultWorkflowSteps?: boolean; invokeTaskCreatedHook?: boolean; ownershipExemption?: IntakeOwnershipExemption; },): Promise<Task> {
|
||||
// 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 {
|
||||
|
||||
239
packages/core/src/tasks/task-intake-owner-resolver.ts
Normal file
239
packages/core/src/tasks/task-intake-owner-resolver.ts
Normal file
@@ -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<RuntimeIntakeOwnershipExemption>;
|
||||
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<string, number>;
|
||||
}
|
||||
|
||||
/*
|
||||
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<string, number>();
|
||||
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<typeof resolveColumnAgentBinding> };
|
||||
|
||||
/*
|
||||
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<string>): 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<string>,
|
||||
): 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<typeof resolveColumnAgentBinding> {
|
||||
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<string, string[]>();
|
||||
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<string>();
|
||||
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;
|
||||
}
|
||||
@@ -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<typeof vi.fn>).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<typeof vi.fn>).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<typeof vi.fn>).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",
|
||||
|
||||
@@ -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<typeof import("../../../../engine/src/agents/agent-session-helpers.js")>();
|
||||
return {
|
||||
...actual,
|
||||
createResolvedAgentSession: mockCreateResolvedAgentSession,
|
||||
};
|
||||
});
|
||||
|
||||
vi.mock("../../../../engine/src/mcp/mcp-resolution.js", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("../../../../engine/src/mcp/mcp-resolution.js")>();
|
||||
return { ...actual, resolveMcpServersForStore: vi.fn(async () => ({ servers: [], errors: [] })) };
|
||||
});
|
||||
|
||||
vi.mock("../../../../engine/src/pi.js", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("../../../../engine/src/pi.js")>();
|
||||
return {
|
||||
...actual,
|
||||
describeModel: vi.fn(() => "mock/intake-owner"),
|
||||
promptWithFallback: vi.fn(async (session: { prompt: (prompt: string) => Promise<void> }, prompt: string) => await session.prompt(prompt)),
|
||||
};
|
||||
});
|
||||
|
||||
vi.mock("../../../../engine/src/worktree/worktree-acquisition.js", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("../../../../engine/src/worktree/worktree-acquisition.js")>();
|
||||
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);
|
||||
});
|
||||
});
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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"]);
|
||||
|
||||
@@ -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<string,
|
||||
* never satisfy a role-pool route, even when their singular compatibility role
|
||||
* matches. A named transient identity is likewise unavailable and holds closed.
|
||||
*/
|
||||
if (!agent || isEphemeralAgent(agent) || agent.state === "paused" || agent.state === "error") return false;
|
||||
if (
|
||||
!agent
|
||||
|| isEphemeralAgent(agent)
|
||||
|| agent.runtimeConfig?.enabled === false
|
||||
|| agent.state === "paused"
|
||||
|| agent.state === "error"
|
||||
) return false;
|
||||
const max = agent.runtimeConfig?.maxWorkflowSessions;
|
||||
return typeof max !== "number" || activeSessions.get(agent.id) === undefined || activeSessions.get(agent.id)! < max;
|
||||
}
|
||||
@@ -160,7 +171,18 @@ export function routeWorkflowPrincipal(input: {
|
||||
const named = (id: string | undefined, authority: WorkflowPrincipalAuthority): WorkflowPrincipalRouteResult | undefined => {
|
||||
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;
|
||||
|
||||
Reference in New Issue
Block a user