diff --git a/apps/bullmq/src/jobs/pr-review-notification.test.ts b/apps/bullmq/src/jobs/pr-review-notification.test.ts index d49cf3ba9..ed6a9a6ba 100644 --- a/apps/bullmq/src/jobs/pr-review-notification.test.ts +++ b/apps/bullmq/src/jobs/pr-review-notification.test.ts @@ -577,6 +577,64 @@ describe('prReviewNotificationJob', () => { expect(mockStickyFooterPost).toHaveBeenCalled(); }); + it('releases deferred feedback exactly once after a live worker heartbeat becomes stale', async () => { + const liveRun = { + id: 1, + payload: { channel: 'C123' }, + slackThreadTs: '111.222', + sourceRunId: null, + status: RunStatus.Idle, + taskPhase: 'running', + workerHeartbeatAt: new Date(), + }; + const deadRun = { + ...liveRun, + workerHeartbeatAt: new Date(Date.now() - WORKER_HEARTBEAT_STALE_MS - 1), + }; + mockFindFirstTaskRun.mockResolvedValue(deadRun); + mockFindFirstTaskRun.mockResolvedValueOnce(liveRun); + mockConsumePending.mockResolvedValueOnce(events).mockResolvedValueOnce([]); + + await prReviewNotificationJob(makeJob() as never); + await prReviewNotificationJob(makeJob({ deferrals: 1 }) as never); + await prReviewNotificationJob(makeJob({ deferrals: 1 }) as never); + + expect(mockSchedule).toHaveBeenCalledTimes(1); + expect(mockConsumePending).toHaveBeenCalledTimes(2); + expect(mockPrepareDelivery).toHaveBeenCalledTimes(1); + expect(mockStickyFooterPost).toHaveBeenCalledTimes(1); + }); + + it('keeps feedback deferred across a worker restart until the replacement run settles', async () => { + const replacementRun = { + id: 2, + payload: { channel: 'C123' }, + slackThreadTs: '111.222', + sourceRunId: 1, + status: RunStatus.Running, + taskPhase: 'running', + workerHeartbeatAt: new Date(), + }; + mockFindFirstTaskRun + .mockResolvedValueOnce(replacementRun) + .mockResolvedValueOnce({ + ...replacementRun, + status: RunStatus.Idle, + taskPhase: 'waiting_for_prompt', + }); + + await prReviewNotificationJob(makeJob() as never); + await prReviewNotificationJob(makeJob({ deferrals: 1 }) as never); + + expect(mockSchedule).toHaveBeenCalledTimes(1); + expect(mockConsumePending).toHaveBeenCalledTimes(1); + expect(mockPrepareDelivery).toHaveBeenCalledTimes(1); + expect(mockStickyFooterPost).toHaveBeenCalledTimes(1); + expect(mockRecordDelivery).toHaveBeenCalledWith( + expect.objectContaining({ runId: 2, taskId: 'task-1' }), + ); + }); + it('drops at the deferral cap when an idle running phase has a fresh heartbeat', async () => { mockFindFirstTaskRun.mockResolvedValue({ id: 1, diff --git a/apps/worker/src/sandbox-server/lib/__tests__/harness-manager.test.ts b/apps/worker/src/sandbox-server/lib/__tests__/harness-manager.test.ts index 113af6d97..044130d6b 100644 --- a/apps/worker/src/sandbox-server/lib/__tests__/harness-manager.test.ts +++ b/apps/worker/src/sandbox-server/lib/__tests__/harness-manager.test.ts @@ -2051,6 +2051,63 @@ describe('HarnessManager touchKeepalive', () => { } }); + it('finalizes each completed turn once when duplicate terminal events arrive', () => { + const onExit = vi.fn(); + const onTaskUpdate = vi.fn(); + const { harness, manager } = createManager({ onExit, onTaskUpdate }); + const completionEvent = { + eventName: TaskEventName.TaskCompleted, + payload: [ + 'task-duplicate-completion', + { + totalTokensIn: 0, + totalTokensOut: 0, + totalCost: 0, + contextTokens: 0, + }, + {}, + { isSubtask: false }, + ], + } as TaskEvent; + + try { + manager.initializeWithoutPrompt(); + manager.startNewTaskFromPrompt({ prompt: 'hello' }); + harness.emitTaskEvent({ + eventName: TaskEventName.TaskStarted, + payload: ['task-duplicate-completion'], + } as TaskEvent); + + harness.emitTaskEvent(completionEvent); + harness.emitTaskEvent(completionEvent); + + expect(onExit).toHaveBeenCalledTimes(1); + expect( + onTaskUpdate.mock.calls.filter( + ([update]) => update.status === 'completed', + ), + ).toHaveLength(1); + + expect( + manager.sendFollowUpPrompt({ prompt: 'run a real follow-up turn' }), + ).toBe(true); + expect(manager.getStatus().phase).toBe('running'); + + harness.emitTaskEvent(completionEvent); + harness.emitTaskEvent(completionEvent); + + expect(onExit).toHaveBeenCalledTimes(2); + expect( + onTaskUpdate.mock.calls.filter( + ([update]) => update.status === 'completed', + ), + ).toHaveLength(2); + } finally { + manager.dispose(); + harness.dispose(); + } + }); + it('is a no-op when in running phase', () => { const { harness, manager } = createManager(); diff --git a/apps/worker/src/sandbox-server/lib/harness-manager.ts b/apps/worker/src/sandbox-server/lib/harness-manager.ts index 84225f383..5b6e9f310 100644 --- a/apps/worker/src/sandbox-server/lib/harness-manager.ts +++ b/apps/worker/src/sandbox-server/lib/harness-manager.ts @@ -1392,6 +1392,16 @@ export class HarnessManager extends EventEmitter { private onTaskCompleted(payload: TaskEventCompletedPayload): void { if (payload[0] === this.state.sessionId) { + if ( + this.state.taskFinishedAt !== undefined && + !isActiveTaskPhase(this.phase) + ) { + this.logger.info( + `[HarnessManager] Ignoring duplicate task completion for settled task ${payload[0]} (phase=${this.phase})`, + ); + return; + } + this.logger.info( `[HarnessManager] Task completed: ${payload[0]} (phase=${this.phase}, queuedRuntimePrompts=${this.runtimeQueuedMessagesCount}, deferredSettlement=${this.deferredTurnSettlement ?? 'none'})`, ); diff --git a/apps/worker/src/sandbox-server/lib/harnesses/__tests__/active-review-follow-up-paired-idle.test.ts b/apps/worker/src/sandbox-server/lib/harnesses/__tests__/active-review-follow-up-paired-idle.test.ts index 4f23a9a11..46c228777 100644 --- a/apps/worker/src/sandbox-server/lib/harnesses/__tests__/active-review-follow-up-paired-idle.test.ts +++ b/apps/worker/src/sandbox-server/lib/harnesses/__tests__/active-review-follow-up-paired-idle.test.ts @@ -191,26 +191,42 @@ function createFixture() { }); const onExit = vi.fn(); + const onTaskUpdate = vi.fn(); const manager = new HarnessManager({ harness, keepaliveMs: 60_000, runId: 100, taskId: 'task-100', logger: { ...createLogger(), log: vi.fn() }, - callbacks: { onExit }, + callbacks: { onExit, onTaskUpdate }, }); const taskEvents: TaskEvent[] = []; harness.subscribe((event) => taskEvents.push(event)); - return { client, harness, manager, onExit, submittedPrompts, taskEvents }; + return { + client, + harness, + manager, + onExit, + onTaskUpdate, + submittedPrompts, + taskEvents, + }; } describe('active PR review follow-up lifecycle (paired session.status idle + session.idle)', () => { it('defers run completion across the paired idle until the drained re-review turn has run', async () => { const fixture = createFixture(); - const { client, harness, manager, onExit, submittedPrompts, taskEvents } = - fixture; + const { + client, + harness, + manager, + onExit, + onTaskUpdate, + submittedPrompts, + taskEvents, + } = fixture; try { await connectHarness(harness, client); @@ -283,6 +299,11 @@ describe('active PR review follow-up lifecycle (paired session.status idle + ses (event) => event.eventName === TaskEventName.TaskCompleted, ), ).toHaveLength(2); + expect( + onTaskUpdate.mock.calls.filter( + ([update]) => update.status === 'completed', + ), + ).toHaveLength(1); } finally { manager.dispose(); harness.dispose();