From 1f9b0e644abb27e19803637803d74e37d7c45ce2 Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Fri, 7 Aug 2026 08:51:05 -0700 Subject: [PATCH] fix(FN-8764): run built-in role provisioning on the locking transaction provisionBuiltinWorkflowRoleAgents took a project-scoped pg_advisory_xact_lock inside transactionImmediate, then ran provision() against the pool. The lock holder therefore needed a SECOND pooled connection to finish while concurrent callers occupied the remaining slots blocking on that same advisory lock. With DEFAULT_POOL_MAX=3 this self-deadlocked: the holder could never complete, the waiters could never take the lock, and every subsequent query -- that is, every DB-backed API route -- queued forever behind an exhausted pool. Observed as a dashboard that booted ("Ready in 6.5s") and then answered no /api request while the event loop sat idle in kevent; pg_stat_activity showed one session idle in transaction holding the lock and two active sessions waiting on it. Thread an optional QueryHandle through listAgents, findAgentByName, createAgent, and writeAgent so provisioning runs on tx and the lock and its work share one connection. Regression test bounds the pool to a single connection, which makes any second checkout unsatisfiable and fails deterministically rather than racing. Verified by reverting the one-line fix: the suite hangs past 300s instead of passing in under 4s. Co-Authored-By: Claude Opus 5 --- .../fix-builtin-role-provisioning-deadlock.md | 7 ++ ...ore-builtin-role-provisioning-pool.test.ts | 89 +++++++++++++++++++ packages/core/src/agents/agent-store.ts | 40 ++++++--- 3 files changed, 123 insertions(+), 13 deletions(-) create mode 100644 .changeset/fix-builtin-role-provisioning-deadlock.md create mode 100644 packages/core/src/__tests__/agent-store-builtin-role-provisioning-pool.test.ts 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; }