diff --git a/doc/execution-semantics.md b/doc/execution-semantics.md index d1130d3181..326da73f0b 100644 --- a/doc/execution-semantics.md +++ b/doc/execution-semantics.md @@ -1156,6 +1156,13 @@ controller, lease, or result). It also checks for contradictory launch/process evidence and verifies local cleanup or exact remote termination receipts. The preparer must have finished or its startup lease must have expired. A missing PID alone does not establish this proof. +For older interrupted preparation rows without a cancellation receipt, the +immutable Paperclip Runner adapter claim, unresolved runtime, and preparing stage +must agree. The old controller must belong to another server boot and its lease +must have expired. No native identity, coordinator, result, adapter invocation, +provider event, or process-launch evidence may exist. Environment cleanup still +requires the same receipts. This historical proof permits explicit Retry or a +newer saved user message; it does not replay the cancelled input. The existing bounded saved-message worker rechecks this proof after restart. Admission atomically settles an unclaimed coordinator and admits one fresh turn, diff --git a/server/src/services/cancelled-native-startup.ts b/server/src/services/cancelled-native-startup.ts index db3f082a88..ef7fb126d5 100644 --- a/server/src/services/cancelled-native-startup.ts +++ b/server/src/services/cancelled-native-startup.ts @@ -4,6 +4,7 @@ import { claimedAdapterType } from "./conversation-continuation.js"; import { PROCESS_IDENTITY_RECORDED, PROCESS_START_REQUESTED } from "./native-local-process-stop.js"; import { hasRemoteTerminationReceipt } from "./remote-execution-termination.js"; import { canContinueCancelledRun } from "./run-cancellation.js"; +import { legacyControllerBootId } from "./legacy-controller-lease.js"; type Run = typeof heartbeatRuns.$inferSelect; type Coordinator = typeof nativeRunFinalizations.$inferSelect; @@ -27,7 +28,16 @@ export async function isCancelledNativeStartup(db: Db, run: Run, coordinator: Co if (run.status !== "cancelled" || !run.finishedAt || run.processPid || run.processGroupId || run.processStartedAt || run.sessionIdAfter) return false; const cancellation = run.resultJson?.startupCancellation as Record | undefined; - const beforeSelection = run.runtimeMode === "legacy" && !run.runtimeModeResolvedAt && + // Older builds did not retain the cancellation fence. Their immutable + // native-adapter claim and unresolved preparation stage still prove that + // provider dispatch did not begin. Require an expired owner from another + // server boot; neither a missing PID nor mutable agent settings is proof. + const historicalBeforeSelection = run.runtimeMode === "legacy" && !run.runtimeModeResolvedAt && + run.executionStage === "preparing" && !run.nativeIssueId && !run.nativeSessionId && !coordinator && + claimedAdapterType(run) === "paperclip_runner" && run.errorCode === "operator_interrupted" && + run.resultJson === null && Boolean(run.controllerBootId && run.controllerBootId !== legacyControllerBootId && + run.controllerLeaseExpiresAt && run.controllerLeaseExpiresAt <= new Date()); + const beforeSelection = historicalBeforeSelection || run.runtimeMode === "legacy" && !run.runtimeModeResolvedAt && !run.nativeSessionId && !coordinator && claimedAdapterType(run) === "paperclip_runner" && cancellation?.beforeNativeSelection === true; const neverClaimed = run.runtimeMode === "native" && coordinator && @@ -42,7 +52,7 @@ export async function isCancelledNativeStartup(db: Db, run: Run, coordinator: Co const leases = await db.select().from(environmentLeases).where(and( eq(environmentLeases.companyId, run.companyId), eq(environmentLeases.heartbeatRunId, run.id), )); - if ((!settled && leases.length === 0) || leases.some(lease => + if ((!settled && leases.length === 0 && !historicalBeforeSelection) || leases.some(lease => lease.provider === "local" ? !lease.releasedAt || lease.status === "pending_cleanup" || lease.cleanupStatus === "failed" : !hasRemoteTerminationReceipt(lease))) return false; @@ -51,7 +61,7 @@ export async function isCancelledNativeStartup(db: Db, run: Run, coordinator: Co const [execution] = await db.select({ id: heartbeatRunEvents.id }).from(heartbeatRunEvents).where(and( eq(heartbeatRunEvents.companyId, run.companyId), eq(heartbeatRunEvents.runId, run.id), or(isNotNull(heartbeatRunEvents.sourceEventId), - inArray(heartbeatRunEvents.eventType, [PROCESS_START_REQUESTED, PROCESS_IDENTITY_RECORDED, + inArray(heartbeatRunEvents.eventType, ["adapter.invoke", PROCESS_START_REQUESTED, PROCESS_IDENTITY_RECORDED, "harness.ready", "session.started", "session.resumed", "session.updated", "turn.started", "provider.event", "provider.rpc_result", "tool.execution.started"])), )).limit(1); diff --git a/server/src/services/explicit-native-continuation.test.ts b/server/src/services/explicit-native-continuation.test.ts index b535771c0f..31b0cb0f4a 100644 --- a/server/src/services/explicit-native-continuation.test.ts +++ b/server/src/services/explicit-native-continuation.test.ts @@ -1,3 +1,4 @@ +import { legacyControllerBootId } from "./legacy-controller-lease.js"; import { instanceSettingsService } from "./instance-settings.js"; import { createPostgresWakeQueueAdapter } from "../modules/wake-queue/adapters/postgres.js"; import { createReleaseIssueExecution } from "../modules/wake-queue/application/use-cases.js"; @@ -597,6 +598,103 @@ const support = await getEmbeddedPostgresTestSupport(); expect(await admit(f)).toMatchObject({ previousRunId: f.sourceRunId }); }); + async function seedHistoricalCancelledPreparation() { + const f = await seedCancelledStartup(); + await db.delete(nativeRunFinalizations).where(eq(nativeRunFinalizations.runId, f.sourceRunId)); + await db.delete(environmentLeases).where(eq(environmentLeases.heartbeatRunId, f.sourceRunId)); + await db.update(heartbeatRuns).set({ runtimeMode: "legacy", runtimeModeResolvedAt: null, + nativeIssueId: null, nativeSessionId: null, executionStage: "preparing", errorCode: "operator_interrupted", + runnerProfileJson: { adapterDispatch: { adapterType: "paperclip_runner" } }, resultJson: null, + }).where(eq(heartbeatRuns.id, f.sourceRunId)); + await db.update(issueRecoveryActions).set({ cause: "legacy_execution_requires_reconciliation" }) + .where(eq(issueRecoveryActions.sourceIssueId, f.issueId)); + return f; + } + + it.each(["message", "retry"])("recovers historical native preparation through an explicit %s", async kind => { + const f = await seedHistoricalCancelledPreparation(); + expect(await getExecutionBlocker(db, f.companyId, f.issueId)).toMatchObject({ canRetry: true }); + await db.insert(heartbeatRuns).values({ companyId: f.companyId, agentId: f.agentId, status: "running" }); + const successor = await heartbeatService(db).wakeup(f.agentId, { + source: "on_demand", triggerDetail: "manual", reason: kind === "retry" ? "retry_failed_run" : "issue_commented", + ...(kind === "retry" ? { failedRunId: f.sourceRunId } : {}), requestedByActorType: "user", requestedByActorId: "board", + payload: { issueId: f.issueId, ...(kind === "message" ? { commentId: f.commentId } : {}) }, + contextSnapshot: { issueId: f.issueId, ...(kind === "message" ? { wakeCommentId: f.commentId } : {}) }, + }); + expect(successor).toMatchObject({ status: "queued", contextSnapshot: { forceFreshSession: true, + previousRunId: f.sourceRunId, explicitUserContinuation: { commentId: kind === "message" ? f.commentId : null } } }); + expect(await getExecutionBlocker(db, f.companyId, f.issueId)).toBeNull(); + const [source] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, f.sourceRunId)); + expect(source).toMatchObject({ status: "cancelled", runtimeMode: "legacy", resultJson: null }); + }); + + it.each(["fresh", "earlier_delivered", "last_delivered"])("resumes saved input after historical native preparation expires across restart exactly once (%s)", async kind => { + const f = await seedHistoricalCancelledPreparation(); + await db.insert(heartbeatRuns).values({ companyId: f.companyId, agentId: f.agentId, status: "running" }); + await db.update(heartbeatRuns).set({ controllerLeaseExpiresAt: new Date(Date.now() + 60_000) }) + .where(eq(heartbeatRuns.id, f.sourceRunId)); + await heartbeatService(db).wakeup(f.agentId, { source: "automation", reason: "issue_commented", + requestedByActorType: "user", requestedByActorId: "board", payload: { issueId: f.issueId, commentId: f.commentId }, + contextSnapshot: { issueId: f.issueId, wakeCommentId: f.commentId } }); + const [waiting] = await db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.companyId, f.companyId)); + expect(waiting.status).toBe("deferred_issue_execution"); + if (kind !== "fresh") { + const deliveredId = randomUUID(); + await db.insert(issueComments).values({ id: deliveredId, companyId: f.companyId, issueId: f.issueId, + authorType: "user", authorUserId: "another-author", body: "Already handled", + createdAt: new Date(kind === "earlier_delivered" ? "2026-09-11T10:30:00Z" : "2026-09-11T11:30:00Z") }); + await db.insert(heartbeatRuns).values({ companyId: f.companyId, agentId: f.agentId, status: "succeeded", + startedAt: new Date("2026-09-11T12:00:00Z"), finishedAt: new Date("2026-09-11T12:01:00Z"), + contextSnapshot: { issueId: f.issueId, wakeCommentIds: [deliveredId] } }); + await db.update(agentWakeupRequests).set({ payload: { ...waiting.payload, + _paperclipWakeContext: { wakeCommentIds: kind === "earlier_delivered" ? [deliveredId, f.commentId] : [f.commentId, deliveredId] } }, + }).where(eq(agentWakeupRequests.id, waiting.id)); + } + await db.update(heartbeatRuns).set({ controllerLeaseExpiresAt: new Date(0) }) + .where(eq(heartbeatRuns.id, f.sourceRunId)); + await db.update(agentWakeupRequests).set({ updatedAt: new Date(0) }).where(eq(agentWakeupRequests.id, waiting.id)); + await Promise.all([heartbeatService(db).resumeExecutionWaitComments(), heartbeatService(db).resumeExecutionWaitComments()]); + const successors = await db.select().from(heartbeatRuns).where(and(eq(heartbeatRuns.companyId, f.companyId), eq(heartbeatRuns.status, "queued"))); + expect(successors).toHaveLength(1); + expect(successors[0].contextSnapshot).toMatchObject({ wakeCommentIds: [f.commentId], + explicitUserContinuation: { commentId: f.commentId }, forceFreshSession: true }); + const [receipt] = await db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, waiting.id)); + expect(receipt).toMatchObject({ status: "coalesced", runId: successors[0].id }); + expect(await getExecutionBlocker(db, f.companyId, f.issueId)).toBeNull(); + }); + + it.each(["resolved", "dispatching", "adapter", "missing_boot", "current_boot", "missing_lease", "live_lease", "receipt", "launch", "invoked", "coordinator", "cleanup", "remote"])( + "holds historical native preparation with contradictory or incomplete %s evidence", async kind => { + const f = await seedHistoricalCancelledPreparation(); + const patch = kind === "resolved" ? { runtimeModeResolvedAt: new Date() } + : kind === "dispatching" ? { executionStage: "dispatching" } + : kind === "adapter" ? { runnerProfileJson: { adapterDispatch: { adapterType: "process" } } } + : kind === "missing_boot" ? { controllerBootId: null } + : kind === "current_boot" ? { controllerBootId: legacyControllerBootId } + : kind === "missing_lease" ? { controllerLeaseExpiresAt: null } + : kind === "live_lease" ? { controllerLeaseExpiresAt: new Date(Date.now() + 60_000) } + : kind === "receipt" ? { resultJson: { startupCancellation: { beforeNativeSelection: false } } } : {}; + if (Object.keys(patch).length) await db.update(heartbeatRuns).set(patch).where(eq(heartbeatRuns.id, f.sourceRunId)); + if (kind === "launch" || kind === "invoked") await db.insert(heartbeatRunEvents).values({ companyId: f.companyId, + agentId: f.agentId, runId: f.sourceRunId, seq: 1, + eventType: kind === "launch" ? PROCESS_START_REQUESTED : "adapter.invoke", payload: { adapterType: "paperclip_runner" } }); + if (kind === "coordinator") { + await db.update(heartbeatRuns).set({ nativeIssueId: f.issueId }).where(eq(heartbeatRuns.id, f.sourceRunId)); + await db.insert(nativeRunFinalizations).values({ companyId: f.companyId, + runId: f.sourceRunId, issueId: f.issueId, phase: "observed", attempt: 0 }); + } + if (kind === "cleanup" || kind === "remote") await db.insert(environmentLeases).values({ companyId: f.companyId, + heartbeatRunId: f.sourceRunId, provider: kind === "remote" ? "daytona" : "local", + providerLeaseId: kind === "remote" ? "unverified" : null, leasePolicy: "ephemeral", + status: "pending_cleanup", cleanupStatus: "failed" }); + expect(await getExecutionBlocker(db, f.companyId, f.issueId)).toMatchObject({ canRetry: false }); + expect(await admit(f)).toBeNull(); + // Later tests exercise the global cleanup sweep with their own retry + // counters. Do not leave this negative fixture as another cleanup target. + await db.delete(environmentLeases).where(eq(environmentLeases.heartbeatRunId, f.sourceRunId)); + }, + ); + it.each(["attempt", "generation", "controller", "lease", "process", "launch", "provider", "cleanup", "remote", "preparing", "closed", "reassigned"])( "retains cancellation safeguards with %s evidence", async kind => { const f = await seedCancelledStartup(); diff --git a/server/src/services/explicit-native-continuation.ts b/server/src/services/explicit-native-continuation.ts index 4dd475b43a..bf1f66530f 100644 --- a/server/src/services/explicit-native-continuation.ts +++ b/server/src/services/explicit-native-continuation.ts @@ -169,14 +169,9 @@ export async function admitExplicitNativeContinuation(input: { const unusedAdmission = run.status === "cancelled" && !run.startedAt && run.errorCode === "execution_reconciliation_required" && !run.processPid && !run.processGroupId && !run.nativeSessionId; - if (partiallyDeliveredQueue && run.runtimeMode !== "native" && !unusedAdmission) return null; const legacyUserTurn = run.runtimeMode === "legacy" && action.cause === "legacy_execution_requires_reconciliation" && isConversationAdapter(agent.adapterType); - if ((queuedInterrupt || queuedRequest) && !legacyUserTurn && !unusedAdmission && - !(queuedRequest && run.runtimeMode === "native" && - (run.status !== "cancelled" || authorizedAt > run.finishedAt)) && - !(queuedInterrupt && response?.source.requiresFreshSession && run.runtimeMode === "native")) return null; // Saved input is a request for a new turn, never permission to undo an // operator Stop or redeliver a message already consumed by this run. if (queuedRequest && !queuedInterrupt && (run.contextSnapshot?.wakeCommentId === commentId || @@ -205,6 +200,12 @@ export async function admitExplicitNativeContinuation(input: { return blocked("workspace_repair_required", "Verify safe workspace staging or repair before continuing. Your message is saved."); } const cancelledStartup = await isCancelledNativeStartup(db, run, coordinator); + if (partiallyDeliveredQueue && run.runtimeMode !== "native" && !unusedAdmission && !cancelledStartup) return null; + if ((queuedInterrupt || queuedRequest) && !legacyUserTurn && !unusedAdmission && + !(queuedRequest && run.runtimeMode === "native" && + (run.status !== "cancelled" || authorizedAt > run.finishedAt!)) && + !(queuedRequest && cancelledStartup && authorizedAt > run.finishedAt!) && + !(queuedInterrupt && response?.source.requiresFreshSession && run.runtimeMode === "native")) return null; if (queuedRequest && !queuedInterrupt && run.status === "cancelled" && !unusedAdmission && !canContinueCancelledRun(run) && !(cancelledStartup && authorizedAt > run.finishedAt!)) return null; if (retry && run.status === "cancelled" && !canContinueCancelledRun(run) && !cancelledStartup) diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index 02162d05c8..88f73ea4c6 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -33,7 +33,7 @@ import { completeTerminatedRemoteNativeSessionCleanup } from "../vendor/papercli import { hasRemoteTerminationReceipt, remoteExecutionHasStopped, remoteTerminationReceipt, stoppedRemoteCleanupScopes } from "./remote-execution-termination.js"; import { applyConnectorSkills, prepareConnectorSkillDelivery, resolveConnectorAssignments } from "./connector-runtime.js"; import { admitExplicitNativeContinuation, undeliveredLegacyUserCommentIds } from "./explicit-native-continuation.js"; -import { canRetryStoppedRun } from "./cancelled-native-startup.js"; +import { canRetryStoppedRun, isCancelledNativeStartup } from "./cancelled-native-startup.js"; import { connectionIntentService } from "./connection-intents.js"; import { managedAiSessionFingerprintConfig, prepareManagedAiRuntime, assertManagedAiProjectAuth, stripAiAuthBindings, isAiConnectionBusy, AI_AUTH_ENV_KEYS } from "./ai-connection-runtime.js"; import { aiConnectionBindingSchema } from "@paperclipai/shared"; @@ -10373,7 +10373,13 @@ export function heartbeatService( async function resumeRemoteStopComments(run: typeof heartbeatRuns.$inferSelect, requestId?: string) { if (!isHeartbeatRunTerminalStatus(run.status) || adapterExecutionControls.has(run.id)) return; - if (run.runtimeMode !== "native" && + const [preparationCoordinator] = run.runtimeMode === "legacy" && run.status === "cancelled" && !run.runtimeModeResolvedAt + ? await db.select().from(nativeRunFinalizations).where(and( + eq(nativeRunFinalizations.companyId, run.companyId), eq(nativeRunFinalizations.runId, run.id), + )) : []; + const cancelledPreparation = run.runtimeMode === "legacy" && + await isCancelledNativeStartup(db, run, preparationCoordinator); + if (run.runtimeMode !== "native" && !cancelledPreparation && parseObject(run.resultJson?.startupCancellation).beforeNativeSelection !== true && !(await remoteExecutionHasStopped(db, run.companyId, run.id))) return; const issueId = run.nativeIssueId ?? (typeof run.contextSnapshot?.issueId === "string" ? run.contextSnapshot.issueId : null); @@ -10414,7 +10420,7 @@ export function heartbeatService( let commentId = deriveCommentId(context, payload); let requestedByActorId = wake.requestedByActorId; const reason = readNonEmptyString(context.wakeReason) ?? wake.reason; - if (run.runtimeMode === "native") { + if (run.runtimeMode === "native" || cancelledPreparation) { const ids = await undeliveredLegacyUserCommentIds(db, run.companyId, issueId, run.agentId, queuedCommentIdsFromWakePayload(payload)); if (!ids.length) continue; @@ -10444,7 +10450,7 @@ export function heartbeatService( const admitted = await admitExplicitNativeContinuation({ db, companyId: run.companyId, issueId, agentId: run.agentId, actorType: wake.requestedByActorType, actorId: requestedByActorId, reason, commentId, successorRunId: randomUUID(), dryRun: true, - ...(run.runtimeMode === "native" ? { queuedCommentRequestId: wake.id } : {}), + ...((run.runtimeMode === "native" || cancelledPreparation) ? { queuedCommentRequestId: wake.id } : {}), onBlocked: (reason, message) => { wait = { reason, message }; }, }); if (!admitted) { @@ -10462,7 +10468,7 @@ export function heartbeatService( await enqueueWakeup(run.agentId, { source: wake.source as WakeupOptions["source"], triggerDetail: (wake.triggerDetail ?? undefined) as WakeupOptions["triggerDetail"], reason, payload, contextSnapshot: context, requestedByActorType: "user", requestedByActorId, - ...(run.runtimeMode === "native" ? { queuedCommentRequestId: wake.id } : {}), + ...((run.runtimeMode === "native" || cancelledPreparation) ? { queuedCommentRequestId: wake.id } : {}), idempotencyKey: `remote-stop-comment:${run.id}:${wake.id}` }, wake.id); break; }