diff --git a/doc/execution-semantics.md b/doc/execution-semantics.md index 46b9e87e7f..6d8ec4af30 100644 --- a/doc/execution-semantics.md +++ b/doc/execution-semantics.md @@ -668,6 +668,8 @@ The same bounded rule applies when the previous heartbeat reported waiting on a A continuation that the staleness gate cancelled with `issue_continuation_waiting_on_review` is a *deliberate park*, not a disappeared execution path. The latest run reported that the issue is waiting for review/approval (for example, an umbrella issue whose work was just decomposed into sub-tasks). Treating that park as a stranded run would retry it, then escalate it to `blocked` with a recovery action and an operator-facing failure notice — even though nothing failed and there is nothing for a human to do. +Execution admission reads a narrow server-owned cancellation-evidence projection. Ordinary run presentation can redact `resultJson` for database encoding or output size; that presentation projection must not decide Retry eligibility or saved-input recovery. The admission projection excludes provider diagnostics and preserves whether the stored result is absent. + Recovery rule for a parked-for-review continuation: - if the issue has a real waiting target — open (non-terminal) sub-tasks or existing unresolved blockers — Paperclip converts the deliberate wait into a first-class dependency wait: it sets the issue `blocked` by those issues, keeps the original assignee, and posts a plain-language comment explaining that the task will resume automatically when its dependencies finish. The issue then self-resumes through the normal `issue_blockers_resolved` path; no recovery action or escalation owner is involved diff --git a/server/src/__tests__/agent-live-run-routes.test.ts b/server/src/__tests__/agent-live-run-routes.test.ts index 28919226bd..55d978e079 100644 --- a/server/src/__tests__/agent-live-run-routes.test.ts +++ b/server/src/__tests__/agent-live-run-routes.test.ts @@ -1163,7 +1163,7 @@ describe("agent live run routes", () => { expect(res.status, JSON.stringify(res.body)).toBe(202); expect(res.body).toEqual(receipt); expect(mockHeartbeatService.getRun).toHaveBeenCalledWith( - failedChatRunId, + failedChatRunId, { includeExecutionEvidence: true }, ); expect( mockChatRunRetries.prepareFailedChatRunRetry, diff --git a/server/src/routes/agents.ts b/server/src/routes/agents.ts index 1fc61ee9e1..094a02a656 100644 --- a/server/src/routes/agents.ts +++ b/server/src/routes/agents.ts @@ -5998,7 +5998,7 @@ export function agentRoutes( "An exact failed-run retry cannot override its execution context.", ); } - const failedRun = await heartbeat.getRun(req.body.failedRunId); + const failedRun = await heartbeat.getRun(req.body.failedRunId, { includeExecutionEvidence: true }); if ( !failedRun || failedRun.companyId !== agent.companyId || diff --git a/server/src/services/explicit-native-continuation.test.ts b/server/src/services/explicit-native-continuation.test.ts index cd0d3ac143..3ecd4596b4 100644 --- a/server/src/services/explicit-native-continuation.test.ts +++ b/server/src/services/explicit-native-continuation.test.ts @@ -623,13 +623,37 @@ const support = await getEmbeddedPostgresTestSupport(); return f; } - it.each(["message", "retry"])("recovers a pre-dispatch review wait through an explicit %s", async kind => { + it("keeps provider diagnostics out of admission evidence and preserves absent results", async () => { const f = await seedCancelledReviewWait(); + const heartbeat = heartbeatService(db); + await db.update(heartbeatRuns).set({ resultJson: { + stopReason: "issue_continuation_waiting_on_review", timeoutSource: "stale_queued_run_gate", + summary: "private provider output", providerPayload: "large provider output".repeat(10000), + } }).where(eq(heartbeatRuns.id, f.sourceRunId)); + const projected = await heartbeat.getRun(f.sourceRunId, { includeExecutionEvidence: true }); + expect(projected?.resultJson).toMatchObject({ stopReason: "issue_continuation_waiting_on_review", timeoutSource: "stale_queued_run_gate" }); + expect(projected?.resultJson).not.toHaveProperty("summary"); + expect(projected?.resultJson).not.toHaveProperty("providerPayload"); + await db.update(heartbeatRuns).set({ resultJson: null }).where(eq(heartbeatRuns.id, f.sourceRunId)); + expect((await heartbeat.getRun(f.sourceRunId, { includeExecutionEvidence: true }))?.resultJson).toBeNull(); + await db.update(heartbeatRuns).set({ resultJson: { unrecognizedReceipt: true } }).where(eq(heartbeatRuns.id, f.sourceRunId)); + expect((await heartbeat.getRun(f.sourceRunId, { includeExecutionEvidence: true }))?.resultJson).not.toBeNull(); + }); + + it.each([false, true].flatMap(ascii => ["message", "retry"].map(kind => ({ ascii, kind }))))( + "recovers a pre-dispatch review wait through an explicit $kind (SQL_ASCII: $ascii)", async ({ ascii, kind }) => { + const f = await seedCancelledReviewWait(); + const heartbeat = heartbeatService(db); + if (ascii) { + const encoding = vi.spyOn(db, "execute").mockResolvedValueOnce([{ server_encoding: "SQL_ASCII" }] as never); + try { expect((await heartbeat.getRun(f.sourceRunId))?.resultJson).toBeNull(); } + finally { encoding.mockRestore(); } + } const notice = await getExecutionBlocker(db, f.companyId, f.issueId); expect(notice).toMatchObject({ canRetry: true, runError: "Waiting for review; this continuation never started." }); expect(notice?.nextAction).not.toContain("Inspect the run before sending a new message"); await db.insert(heartbeatRuns).values({ companyId: f.companyId, agentId: f.agentId, status: "running" }); - const successor = await heartbeatService(db).wakeup(f.agentId, { source: kind === "retry" ? "on_demand" : "automation", triggerDetail: "manual", + const successor = await heartbeat.wakeup(f.agentId, { source: kind === "retry" ? "on_demand" : "automation", 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 } : {}) }, @@ -641,13 +665,19 @@ const support = await getEmbeddedPostgresTestSupport(); expect((await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, f.sourceRunId)))[0].status).toBe("cancelled"); }); - it("reconsiders saved input after a pre-dispatch review wait exactly once", async () => { + it.each([false, true])("reconsiders saved input after a pre-dispatch review wait exactly once (SQL_ASCII: %s)", async ascii => { const f = await seedCancelledReviewWait(); + const heartbeat = heartbeatService(db); + if (ascii) { + const encoding = vi.spyOn(db, "execute").mockResolvedValueOnce([{ server_encoding: "SQL_ASCII" }] as never); + try { expect((await heartbeat.getRun(f.sourceRunId))?.resultJson).toBeNull(); } + finally { encoding.mockRestore(); } + } await db.insert(heartbeatRuns).values({ companyId: f.companyId, agentId: f.agentId, status: "running" }); const leaseId = randomUUID(); await db.insert(environmentLeases).values({ id: leaseId, companyId: f.companyId, heartbeatRunId: f.sourceRunId, provider: "local", status: "pending_cleanup", cleanupStatus: "failed" }); - await heartbeatService(db).wakeup(f.agentId, { source: "automation", reason: "issue_commented", + await heartbeat.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)); @@ -655,7 +685,7 @@ const support = await getEmbeddedPostgresTestSupport(); await db.update(environmentLeases).set({ status: "released", releasedAt: new Date(), cleanupStatus: "succeeded" }) .where(eq(environmentLeases.id, leaseId)); await db.update(agentWakeupRequests).set({ updatedAt: new Date(0) }).where(eq(agentWakeupRequests.id, waiting.id)); - await Promise.all([heartbeatService(db).resumeExecutionWaitComments(), heartbeatService(db).resumeExecutionWaitComments()]); + await Promise.all([heartbeat.resumeExecutionWaitComments(), heartbeat.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], diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index 9ca39eb35c..4c7f3a281c 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -3558,6 +3558,27 @@ const heartbeatRunSafeResultJsonColumn = sql | null>` end `.as("resultJson"); +// Execution admission needs retained server receipts even when presentation +// projection omits resultJson (SQL_ASCII or oversized provider output). Select +// only the evidence used by eligibility; never retrieve provider diagnostics. +const heartbeatRunExecutionEvidenceColumn = sql | null>` + case when ${heartbeatRuns.resultJson} is null then null else jsonb_build_object( + 'startupCancellation', ${heartbeatRuns.resultJson} -> 'startupCancellation', + 'startupPreparationSettledAt', ${heartbeatRuns.resultJson} -> 'startupPreparationSettledAt', + 'stopReason', ${heartbeatRuns.resultJson} -> 'stopReason', + 'timeoutSource', ${heartbeatRuns.resultJson} -> 'timeoutSource', + 'workspaceRestoreFailure', ${heartbeatRuns.resultJson} -> 'workspaceRestoreFailure', + 'executionCancellation', ${heartbeatRuns.resultJson} -> 'executionCancellation', + 'nativeCancellation', ${heartbeatRuns.resultJson} -> 'nativeCancellation', + 'cancelledByActorType', ${heartbeatRuns.resultJson} -> 'cancelledByActorType', + 'cancelledByUserId', ${heartbeatRuns.resultJson} -> 'cancelledByUserId', + 'conversationContinuation', ${heartbeatRuns.resultJson} -> 'conversationContinuation', + 'cancellation', ${heartbeatRuns.resultJson} -> 'cancellation', + 'acpToolInventoryComplete', ${heartbeatRuns.resultJson} -> 'acpToolInventoryComplete', + 'acpPendingToolCount', ${heartbeatRuns.resultJson} -> 'acpPendingToolCount' + ) end +`.as("resultJson"); + const heartbeatRunSafeColumns = { ...getTableColumns(heartbeatRuns), processGroupId: heartbeatRunProcessGroupIdColumn, @@ -10384,7 +10405,7 @@ export function heartbeatService( !(await remoteExecutionHasStopped(db, run.companyId, run.id))) return; const issueId = run.nativeIssueId ?? (typeof run.contextSnapshot?.issueId === "string" ? run.contextSnapshot.issueId : null); if (!issueId) return; - const currentRun = run.runtimeMode === "native" ? await getRun(run.id) : null; + const currentRun = run.runtimeMode === "native" ? await getRun(run.id, { includeExecutionEvidence: true }) : null; const [coordinator] = currentRun ? await db.select({ phase: nativeRunFinalizations.phase, leaseOwner: nativeRunFinalizations.leaseOwner }).from(nativeRunFinalizations).where(and( eq(nativeRunFinalizations.companyId, run.companyId), eq(nativeRunFinalizations.runId, run.id), @@ -10399,7 +10420,7 @@ export function heartbeatService( await acknowledgedNativeStopExecutionHasStopped(db, currentRun) && !(await getExecutionBlocker(db, run.companyId, issueId)); const legacyContinuation = run.runtimeMode === "legacy" && - hasConversationContinuationPolicy((await getRun(run.id))?.resultJson) && + hasConversationContinuationPolicy((await getRun(run.id, { includeExecutionEvidence: true }))?.resultJson) && !(await getExecutionBlocker(db, run.companyId, issueId)); if (run.runtimeMode !== "native" && run.runtimeMode !== "legacy") return; const pending = await db.select().from(agentWakeupRequests).where(and( @@ -10630,7 +10651,7 @@ export function heartbeatService( const blocker = await getExecutionBlocker(db, wake.companyId, issueId); const sourceId = blocker?.runId; if (!sourceId || !isUuidLike(sourceId)) continue; - const run = await getRun(sourceId); + const run = await getRun(sourceId, { includeExecutionEvidence: true }); if (!run || run.companyId !== wake.companyId || run.agentId !== wake.agentId) continue; if (canContinueCancelledRun(run)) { await resumeSavedLegacyComments(wake.companyId, wake.id).catch(err => { @@ -10731,18 +10752,17 @@ export function heartbeatService( async function getRun( runId: string, - opts?: { unsafeFullResultJson?: boolean }, + opts?: { unsafeFullResultJson?: boolean; includeExecutionEvidence?: boolean }, ) { const safeForLegacyEncoding = !opts?.unsafeFullResultJson && (await hasUnsafeTextProjectionDatabase()); + const columns = opts?.unsafeFullResultJson + ? getTableColumns(heartbeatRuns) + : safeForLegacyEncoding ? heartbeatRunSqlAsciiSafeColumns : heartbeatRunSafeColumns; return db - .select( - opts?.unsafeFullResultJson - ? getTableColumns(heartbeatRuns) - : safeForLegacyEncoding - ? heartbeatRunSqlAsciiSafeColumns - : heartbeatRunSafeColumns, - ) + .select(opts?.includeExecutionEvidence + ? { ...columns, resultJson: heartbeatRunExecutionEvidenceColumn } + : columns) .from(heartbeatRuns) .where(eq(heartbeatRuns.id, runId)) .then((rows) => rows[0] ?? null); @@ -26905,7 +26925,7 @@ export function heartbeatService( } if (opts.failedRunId) { - const failed = await getRun(opts.failedRunId); + const failed = await getRun(opts.failedRunId, { includeExecutionEvidence: true }); if (opts.requestedByActorType !== "user" || !opts.requestedByActorId || reason !== "retry_failed_run" || source !== "on_demand" || triggerDetail !== "manual" || !failed || failed.companyId !== agent.companyId || failed.agentId !== agentId ||