From 021dc0ad54907c738ecba6e45707a7227e7a2b3f Mon Sep 17 00:00:00 2001 From: nickyleach <331803+nickyleach@users.noreply.github.com> Date: Thu, 10 Sep 2026 22:55:37 +0000 Subject: [PATCH] fix(server): bound the live steering-state probe so it cannot hold row locks open buildQueueSnapshot awaited getNativeSessionSteeringState with no time limit, inside a transaction that holds row locks on the issue, wake, and run rows with `for("update")`. A stalled native-runtime call would have kept those locks open for as long as the provider took to answer. Add a fixed timeout around the probe so it always falls back to "temporarily_unavailable" within a bounded time. Co-authored-by: Paperclip --- .../adapters/queued-comment-postgres.test.ts | 54 +++++++++++++++++++ .../adapters/queued-comment-postgres.ts | 26 ++++++++- 2 files changed, 79 insertions(+), 1 deletion(-) diff --git a/server/src/modules/wake-queue/adapters/queued-comment-postgres.test.ts b/server/src/modules/wake-queue/adapters/queued-comment-postgres.test.ts index f29482b826..c2f4741ced 100644 --- a/server/src/modules/wake-queue/adapters/queued-comment-postgres.test.ts +++ b/server/src/modules/wake-queue/adapters/queued-comment-postgres.test.ts @@ -505,4 +505,58 @@ describeEmbeddedPostgres("queued-comment postgres adapter", () => { expect(second.queue.entries).toHaveLength(0); expect(second.queue.steeringDisposition).toBe("temporarily_unavailable"); }); + + it("bounds the live steering probe so a stalled provider answer cannot hold the transaction's row locks open", async () => { + const companyId = await seedCompany(); + const agentId = await seedAgent({ companyId, adapterType: "paperclip_runner" }); + const issueId = await seedIssue({ companyId, assigneeAgentId: agentId }); + const firstCommentId = await seedComment({ companyId, issueId, authorUserId: "user-1" }); + const secondCommentId = await seedComment({ companyId, issueId, authorUserId: "user-1" }); + const wakeId = await seedDeferredWake({ + companyId, + agentId, + issueId, + commentIds: [firstCommentId, secondCommentId], + }); + const targetRunId = await seedRunningNativeRun({ companyId, agentId, issueId }); + + steerNativeSessionMock.mockResolvedValueOnce({ turnId: "turn-1" }); + // The live provider never answers, standing in for a stalled + // native-runtime call. The probe must resolve on its own bound instead + // of leaving the steer's transaction, and its row locks, open forever. + // A second queued comment stays behind after this steer, so the shared + // rule still asks for a live probe instead of short-circuiting to + // "temporarily_unavailable" for an empty queue. + getNativeSessionSteeringStateMock.mockImplementationOnce(() => new Promise(() => {})); + + const issueLock = createQueuedCommentIssueLockWriter(db, noopDeps); + const peeked = await issueLock.withLockedQueue( + { + issue: { id: issueId, companyId, assigneeAgentId: agentId, executionRunId: null }, + actor: { actorType: "user", actorId: "user-1", agentId: null, runId: null, agentApiKeyId: null }, + queueId: wakeId, + }, + async (locked) => locked.queue, + ); + + vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] }); + try { + const resultPromise = issueLock.steerQueuedWakeComment({ + issue: { id: issueId, companyId, assigneeAgentId: agentId, executionRunId: null }, + actor: { actorType: "user", actorId: "user-1", agentId: null, runId: null, agentApiKeyId: null }, + commentId: firstCommentId, + queueId: wakeId, + targetRunId, + revision: peeked.revision, + }); + // Advance past the probe's own bound. A resolved promise here proves + // the bound fired; an unbounded wait would leave this promise pending. + await vi.advanceTimersByTimeAsync(5_000); + const result = await resultPromise; + expect(getNativeSessionSteeringStateMock).toHaveBeenCalledWith(targetRunId); + expect(result.queue.steeringDisposition).toBe("temporarily_unavailable"); + } finally { + vi.useRealTimers(); + } + }); }); diff --git a/server/src/modules/wake-queue/adapters/queued-comment-postgres.ts b/server/src/modules/wake-queue/adapters/queued-comment-postgres.ts index fb9c38c19c..17ef9b33e5 100644 --- a/server/src/modules/wake-queue/adapters/queued-comment-postgres.ts +++ b/server/src/modules/wake-queue/adapters/queued-comment-postgres.ts @@ -40,6 +40,30 @@ import type { type WakeRow = typeof agentWakeupRequests.$inferSelect; type RunRow = typeof heartbeatRuns.$inferSelect; +// `buildQueueSnapshot` awaits the live steering probe below while it runs +// inside a lock-holding transaction (`for("update")` on the issue, wake, and +// run rows). The native runtime call has no timeout of its own, so an +// unbounded wait would hold those row locks for as long as the provider +// takes to answer. This bound keeps the wait, and so the lock hold, short. +const LIVE_STEERING_PROBE_TIMEOUT_MS = 2_000; + +/** Rejects after `timeoutMs` if `promise` has not settled yet, so a caller can bound how long it waits. */ +function rejectAfterTimeout(promise: Promise, timeoutMs: number): Promise { + return new Promise((resolve, reject) => { + const timer = setTimeout(() => reject(new Error("live steering probe timed out")), timeoutMs); + promise.then( + (value) => { + clearTimeout(timer); + resolve(value); + }, + (error) => { + clearTimeout(timer); + reject(error); + }, + ); + }); +} + function toWakeRow(row: WakeRow): QueuedCommentWakeRow { return { id: row.id, agentId: row.agentId, status: row.status, runId: row.runId, payload: parseObject(row.payload) }; } @@ -174,7 +198,7 @@ function buildTransaction(tx: Db, companyId: string, deps: QueuedCommentQueuePos steering.kind !== "probe" ? steering.kind : probeLiveSteering - ? await getNativeSessionSteeringState(steering.steeringRunId) + ? await rejectAfterTimeout(getNativeSessionSteeringState(steering.steeringRunId), LIVE_STEERING_PROBE_TIMEOUT_MS) .then((liveState) => liveState.disposition) .catch(() => "temporarily_unavailable" as const) : "temporarily_unavailable";