mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-06 10:48:12 +02:00
fix: prevent run identity locks from blocking audit checks (#14478)
Use NO KEY UPDATE for identity locks so audit foreign-key checks can proceed while identity writers remain serialized. Preserve task-before-run ordering, company scoping, and foreign keys. Verified with PostgreSQL concurrency regressions, focused tests, typecheck, build, and full CI. Co-Authored-By: Paperclip <noreply@paperclip.ing>
This commit is contained in:
1 parent
faf72cb1d5
commit
4f6cf5b3ff
3 files changed
+65
-9
No files matched your search
@@ -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`:
|
||||
|
||||
@@ -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<typeof initializeRunIdentity>;
|
||||
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" });
|
||||
|
||||
@@ -57,7 +57,13 @@ export async function explicitOperatorRunIdentity(
|
||||
export type RunIdentityContext = typeof runIdentityContexts.$inferSelect;
|
||||
type Executor = Pick<Db, "select" | "insert" | "update">;
|
||||
|
||||
/** 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<Db, "select">,
|
||||
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);
|
||||
});
|
||||
|
||||
Reference in new issue
Block a user