diff --git a/server/src/__tests__/native-cancellation-request.integration.test.ts b/server/src/__tests__/native-cancellation-request.integration.test.ts index 9ada636e8f..16ae43c022 100644 --- a/server/src/__tests__/native-cancellation-request.integration.test.ts +++ b/server/src/__tests__/native-cancellation-request.integration.test.ts @@ -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>; @@ -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(resolve => { locked = resolve; }); + const hold = new Promise(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" }]); + }); + }); diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index c4281a4367..dfbe950d60 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -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, + cancellationCondition?: ReturnType, ) { // 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) { diff --git a/server/src/services/native-runtime/native-cancellation-request.ts b/server/src/services/native-runtime/native-cancellation-request.ts index 54d11f190c..4d588ad9b1 100644 --- a/server/src/services/native-runtime/native-cancellation-request.ts +++ b/server/src/services/native-runtime/native-cancellation-request.ts @@ -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. */ diff --git a/server/src/services/native-runtime/native-session-executor.test.ts b/server/src/services/native-runtime/native-session-executor.test.ts index 66233a6737..2c5f961c18 100644 --- a/server/src/services/native-runtime/native-session-executor.test.ts +++ b/server/src/services/native-runtime/native-session-executor.test.ts @@ -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; @@ -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) => ({ 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 } } }); diff --git a/server/src/services/native-runtime/native-session-executor.ts b/server/src/services/native-runtime/native-session-executor.ts index bd7d2b8416..dd94378756 100644 --- a/server/src/services/native-runtime/native-session-executor.ts +++ b/server/src/services/native-runtime/native-session-executor.ts @@ -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;