fix(server): fence native retry cancellation through final commit

Recheck retry eligibility under the coordinator lock before creating a cancellation intent. Preserve later run and coordinator outcomes with a failed-only acknowledged-result CAS, while retaining same-intent recovery and NOWAIT conflict handling.

Co-Authored-By: Paperclip <noreply@paperclip.ing>
(cherry picked from commit 29230e2000)
This commit is contained in:
Dotta committed 2026-09-30 13:54:30 -05:00
1 parent 45b977bca0
commit e4bf459e31
5 files changed
+186 -34

No files matched your search

@@ -1,9 +1,9 @@
import { randomUUID } from "node:crypto";
import { afterAll, beforeAll, describe, expect, it } from "vitest";
import { eq, sql } from "drizzle-orm";
import { and, eq, inArray, sql } from "drizzle-orm";
import { agents, companies, createDb, heartbeatRuns, issues, nativeRunFinalizations } from "@paperclipai/db";
import { startEmbeddedPostgresTestDatabase } from "./helpers/embedded-postgres.js";
import { cancellationIntentId, claimCancellationRequest, startupCancellationFence } from "../services/native-runtime/native-cancellation-request.js";
import { rethrowNativeCancellationLockConflict, nativeRetryCancellationCommitCondition, cancellationIntentId, claimCancellationRequest, startupCancellationFence } from "../services/native-runtime/native-cancellation-request.js";
describe("atomic caller cancellation request ownership", () => {
let database: Awaited<ReturnType<typeof startEmbeddedPostgresTestDatabase>>;
@@ -101,4 +101,63 @@ describe("atomic caller cancellation request ownership", () => {
expect(run.resultJson?.startupCancellation).toBeUndefined();
});
it.each(["failed", "running"])("preserves a %s terminal/recovery outcome that wins after reservation", async status => {
const f = await retryFixture(), id = randomUUID();
await claimCancellationRequest(db, f.runId, f.companyId, id, "board-user");
await db.update(nativeRunFinalizations).set({ phase: status === "failed" ? "terminal_failure" : "observed",
failureCode: "board_recovery", nextAttemptAt: null }).where(eq(nativeRunFinalizations.runId, f.runId));
await db.update(heartbeatRuns).set({ status, errorCode: "board_recovery" }).where(eq(heartbeatRuns.id, f.runId));
await expect(cancelNativeSession(f.runId, "Stop", { db, scope: "run", cancellationRequestId: id })).rejects.toMatchObject({ status: 409 });
const [run] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, f.runId));
expect(run).toMatchObject({ status, errorCode: "board_recovery" });
expect(run.resultJson?.nativeCancellation).toBeUndefined();
});
it("does not overwrite a coordinator-only terminal change after acknowledgement", async () => {
const f = await retryFixture(), id = randomUUID();
await claimCancellationRequest(db, f.runId, f.companyId, id, "board-user");
await cancelNativeSession(f.runId, "Stop", { db, scope: "run", cancellationRequestId: id });
const [acknowledged] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, f.runId));
const commit = () => db.update(heartbeatRuns).set({ status: "cancelled" }).where(and(eq(heartbeatRuns.id, f.runId),
inArray(heartbeatRuns.status, ["queued", "running", "scheduled_retry", "failed"]), nativeRetryCancellationCommitCondition(acknowledged.resultJson))).returning();
let release!: () => void, locked!: () => void;
const acquired = new Promise<void>(resolve => { locked = resolve; });
const hold = new Promise<void>(resolve => { release = resolve; });
const advance = db.transaction(async tx => {
await tx.update(nativeRunFinalizations).set({ failureCode: "board_recovery" }).where(eq(nativeRunFinalizations.runId, f.runId));
locked(); await hold;
});
try { await acquired; await expect(commit().catch(rethrowNativeCancellationLockConflict)).rejects.toMatchObject({ status: 409, message: "Native retry cancellation is busy; retry the same request" }); }
finally { release(); await advance; }
expect(await commit()).toEqual([]);
const [preserved] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, f.runId));
expect(preserved).toMatchObject({ status: "failed", resultJson: acknowledged.resultJson });
});
it.each(["running", "queued", "scheduled_retry"])("preserves a status-only advance to %s after acknowledgement", async status => {
const f = await retryFixture(), id = randomUUID();
await claimCancellationRequest(db, f.runId, f.companyId, id, "board-user");
await cancelNativeSession(f.runId, "Stop", { db, scope: "run", cancellationRequestId: id });
const [acknowledged] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, f.runId));
await db.update(heartbeatRuns).set({ status }).where(eq(heartbeatRuns.id, f.runId));
// Match heartbeat's broader candidate list: the production retry predicate,
// not a test-only failed-status filter, must preserve the newer state.
const committed = await db.update(heartbeatRuns).set({ status: "cancelled" }).where(and(eq(heartbeatRuns.id, f.runId),
inArray(heartbeatRuns.status, ["queued", "running", "scheduled_retry", "failed"]),
nativeRetryCancellationCommitCondition(acknowledged.resultJson))).returning();
expect(committed).toEqual([]);
const [preserved] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, f.runId));
expect(preserved).toMatchObject({ status, resultJson: acknowledged.resultJson });
});
it("commits an unchanged acknowledged retry but preserves a newer run result", async () => {
const f = await retryFixture(), id = randomUUID();
await claimCancellationRequest(db, f.runId, f.companyId, id, "board-user");
await cancelNativeSession(f.runId, "Stop", { db, scope: "run", cancellationRequestId: id });
const [acknowledged] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, f.runId));
const commit = () => db.update(heartbeatRuns).set({ status: "cancelled" }).where(and(eq(heartbeatRuns.id, f.runId),
inArray(heartbeatRuns.status, ["queued", "running", "scheduled_retry", "failed"]), nativeRetryCancellationCommitCondition(acknowledged.resultJson))).returning();
await db.update(heartbeatRuns).set({ resultJson: { ...acknowledged.resultJson, boardRecovery: "new-outcome" } }).where(eq(heartbeatRuns.id, f.runId));
expect(await commit()).toEqual([]);
await db.update(heartbeatRuns).set({ resultJson: acknowledged.resultJson }).where(eq(heartbeatRuns.id, f.runId));
expect(await commit()).toMatchObject([{ status: "cancelled" }]);
});
});
+12 -4
View File
@@ -1,4 +1,4 @@
import { claimCancellationRequest, startupCancellationFence } from "./native-runtime/native-cancellation-request.js";
import { nativeRetryCancellationCommitCondition, rethrowNativeCancellationLockConflict, claimCancellationRequest, startupCancellationFence } from "./native-runtime/native-cancellation-request.js";
import { CHAT_COMPLETION_WAKE_REASON, prepareChatCompletionTurn, chatCompletionInstruction, isCompletedOnboardingHandoffWake } from "./chat-completion-delivery.js";
import { isAgentDirectoryCopy } from "./agent-directory-working-copies.js";
@@ -12883,6 +12883,7 @@ export function heartbeatService(
status: string,
fromStatuses: string[],
patch?: Partial<typeof heartbeatRuns.$inferInsert>,
cancellationCondition?: ReturnType<typeof nativeRetryCancellationCommitCondition>,
) {
// fromStatuses can name a terminal status as its own source (for example,
// an idempotent "still failed" patch), so the write below is not always a
@@ -12937,13 +12938,18 @@ export function heartbeatService(
and(
eq(heartbeatRuns.id, runId),
inArray(heartbeatRuns.status, fromStatuses),
...(cancellationCondition ? [cancellationCondition] : []),
...(isHeartbeatRunTerminalStatus(status)
? [nativeRunnerOwnershipNotHeldCondition()]
: []),
),
)
.returning()
.then((rows) => rows[0] ?? null);
.then((rows) => rows[0] ?? null)
.catch((error: unknown) => {
if (cancellationCondition) rethrowNativeCancellationLockConflict(error);
throw error;
});
if (updated) {
publishLiveEvent({
@@ -29173,8 +29179,9 @@ export function heartbeatService(
run = await claimCancellationRequest(db, runId, run.companyId, options.cancellationRequestId, options.cancellationRequestedByUserId ?? null);
}
// The caller claim checked retry eligibility under both durable row locks.
// Its own native intent may already have disabled the retry; re-reading only
// retryable_failure would strand that same-ID recovery.
// This is only a cancellation candidate: dispatch rechecks the coordinator
// under lock, and the final CAS requires its own unchanged acknowledged
// retry fence. Re-reading only retryable_failure here would strand replay.
const pendingNativeRetry =
run.runtimeMode === "native" && run.status === "failed" && (
Boolean(options.cancellationRequestId) || await db
@@ -29401,6 +29408,7 @@ export function heartbeatService(
}
: {}),
},
pendingNativeRetry ? nativeRetryCancellationCommitCondition(persistedCancellationResult) : undefined,
);
} catch (error) {
if (processCancellationSettlement) {
@@ -22,6 +22,56 @@ export function assertCancellationRequest(result: unknown, requestId: string, al
}
}
/** Evaluate only while both the run and coordinator are locked. A UUID alone
* is not retry authority; a resumed Stop needs its exact durable native intent. */
export function nativeRetryCancellationEligible(input: {
runId: string; companyId: string; issueId: string; resultJson: unknown; requestId?: string; scope?: string;
}, coordinator: { phase?: string; failureCode?: string | null } | null | undefined) {
if (coordinator?.phase === "retryable_failure") return true;
if (coordinator?.phase !== "terminal_failure" || coordinator.failureCode !== "native_retry_cancelled") return false;
const result = record(input.resultJson), native = record(result.nativeCancellation);
return (!input.requestId || (record(result.startupCancellation).cancellationRequestId === input.requestId
&& native.intentId === cancellationIntentId(input.requestId)))
&& native.schema === "paperclip.native-cancellation.v1" && typeof native.intentId === "string" && !!native.intentId
&& native.runId === input.runId && native.companyId === input.companyId && native.issueId === input.issueId
&& native.scope === (input.scope ?? "run") && ["pending", "acknowledged"].includes(String(native.dispatchState))
&& typeof native.dispatched === "boolean" && typeof native.intentAuditId === "string" && !!native.intentAuditId
&& (native.dispatchState === "pending" || (typeof native.acknowledgementAuditId === "string"
&& !!native.acknowledgementAuditId && native.acknowledgementAuditId !== native.intentAuditId));
}
/** Final failed -> cancelled CAS: a later board/recovery outcome must win.
* The row-locking subquery prevents a coordinator-only writer from changing
* the retry disposition between this check and the status update. */
export function nativeRetryCancellationCommitCondition(expectedResult: unknown) {
const native = record(record(expectedResult).nativeCancellation);
if (native.dispatchState !== "acknowledged" || typeof native.runId !== "string"
|| typeof native.companyId !== "string" || typeof native.issueId !== "string"
|| !nativeRetryCancellationEligible({ runId: native.runId, companyId: native.companyId,
issueId: native.issueId, resultJson: expectedResult }, { phase: "terminal_failure", failureCode: "native_retry_cancelled" })) {
return sql`false`;
}
return sql`${heartbeatRuns.status} = 'failed' and ${heartbeatRuns.runtimeMode} = 'native'
and ${heartbeatRuns.id} = ${native.runId} and ${heartbeatRuns.companyId} = ${native.companyId}
and ${heartbeatRuns.nativeIssueId} = ${native.issueId}
and ${heartbeatRuns.resultJson} = ${JSON.stringify(expectedResult)}::jsonb
and exists (select 1 from ${nativeRunFinalizations}
where ${nativeRunFinalizations.runId} = ${heartbeatRuns.id}
and ${nativeRunFinalizations.companyId} = ${heartbeatRuns.companyId}
and ${nativeRunFinalizations.issueId} = ${heartbeatRuns.nativeIssueId}
and ${nativeRunFinalizations.phase} = 'terminal_failure'
and ${nativeRunFinalizations.failureCode} = 'native_retry_cancelled'
for update nowait)`;
}
export function rethrowNativeCancellationLockConflict(error: unknown): never {
const value = record(error);
if (value.code === "55P03" || record(value.cause).code === "55P03") {
throw conflict("Native retry cancellation is busy; retry the same request");
}
throw error;
}
/** Reserve the caller's identity before handoff or dispatch. No schema change:
* the existing startup fence and native intent carry the correlation. */
export async function claimCancellationRequest(db: Db, runId: string, companyId: string, requestId: string, userId: string | null) {
@@ -43,16 +93,8 @@ export async function claimCancellationRequest(db: Db, runId: string, companyId:
.where(and(eq(nativeRunFinalizations.runId, runId), eq(nativeRunFinalizations.companyId, companyId),
eq(nativeRunFinalizations.issueId, run.nativeIssueId)))
.for("update", { noWait: true }).limit(1).then(rows => rows[0]);
const native = record(result.nativeCancellation);
const ownIntent = record(result.startupCancellation).cancellationRequestId === requestId
&& native.schema === "paperclip.native-cancellation.v1" && native.intentId === cancellationIntentId(requestId)
&& native.runId === runId && native.companyId === companyId && native.issueId === run.nativeIssueId
&& native.scope === "run" && ["pending", "acknowledged"].includes(String(native.dispatchState))
&& typeof native.dispatched === "boolean" && typeof native.intentAuditId === "string" && !!native.intentAuditId
&& (native.dispatchState === "pending" || (typeof native.acknowledgementAuditId === "string"
&& !!native.acknowledgementAuditId && native.acknowledgementAuditId !== native.intentAuditId));
pendingRetry = coordinator?.phase === "retryable_failure"
|| (coordinator?.phase === "terminal_failure" && coordinator.failureCode === "native_retry_cancelled" && ownIntent);
pendingRetry = nativeRetryCancellationEligible({ runId, companyId, issueId: run.nativeIssueId,
resultJson: result, requestId }, coordinator);
}
if (record(result.startupCancellation).cancellationRequestId === requestId) {
const actor = record(record(result.startupCancellation).requestedBy);
@@ -70,19 +112,11 @@ export async function claimCancellationRequest(db: Db, runId: string, companyId:
}
if (!pendingRetry && !["queued", "running", "scheduled_retry"].includes(run.status)) throw conflict("Correlated Stop run is already terminal");
const [claimed] = await tx.update(heartbeatRuns).set({ resultJson: { ...result,
startupCancellation: { requestedAt: new Date().toISOString(), beforeNativeSelection: false, cancellationRequestId: requestId, requestedBy: { type: "board", userId } },
startupCancellation: { requestedAt: new Date().toISOString(), beforeNativeSelection: false, cancellationRequestId: requestId, requestedBy: { type: "board", userId }, ...(run.status === "failed" ? { retryCancellation: true } : {}) },
} }).where(and(eq(heartbeatRuns.id, runId), eq(heartbeatRuns.companyId, companyId))).returning();
if (!claimed) throw conflict("Correlated Stop run changed");
return claimed;
}).catch((error: unknown) => {
// postgres-js can be wrapped by Drizzle; retain no raw database detail in
// the public conflict while allowing the exact request to retry later.
const value = record(error);
if (value.code === "55P03" || record(value.cause).code === "55P03") {
throw conflict("Native retry cancellation is busy; retry the same request");
}
throw error;
});
}).catch(rethrowNativeCancellationLockConflict);
}
/** Used by the default Stop path too: never erase a correlated reservation. */
@@ -4925,7 +4925,10 @@ function cancellationDb(options?: {
runId: string;
assessmentId: string | null;
decisionId?: string | null;
phase?: string;
failureCode?: string | null;
} | null;
status?: string;
failResultJsonUpdateAt?: number;
ownershipHeld?: boolean;
resultJson?: Record<string, unknown>;
@@ -4936,6 +4939,7 @@ function cancellationDb(options?: {
companyId: execution.binding.companyId,
nativeIssueId: execution.binding.issueId,
runtimeMode: "native",
status: options?.status ?? "running",
...(options?.ownershipHeld
? {
status: "running",
@@ -4995,6 +4999,7 @@ function cancellationDb(options?: {
set: (values: Record<string, unknown>) => ({
where: () => {
updates.push({ table, values });
if (table === nativeRunFinalizations && coordinator) Object.assign(coordinator, values);
const updatesResultJson = "resultJson" in values;
if (updatesResultJson) resultJsonUpdateCount += 1;
const shouldFail =
@@ -5516,6 +5521,43 @@ describe("native session cancellation", () => {
expect(state.publishActivity).toHaveBeenCalledTimes(2);
});
it.each(["terminal_failure", "observed", "committed"])("rejects retry phase %s that changed after caller reservation, before any native effect", async phase => {
const id = "11111111-1111-4111-8111-111111111111";
const persistence = cancellationDb({ status: "failed",
coordinator: { runId: execution.binding.runId, assessmentId: null, phase, failureCode: "board_recovery" },
resultJson: { startupCancellation: { cancellationRequestId: id, retryCancellation: true } } });
await expect(cancelNativeSession(execution.binding.runId, "operator Stop", { db: persistence.db,
scope: "run", cancellationRequestId: id })).rejects.toMatchObject({ status: 409, message: "Native retry is no longer cancellable" });
expect(persistence.updates).toEqual([]);
expect(state.cancel).not.toHaveBeenCalled();
expect(state.persistActivity).not.toHaveBeenCalled();
});
it("does not mistake a recovered running attempt for the failed retry that reserved Stop", async () => {
const id = "11111111-1111-4111-8111-111111111111";
const persistence = cancellationDb({ status: "running",
coordinator: { runId: execution.binding.runId, assessmentId: null, phase: "observed" },
resultJson: { startupCancellation: { cancellationRequestId: id, retryCancellation: true } } });
await expect(cancelNativeSession(execution.binding.runId, "operator Stop", { db: persistence.db,
scope: "run", cancellationRequestId: id })).rejects.toMatchObject({ status: 409 });
expect(persistence.updates).toEqual([]);
expect(state.cancel).not.toHaveBeenCalled();
});
it("fences a failed retry at dispatch and preserves the same audited intent on replay", async () => {
const id = "11111111-1111-4111-8111-111111111111";
const persistence = cancellationDb({ status: "failed",
coordinator: { runId: execution.binding.runId, assessmentId: null, phase: "retryable_failure" },
resultJson: { startupCancellation: { cancellationRequestId: id, retryCancellation: true } } });
const first = await cancelNativeSession(execution.binding.runId, "operator Stop", { db: persistence.db, scope: "run", cancellationRequestId: id });
expect(persistence.updates).toContainEqual(expect.objectContaining({ table: nativeRunFinalizations,
values: expect.objectContaining({ phase: "terminal_failure", failureCode: "native_retry_cancelled", nextAttemptAt: null }) }));
const writes = persistence.getResultJsonUpdateCount();
await expect(cancelNativeSession(execution.binding.runId, "retry Stop", { db: persistence.db,
scope: "run", cancellationRequestId: id })).resolves.toMatchObject({ auditId: first.auditId });
expect(persistence.getResultJsonUpdateCount()).toBe(writes);
});
it("uses the reserved caller UUID and permits same-intent and default retries", async () => {
const id = "11111111-1111-4111-8111-111111111111";
const persistence = cancellationDb({ resultJson: { startupCancellation: { cancellationRequestId: id } } });
@@ -1,4 +1,4 @@
import { assertCancellationRequest, cancellationIntentId as callerCancellationIntentId, cancellationRequestId } from "./native-cancellation-request.js";
import { nativeRetryCancellationEligible, rethrowNativeCancellationLockConflict, assertCancellationRequest, cancellationIntentId as callerCancellationIntentId, cancellationRequestId } from "./native-cancellation-request.js";
import { readNativeCursorPlanWait } from "./native-cursor-plan-wait.js";
import { resolveAcpxQualification } from "./acpx-qualification.js";
import { readLocalAiCredentialFile } from "../local-ai-credential-file.js";
@@ -157,7 +157,7 @@ import {
type NativeAuthoritativeIssueStatus,
type NativeStatusDecision,
} from "./status-arbiter.js";
import { HttpError } from "../../errors.js";
import { conflict, HttpError } from "../../errors.js";
import { redactSensitiveText } from "../../redaction.js";
import { resolvePaperclipRunnerBinary } from "./native-codex-runner.js";
import {
@@ -6757,8 +6757,10 @@ export async function cancelNativeSession(
}
if (isNativeRunnerOwnershipHeld(lockedRun))
throw new NativeRunnerOwnershipUnverifiedError();
const coordinator = await tx
.select({ runId: nativeRunFinalizations.runId })
const retryCancellation = lockedRun.status === "failed"
|| record(record(lockedRun.resultJson).startupCancellation).retryCancellation === true;
const coordinatorQuery = tx
.select({ runId: nativeRunFinalizations.runId, phase: nativeRunFinalizations.phase, failureCode: nativeRunFinalizations.failureCode })
.from(nativeRunFinalizations)
.where(
and(
@@ -6766,9 +6768,12 @@ export async function cancelNativeSession(
eq(nativeRunFinalizations.companyId, cancellationContext.companyId),
eq(nativeRunFinalizations.issueId, cancellationContext.issueId),
),
)
.limit(1)
.then((rows) => rows[0] ?? null);
);
// Recheck after reservation, at the same transaction that records intent
// and disables the retry. NOWAIT avoids coordinator -> run lock inversion.
const coordinator = await (retryCancellation
? coordinatorQuery.for("update", { noWait: true }) : coordinatorQuery)
.limit(1).then((rows) => rows[0] ?? null);
if (!coordinator)
throw new Error("native_cancellation_coordinator_missing");
@@ -6778,6 +6783,10 @@ export async function cancelNativeSession(
const requestedId = options.cancellationRequestId
?? cancellationRequestId(record(resultJson.startupCancellation).cancellationRequestId);
if (requestedId) assertCancellationRequest(resultJson, requestedId);
if (retryCancellation && !nativeRetryCancellationEligible({
runId, companyId: cancellationContext.companyId, issueId: cancellationContext.issueId,
resultJson, requestId: requestedId, scope: options.scope ?? "run",
}, coordinator)) throw conflict("Native retry is no longer cancellable");
const existing = record(resultJson.nativeCancellation);
const existingIntentId =
typeof existing.intentId === "string" && existing.intentId.length > 0
@@ -6907,7 +6916,7 @@ export async function cancelNativeSession(
priorCoordinatorDecisionId: cancellationContext.coordinatorDecisionId,
existing: false,
};
});
}).catch(rethrowNativeCancellationLockConflict);
if (intentPublication) publishActivity(intentPublication);
cancellationIntentId = intent.intentId;
auditId = intent.auditId;