diff --git a/apps/bullmq/src/scheduled-jobs/__tests__/sleep-check.test.ts b/apps/bullmq/src/scheduled-jobs/__tests__/sleep-check.test.ts index 89d4016cb..387017aee 100644 --- a/apps/bullmq/src/scheduled-jobs/__tests__/sleep-check.test.ts +++ b/apps/bullmq/src/scheduled-jobs/__tests__/sleep-check.test.ts @@ -496,10 +496,48 @@ describe('sleepCheckJob', () => { status: RunStatus.Completed, error: 'Auto-snapshot could not run because instance sb-1 was stopped.', }); - expect(mockRecordTaskRunEvent).toHaveBeenCalled(); + expect(mockRecordTaskRunEvent).toHaveBeenCalledWith( + expect.anything(), + expect.objectContaining({ + runId: 42, + eventType: 'completed', + message: 'Completed idle task run #42 without a snapshot.', + details: expect.objectContaining({ + decision: 'complete_without_snapshot', + }), + }), + ); expect(warnSpy).not.toHaveBeenCalled(); }); + it('cancels an active unfinished review when its due instance is already stopped', async () => { + mockJobQueries({ + dueJobs: [ + { + id: 105, + machineId: 'sb-review-stopped', + payloadKind: TaskPayloadKind.GithubPrReviewSync, + status: RunStatus.Running, + taskPhase: 'waiting_for_prompt', + vendor: 'modal', + snapshotId: null, + sleepRequestedAt: null, + snapshotRequestedAt: null, + }, + ], + }); + mockGetInstanceStatus.mockResolvedValue({ status: 'stopped' }); + + await sleepCheckJob(); + + expect(mockFinishRun).toHaveBeenCalledWith({ + id: 105, + status: RunStatus.Canceled, + error: + 'Due sleep handling found active instance sb-review-stopped in status stopped', + }); + }); + it('skips idle completion when another process already claimed the row', async () => { mockJobQueries({ dueJobs: [ @@ -928,7 +966,7 @@ describe('sleepCheckJob', () => { ); }); - it('shuts down due non-resumable jobs', async () => { + it('cancels an unfinished review when its due sandbox is shut down', async () => { const mockJob = { id: 99, machineId: 'sb-2', @@ -964,15 +1002,22 @@ describe('sleepCheckJob', () => { expect(setFn).toHaveBeenNthCalledWith(1, { sleepRequestedAt: expect.any(Date), }); - expect(setFn).toHaveBeenNthCalledWith(2, { - sleepAt: null, - taskPhase: null, - status: RunStatus.Completed, - completedAt: expect.any(Date), + expect(mockFinishRun).toHaveBeenCalledWith({ + id: 99, + status: RunStatus.Canceled, + error: + 'Due sleep handling terminated the review sandbox before the review completed', }); - expect(mockMaybeEnqueueBrainMemoryForCompletedRun).toHaveBeenCalledWith( + expect(mockMaybeEnqueueBrainMemoryForCompletedRun).not.toHaveBeenCalled(); + expect(mockRecordTaskRunEvent).toHaveBeenCalledWith( expect.anything(), - 99, + expect.objectContaining({ + runId: 99, + eventType: 'decision', + details: expect.objectContaining({ + decision: 'cancel_unfinished_review_on_sandbox_termination', + }), + }), ); expect(mockRecordComputeProviderUsage).toHaveBeenCalledWith({ runId: 99, @@ -989,6 +1034,33 @@ describe('sleepCheckJob', () => { expect(captureBullMqMessageMock).not.toHaveBeenCalled(); }); + it('releases the shutdown claim when review finalization fails', async () => { + const mockJob = { + id: 104, + machineId: 'sb-review-finalize-retry', + payloadKind: TaskPayloadKind.GithubPrReview, + status: RunStatus.Idle, + taskPhase: 'waiting_for_prompt', + vendor: 'modal', + taskId: 'task-review-finalize-retry', + snapshotId: null, + sleepRequestedAt: null, + snapshotRequestedAt: null, + }; + + mockJobQueries({ dueJobs: [mockJob] }); + mockGetInstanceStatus.mockResolvedValue({ + status: 'running', + timeoutRemainingMs: 5 * 60 * 60 * 1_000, + }); + returningFn.mockResolvedValue([{ id: 104 }]); + mockFinishRun.mockRejectedValueOnce(new Error('finalization failed')); + + await sleepCheckJob(); + + expect(setFn).toHaveBeenCalledWith({ sleepRequestedAt: null }); + }); + it('reports provider-timeout backstop shutdowns to Sentry', async () => { const mockJob = { id: 100, @@ -1072,6 +1144,38 @@ describe('sleepCheckJob', () => { expect(mockEnterStandby).not.toHaveBeenCalled(); }); + it('cancels an unfinished review whose provider instance no longer exists', async () => { + mockJobQueries({ + hardLimitJobs: [ + { + id: 106, + machineId: 'sb-review-missing', + payloadKind: TaskPayloadKind.GithubPrReview, + status: RunStatus.Running, + taskPhase: 'waiting_for_prompt', + vendor: 'azure', + snapshotId: null, + sleepRequestedAt: null, + snapshotRequestedAt: null, + sleepAt: new Date(Date.now() + 30 * 60 * 1_000), + }, + ], + }); + const { AzureDataPlaneError } = await import('@roomote/compute-providers'); + mockGetInstanceStatus.mockRejectedValue( + new AzureDataPlaneError('Requested document not found.', 404), + ); + + await sleepCheckJob(); + + expect(mockFinishRun).toHaveBeenCalledWith({ + id: 106, + status: RunStatus.Canceled, + error: + 'Provider-timeout backstop found that instance sb-review-missing no longer exists', + }); + }); + it('completes an idle run whose azure instance no longer exists', async () => { mockJobQueries({ dueJobs: [ @@ -1696,6 +1800,110 @@ describe('sleepCheckJob', () => { ); }); + it('completes an idle review when its stale sandbox is already gone', async () => { + const mockJob = { + id: 103, + machineId: 'sb-review-gone', + payloadKind: TaskPayloadKind.GithubPrReview, + status: RunStatus.Idle, + taskPhase: 'waiting_for_prompt', + vendor: 'modal', + workerHeartbeatAt: new Date(Date.now() - 5 * 60 * 1_000), + snapshotId: null, + sleepRequestedAt: null, + snapshotRequestedAt: null, + sleepAt: null, + }; + + mockJobQueries({ staleJobs: [mockJob] }); + mockGetInstanceStatus.mockResolvedValue({ status: 'stopped' }); + returningFn.mockResolvedValueOnce([{ id: 103 }]); + + await sleepCheckJob(); + + expect(mockFinishRun).toHaveBeenCalledWith({ + id: 103, + status: RunStatus.Completed, + error: + 'Idle session could not be snapshotted because instance sb-review-gone was stopped.', + }); + expect(mockRecordTaskRunEvent).toHaveBeenCalledWith( + expect.anything(), + expect.objectContaining({ + runId: 103, + eventType: 'completed', + message: 'Completed idle task run #103 without a snapshot.', + details: expect.objectContaining({ + decision: 'complete_without_snapshot', + }), + }), + ); + }); + + it('completes an idle review when its due sandbox is shut down', async () => { + const mockJob = { + id: 107, + machineId: 'sb-idle-review', + payloadKind: TaskPayloadKind.GithubPrReview, + status: RunStatus.Idle, + taskPhase: 'waiting_for_prompt', + vendor: 'modal', + taskId: 'task-idle-review', + snapshotId: null, + sleepRequestedAt: null, + snapshotRequestedAt: null, + }; + + mockJobQueries({ dueJobs: [mockJob] }); + mockGetInstanceStatus.mockResolvedValue({ + status: 'running', + timeoutRemainingMs: 5 * 60 * 60 * 1_000, + }); + returningFn.mockResolvedValue([{ id: 107 }]); + + await sleepCheckJob(); + + expect(mockFinishRun).toHaveBeenCalledWith({ + id: 107, + status: RunStatus.Completed, + }); + expect(mockMaybeEnqueueBrainMemoryForCompletedRun).not.toHaveBeenCalled(); + }); + + it('cancels an idle review when its due sandbox is shut down after a stop request', async () => { + const mockJob = { + id: 108, + machineId: 'sb-stopped-idle-review', + payloadKind: TaskPayloadKind.GithubPrReview, + status: RunStatus.Idle, + taskPhase: 'waiting_for_prompt', + vendor: 'modal', + taskId: 'task-stopped-idle-review', + snapshotId: null, + sleepRequestedAt: null, + snapshotRequestedAt: null, + }; + + mockJobQueries({ dueJobs: [mockJob] }); + mockGetInstanceStatus.mockResolvedValue({ + status: 'running', + timeoutRemainingMs: 5 * 60 * 60 * 1_000, + }); + mockDbQueryTaskRunsFindFirst.mockResolvedValue({ + cancelRequestedAt: new Date(), + }); + returningFn.mockResolvedValue([{ id: 108 }]); + + await sleepCheckJob(); + + expect(mockFinishRun).toHaveBeenCalledWith({ + id: 108, + status: RunStatus.Canceled, + error: + 'Due sleep handling terminated the review sandbox before the review completed', + }); + }); + it('cancels preparing jobs with a stale heartbeat and stop request when the sandbox is gone', async () => { const mockJob = { id: 101, @@ -1813,7 +2021,7 @@ describe('sleepCheckJob', () => { ); }); - it('destroys and fails stale non-resumable jobs', async () => { + it('destroys and cancels stale unfinished reviews', async () => { const mockJob = { id: 92, machineId: 'sb-non-resumable-stale', @@ -1846,7 +2054,7 @@ describe('sleepCheckJob', () => { }); expect(mockFinishRun).toHaveBeenCalledWith({ id: 92, - status: RunStatus.Failed, + status: RunStatus.Canceled, error: 'Worker heartbeat stale for instance sb-non-resumable-stale', }); expect(mockCreateComputeProviderMutationEventRecorder).toHaveBeenCalledWith( @@ -1948,7 +2156,7 @@ describe('sleepCheckJob', () => { ); }); - it('destroys and fails non-resumable booting jobs that miss their initial heartbeat', async () => { + it('destroys and cancels booting reviews that miss their initial heartbeat', async () => { const mockJob = { id: 98, machineId: 'sb-non-resumable-booting', @@ -1982,7 +2190,7 @@ describe('sleepCheckJob', () => { }); expect(mockFinishRun).toHaveBeenCalledWith({ id: 98, - status: RunStatus.Failed, + status: RunStatus.Canceled, error: 'Initial worker heartbeat missing for instance sb-non-resumable-booting', }); diff --git a/apps/bullmq/src/scheduled-jobs/sleep-check.ts b/apps/bullmq/src/scheduled-jobs/sleep-check.ts index 6828ff1bf..88dc1ac85 100644 --- a/apps/bullmq/src/scheduled-jobs/sleep-check.ts +++ b/apps/bullmq/src/scheduled-jobs/sleep-check.ts @@ -2,6 +2,7 @@ import { type ComputeProvider, type TaskPhase, RunStatus, + TaskPayloadKind, sleepCheckManagedComputeProviders, isStandbyResumeCapableComputeProvider, isTaskResumeCapableComputeProvider, @@ -91,6 +92,30 @@ type SleepCheckJob = Pick< | 'workerHeartbeatAt' >; +function isGithubPrReviewJob(job: SleepCheckJob): boolean { + return ( + job.payloadKind === TaskPayloadKind.GithubPrReview || + job.payloadKind === TaskPayloadKind.GithubPrReviewSync + ); +} + +function resolveSandboxTerminationStatus( + job: SleepCheckJob, + fallback: RunStatus.Completed | RunStatus.Failed | RunStatus.Canceled, +): RunStatus.Completed | RunStatus.Failed | RunStatus.Canceled { + if (!isGithubPrReviewJob(job)) { + return fallback; + } + + if (fallback === RunStatus.Canceled) { + return RunStatus.Canceled; + } + + return job.status === RunStatus.Idle + ? RunStatus.Completed + : RunStatus.Canceled; +} + const SENTRY_DESTROY_INSTANCE_REASONS = new Set([ 'provider_timeout_backstop', 'worker_heartbeat_stale', @@ -1005,7 +1030,10 @@ async function finalizeRunForMissingInstance( return; } - const finalStatus = await resolveSweptJobFinalStatus(job.id); + const finalStatus = resolveSandboxTerminationStatus( + job, + await resolveSweptJobFinalStatus(job.id), + ); await recordSleepCheckEvent( job, @@ -1065,13 +1093,16 @@ async function handleTimedSleepCandidate(params: { return { snapshotted: 0, shutDown: 0, failed: 0 }; } - const finalStatus = await resolveSweptJobFinalStatus(job.id); + const finalStatus = resolveSandboxTerminationStatus( + job, + await resolveSweptJobFinalStatus(job.id), + ); await recordSleepCheckEvent( job, finalStatus === RunStatus.Canceled ? 'decision' : 'failed', finalStatus === RunStatus.Canceled - ? `${describeSleepCheckPath(path)} found active instance ${job.machineId} in status ${status}; finalizing task run #${job.id} as canceled after its stop request.` + ? `${describeSleepCheckPath(path)} found active instance ${job.machineId} in status ${status}; finalizing task run #${job.id} as canceled.` : `${describeSleepCheckPath(path)} found active instance ${job.machineId} in status ${status}; failing the task run.`, details, ); @@ -1233,29 +1264,61 @@ async function handleTimedSleepCandidate(params: { 'sleepCheck', ); - const endedAt = new Date(); + const reviewJob = isGithubPrReviewJob(job); + const reviewFallback = reviewJob + ? await resolveSweptJobFinalStatus(job.id) + : RunStatus.Completed; + const reviewStatus = resolveSandboxTerminationStatus( + job, + reviewFallback === RunStatus.Canceled + ? RunStatus.Canceled + : RunStatus.Completed, + ); + const reviewCanceled = reviewStatus === RunStatus.Canceled; + if (reviewJob) { + try { + await finishRun({ + id: job.id, + status: reviewStatus, + ...(reviewCanceled + ? { + error: `${describeSleepCheckPath(path)} terminated the review sandbox before the review completed`, + } + : {}), + }); + } catch (error) { + await db + .update(taskRuns) + .set({ sleepRequestedAt: null }) + .where(eq(taskRuns.id, job.id)) + .catch(() => {}); + throw error; + } + } else { + const endedAt = new Date(); - await db.transaction(async (tx) => { - await tx - .update(taskRuns) - .set({ - sleepAt: null, - taskPhase: null, - status: RunStatus.Completed, - completedAt: endedAt, - }) - .where(eq(taskRuns.id, job.id)); + await db.transaction(async (tx) => { + await tx + .update(taskRuns) + .set({ + sleepAt: null, + taskPhase: null, + status: RunStatus.Completed, + completedAt: endedAt, + }) + .where(eq(taskRuns.id, job.id)); - // Direct-completion path (not via finishRun): derive the task state - // from all its runs now that this run is completed. - await syncTaskStateFromRuns(tx, job.taskId); - await maybeEnqueueBrainMemoryForCompletedRun(tx, job.id); + // Direct-completion path (not via finishRun): derive the task state + // from all its runs now that this run is completed. + await syncTaskStateFromRuns(tx, job.taskId); + await maybeEnqueueBrainMemoryForCompletedRun(tx, job.id); - await markTaskStartParallelCountEndedAt(tx, { - runId: job.id, - endedAt, + await markTaskStartParallelCountEndedAt(tx, { + runId: job.id, + endedAt, + }); }); - }); + } if (job.taskId) { try { @@ -1274,12 +1337,15 @@ async function handleTimedSleepCandidate(params: { await recordSleepCheckEvent( job, - 'completed', - `Shut down instance ${job.machineId} for non-resumable task run #${job.id}${path === 'hard_limit' ? ' due to provider-timeout backstop' : ''}.`, + reviewCanceled ? 'decision' : 'completed', + reviewCanceled + ? `Shut down instance ${job.machineId} and canceled unfinished review task run #${job.id}.` + : `Shut down instance ${job.machineId} for non-resumable task run #${job.id}${path === 'hard_limit' ? ' due to provider-timeout backstop' : ''}.`, { path, - decision: - path === 'hard_limit' + decision: reviewCanceled + ? 'cancel_unfinished_review_on_sandbox_termination' + : path === 'hard_limit' ? 'complete_hard_limit_non_resumable_shutdown' : 'complete_non_resumable_shutdown', ...buildSleepCheckDetails(job), @@ -1416,7 +1482,10 @@ async function handleHeartbeatRecoveryCandidate(params: { return { snapshotted: 0, failed: 0 }; } - const finalStatus = await resolveSweptJobFinalStatus(job.id); + const finalStatus = resolveSandboxTerminationStatus( + job, + await resolveSweptJobFinalStatus(job.id), + ); await finishRun({ id: job.id, @@ -1509,7 +1578,10 @@ async function handleHeartbeatRecoveryCandidate(params: { 'sleepCheck', ); - const finalStatus = await resolveSweptJobFinalStatus(job.id); + const finalStatus = resolveSandboxTerminationStatus( + job, + await resolveSweptJobFinalStatus(job.id), + ); await finishRun({ id: job.id, @@ -1776,10 +1848,20 @@ async function completeIdleJobWithoutSnapshot( }, ); + const reviewFallback = isGithubPrReviewJob(job) + ? await resolveSweptJobFinalStatus(job.id) + : RunStatus.Completed; + const finalStatus = resolveSandboxTerminationStatus( + job, + reviewFallback === RunStatus.Canceled + ? RunStatus.Canceled + : RunStatus.Completed, + ); + try { await finishRun({ id: job.id, - status: RunStatus.Completed, + status: finalStatus, error: errorMessage, }); } catch (error) { @@ -1797,11 +1879,16 @@ async function completeIdleJobWithoutSnapshot( await recordSleepCheckEvent( job, - 'completed', - `Completed idle task run #${job.id} without a snapshot.`, + finalStatus === RunStatus.Canceled ? 'decision' : 'completed', + finalStatus === RunStatus.Canceled + ? `Canceled idle review task run #${job.id} without a snapshot.` + : `Completed idle task run #${job.id} without a snapshot.`, { ...details, - decision: 'complete_without_snapshot', + decision: + finalStatus === RunStatus.Canceled + ? 'cancel_unfinished_review_without_snapshot' + : 'complete_without_snapshot', snapshotFailedAt: snapshotFailedAt.toISOString(), error: errorMessage, }, diff --git a/packages/github/src/api.ts b/packages/github/src/api.ts index 310d2d537..c258431aa 100644 --- a/packages/github/src/api.ts +++ b/packages/github/src/api.ts @@ -2231,6 +2231,15 @@ export async function createCheckRun( type UpdateCheck = Checks['update']; +type GetCheck = Checks['get']; + +export async function getCheckRun( + token: string, + params: GetCheck['parameters'], +): Promise { + return getOctokit(token).rest.checks.get(params); +} + export async function updateCheckRun( token: string, params: UpdateCheck['parameters'], diff --git a/packages/sdk/src/server/lib/task-runs/__tests__/dequeue-helpers.test.ts b/packages/sdk/src/server/lib/task-runs/__tests__/dequeue-helpers.test.ts index 7e924e66b..693f840e4 100644 --- a/packages/sdk/src/server/lib/task-runs/__tests__/dequeue-helpers.test.ts +++ b/packages/sdk/src/server/lib/task-runs/__tests__/dequeue-helpers.test.ts @@ -14,6 +14,7 @@ const { mockTaskRunsFindFirst, mockNotifySourceRunOnSettle, mockCaptureTaskSettled, + mockFinishRun, } = vi.hoisted(() => ({ mockDecryptSecrets: vi.fn(), mockEnvironmentVariablesFindMany: vi.fn(), @@ -27,6 +28,7 @@ const { mockTaskRunsFindFirst: vi.fn(), mockNotifySourceRunOnSettle: vi.fn(), mockCaptureTaskSettled: vi.fn(), + mockFinishRun: vi.fn(), })); vi.mock('@roomote/db/encryption', () => ({ @@ -114,9 +116,14 @@ vi.mock('../notify-fast-agent-parent-on-settle', () => ({ notifyFastAgentParentOnSettle: vi.fn().mockResolvedValue(undefined), })); +vi.mock('../finish-run', () => ({ + finishRun: (...args: unknown[]) => mockFinishRun(...args), +})); + import { resolveWorkspaceSourceControlProvider } from '@roomote/db/server'; import { + cancelAndReleaseTaskRun, createSourceControlTokenForTaskRun, fetchResolvedRuntimeEnvVars, notifyCanceledTaskRunOnSettle, @@ -124,6 +131,26 @@ import { redactSourceControlProviderEnvVars, } from '../dequeue-helpers'; +describe('cancelAndReleaseTaskRun', () => { + it('uses the terminal finalizer so linked review artifacts cannot remain pending', async () => { + const taskRun = { + ...makeTaskRun({ repo: 'owner/repo' }), + payloadKind: TaskPayloadKind.GithubPrReview, + } as TaskRun; + + await cancelAndReleaseTaskRun( + taskRun, + 'Failed to create source control token.', + ); + + expect(mockFinishRun).toHaveBeenCalledWith({ + id: taskRun.id, + status: RunStatus.Canceled, + error: 'Failed to create source control token.', + }); + }); +}); + function makeTaskRun(payload: TaskRun['payload']): TaskRun { return { id: 123, diff --git a/packages/sdk/src/server/lib/task-runs/__tests__/dequeue-task-run.test.ts b/packages/sdk/src/server/lib/task-runs/__tests__/dequeue-task-run.test.ts index c294261c3..604c2628f 100644 --- a/packages/sdk/src/server/lib/task-runs/__tests__/dequeue-task-run.test.ts +++ b/packages/sdk/src/server/lib/task-runs/__tests__/dequeue-task-run.test.ts @@ -605,7 +605,6 @@ describe('dequeueTaskRun', () => { expect(mockCancelAndReleaseTaskRun).toHaveBeenCalledWith( taskRun, 'Failed to create source control token.', - expect.any(String), ); }); @@ -742,7 +741,10 @@ describe('dequeueTaskRun', () => { existingArtifacts: taskRun.artifacts, }), ); - expect(mockNotifyCanceledTaskRunOnSettle).toHaveBeenCalledWith(taskRun); + expect(mockCancelAndReleaseTaskRun).toHaveBeenCalledWith( + taskRun, + 'Task run is not valid.', + ); }); it('cancels the task run when launch metadata persistence fails for PR review runs', async () => { @@ -771,7 +773,6 @@ describe('dequeueTaskRun', () => { expect(mockCancelAndReleaseTaskRun).toHaveBeenCalledWith( taskRun, expect.stringContaining('Failed to persist launch metadata'), - '[dequeueTaskRun]', ); }); }); diff --git a/packages/sdk/src/server/lib/task-runs/__tests__/finish-run.test.ts b/packages/sdk/src/server/lib/task-runs/__tests__/finish-run.test.ts index 2fe9e11c1..d59ab73b7 100644 --- a/packages/sdk/src/server/lib/task-runs/__tests__/finish-run.test.ts +++ b/packages/sdk/src/server/lib/task-runs/__tests__/finish-run.test.ts @@ -195,6 +195,12 @@ vi.mock('@roomote/cloud-agents/server', () => ({ /\s*(?:Self-reviewing the PR(?: with fresh eyes)? now\.|Reviewing the PR now\.|Re-reviewing new commits now\.)/i.test( body, ), + isSafetyNetReviewStatusLine: (line: string) => + [ + 'Review complete.', + 'Review could not be completed.', + 'Review was canceled.', + ].some((message) => line.trim().startsWith(message)), parseReviewSummaryMarkerSha: (body: string) => body.match(/roomote-review-summary\s+sha=([0-9a-f]+)/i)?.[1], getMarkedSection: ({ @@ -226,12 +232,14 @@ vi.mock('@roomote/telemetry/server', () => ({ })); const mockCreateIssueComment = vi.fn().mockResolvedValue(undefined); +const mockGetCheckRun = vi.fn(); const mockUpdateCheckRun = vi.fn().mockResolvedValue(undefined); vi.mock('@roomote/github', () => ({ createTaskRunGitHubToken: vi.fn().mockResolvedValue('github-token'), createIssueComment: (...args: unknown[]) => mockCreateIssueComment(...args), deleteReaction: vi.fn(), + getCheckRun: (...args: unknown[]) => mockGetCheckRun(...args), updateCheckRun: (...args: unknown[]) => mockUpdateCheckRun(...args), })); @@ -438,6 +446,7 @@ describe('finishRun', () => { mockUpdateMessage.mockResolvedValue(true); mockDbExecute.mockResolvedValue([]); mockResolveDefaultComputeProvider.mockResolvedValue('modal'); + mockGetCheckRun.mockResolvedValue({ data: { status: 'in_progress' } }); mockUpdatePendingEnvironmentSnapshot.mockResolvedValue(true); syncRunRows = []; mockDbTransaction.mockImplementation( @@ -2330,6 +2339,106 @@ describe('finishRun', () => { ); }); + it.each(['queued', 'in_progress'] as const)( + 'cancels a %s check and finalizes a pending summary when the review sandbox terminates', + async (checkStatus) => { + mockFindFirstRun.mockResolvedValue( + makeRun( + { + payloadKind: TaskPayloadKind.GithubPrReview, + payload: { repo: 'owner/repo', headSha: 'abc1234' }, + }, + { workflow: 'pr_review', surface: 'github' }, + ), + ); + mockFindManyTaskPullRequests.mockResolvedValue([reviewPrRow]); + mockFinalizeGithubPrReviewComment.mockResolvedValueOnce({ + finalized: true, + body: '', + }); + mockGetCheckRun.mockResolvedValueOnce({ + data: { status: checkStatus }, + }); + + await finishRun({ + id: 1, + status: RunStatus.Canceled, + error: 'Review sandbox terminated before the review completed.', + }); + + expect(mockBuildTerminalReviewStatus).toHaveBeenCalledWith({ + outcome: 'canceled', + taskUrl: 'https://example.com/task', + }); + expect(mockFinalizeGithubPrReviewComment).toHaveBeenCalledWith( + expect.objectContaining({ + commentId: 456, + terminalStatus: 'terminal-status', + }), + ); + expect(mockUpdateCheckRun).toHaveBeenCalledWith( + 'github-token', + expect.objectContaining({ + check_run_id: 123, + status: 'completed', + conclusion: 'cancelled', + }), + ); + }, + ); + + it('preserves a canonical successful review when the sandbox later terminates', async () => { + mockFindFirstRun.mockResolvedValue( + makeRun( + { + payloadKind: TaskPayloadKind.GithubPrReview, + payload: { repo: 'owner/repo', headSha: 'abc1234' }, + }, + { workflow: 'pr_review', surface: 'github' }, + ), + ); + mockFindManyTaskPullRequests.mockResolvedValue([reviewPrRow]); + mockFinalizeGithubPrReviewComment.mockResolvedValueOnce({ + finalized: false, + body: '\n\nNo issues found.\n\n\n', + }); + + await finishRun({ id: 1, status: RunStatus.Canceled }); + + expect(mockUpdateCheckRun).toHaveBeenCalledWith( + 'github-token', + expect.objectContaining({ + check_run_id: 123, + status: 'completed', + conclusion: 'success', + }), + ); + }); + + it('does not overwrite an already completed review check during sandbox termination', async () => { + mockFindFirstRun.mockResolvedValue( + makeRun( + { + payloadKind: TaskPayloadKind.GithubPrReview, + payload: { repo: 'owner/repo', headSha: 'abc1234' }, + }, + { workflow: 'pr_review', surface: 'github' }, + ), + ); + mockFindManyTaskPullRequests.mockResolvedValue([reviewPrRow]); + mockFinalizeGithubPrReviewComment.mockResolvedValueOnce({ + finalized: false, + body: '\n\nNo issues found.\n', + }); + mockGetCheckRun.mockResolvedValueOnce({ + data: { status: 'completed', conclusion: 'success' }, + }); + + await finishRun({ id: 1, status: RunStatus.Canceled }); + + expect(mockUpdateCheckRun).not.toHaveBeenCalled(); + }); + it('refreshes a missing check id when publication races with finalization', async () => { mockFindFirstRun.mockResolvedValue( makeRun( diff --git a/packages/sdk/src/server/lib/task-runs/dequeue-helpers.ts b/packages/sdk/src/server/lib/task-runs/dequeue-helpers.ts index 8e5cb1cba..c8dff4c0a 100644 --- a/packages/sdk/src/server/lib/task-runs/dequeue-helpers.ts +++ b/packages/sdk/src/server/lib/task-runs/dequeue-helpers.ts @@ -43,7 +43,6 @@ import { createTaskRunGiteaCredentials } from '@roomote/gitea'; import { createTaskRunAdoCredentials } from '@roomote/ado'; import { DEFAULT_ROOMOTE_COMMIT_AUTHOR, - releaseTaskRun, resolvePublicGitAuthor, resolveRunCommitAuthor, } from '@roomote/cloud-agents/server'; @@ -52,6 +51,7 @@ import { withBootstrapFailureSignal } from '../../../bootstrap-failure-signal'; import { notifySourceRunOnSettle } from './notify-source-run-on-settle'; import { notifyFastAgentParentOnSettle } from './notify-fast-agent-parent-on-settle'; import { settleSlackLiveTaskCardOnExit } from './settle-slack-live-task-card-on-exit'; +import { finishRun } from './finish-run'; /** * Resolved git author identity for commits made by the worker. @@ -342,43 +342,12 @@ async function loadPersistedDeploymentEnvVarsFromDb(): Promise< export async function cancelAndReleaseTaskRun( taskRun: TaskRun, errorMessage: string, - logPrefix: string, ): Promise { - const endedAt = new Date(); - - await db.transaction(async (tx) => { - await tx - .update(taskRuns) - .set({ - status: RunStatus.Canceled, - canceledAt: endedAt, - error: errorMessage, - }) - .where(eq(taskRuns.id, taskRun.id)); - - // Derive the owning task's state from all its runs. finishRun never - // runs for runs canceled before/at dequeue, so without this sync the task - // would stay 'active' forever. The shared helper deprioritizes this - // never-started cancel, so an earlier completed sibling still wins. - await syncTaskStateFromRuns(tx, taskRun.taskId); - - await markTaskStartParallelCountEndedAt(tx, { - runId: taskRun.id, - endedAt, - }); + await finishRun({ + id: taskRun.id, + status: RunStatus.Canceled, + error: errorMessage, }); - - try { - await releaseTaskRun(taskRun); - } catch (error) { - console.error( - `${logPrefix} Failed to release lock for run ${taskRun.id}: ${ - error instanceof Error ? error.message : String(error) - }`, - ); - } - - await notifyCanceledTaskRunOnSettle(taskRun, errorMessage); } /** diff --git a/packages/sdk/src/server/lib/task-runs/dequeue-resume-task-run.ts b/packages/sdk/src/server/lib/task-runs/dequeue-resume-task-run.ts index ecbdf7099..a14675c77 100644 --- a/packages/sdk/src/server/lib/task-runs/dequeue-resume-task-run.ts +++ b/packages/sdk/src/server/lib/task-runs/dequeue-resume-task-run.ts @@ -412,7 +412,6 @@ export const dequeueResumeTaskRun = async ( await cancelAndReleaseTaskRun( result.taskRun, 'Failed to create source control token.', - tag, ); return undefined; @@ -447,7 +446,7 @@ export const dequeueResumeTaskRun = async ( }, }); - await cancelAndReleaseTaskRun(result.taskRun, message, tag); + await cancelAndReleaseTaskRun(result.taskRun, message); return undefined; } diff --git a/packages/sdk/src/server/lib/task-runs/dequeue-task-run.ts b/packages/sdk/src/server/lib/task-runs/dequeue-task-run.ts index cf3ba6be3..ce722080b 100644 --- a/packages/sdk/src/server/lib/task-runs/dequeue-task-run.ts +++ b/packages/sdk/src/server/lib/task-runs/dequeue-task-run.ts @@ -20,7 +20,7 @@ import { and, eq, } from '@roomote/db/server'; -import { releaseTaskRun, generatePrompt } from '@roomote/cloud-agents/server'; +import { generatePrompt } from '@roomote/cloud-agents/server'; import { type GitAuthor, @@ -31,7 +31,6 @@ import { createSourceControlTokenForTaskRun, type SourceControlRuntimeToken, cancelTaskRun, - notifyCanceledTaskRunOnSettle, reportBootstrapFailure, resolveGitAuthor, claimJobById, @@ -291,7 +290,7 @@ export const dequeueTaskRun = async ( const tag = '[dequeueTaskRun]'; type TransactionResult = - | { error: true; taskRun?: TaskRun } + | { error: true; taskRun?: TaskRun; errorMessage?: string } | { error: false; taskRun: TaskRun; @@ -326,7 +325,11 @@ export const dequeueTaskRun = async ( if (!taskRun) { console.error(`${tag} Task run not found: ${JSON.stringify(dequeued)}`); - return { error: true, taskRun }; + return { + error: true, + taskRun, + errorMessage: 'Task run is not valid.', + }; } const task = taskRun.task; @@ -364,7 +367,11 @@ export const dequeueTaskRun = async ( bootstrapFailureReason: 'schema_validation_failed', existingArtifacts: taskRun.artifacts, }); - return { error: true, taskRun }; + return { + error: true, + taskRun, + errorMessage: 'Task run is not valid.', + }; } const gitAuthor = await resolveGitAuthor(tx, taskRun); @@ -418,17 +425,10 @@ export const dequeueTaskRun = async ( const { taskRun } = txResult; if (taskRun) { - try { - await releaseTaskRun(taskRun); - } catch (error) { - console.error( - `${tag} Failed to release lock for run ${taskRun.id}: ${ - error instanceof Error ? error.message : String(error) - }`, - ); - } - - await notifyCanceledTaskRunOnSettle(taskRun); + await cancelAndReleaseTaskRun( + taskRun, + txResult.errorMessage ?? 'Task run was canceled during dequeue.', + ); } return undefined; @@ -475,7 +475,6 @@ export const dequeueTaskRun = async ( await cancelAndReleaseTaskRun( txResult.taskRun, 'Failed to create source control token.', - tag, ); return undefined; } @@ -548,7 +547,7 @@ export const dequeueTaskRun = async ( } catch (error) { const message = `${tag} Failed to generate prompt for task run ${txResult.taskRun.id}: ${error instanceof Error ? error.message : String(error)}`; console.error(message); - await cancelAndReleaseTaskRun(txResult.taskRun, message, tag); + await cancelAndReleaseTaskRun(txResult.taskRun, message); return undefined; } } @@ -574,7 +573,7 @@ export const dequeueTaskRun = async ( error instanceof Error ? error.message : 'Failed to resolve harness runtime credentials.'; - await cancelAndReleaseTaskRun(txResult.taskRun, message, tag); + await cancelAndReleaseTaskRun(txResult.taskRun, message); return undefined; } @@ -633,7 +632,7 @@ export const dequeueTaskRun = async ( error instanceof Error ? error.message : String(error) }`; console.error(message); - await cancelAndReleaseTaskRun(result.taskRun, message, tag); + await cancelAndReleaseTaskRun(result.taskRun, message); return undefined; } diff --git a/packages/sdk/src/server/lib/task-runs/finish-run.ts b/packages/sdk/src/server/lib/task-runs/finish-run.ts index e8c0cbe94..ec56cc091 100644 --- a/packages/sdk/src/server/lib/task-runs/finish-run.ts +++ b/packages/sdk/src/server/lib/task-runs/finish-run.ts @@ -61,6 +61,7 @@ import { createTaskRunGitHubToken, createIssueComment, deleteReaction, + getCheckRun, updateCheckRun, } from '@roomote/github'; import { revokeTaskRunScopedGitLabTokens } from '@roomote/gitlab'; @@ -77,7 +78,10 @@ import { settleSlackLiveTaskCardOnExit } from './settle-slack-live-task-card-on- import { refreshTaskTitleOnCompletion } from './record-task-message-envelope'; import { getRedis } from '@roomote/redis'; import { resolveSlackTaskRunRouting } from './slack-task-run-routing'; -import { getGithubPrReviewCheckResult } from './github-pr-review-check'; +import { + getGithubPrReviewCheckResult, + getTerminalReviewSummaryResult, +} from './github-pr-review-check'; import { SlackNotifier, buildTaskFailedMessage, @@ -818,19 +822,40 @@ async function cleanupGithubPrReviewArtifacts( if (checkRunId) { try { token ??= await createTaskRunGitHubToken(run); - const checkResult = getGithubPrReviewCheckResult({ - runStatus: status, - reviewSummaryBody: reviewSummary.body, - safetyNetFinalized: reviewSummary.finalized, - expectedHeadSha: - 'latestObservedHeadSha' in run.payload && - typeof run.payload.latestObservedHeadSha === 'string' - ? run.payload.latestObservedHeadSha - : 'headSha' in run.payload && - typeof run.payload.headSha === 'string' - ? run.payload.headSha - : undefined, - }); + if (status === RunStatus.Canceled) { + const { data: currentCheckRun } = await getCheckRun(token, { + owner, + repo, + check_run_id: checkRunId, + }); + if (currentCheckRun.status === 'completed') { + continue; + } + } + + const expectedHeadSha = + 'latestObservedHeadSha' in run.payload && + typeof run.payload.latestObservedHeadSha === 'string' + ? run.payload.latestObservedHeadSha + : 'headSha' in run.payload && + typeof run.payload.headSha === 'string' + ? run.payload.headSha + : undefined; + const terminalSummaryResult = + status === RunStatus.Canceled && expectedHeadSha + ? getTerminalReviewSummaryResult({ + reviewSummaryBody: reviewSummary.body, + expectedHeadSha, + }) + : null; + const checkResult = + terminalSummaryResult ?? + getGithubPrReviewCheckResult({ + runStatus: status, + reviewSummaryBody: reviewSummary.body, + safetyNetFinalized: reviewSummary.finalized, + expectedHeadSha, + }); const taskUrl = getTaskUrl({ taskId: run.taskId, utm: { diff --git a/packages/sdk/src/server/lib/task-runs/github-pr-review-check.ts b/packages/sdk/src/server/lib/task-runs/github-pr-review-check.ts index 65d405abc..38ef28de3 100644 --- a/packages/sdk/src/server/lib/task-runs/github-pr-review-check.ts +++ b/packages/sdk/src/server/lib/task-runs/github-pr-review-check.ts @@ -77,7 +77,7 @@ function classifyReviewSummary(input: { }; } -function getTerminalReviewSummaryResult(input: { +export function getTerminalReviewSummaryResult(input: { reviewSummaryBody?: string; expectedHeadSha: string; }) {