diff --git a/doc/DATABASE.md b/doc/DATABASE.md index 3c82a93a2b..151581873d 100644 --- a/doc/DATABASE.md +++ b/doc/DATABASE.md @@ -165,6 +165,15 @@ idempotent actor synchronization operations, not arbitrary transactions. A persistent outage still fails the request after the bounded retries; each connection attempt remains subject to the configured database connect timeout. +## Execution identity row locks + +Identity initialization, credential acquisition, and steering reconciliation lock +the task before its run. These operations use `FOR NO KEY UPDATE`: they change +identity state, not parent keys. The lock still serializes identity writers and +blocks concurrent task or run updates. It allows audit inserts to retain their +foreign-key `KEY SHARE` locks without waiting on identity acquisition. The audit +foreign keys and their deletion behavior remain enforced. + ## Switching between modes The database mode is controlled by `DATABASE_URL`: diff --git a/server/src/__tests__/run-identity.test.ts b/server/src/__tests__/run-identity.test.ts index 9155a9b9fa..d6dbbbfc3f 100644 --- a/server/src/__tests__/run-identity.test.ts +++ b/server/src/__tests__/run-identity.test.ts @@ -1,7 +1,7 @@ import { randomUUID } from "node:crypto"; import { eq, sql } from "drizzle-orm"; import { afterAll, beforeAll, describe, expect, it } from "vitest"; -import { agentWakeupRequests, agents, companies, createDb, heartbeatRuns, heartbeatRunEvents, issueComments, issueThreadInteractions, issues } from "@paperclipai/db"; +import { agentWakeupRequests, agents, companies, createDb, heartbeatRuns, heartbeatRunEvents, issueComments, issueThreadInteractions, issues, secretAccessEvents } from "@paperclipai/db"; import { getEmbeddedPostgresTestSupport, startEmbeddedPostgresTestDatabase } from "./helpers/embedded-postgres.js"; import { acceptSteeredIdentity, captureRunIdentity, initializeRunIdentity, listRunIdentityContexts, rejectSteeredIdentity, reserveSteeredIdentity } from "../services/run-identity.js"; @@ -217,11 +217,11 @@ const support = await getEmbeddedPostgresTestSupport(); } }); - it("does not deadlock identity initialization against a task mutation that also updates the run", async () => { + it.each(["update", "no key update"] as const)("serializes identity initialization behind a task's %s lock", async (lockMode) => { const input = await seed(); let initialization!: ReturnType; await db.transaction(async (tx) => { - await tx.select().from(issues).where(eq(issues.id, input.issueId)).for("update"); + await tx.select().from(issues).where(eq(issues.id, input.issueId)).for(lockMode); const [backend] = await tx.execute(sql`select pg_backend_pid() as pid`) as unknown as Array<{ pid: number }>; initialization = initializeRunIdentity(db, { ...input, messageIds: [], responsibleUserId: "A", cause: "instruction" }); // Wait until initialization is blocked by this task mutation, rather than @@ -241,6 +241,47 @@ const support = await getEmbeddedPostgresTestSupport(); await expect(initialization).resolves.toMatchObject({ responsibleUserId: "A" }); }); + it.each(["initialize", "capture"] as const)( + "can %s identity while an audit append holds foreign-key locks", + async (operation) => { + const input = await seed(); + if (operation === "capture") { + await initializeRunIdentity(db, { ...input, messageIds: [], responsibleUserId: "A", cause: "instruction" }); + } + // A real append holds KEY SHARE on both parent rows until it commits. + // Bound the other connection's wait so a conflicting lock fails this + // regression instead of leaving both transactions waiting for each other. + const identityDb = createDb(`${database.connectionString}?options=-c%20lock_timeout%3D1000`, { + maxConnections: 1, + }); + const [settings] = await identityDb.execute(sql`show lock_timeout`); + expect(settings?.lock_timeout).toBe("1s"); + await db.transaction(async (audit) => { + await audit.insert(secretAccessEvents).values({ + companyId: input.companyId, + heartbeatRunId: input.runId, + issueId: input.issueId, + provider: "local_encrypted", + actorType: "agent", + actorId: input.agentId, + consumerType: "agent", + consumerId: input.agentId, + outcome: "granted", + }); + if (operation === "initialize") { + await expect(initializeRunIdentity(identityDb, { + ...input, messageIds: [], responsibleUserId: "A", cause: "instruction", + })).resolves.toMatchObject({ responsibleUserId: "A" }); + } else { + await expect(captureRunIdentity(identityDb, input)).resolves.toMatchObject({ + run: { responsibleUserId: "A" }, + context: { responsibleUserId: "A" }, + }); + } + }); + }, + ); + it("does not turn a company-default fallback into personal consent on continuation", async () => { const input = await seed(); await initializeRunIdentity(db, { ...input, messageIds: [], responsibleUserId: "A", cause: "company_default" }); diff --git a/server/src/services/run-identity.ts b/server/src/services/run-identity.ts index dfc5f5bb89..496b58719d 100644 --- a/server/src/services/run-identity.ts +++ b/server/src/services/run-identity.ts @@ -57,7 +57,13 @@ export async function explicitOperatorRunIdentity( export type RunIdentityContext = typeof runIdentityContexts.$inferSelect; type Executor = Pick; -/** Match task mutation ordering: lock the task before the run, never the reverse. */ +/** + * Lock the task before the run, matching task mutation ordering. Identity + * operations do not change parent keys: NO KEY UPDATE still serializes writers + * and steering, while allowing audit inserts to check their foreign keys. + * FOR UPDATE can deadlock with an append that holds KEY SHARE on the run and + * then checks the task while identity capture holds the task and waits on the run. + */ async function lockIdentityTask( executor: Pick, companyId: string, @@ -85,7 +91,7 @@ async function lockIdentityTask( .select({ id: issues.id }) .from(issues) .where(and(eq(issues.id, issueId), eq(issues.companyId, companyId))) - .for("update"); + .for("no key update"); } async function append( @@ -181,7 +187,7 @@ export async function initializeRunIdentity( eq(heartbeatRuns.companyId, input.companyId), ), ) - .for("update"); + .for("no key update"); if (!run) throw forbidden("Run identity does not belong to this company"); if (run.activeIdentityContextId) { const [current] = await tx @@ -345,7 +351,7 @@ export async function reserveSteeredIdentity( eq(heartbeatRuns.companyId, input.companyId), ), ) - .for("update"); + .for("no key update"); // Processes started before the broker rollout keep their original environment. if (!run?.activeIdentityContextId) return null; const [pending] = await tx @@ -439,7 +445,7 @@ export async function captureRunIdentity( eq(heartbeatRuns.agentId, input.agentId), ), ) - .for("update"); + .for("no key update"); if (!run || run.status !== "running") throw forbidden( "Credential acquisition requires this agent's active run", @@ -512,7 +518,7 @@ export async function reconcileSteeredIdentity( eq(heartbeatRuns.companyId, context.companyId), ), ) - .for("update"); + .for("no key update"); if (!run) return; await acceptSteeredIdentity(tx, context); });