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 <noreply@anthropic.com>
This commit is contained in:
gsxdsm
2026-08-07 08:51:05 -07:00
parent 039db99307
commit 1f9b0e644a
3 changed files with 123 additions and 13 deletions

View File

@@ -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.

View File

@@ -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);
});

View File

@@ -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 * Each helper targets the project-schema tables via Drizzle and is the async
* equivalent of the sync this.db.prepare() call sites below. * equivalent of the sync this.db.prepare() call sites below.
*/ */
import type { QueryHandle } from "../async-stores/async-mission-store-queries.js";
import { import {
writeAgent as writeAgentAsync, writeAgent as writeAgentAsync,
readAgent as readAgentAsync, readAgent as readAgentAsync,
@@ -731,10 +732,10 @@ export class AgentStore extends EventEmitter {
* @param name - Agent name to match exactly * @param name - Agent name to match exactly
* @returns Matching non-ephemeral agent, or null when none exists * @returns Matching non-ephemeral agent, or null when none exists
*/ */
async findAgentByName(name: string): Promise<Agent | null> { async findAgentByName(name: string, executor?: QueryHandle): Promise<Agent | null> {
// FNXC:SqliteFinalRemoval 2026-06-25-23:45: // FNXC:SqliteFinalRemoval 2026-06-25-23:45:
// Backend mode: read via async Drizzle helper, filter ephemeral in-memory. // 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) { for (const agent of agents) {
if (!isEphemeralAgent(agent)) { if (!isEphemeralAgent(agent)) {
return this.parseAgent(agent as unknown as AgentData); return this.parseAgent(agent as unknown as AgentData);
@@ -772,7 +773,7 @@ export class AgentStore extends EventEmitter {
* @returns The created agent * @returns The created agent
* @throws Error if input is invalid or a duplicate non-ephemeral name exists * @throws Error if input is invalid or a duplicate non-ephemeral name exists
*/ */
async createAgent(input: AgentCreateInput): Promise<Agent> { async createAgent(input: AgentCreateInput, executor?: QueryHandle): Promise<Agent> {
if (!input.name?.trim()) { if (!input.name?.trim()) {
throw new Error("Agent name is required"); 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 }); const ephemeral = isEphemeralAgent({ metadata, name: input.name, role: roles[0], reportsTo: input.reportsTo });
if (!ephemeral) { if (!ephemeral) {
const existing = await this.findAgentByName(normalizedName); const existing = await this.findAgentByName(normalizedName, executor);
if (existing) { if (existing) {
throw new Error(`Agent with name "${normalizedName}" already exists (agentId: ${existing.id})`); throw new Error(`Agent with name "${normalizedName}" already exists (agentId: ${existing.id})`);
} }
@@ -838,7 +839,7 @@ export class AgentStore extends EventEmitter {
...(resolvedHeartbeatProcedurePath && { heartbeatProcedurePath: resolvedHeartbeatProcedurePath }), ...(resolvedHeartbeatProcedurePath && { heartbeatProcedurePath: resolvedHeartbeatProcedurePath }),
}; };
await this.writeAgent(agent); await this.writeAgent(agent, executor);
this.emit("agent:created", agent); this.emit("agent:created", agent);
return agent; return agent;
@@ -1960,11 +1961,14 @@ export class AgentStore extends EventEmitter {
* @param filter - Optional filter criteria * @param filter - Optional filter criteria
* @returns Array of agents * @returns Array of agents
*/ */
async listAgents(filter?: { state?: AgentState; role?: AgentCapability; includeEphemeral?: boolean }): Promise<Agent[]> { async listAgents(
filter?: { state?: AgentState; role?: AgentCapability; includeEphemeral?: boolean },
executor?: QueryHandle,
): Promise<Agent[]> {
// FNXC:WorkflowAgentRouting 2026-08-07-03:12: // FNXC:WorkflowAgentRouting 2026-08-07-03:12:
// Role-pool membership is canonical multi-tag state, so SQL must not use the // Role-pool membership is canonical multi-tag state, so SQL must not use the
// deprecated singular projection to exclude a matching durable principal. // 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 return agents
.map((a) => this.parseAgent(a as unknown as AgentData)) .map((a) => this.parseAgent(a as unknown as AgentData))
.filter((agent) => !filter?.role || agent.roles.includes(filter.role)) .filter((agent) => !filter?.role || agent.roles.includes(filter.role))
@@ -1978,14 +1982,14 @@ export class AgentStore extends EventEmitter {
* members with the same tag. * members with the same tag.
*/ */
async provisionBuiltinWorkflowRoleAgents(): Promise<Agent[]> { async provisionBuiltinWorkflowRoleAgents(): Promise<Agent[]> {
const provision = async (): Promise<Agent[]> => { const provision = async (executor?: QueryHandle): Promise<Agent[]> => {
const definitions: ReadonlyArray<{ role: AgentCapability; name: string; title: string }> = [ const definitions: ReadonlyArray<{ role: AgentCapability; name: string; title: string }> = [
{ role: "triage", name: "Workflow Planner", title: "Built-in workflow planning owner" }, { role: "triage", name: "Workflow Planner", title: "Built-in workflow planning owner" },
{ role: "executor", name: "Workflow Executor", title: "Built-in workflow execution owner" }, { role: "executor", name: "Workflow Executor", title: "Built-in workflow execution owner" },
{ role: "reviewer", name: "Workflow Reviewer", title: "Built-in workflow review owner" }, { role: "reviewer", name: "Workflow Reviewer", title: "Built-in workflow review owner" },
{ role: "merger", name: "Workflow Merger", title: "Built-in workflow merge 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( const builtins = new Map(
existing existing
.filter((agent) => agent.metadata?.builtInWorkflowRole === true) .filter((agent) => agent.metadata?.builtInWorkflowRole === true)
@@ -2005,7 +2009,7 @@ export class AgentStore extends EventEmitter {
metadata: { builtInWorkflowRole: true, workflowRole: definition.role }, metadata: { builtInWorkflowRole: true, workflowRole: definition.role },
// Disabled scheduling does not make the agent unavailable for graph sessions. // Disabled scheduling does not make the agent unavailable for graph sessions.
runtimeConfig: { enabled: false }, runtimeConfig: { enabled: false },
})); }, executor));
} }
return result; return result;
}; };
@@ -2018,9 +2022,19 @@ export class AgentStore extends EventEmitter {
* duplicate owners. User-created same-role agents are intentionally outside * duplicate owners. User-created same-role agents are intentionally outside
* this provenance lock and remain valid pool members. * 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) => { return this.asyncLayer.transactionImmediate(async (tx) => {
await tx.execute(sql`SELECT pg_advisory_xact_lock(hashtext(${this.backendProjectId}), hashtext('builtin-workflow-role-provisioning'))`); 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<void> { private async writeAgent(agent: Agent, executor?: QueryHandle): Promise<void> {
// FNXC:SqliteFinalRemoval 2026-06-25-23:40: // FNXC:SqliteFinalRemoval 2026-06-25-23:40:
// Backend mode: delegate to async Drizzle writeAgent helper. // 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; return;
} }