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 <noreply@paperclip.ing>
This commit is contained in:
nickyleach 2026-09-10 22:55:37 +00:00
parent 92adb87d5b
commit 021dc0ad54
2 changed files with 79 additions and 1 deletions

View File

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

View File

@ -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<T>(promise: Promise<T>, timeoutMs: number): Promise<T> {
return new Promise<T>((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";