refactor(server): move the steering queue mutation into the wake-queue module
The steering route kept the identity lifecycle, the row locks, the acknowledgement replay checks, the provider steering call, the queue removal, the wake cancellation, and the run acknowledgement write inline. Move that logic into createQueuedCommentQueue as steerQueuedWakeComment, the fourth mutation on the same queue adapter. The route now keeps only authorization, validation, issue access, error mapping, the activity log, and the response shape. The three identity calls that must run before or after the module's own transaction (reserveSteeredIdentity, reconcileSteeredIdentity, rejectSteeredIdentity) stay on the root database handle, so a lost response after the final queued message still resolves to a durable idempotency record instead of a second steer. The module reserves the identity itself; no caller-supplied pending identity id crosses the port. Delete lockQueuedCommentState and assertQueueMutationTarget from the route now that steering is their only caller. buildQueuedCommentQueue stays, since the read-only queued-comments route still needs it. Co-authored-by: Paperclip <noreply@paperclip.ing>
This commit is contained in:
parent
d1ba17eeca
commit
c66868abb9
|
|
@ -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<SteerQueuedWakeCommentResult> {
|
||||
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;
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -157,6 +157,24 @@ export interface QueuedCommentQueueTransaction {
|
|||
logActivity(input: QueuedCommentActivityLogInput): Promise<QueuedCommentActivityPublication>;
|
||||
}
|
||||
|
||||
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<T>,
|
||||
): Promise<T>;
|
||||
|
||||
/**
|
||||
* 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<SteerQueuedWakeCommentResult>;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -128,6 +128,7 @@ function createFakeTransaction(overrides: Partial<QueuedCommentQueueTransaction>
|
|||
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 })),
|
||||
};
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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<SteerQueuedWakeCommentResult> {
|
||||
return deps.issueLock.steerQueuedWakeComment(input);
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 }),
|
||||
};
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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<Parameters<Db["transaction"]>[0]>[0];
|
||||
type IssueQueueTx = Parameters<Parameters<Db["transaction"]>[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<typeof getActorInfo>;
|
||||
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<typeof getActorInfo> }) {
|
||||
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),
|
||||
);
|
||||
},
|
||||
);
|
||||
|
|
|
|||
Loading…
Reference in New Issue