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 375a99eb63..fa53a21b7f 100644 --- a/server/src/modules/wake-queue/adapters/queued-comment-postgres.ts +++ b/server/src/modules/wake-queue/adapters/queued-comment-postgres.ts @@ -1,4 +1,4 @@ -import { and, eq, inArray } from "drizzle-orm"; +import { and, asc, eq, inArray } from "drizzle-orm"; import type { Db } from "@paperclipai/db"; import { agentWakeupRequests, agents, heartbeatRuns, issueComments, issues } from "@paperclipai/db"; import type { IssueComment, IssueQueuedCommentQueue } from "@paperclipai/shared"; @@ -10,9 +10,20 @@ import { withQueuedCommentIdsInWakePayload, } from "../../../services/issue-queued-comment-queue.js"; import { logActivity as persistActivityLogRow, type ActivityPublication } from "../../../services/activity-log.js"; +import { + NativeSessionSteeringError, + steerNativeSession, +} from "../../../services/native-runtime/native-session-executor.js"; +import { + acceptSteeredIdentity, + reconcileSteeredIdentity, + rejectSteeredIdentity, + reserveSteeredIdentity, + storedSteeringAcknowledgement, +} from "../../../services/run-identity.js"; import { decideQueuedCommentWakeLookup } from "../domain/policy.js"; import { parseObject, readNonEmptyString } from "../domain/values.js"; -import { QueuedCommentMutationError } from "../application/queued-comment-use-cases.js"; +import { QueuedCommentMutationError, requireMutationTarget } from "../application/queued-comment-use-cases.js"; import type { LockedQueuedCommentState, QueuedCommentActivityLogInput, @@ -21,6 +32,8 @@ import type { QueuedCommentQueueTransaction, QueuedCommentRunRow, QueuedCommentWakeRow, + SteerQueuedWakeCommentInput, + SteerQueuedWakeCommentResult, } from "../application/queued-comment-ports.js"; type WakeRow = typeof agentWakeupRequests.$inferSelect; @@ -202,6 +215,58 @@ function buildTransaction(tx: Db, companyId: string, deps: QueuedCommentQueuePos }; } +/** + * Finds the issue's current live queue. It scans the assigned agent's own + * pending wakes, the same lookup the read-only queued-comments route runs. + * The steering replay branch needs this fresh read for one reason: by the + * time a retry arrives, the wake it named can already be cancelled. The + * response must then show whatever queue is live now, not the old one. + */ +async function findCurrentQueuedCommentWake( + tx: Db, + companyId: string, + issue: { id: string; assigneeAgentId: string | null }, +): Promise<{ wake: WakeRow; state: "deferred" | "queued"; queueRun: RunRow | null } | null> { + if (!issue.assigneeAgentId) return null; + const rows = await tx + .select() + .from(agentWakeupRequests) + .where( + and( + eq(agentWakeupRequests.companyId, companyId), + eq(agentWakeupRequests.agentId, issue.assigneeAgentId), + inArray(agentWakeupRequests.status, ["deferred_issue_execution", "queued"]), + ), + ) + .orderBy(asc(agentWakeupRequests.requestedAt)); + + for (const wake of rows) { + if (parseObject(wake.payload).issueId !== issue.id || queuedCommentIdsFromWakePayload(wake.payload).length === 0) { + continue; + } + if (wake.status === "deferred_issue_execution") { + return { wake, state: "deferred", queueRun: null }; + } + if (!wake.runId) continue; + const queueRun = await tx + .select() + .from(heartbeatRuns) + .where( + and( + eq(heartbeatRuns.id, wake.runId), + eq(heartbeatRuns.companyId, companyId), + eq(heartbeatRuns.agentId, issue.assigneeAgentId), + eq(heartbeatRuns.wakeupRequestId, wake.id), + eq(heartbeatRuns.status, "queued"), + ), + ) + .limit(1) + .then((queueRunRows) => queueRunRows[0] ?? null); + if (queueRun) return { wake, state: "queued", queueRun }; + } + return null; +} + export function createQueuedCommentIssueLockWriter(db: Db, deps: QueuedCommentQueuePostgresAdapterDeps): QueuedCommentIssueLockWriter { return { async withLockedQueue(input, fn) { @@ -307,5 +372,264 @@ export function createQueuedCommentIssueLockWriter(db: Db, deps: QueuedCommentQu return fn(locked, transaction); }); }, + + async steerQueuedWakeComment(input: SteerQueuedWakeCommentInput): Promise { + const { issue, actor, commentId, queueId, targetRunId, revision } = input; + const companyId = issue.companyId; + + // Reserve the run's pending steering identity on the root handle. + // This call opens its own transaction and runs before this method's + // own transaction opens. So the reservation survives a rollback of + // the steer that follows it. + const steeringIdentity = await reserveSteeredIdentity(db, { + companyId, + runId: targetRunId, + issueId: issue.id, + messageId: commentId, + }); + + let steeringDeliveryAttempted = false; + let turnId: string | null = null; + let duplicate = false; + + try { + const queue = await db.transaction(async (rawTx) => { + const tx = rawTx as unknown as Db; + const transaction = buildTransaction(tx, companyId, deps); + const now = new Date(); + + // A client can lose the successful response after the final + // queued message cancels its wake. Lock the issue row first. This + // keeps the persisted acknowledgement a durable idempotency + // record, even when no pending queue remains by the time a retry + // arrives. + await tx + .select({ id: issues.id }) + .from(issues) + .where(and(eq(issues.id, issue.id), eq(issues.companyId, companyId))) + .for("update"); + + const wakeRow = await tx + .select() + .from(agentWakeupRequests) + .where( + and( + eq(agentWakeupRequests.id, queueId), + eq(agentWakeupRequests.companyId, companyId), + issue.assigneeAgentId ? eq(agentWakeupRequests.agentId, issue.assigneeAgentId) : undefined, + ), + ) + .for("update") + .limit(1) + .then((rows) => rows[0] ?? null); + + const retryRunRow = + wakeRow && parseObject(wakeRow.payload).issueId === issue.id + ? await tx + .select() + .from(heartbeatRuns) + .where( + and( + eq(heartbeatRuns.id, targetRunId), + eq(heartbeatRuns.companyId, companyId), + eq(heartbeatRuns.agentId, wakeRow.agentId), + ), + ) + .for("update") + .limit(1) + .then((rows) => rows[0] ?? null) + : null; + + const retryRunContext = parseObject(retryRunRow?.contextSnapshot); + const retryRunResult = parseObject(retryRunRow?.resultJson); + const retryAcknowledgements = parseObject(retryRunResult.queuedSteeringAcknowledgements); + const retryAcknowledgement = parseObject(retryAcknowledgements[commentId]); + + if ( + retryRunRow && + (retryRunContext.issueId === issue.id || retryRunContext.taskId === issue.id) && + retryAcknowledgement.status === "acknowledged" && + retryAcknowledgement.queueId === queueId + ) { + duplicate = true; + turnId = typeof retryAcknowledgement.turnId === "string" ? retryAcknowledgement.turnId : null; + const current = await findCurrentQueuedCommentWake(tx, companyId, issue); + return transaction.buildQueueSnapshot({ + issue, + actor, + wake: current ? toWakeRow(current.wake) : null, + state: current?.state ?? null, + queueRun: current?.queueRun ? toRunRow(current.queueRun) : null, + activeRun: retryRunRow.status === "running" ? toRunRow(retryRunRow) : null, + }); + } + + if ( + !wakeRow || + parseObject(wakeRow.payload).issueId !== issue.id || + queuedCommentIdsFromWakePayload(wakeRow.payload).length === 0 + ) { + throw new QueuedCommentMutationError("queued_comment_not_pending", "The queued message is no longer pending"); + } + + let state: "deferred" | "queued"; + let queueRunRow: RunRow | null = null; + if (wakeRow.status === "deferred_issue_execution") { + state = "deferred"; + } else if (wakeRow.status === "queued" && wakeRow.runId) { + queueRunRow = await tx + .select() + .from(heartbeatRuns) + .where( + and( + eq(heartbeatRuns.id, wakeRow.runId), + eq(heartbeatRuns.companyId, companyId), + eq(heartbeatRuns.agentId, wakeRow.agentId), + eq(heartbeatRuns.wakeupRequestId, wakeRow.id), + ), + ) + .for("update") + .limit(1) + .then((rows) => rows[0] ?? null); + if (!queueRunRow || queueRunRow.status !== "queued") { + throw new QueuedCommentMutationError( + "queued_comment_already_dispatching", + "The queued message is already being dispatched", + ); + } + state = "queued"; + } else if ( + wakeRow.status === "claimed" || + wakeRow.status === "running" || + (wakeRow.runId && (wakeRow.status === "succeeded" || wakeRow.status === "failed")) + ) { + throw new QueuedCommentMutationError( + "queued_comment_already_dispatching", + "The queued message is already being dispatched", + ); + } else { + throw new QueuedCommentMutationError("queued_comment_not_pending", "The queued message is no longer pending"); + } + + const activeRunId = state === "deferred" ? targetRunId : null; + const activeRunRow = activeRunId + ? await tx + .select() + .from(heartbeatRuns) + .where( + and( + eq(heartbeatRuns.id, activeRunId), + eq(heartbeatRuns.companyId, companyId), + eq(heartbeatRuns.status, "running"), + ), + ) + .for("update") + .limit(1) + .then((rows) => rows[0] ?? null) + : null; + const activeRunContext = parseObject(activeRunRow?.contextSnapshot); + if (!activeRunRow || (activeRunContext.issueId !== issue.id && activeRunContext.taskId !== issue.id)) { + throw new QueuedCommentMutationError("queued_comment_stale_target", "The queued message targets a stale run"); + } + + const lockedQueue = await transaction.buildQueueSnapshot({ + issue, + actor, + wake: toWakeRow(wakeRow), + state, + queueRun: queueRunRow ? toRunRow(queueRunRow) : null, + activeRun: toRunRow(activeRunRow), + }); + + const runResult = parseObject(activeRunRow.resultJson); + const acknowledgements = parseObject(runResult.queuedSteeringAcknowledgements); + const priorAcknowledgement = parseObject(acknowledgements[commentId]); + if (priorAcknowledgement.status === "acknowledged" && priorAcknowledgement.queueId === queueId) { + duplicate = true; + turnId = typeof priorAcknowledgement.turnId === "string" ? priorAcknowledgement.turnId : null; + return lockedQueue; + } + + requireMutationTarget(lockedQueue, queueId, revision); + if (lockedQueue.protocol !== "paperclip_runner_v1") { + throw new QueuedCommentMutationError("steering_unsupported", "This runner does not support same-turn steering"); + } + const entry = lockedQueue.entries.find((candidate) => candidate.comment.id === commentId); + if (!entry) { + throw new QueuedCommentMutationError("queued_comment_not_pending", "The queued message is no longer pending"); + } + + steeringDeliveryAttempted = true; + const acknowledgement = + (steeringIdentity ? await storedSteeringAcknowledgement(tx, steeringIdentity) : null) ?? + (await steerNativeSession({ + runId: activeRunRow.id, + message: entry.comment.body, + correlationId: commentId, + onAcknowledged: steeringIdentity ? () => reconcileSteeredIdentity(db, steeringIdentity) : undefined, + })); + if (steeringIdentity) await acceptSteeredIdentity(tx, steeringIdentity); + turnId = acknowledgement.turnId; + + const remainingIds = lockedQueue.entries.map((candidate) => candidate.comment.id).filter((candidateId) => candidateId !== commentId); + let nextWakeRow: WakeRow | null; + if (remainingIds.length === 0) { + await tx + .update(agentWakeupRequests) + .set({ status: "cancelled", finishedAt: now, updatedAt: now }) + .where(and(eq(agentWakeupRequests.id, wakeRow.id), eq(agentWakeupRequests.companyId, companyId))); + nextWakeRow = null; + } else { + nextWakeRow = await tx + .update(agentWakeupRequests) + .set({ payload: withQueuedCommentIdsInWakePayload(wakeRow.payload, remainingIds), updatedAt: now }) + .where(and(eq(agentWakeupRequests.id, wakeRow.id), eq(agentWakeupRequests.companyId, companyId))) + .returning() + .then((rows) => rows[0] ?? wakeRow); + } + + await tx + .update(heartbeatRuns) + .set({ + resultJson: { + ...runResult, + queuedSteeringAcknowledgements: { + ...acknowledgements, + [commentId]: { + status: "acknowledged", + queueId, + turnId: acknowledgement.turnId, + acknowledgedAt: now.toISOString(), + }, + }, + }, + updatedAt: now, + }) + .where(and(eq(heartbeatRuns.id, activeRunRow.id), eq(heartbeatRuns.companyId, companyId))); + + return transaction.buildQueueSnapshot({ + issue, + actor, + wake: nextWakeRow ? toWakeRow(nextWakeRow) : null, + state: nextWakeRow ? "deferred" : null, + queueRun: null, + activeRun: toRunRow(activeRunRow), + }); + }); + return { queue, turnId, duplicate }; + } catch (error) { + // A late native acknowledgement can arrive after this transaction + // rolls back. Only a delivery attempt with a still-unknown outcome + // keeps the identity reservation pending for later reconciliation. + // A definite provider answer never keeps it pending. + const uncertain = + steeringDeliveryAttempted && + (!(error instanceof NativeSessionSteeringError) || error.code === "steering_timeout"); + if (steeringIdentity && !uncertain) { + await rejectSteeredIdentity(db, steeringIdentity); + } + throw error; + } + }, }; } diff --git a/server/src/modules/wake-queue/application/queued-comment-ports.ts b/server/src/modules/wake-queue/application/queued-comment-ports.ts index a30d0c6bc0..dbf3e34f65 100644 --- a/server/src/modules/wake-queue/application/queued-comment-ports.ts +++ b/server/src/modules/wake-queue/application/queued-comment-ports.ts @@ -157,6 +157,24 @@ export interface QueuedCommentQueueTransaction { logActivity(input: QueuedCommentActivityLogInput): Promise; } +export type SteerQueuedWakeCommentInput = { + /** Also carries the company id. Every locked read and write binds its `companyId` predicate to `issue.companyId`. */ + issue: QueuedCommentIssueContext; + actor: QueuedCommentActor; + commentId: string; + queueId: string; + targetRunId: string; + revision: string; +}; + +export type SteerQueuedWakeCommentResult = { + queue: QueuedCommentQueueSnapshot; + /** Null when the provider acknowledgement carried no turn id. */ + turnId: string | null; + /** True when the caller retried a steer whose acknowledgement was already recorded. */ + duplicate: boolean; +}; + export interface QueuedCommentIssueLockWriter { /** * Opens the one transaction a mutation runs in: locks the issue row, locks @@ -176,4 +194,18 @@ export interface QueuedCommentIssueLockWriter { }, fn: (locked: LockedQueuedCommentState, transaction: QueuedCommentQueueTransaction) => Promise, ): Promise; + + /** + * Runs the fourth queue mutation: same-turn steering. First it reserves + * the run's pending steering identity on the root database handle. This + * reservation survives a rollback of the transaction that follows it. + * Then it opens the one transaction the steer itself runs in. It replays + * a durably recorded acknowledgement when the caller retries a lost + * response. Otherwise it delivers the message and writes the queue, the + * wake, and the run rows. On a failed delivery with an uncertain outcome, + * it leaves the identity reservation pending for later reconciliation. On + * every other failure, it rejects the reservation and rethrows the + * original error unchanged. + */ + steerQueuedWakeComment(input: SteerQueuedWakeCommentInput): Promise; } diff --git a/server/src/modules/wake-queue/application/queued-comment-use-cases.test.ts b/server/src/modules/wake-queue/application/queued-comment-use-cases.test.ts index d7baf86438..d3fe4422b2 100644 --- a/server/src/modules/wake-queue/application/queued-comment-use-cases.test.ts +++ b/server/src/modules/wake-queue/application/queued-comment-use-cases.test.ts @@ -128,6 +128,7 @@ function createFakeTransaction(overrides: Partial function createFakeIssueLock(locked: LockedQueuedCommentState, transaction: QueuedCommentQueueTransaction): QueuedCommentIssueLockWriter { return { withLockedQueue: vi.fn(async (_input, fn) => fn(locked, transaction)), + steerQueuedWakeComment: vi.fn(async () => ({ queue: queueSnapshot(), turnId: null, duplicate: false })), }; } diff --git a/server/src/modules/wake-queue/application/queued-comment-use-cases.ts b/server/src/modules/wake-queue/application/queued-comment-use-cases.ts index 8fc1624559..743befa0bf 100644 --- a/server/src/modules/wake-queue/application/queued-comment-use-cases.ts +++ b/server/src/modules/wake-queue/application/queued-comment-use-cases.ts @@ -11,6 +11,8 @@ import type { QueuedCommentQueueSnapshot, QueuedCommentQueueTransaction, QueuedCommentRunRow, + SteerQueuedWakeCommentInput, + SteerQueuedWakeCommentResult, } from "./queued-comment-ports.js"; export type QueuedCommentMutationErrorCode = @@ -18,7 +20,9 @@ export type QueuedCommentMutationErrorCode = | "queued_comment_already_dispatching" | "queued_comment_stale_queue" | "queued_comment_revision_conflict" - | "queued_comment_order_mismatch"; + | "queued_comment_order_mismatch" + | "queued_comment_stale_target" + | "steering_unsupported"; /** The route maps this 1:1 onto the `conflict(...)` HTTP error it threw before this move, using `code` and `message` unchanged. */ export class QueuedCommentMutationError extends Error { @@ -39,7 +43,8 @@ export class QueuedCommentMutationForbiddenError extends Error { } } -function requireMutationTarget(queue: QueuedCommentQueueSnapshot, queueId: string, revision: string): void { +/** The steering adapter reuses this check below. It applies the same rule. */ +export function requireMutationTarget(queue: QueuedCommentQueueSnapshot, queueId: string, revision: string): void { if (queue.queueId !== queueId) { throw new QueuedCommentMutationError("queued_comment_stale_queue", "The queued message targets a stale queue"); } @@ -355,3 +360,17 @@ export function createDiscardQueuedComment(deps: { issueLock: QueuedCommentIssue ); }; } + +/** + * Wires the fourth queue mutation, same-turn steering, onto the port. The + * adapter owns the whole flow: the identity reservation on the root handle, + * the one transaction, and the rejection on failure. None of it decomposes + * into a pure decision that this layer could hold instead. This factory + * only keeps the composition symmetric with the other three mutations + * above. + */ +export function createSteerQueuedWakeComment(deps: { issueLock: QueuedCommentIssueLockWriter }) { + return function steerQueuedWakeComment(input: SteerQueuedWakeCommentInput): Promise { + return deps.issueLock.steerQueuedWakeComment(input); + }; +} diff --git a/server/src/modules/wake-queue/index.ts b/server/src/modules/wake-queue/index.ts index b725f6a805..65b33f1985 100644 --- a/server/src/modules/wake-queue/index.ts +++ b/server/src/modules/wake-queue/index.ts @@ -12,6 +12,7 @@ import { createDiscardQueuedComment, createEditQueuedComment, createReorderQueuedComments, + createSteerQueuedWakeComment, } from "./application/queued-comment-use-cases.js"; import type { IssueSnapshot, @@ -54,6 +55,8 @@ export type { QueuedCommentActor, QueuedCommentIssueContext, QueuedCommentQueueSnapshot, + SteerQueuedWakeCommentInput, + SteerQueuedWakeCommentResult, } from "./application/queued-comment-ports.js"; export type { QueuedCommentQueuePostgresAdapterDeps } from "./adapters/queued-comment-postgres.js"; @@ -123,6 +126,7 @@ export function createQueuedCommentQueue(db: Db, deps: QueuedCommentQueuePostgre editQueuedComment: createEditQueuedComment({ issueLock }), reorderQueuedComments: createReorderQueuedComments({ issueLock }), discardQueuedComment: createDiscardQueuedComment({ issueLock }), + steerQueuedWakeComment: createSteerQueuedWakeComment({ issueLock }), }; } diff --git a/server/src/routes/issues.ts b/server/src/routes/issues.ts index 0cc963a4e0..9a5ccd46fe 100644 --- a/server/src/routes/issues.ts +++ b/server/src/routes/issues.ts @@ -5,13 +5,6 @@ import { validateExecutionReconciliation, markExecutionReconciliation, } from "../services/execution-recovery-resolution.js"; -import { - storedSteeringAcknowledgement, - reconcileSteeredIdentity, - reserveSteeredIdentity, - acceptSteeredIdentity, - rejectSteeredIdentity, -} from "../services/run-identity.js"; import { createHash, randomUUID } from "node:crypto"; import { Router, type Request, type Response } from "express"; import multer from "multer"; @@ -334,7 +327,6 @@ import { import { getNativeSessionSteeringState, NativeSessionSteeringError, - steerNativeSession, } from "../services/native-runtime/native-session-executor.js"; import { buildQueuedCommentQueueSnapshot, @@ -6751,7 +6743,6 @@ export function issueRoutes( } type IssueQueueDb = Db | Parameters[0]>[0]; - type IssueQueueTx = Parameters[0]>[0]; type IssueQueueWake = typeof agentWakeupRequests.$inferSelect; type IssueQueueRun = typeof heartbeatRuns.$inferSelect; type IssueQueueState = { @@ -6887,155 +6878,6 @@ export function issueRoutes( }); } - function assertQueueMutationTarget(input: { - queue: IssueQueuedCommentQueue; - queueId: string; - revision: string; - }) { - if (input.queue.queueId !== input.queueId) { - throw conflict("The queued message targets a stale queue", { - code: "queued_comment_stale_queue", - }); - } - if (input.queue.revision !== input.revision) { - throw conflict("The queued messages changed in another session", { - code: "queued_comment_revision_conflict", - }); - } - } - - async function lockQueuedCommentState(input: { - tx: IssueQueueTx; - issue: { - id: string; - companyId: string; - assigneeAgentId: string | null; - executionRunId?: string | null; - }; - actor: ReturnType; - queueId: string; - targetRunId?: string; - }) { - await input.tx - .select({ id: issueRows.id }) - .from(issueRows) - .where( - and( - eq(issueRows.id, input.issue.id), - eq(issueRows.companyId, input.issue.companyId), - ), - ) - .for("update"); - const wake = await input.tx - .select() - .from(agentWakeupRequests) - .where( - and( - eq(agentWakeupRequests.id, input.queueId), - eq(agentWakeupRequests.companyId, input.issue.companyId), - input.issue.assigneeAgentId - ? eq(agentWakeupRequests.agentId, input.issue.assigneeAgentId) - : undefined, - ), - ) - .for("update") - .limit(1) - .then((rows) => rows[0] ?? null); - if ( - !wake || - readObject(wake.payload).issueId !== input.issue.id || - queuedCommentIdsFromWakePayload(wake.payload).length === 0 - ) { - throw conflict("The queued message is no longer pending", { - code: "queued_comment_not_pending", - }); - } - - let state: IssueQueueState["state"]; - let queueRun: IssueQueueRun | null = null; - if (wake.status === "deferred_issue_execution") { - state = "deferred"; - } else if (wake.status === "queued" && wake.runId) { - queueRun = await input.tx - .select() - .from(heartbeatRuns) - .where( - and( - eq(heartbeatRuns.id, wake.runId), - eq(heartbeatRuns.companyId, input.issue.companyId), - eq(heartbeatRuns.agentId, wake.agentId), - eq(heartbeatRuns.wakeupRequestId, wake.id), - ), - ) - .for("update") - .limit(1) - .then((rows) => rows[0] ?? null); - if (!queueRun || queueRun.status !== "queued") { - throw conflict("The queued message is already being dispatched", { - code: "queued_comment_already_dispatching", - }); - } - state = "queued"; - } else if ( - wake.status === "claimed" || - wake.status === "running" || - (wake.runId && (wake.status === "succeeded" || wake.status === "failed")) - ) { - throw conflict("The queued message is already being dispatched", { - code: "queued_comment_already_dispatching", - }); - } else { - throw conflict("The queued message is no longer pending", { - code: "queued_comment_not_pending", - }); - } - - const activeRunId = - state === "deferred" - ? (input.targetRunId ?? input.issue.executionRunId ?? null) - : null; - const activeRun = activeRunId - ? await input.tx - .select() - .from(heartbeatRuns) - .where( - and( - eq(heartbeatRuns.id, activeRunId), - eq(heartbeatRuns.companyId, input.issue.companyId), - eq(heartbeatRuns.status, "running"), - ), - ) - .for("update") - .limit(1) - .then((rows) => rows[0] ?? null) - : null; - if (input.targetRunId) { - const runContext = readObject(activeRun?.contextSnapshot); - if ( - !activeRun || - (runContext.issueId !== input.issue.id && - runContext.taskId !== input.issue.id) - ) { - throw conflict("The queued message targets a stale run", { - code: "queued_comment_stale_target", - }); - } - } - const queueState = { wake, state, queueRun } satisfies IssueQueueState; - const queue = await buildQueuedCommentQueue({ - executor: input.tx, - issue: input.issue, - activeRun, - actor: input.actor, - queueState, - steeringDisposition: - activeRun?.runtimeMode === "native" - ? "temporarily_unavailable" - : "unsupported", - }); - return { activeRun, wake, queueRun, state, queue, queueState }; - } - function operatorInterruptCancelOptions(input: { issueId: string; actor: ReturnType }) { return { errorCode: "operator_interrupted", @@ -15203,224 +15045,21 @@ export function issueRoutes( ); if (!issue) return; const actor = getActorInfo(req); - const steeringIdentity = await reserveSteeredIdentity(db, { - companyId: issue.companyId, - runId: req.body.targetRunId, - issueId: issue.id, - messageId: commentId, - }); - let steeringDeliveryAttempted = false; - let acknowledgedTurnId: string | null = null; - let duplicate = false; - let queue: IssueQueuedCommentQueue; + let result; try { - queue = await db.transaction(async (tx) => { - // A client can lose the successful response after the final queued - // message cancels its wake. Lock the original queue and target run - // first so that the persisted acknowledgement remains a durable - // idempotency record even when no pending queue remains. - await tx - .select({ id: issueRows.id }) - .from(issueRows) - .where( - and( - eq(issueRows.id, issue.id), - eq(issueRows.companyId, issue.companyId), - ), - ) - .for("update"); - const retryWake = await tx - .select() - .from(agentWakeupRequests) - .where( - and( - eq(agentWakeupRequests.id, req.body.queueId), - eq(agentWakeupRequests.companyId, issue.companyId), - issue.assigneeAgentId - ? eq(agentWakeupRequests.agentId, issue.assigneeAgentId) - : undefined, - ), - ) - .for("update") - .limit(1) - .then((rows) => rows[0] ?? null); - const retryRun = - retryWake && readObject(retryWake.payload).issueId === issue.id - ? await tx - .select() - .from(heartbeatRuns) - .where( - and( - eq(heartbeatRuns.id, req.body.targetRunId), - eq(heartbeatRuns.companyId, issue.companyId), - eq(heartbeatRuns.agentId, retryWake.agentId), - ), - ) - .for("update") - .limit(1) - .then((rows) => rows[0] ?? null) - : null; - const retryRunContext = readObject(retryRun?.contextSnapshot); - const retryRunResult = readObject(retryRun?.resultJson); - const retryAcknowledgements = readObject( - retryRunResult.queuedSteeringAcknowledgements, - ); - const retryAcknowledgement = readObject( - retryAcknowledgements[commentId], - ); - if ( - retryRun && - (retryRunContext.issueId === issue.id || - retryRunContext.taskId === issue.id) && - retryAcknowledgement.status === "acknowledged" && - retryAcknowledgement.queueId === req.body.queueId - ) { - duplicate = true; - acknowledgedTurnId = - typeof retryAcknowledgement.turnId === "string" - ? retryAcknowledgement.turnId - : null; - return buildQueuedCommentQueue({ - executor: tx, - issue, - activeRun: retryRun.status === "running" ? retryRun : null, - actor, - }); - } - - const locked = await lockQueuedCommentState({ - tx, - issue, - actor, - queueId: req.body.queueId, - targetRunId: req.body.targetRunId, - }); - if (!locked.activeRun) { - throw conflict("The queued message targets a stale run", { - code: "queued_comment_stale_target", - }); - } - const runResult = readObject(locked.activeRun.resultJson); - const acknowledgements = readObject( - runResult.queuedSteeringAcknowledgements, - ); - const priorAcknowledgement = readObject(acknowledgements[commentId]); - if ( - priorAcknowledgement.status === "acknowledged" && - priorAcknowledgement.queueId === req.body.queueId - ) { - duplicate = true; - acknowledgedTurnId = - typeof priorAcknowledgement.turnId === "string" - ? priorAcknowledgement.turnId - : null; - return buildQueuedCommentQueue({ - executor: tx, - issue, - activeRun: locked.activeRun, - actor, - queueState: locked.queueState, - }); - } - assertQueueMutationTarget({ - queue: locked.queue, - queueId: req.body.queueId, - revision: req.body.revision, - }); - if (locked.queue.protocol !== "paperclip_runner_v1") { - throw conflict("This runner does not support same-turn steering", { - code: "steering_unsupported", - }); - } - const entry = locked.queue.entries.find( - (candidate) => candidate.comment.id === commentId, - ); - if (!entry) { - throw conflict("The queued message is no longer pending", { - code: "queued_comment_not_pending", - }); - } - - steeringDeliveryAttempted = true; - const acknowledgement = - (steeringIdentity - ? await storedSteeringAcknowledgement(tx, steeringIdentity) - : null) ?? - (await steerNativeSession({ - runId: locked.activeRun.id, - message: entry.comment.body, - correlationId: commentId, - onAcknowledged: steeringIdentity - ? () => reconcileSteeredIdentity(db, steeringIdentity) - : undefined, - })); - if (steeringIdentity) - await acceptSteeredIdentity(tx, steeringIdentity); - acknowledgedTurnId = acknowledgement.turnId; - const remainingIds = locked.queue.entries - .map((candidate) => candidate.comment.id) - .filter((candidateId) => candidateId !== commentId); - const now = new Date(); - const nextWake = - remainingIds.length === 0 - ? await tx - .update(agentWakeupRequests) - .set({ status: "cancelled", finishedAt: now, updatedAt: now }) - .where(eq(agentWakeupRequests.id, locked.wake.id)) - .returning() - .then(() => null) - : await tx - .update(agentWakeupRequests) - .set({ - payload: withQueuedCommentIdsInWakePayload( - locked.wake.payload, - remainingIds, - ), - updatedAt: now, - }) - .where(eq(agentWakeupRequests.id, locked.wake.id)) - .returning() - .then((rows) => rows[0] ?? locked.wake); - await tx - .update(heartbeatRuns) - .set({ - resultJson: { - ...runResult, - queuedSteeringAcknowledgements: { - ...acknowledgements, - [commentId]: { - status: "acknowledged", - queueId: req.body.queueId, - turnId: acknowledgement.turnId, - acknowledgedAt: now.toISOString(), - }, - }, - }, - updatedAt: now, - }) - .where(eq(heartbeatRuns.id, locked.activeRun.id)); - return buildQueuedCommentQueue({ - executor: tx, - issue, - activeRun: locked.activeRun, - actor, - queueState: nextWake - ? { wake: nextWake, state: "deferred", queueRun: null } - : null, - }); + result = await queuedCommentQueue.steerQueuedWakeComment({ + issue: buildQueuedCommentIssueContext(issue), + actor, + commentId, + queueId: req.body.queueId, + targetRunId: req.body.targetRunId, + revision: req.body.revision, }); } catch (error) { - const uncertain = - steeringDeliveryAttempted && - (!(error instanceof NativeSessionSteeringError) || - error.code === "steering_timeout"); - if (steeringIdentity && !uncertain) - await rejectSteeredIdentity(db, steeringIdentity); - if (error instanceof NativeSessionSteeringError) { throw conflict(error.message, { code: error.code, retryable: true }); } - throw error; + throwForQueuedCommentMutationError(error); } await logActivity(db, { companyId: issue.companyId, @@ -15435,12 +15074,12 @@ export function issueRoutes( details: { commentId, targetRunId: req.body.targetRunId, - turnId: acknowledgedTurnId, - duplicate, + turnId: result.turnId, + duplicate: result.duplicate, }, }); res.json( - await runRedactions.redactForIssue(issue.companyId, issue.id, queue), + await runRedactions.redactForIssue(issue.companyId, issue.id, result.queue), ); }, );