diff --git a/.changeset/fix-builtin-role-provisioning-deadlock.md b/.changeset/fix-builtin-role-provisioning-deadlock.md new file mode 100644 index 0000000000..493ebf0e57 --- /dev/null +++ b/.changeset/fix-builtin-role-provisioning-deadlock.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Fix a startup deadlock that made the dashboard stop responding to all requests. +category: fix +dev: `provisionBuiltinWorkflowRoleAgents` (FN-8764) held a `pg_advisory_xact_lock` transaction while running its reads/writes on the pool, requiring a second connection. With concurrent callers blocking on the same lock and `DEFAULT_POOL_MAX=3`, the pool self-deadlocked and every DB-backed API route queued forever. `listAgents`/`findAgentByName`/`createAgent`/`writeAgent` now accept an optional `QueryHandle` so the provisioning work runs on the locking transaction. diff --git a/packages/core/src/__tests__/agent-store-builtin-role-provisioning-pool.test.ts b/packages/core/src/__tests__/agent-store-builtin-role-provisioning-pool.test.ts new file mode 100644 index 0000000000..d924a84d79 --- /dev/null +++ b/packages/core/src/__tests__/agent-store-builtin-role-provisioning-pool.test.ts @@ -0,0 +1,89 @@ +import { afterEach, beforeEach, expect, it } from "vitest"; +import { AgentStore } from "../agents/agent-store.js"; +import { createTaskStoreForTest, pgDescribe, type PgTestHarness } from "../__test-utils__/pg-test-harness.js"; + +/* +FNXC:WorkflowAgentRouting 2026-08-07-16:40: +Regression guard for the FN-8764 provisioning deadlock. `provisionBuiltinWorkflowRoleAgents` +takes a project-scoped `pg_advisory_xact_lock` inside `transactionImmediate`, so the lock +holder occupies one pooled connection for the whole transaction. When the provisioning work +ran on the POOL instead of on `tx`, the holder needed a SECOND connection to finish while +concurrent callers occupied the remaining slots blocking on that same lock — with +DEFAULT_POOL_MAX=3 that self-deadlocked and every later DB-backed query (i.e. every API +route) queued forever behind an exhausted pool. + +The invariant is "the lock and its work share one connection", so these tests bound the pool +rather than reproducing the original three-caller race: `poolMax: 1` makes ANY second +connection checkout unsatisfiable, which fails the pre-fix code deterministically instead of +depending on scheduling. Both the create path (no built-ins yet) and the idempotent re-entry +path (built-ins already present) are covered, since only the former exercises writes under +the lock. Each assertion carries its own timeout so a regression surfaces as a failure rather +than a hung suite. +*/ +pgDescribe("AgentStore built-in workflow role provisioning under a saturated pool", () => { + let harness: PgTestHarness; + let agentStore: AgentStore; + + beforeEach(async () => { + harness = await createTaskStoreForTest({ poolMax: 1, prefix: "fusion_test_provision_pool" }); + agentStore = new AgentStore({ + rootDir: harness.rootDir, + // The advisory lock is project-scoped, so the layer must be project-bound; + // the shared harness layer is deliberately project-agnostic (projectId ""). + asyncLayer: { ...harness.layer, projectId: "proj_provision_pool" }, + taskStore: harness.store, + }); + }); + + afterEach(async () => { + await harness?.teardown(); + }); + + it("completes the initial create path on a single-connection pool", async () => { + const agents = await agentStore.provisionBuiltinWorkflowRoleAgents(); + + expect(agents).toHaveLength(4); + expect(agents.map((a) => a.metadata?.workflowRole).sort()).toEqual([ + "executor", + "merger", + "reviewer", + "triage", + ]); + for (const agent of agents) { + expect(agent.metadata?.builtInWorkflowRole).toBe(true); + } + }, 20_000); + + it("stays idempotent on re-entry without checking out a second connection", async () => { + const first = await agentStore.provisionBuiltinWorkflowRoleAgents(); + const second = await agentStore.provisionBuiltinWorkflowRoleAgents(); + + expect(second).toHaveLength(4); + // Re-entry must reuse the same durable owners, not add duplicates. + expect(second.map((a) => a.id).sort()).toEqual(first.map((a) => a.id).sort()); + + const durable = (await agentStore.listAgents({ includeEphemeral: true })).filter( + (a) => a.metadata?.builtInWorkflowRole === true, + ); + expect(durable).toHaveLength(4); + }, 20_000); + + it("serializes concurrent callers without deadlocking or duplicating owners", async () => { + // The original failure needed >1 in-flight caller. Even serialized behind a + // single connection, concurrent callers must converge on one set of owners. + const [a, b, c] = await Promise.all([ + agentStore.provisionBuiltinWorkflowRoleAgents(), + agentStore.provisionBuiltinWorkflowRoleAgents(), + agentStore.provisionBuiltinWorkflowRoleAgents(), + ]); + + for (const result of [a, b, c]) { + expect(result).toHaveLength(4); + } + + const durable = (await agentStore.listAgents({ includeEphemeral: true })).filter( + (agent) => agent.metadata?.builtInWorkflowRole === true, + ); + expect(durable).toHaveLength(4); + }, 30_000); +}); diff --git a/packages/core/src/agents/agent-store.ts b/packages/core/src/agents/agent-store.ts index 4b10bb2c0a..eca7ca3614 100644 --- a/packages/core/src/agents/agent-store.ts +++ b/packages/core/src/agents/agent-store.ts @@ -70,6 +70,7 @@ import { and, eq, gt, lte, sql } from "drizzle-orm"; * Each helper targets the project-schema tables via Drizzle and is the async * equivalent of the sync this.db.prepare() call sites below. */ +import type { QueryHandle } from "../async-stores/async-mission-store-queries.js"; import { writeAgent as writeAgentAsync, readAgent as readAgentAsync, @@ -731,10 +732,10 @@ export class AgentStore extends EventEmitter { * @param name - Agent name to match exactly * @returns Matching non-ephemeral agent, or null when none exists */ - async findAgentByName(name: string): Promise { + async findAgentByName(name: string, executor?: QueryHandle): Promise { // FNXC:SqliteFinalRemoval 2026-06-25-23:45: // Backend mode: read via async Drizzle helper, filter ephemeral in-memory. - const agents = await findAgentRowsByNameAsync(this.asyncLayer!.db, name); + const agents = await findAgentRowsByNameAsync(executor ?? this.asyncLayer!.db, name); for (const agent of agents) { if (!isEphemeralAgent(agent)) { return this.parseAgent(agent as unknown as AgentData); @@ -772,7 +773,7 @@ export class AgentStore extends EventEmitter { * @returns The created agent * @throws Error if input is invalid or a duplicate non-ephemeral name exists */ - async createAgent(input: AgentCreateInput): Promise { + async createAgent(input: AgentCreateInput, executor?: QueryHandle): Promise { if (!input.name?.trim()) { throw new Error("Agent name is required"); } @@ -783,7 +784,7 @@ export class AgentStore extends EventEmitter { const ephemeral = isEphemeralAgent({ metadata, name: input.name, role: roles[0], reportsTo: input.reportsTo }); if (!ephemeral) { - const existing = await this.findAgentByName(normalizedName); + const existing = await this.findAgentByName(normalizedName, executor); if (existing) { throw new Error(`Agent with name "${normalizedName}" already exists (agentId: ${existing.id})`); } @@ -838,7 +839,7 @@ export class AgentStore extends EventEmitter { ...(resolvedHeartbeatProcedurePath && { heartbeatProcedurePath: resolvedHeartbeatProcedurePath }), }; - await this.writeAgent(agent); + await this.writeAgent(agent, executor); this.emit("agent:created", agent); return agent; @@ -1960,11 +1961,14 @@ export class AgentStore extends EventEmitter { * @param filter - Optional filter criteria * @returns Array of agents */ - async listAgents(filter?: { state?: AgentState; role?: AgentCapability; includeEphemeral?: boolean }): Promise { + async listAgents( + filter?: { state?: AgentState; role?: AgentCapability; includeEphemeral?: boolean }, + executor?: QueryHandle, + ): Promise { // FNXC:WorkflowAgentRouting 2026-08-07-03:12: // Role-pool membership is canonical multi-tag state, so SQL must not use the // deprecated singular projection to exclude a matching durable principal. - const agents = await listAgentRowsAsync(this.asyncLayer!.db, { state: filter?.state }); + const agents = await listAgentRowsAsync(executor ?? this.asyncLayer!.db, { state: filter?.state }); return agents .map((a) => this.parseAgent(a as unknown as AgentData)) .filter((agent) => !filter?.role || agent.roles.includes(filter.role)) @@ -1978,14 +1982,14 @@ export class AgentStore extends EventEmitter { * members with the same tag. */ async provisionBuiltinWorkflowRoleAgents(): Promise { - const provision = async (): Promise => { + const provision = async (executor?: QueryHandle): Promise => { const definitions: ReadonlyArray<{ role: AgentCapability; name: string; title: string }> = [ { role: "triage", name: "Workflow Planner", title: "Built-in workflow planning owner" }, { role: "executor", name: "Workflow Executor", title: "Built-in workflow execution owner" }, { role: "reviewer", name: "Workflow Reviewer", title: "Built-in workflow review owner" }, { role: "merger", name: "Workflow Merger", title: "Built-in workflow merge owner" }, ]; - const existing = await this.listAgents({ includeEphemeral: true }); + const existing = await this.listAgents({ includeEphemeral: true }, executor); const builtins = new Map( existing .filter((agent) => agent.metadata?.builtInWorkflowRole === true) @@ -2005,7 +2009,7 @@ export class AgentStore extends EventEmitter { metadata: { builtInWorkflowRole: true, workflowRole: definition.role }, // Disabled scheduling does not make the agent unavailable for graph sessions. runtimeConfig: { enabled: false }, - })); + }, executor)); } return result; }; @@ -2018,9 +2022,19 @@ export class AgentStore extends EventEmitter { * duplicate owners. User-created same-role agents are intentionally outside * this provenance lock and remain valid pool members. */ + /* + * FNXC:WorkflowAgentRouting 2026-08-07-16:22: + * The provisioning work MUST run on `tx`, not on the pool. The lock holder + * previously called provision() with no executor, so its reads/writes asked + * the pool for a SECOND connection while concurrent callers occupied the + * remaining slots blocking on this same advisory lock. With DEFAULT_POOL_MAX=3 + * that self-deadlocks: the holder can never finish, the waiters can never take + * the lock, and every later query — i.e. every DB-backed API route — queues + * forever behind an exhausted pool. Keep the lock and the work on one connection. + */ return this.asyncLayer.transactionImmediate(async (tx) => { await tx.execute(sql`SELECT pg_advisory_xact_lock(hashtext(${this.backendProjectId}), hashtext('builtin-workflow-role-provisioning'))`); - return provision(); + return provision(tx); }); } @@ -3047,10 +3061,10 @@ export class AgentStore extends EventEmitter { }; } - private async writeAgent(agent: Agent): Promise { + private async writeAgent(agent: Agent, executor?: QueryHandle): Promise { // FNXC:SqliteFinalRemoval 2026-06-25-23:40: // Backend mode: delegate to async Drizzle writeAgent helper. - await writeAgentAsync(this.asyncLayer!.db, agent, this.asyncLayer!.projectId); + await writeAgentAsync(executor ?? this.asyncLayer!.db, agent, this.asyncLayer!.projectId); return; }