From a2e2d7d87a50f248e791e62b6a8c3b09388934af Mon Sep 17 00:00:00 2001 From: "@daniel-lxs" <57051444+daniel-lxs@users.noreply.github.com> Date: Sat, 15 Aug 2026 16:44:29 +0000 Subject: [PATCH 1/2] fix: suppress premature chat closeout fallbacks --- ...ng-chat-closeout-fallback-delivery.test.ts | 116 ++++++++++++++++++ .../runtime-envelope-subscription.test.ts | 46 ++++++- ...ssing-chat-closeout-fallback-settlement.ts | 61 ++++++++- .../run-task/subscribe-harness-callbacks.ts | 10 ++ 4 files changed, 229 insertions(+), 4 deletions(-) diff --git a/apps/worker/src/run-task/__tests__/missing-chat-closeout-fallback-delivery.test.ts b/apps/worker/src/run-task/__tests__/missing-chat-closeout-fallback-delivery.test.ts index 11ba476a6..73c108d42 100644 --- a/apps/worker/src/run-task/__tests__/missing-chat-closeout-fallback-delivery.test.ts +++ b/apps/worker/src/run-task/__tests__/missing-chat-closeout-fallback-delivery.test.ts @@ -1,3 +1,5 @@ +import { TaskEventName } from '@roomote/types'; + const { claimDelivery, releaseDelivery, replyToChatThread } = vi.hoisted( () => ({ claimDelivery: vi.fn(), @@ -10,7 +12,10 @@ vi.mock('@roomote/sdk/client', () => ({ sdk: { taskRuns: { claimMissingChatCloseoutFallbackDelivery: claimDelivery, + recordInferenceUsage: vi.fn().mockResolvedValue({ recorded: true }), + recordMessageEnvelope: vi.fn().mockResolvedValue(null), releaseMissingChatCloseoutFallbackDelivery: releaseDelivery, + stampMilestone: vi.fn().mockResolvedValue(undefined), }, }, })); @@ -20,7 +25,9 @@ vi.mock('../../mcp/roomote-mcp-server/chat-api-client', () => ({ })); import { deliverMissingChatCloseoutFallback } from '../missing-chat-closeout-fallback-delivery'; +import { subscribeHarnessCallbacks } from '../subscribe-harness-callbacks'; import { + cancelPendingMissingChatCloseoutFallback, recordMissingChatCloseoutFallback, settleMissingChatCloseoutFallback, waitForMissingChatCloseoutFallbackDelivery, @@ -50,7 +57,12 @@ describe('deliverMissingChatCloseoutFallback', () => { replyToChatThread.mockResolvedValue({ messageTs: '123.456' }); }); + afterEach(() => { + vi.useRealTimers(); + }); + it('holds the fallback until the matching completion is settled', async () => { + vi.useFakeTimers(); const context = {}; recordMissingChatCloseoutFallback(context, { runId: 42, @@ -66,9 +78,14 @@ describe('deliverMissingChatCloseoutFallback', () => { expect(replyToChatThread).not.toHaveBeenCalled(); await settleMissingChatCloseoutFallback(context, 'completion-settled'); + expect(replyToChatThread).not.toHaveBeenCalled(); + + await vi.runAllTimersAsync(); + await waitForMissingChatCloseoutFallbackDelivery(context); expect(replyToChatThread).toHaveBeenCalledWith(expect.any(Object), { text: 'Final answer after the goal settles.', }); + vi.useRealTimers(); }); it('drops an exhausted closeout when a later completion supersedes it', async () => { @@ -89,6 +106,7 @@ describe('deliverMissingChatCloseoutFallback', () => { }); it('delivers when settlement wins the event-ordering race', async () => { + vi.useFakeTimers(); const context = {}; await settleMissingChatCloseoutFallback(context, 'completion-late-record'); @@ -99,11 +117,109 @@ describe('deliverMissingChatCloseoutFallback', () => { mcpTaskEnv, logger, }); + await vi.runAllTimersAsync(); await waitForMissingChatCloseoutFallbackDelivery(context); expect(replyToChatThread).toHaveBeenCalledWith(expect.any(Object), { text: 'Final answer recorded after settlement.', }); + vi.useRealTimers(); + }); + + it('cancels a settled fallback when later runtime activity arrives', async () => { + vi.useFakeTimers(); + const context = {}; + recordMissingChatCloseoutFallback(context, { + runId: 42, + completionId: 'completion-with-late-activity', + text: null, + mcpTaskEnv, + logger, + }); + + await settleMissingChatCloseoutFallback( + context, + 'completion-with-late-activity', + ); + cancelPendingMissingChatCloseoutFallback(context); + await vi.runAllTimersAsync(); + + expect(replyToChatThread).not.toHaveBeenCalled(); + vi.useRealTimers(); + }); + + it('cancels a settled fallback through the runtime subscription path', async () => { + vi.useFakeTimers(); + const context = {}; + let outputListener: ((event: unknown) => void) | undefined; + let taskListener: ((event: unknown) => void) | undefined; + const unsubscribe = subscribeHarnessCallbacks({ + harness: { + subscribe: (listener: (event: unknown) => void) => { + taskListener = listener; + return () => {}; + }, + subscribeRuntimeInferenceUsage: () => () => {}, + subscribeRuntimeOutput: (listener: (event: unknown) => void) => { + outputListener = listener; + return () => {}; + }, + subscribeRuntimePersistedEnvelope: () => () => {}, + subscribeRuntimeTurnCompleted: () => () => {}, + } as never, + taskRun: { id: 42, taskId: 'task-with-late-tool-activity' } as never, + callbacks: {}, + context, + logger, + mcpTaskEnv, + }); + taskListener?.({ + eventName: TaskEventName.TaskCompleted, + payload: [ + 'task-with-late-tool-activity', + {}, + {}, + { + completionId: 'completion-with-late-tool-activity', + isSubtask: false, + missingChatCloseout: { reminderCount: 3 }, + }, + ], + }); + await settleMissingChatCloseoutFallback( + context, + 'completion-with-late-tool-activity', + ); + + outputListener?.({ + eventType: 'roomote_runtime.tool_call_update', + role: 'assistant', + }); + await vi.runAllTimersAsync(); + await unsubscribe(); + + expect(replyToChatThread).not.toHaveBeenCalled(); + vi.useRealTimers(); + }); + + it('flushes a settled fallback immediately during harness teardown', async () => { + vi.useFakeTimers(); + const context = {}; + recordMissingChatCloseoutFallback(context, { + runId: 42, + completionId: 'completion-at-teardown', + text: null, + mcpTaskEnv, + logger, + }); + await settleMissingChatCloseoutFallback(context, 'completion-at-teardown'); + + await waitForMissingChatCloseoutFallbackDelivery(context); + + expect(replyToChatThread).toHaveBeenCalledWith(expect.any(Object), { + text: "I'm finished. Is there anything else you'd like me to do?", + }); + vi.useRealTimers(); }); it('posts the last finalized assistant message through the shared chat path', async () => { diff --git a/apps/worker/src/run-task/__tests__/runtime-envelope-subscription.test.ts b/apps/worker/src/run-task/__tests__/runtime-envelope-subscription.test.ts index c4fc56f5d..f34b142d3 100644 --- a/apps/worker/src/run-task/__tests__/runtime-envelope-subscription.test.ts +++ b/apps/worker/src/run-task/__tests__/runtime-envelope-subscription.test.ts @@ -1,10 +1,12 @@ import { TaskEventName, type TaskEvent } from '@roomote/types'; const { + mockCancelPendingMissingChatCloseoutFallback, mockRecordMissingChatCloseoutFallback, mockDeliverShowWidgetFallback, mockWaitForMissingChatCloseoutFallbackDelivery, } = vi.hoisted(() => ({ + mockCancelPendingMissingChatCloseoutFallback: vi.fn(), mockRecordMissingChatCloseoutFallback: vi.fn(), mockDeliverShowWidgetFallback: vi.fn().mockResolvedValue(undefined), mockWaitForMissingChatCloseoutFallbackDelivery: vi @@ -27,6 +29,8 @@ vi.mock('../show-widget-fallback-delivery', () => ({ })); vi.mock('../missing-chat-closeout-fallback-settlement', () => ({ + cancelPendingMissingChatCloseoutFallback: + mockCancelPendingMissingChatCloseoutFallback, recordMissingChatCloseoutFallback: mockRecordMissingChatCloseoutFallback, waitForMissingChatCloseoutFallbackDelivery: mockWaitForMissingChatCloseoutFallbackDelivery, @@ -41,6 +45,7 @@ import { captureWorkerException } from '../../monitoring/sentry'; import { ACP_ENVELOPE_EVENT_TYPES, + ACP_LIVE_EVENT_TYPES, type AcpMessage, type AcpPersistedEnvelope, type AcpTurnCompletedEvent, @@ -133,6 +138,7 @@ describe('subscribeHarnessCallbacks', () => { recordInferenceUsageMock.mockClear(); recordInferenceUsageMock.mockResolvedValue({ recorded: true }); captureWorkerExceptionMock.mockClear(); + mockCancelPendingMissingChatCloseoutFallback.mockClear(); mockRecordMissingChatCloseoutFallback.mockClear(); mockWaitForMissingChatCloseoutFallbackDelivery.mockClear(); mockDeliverShowWidgetFallback.mockClear(); @@ -670,12 +676,13 @@ describe('subscribeHarnessCallbacks', () => { it('does not persist raw Roomote runtime output events to task_messages', async () => { const { harness, emitOutput } = createRuntimeHarness(); + const context = {}; const unsubscribe = subscribeHarnessCallbacks({ harness: harness as never, taskRun: { id: 49, taskId: 'task-no-raw-output' } as never, callbacks: { onMessage: vi.fn().mockResolvedValue(undefined) }, - context: {}, + context, logger: { runId: 49, filePath: '/tmp/test.log', @@ -702,11 +709,48 @@ describe('subscribeHarnessCallbacks', () => { }, }); + expect(mockCancelPendingMissingChatCloseoutFallback).toHaveBeenCalledWith( + context, + ); expect(recordMessageEnvelopeMock).not.toHaveBeenCalled(); await unsubscribe(); }); + it('does not cancel a pending closeout fallback for usage-only activity', async () => { + const { harness, emitOutput } = createRuntimeHarness(); + const context = {}; + const unsubscribe = subscribeHarnessCallbacks({ + harness: harness as never, + taskRun: { id: 51, taskId: 'task-usage-output' } as never, + callbacks: {}, + context, + logger: { + runId: 51, + filePath: '/tmp/test.log', + info: vi.fn(), + warn: vi.fn(), + error: vi.fn(), + log: vi.fn(), + }, + }); + + emitOutput({ + id: 'runtime-session-usage:1', + ts: 1772823379100, + eventType: ACP_LIVE_EVENT_TYPES.UsageUpdate, + role: 'assistant', + kind: 'unknown', + contentBlocks: [], + metadata: { sessionId: 'runtime-session-usage' }, + payload: {}, + }); + + expect(mockCancelPendingMissingChatCloseoutFallback).not.toHaveBeenCalled(); + + await unsubscribe(); + }); + it('forwards request_user_input envelopes as callback events', async () => { const { harness, emitEnvelope } = createRuntimeHarness(); const callbacks = { onMessage: vi.fn().mockResolvedValue(undefined) }; diff --git a/apps/worker/src/run-task/missing-chat-closeout-fallback-settlement.ts b/apps/worker/src/run-task/missing-chat-closeout-fallback-settlement.ts index 246e06477..a9246ec5a 100644 --- a/apps/worker/src/run-task/missing-chat-closeout-fallback-settlement.ts +++ b/apps/worker/src/run-task/missing-chat-closeout-fallback-settlement.ts @@ -13,9 +13,14 @@ interface PendingMissingChatCloseoutFallback { interface MissingChatCloseoutFallbackState { pending: PendingMissingChatCloseoutFallback | null; settledCompletionIds: Set; + deliveryTimer: ReturnType | null; deliveryWork: Set>; } +// OpenCode can emit a stale idle before trailing tool activity reaches Roomote. +// Give that activity a chance to cancel the otherwise-terminal fallback. +const MISSING_CHAT_CLOSEOUT_FALLBACK_GRACE_MS = 5_000; + const stateByContext = new WeakMap< RunTaskContext, MissingChatCloseoutFallbackState @@ -30,13 +35,23 @@ function getState(context: RunTaskContext): MissingChatCloseoutFallbackState { const state: MissingChatCloseoutFallbackState = { pending: null, settledCompletionIds: new Set(), + deliveryTimer: null, deliveryWork: new Set(), }; stateByContext.set(context, state); return state; } -function startDeliveryIfSettled( +function clearDeliveryTimer(state: MissingChatCloseoutFallbackState): void { + if (!state.deliveryTimer) { + return; + } + + clearTimeout(state.deliveryTimer); + state.deliveryTimer = null; +} + +function startDeliveryNowIfSettled( state: MissingChatCloseoutFallbackState, ): Promise | null { const pending = state.pending; @@ -54,13 +69,50 @@ function startDeliveryIfSettled( return work; } +function scheduleDeliveryIfSettled( + state: MissingChatCloseoutFallbackState, +): void { + const pending = state.pending; + if ( + state.deliveryTimer || + !pending || + !state.settledCompletionIds.has(pending.completionId) + ) { + return; + } + + state.deliveryTimer = setTimeout(() => { + state.deliveryTimer = null; + startDeliveryNowIfSettled(state); + }, MISSING_CHAT_CLOSEOUT_FALLBACK_GRACE_MS); + state.deliveryTimer.unref?.(); +} + export function recordMissingChatCloseoutFallback( context: RunTaskContext, pending: PendingMissingChatCloseoutFallback | null, ): void { const state = getState(context); + clearDeliveryTimer(state); state.pending = pending; - void startDeliveryIfSettled(state); + if (!pending) { + state.settledCompletionIds.clear(); + return; + } + scheduleDeliveryIfSettled(state); +} + +export function cancelPendingMissingChatCloseoutFallback( + context: RunTaskContext, +): void { + const state = stateByContext.get(context); + if (!state?.pending) { + return; + } + + clearDeliveryTimer(state); + state.pending = null; + state.settledCompletionIds.clear(); } export async function settleMissingChatCloseoutFallback( @@ -69,7 +121,7 @@ export async function settleMissingChatCloseoutFallback( ): Promise { const state = getState(context); state.settledCompletionIds.add(completionId); - await startDeliveryIfSettled(state); + scheduleDeliveryIfSettled(state); } export async function waitForMissingChatCloseoutFallbackDelivery( @@ -80,6 +132,9 @@ export async function waitForMissingChatCloseoutFallbackDelivery( return; } + clearDeliveryTimer(state); + await startDeliveryNowIfSettled(state); + while (state.deliveryWork.size > 0) { await Promise.allSettled([...state.deliveryWork]); } diff --git a/apps/worker/src/run-task/subscribe-harness-callbacks.ts b/apps/worker/src/run-task/subscribe-harness-callbacks.ts index 30c9a601f..ed244cb2f 100644 --- a/apps/worker/src/run-task/subscribe-harness-callbacks.ts +++ b/apps/worker/src/run-task/subscribe-harness-callbacks.ts @@ -1,4 +1,5 @@ import { + ACP_LIVE_EVENT_TYPES, type AcpPersistedEnvelope, type AcpTurnCompletedEvent, TaskEventName, @@ -19,11 +20,16 @@ import { captureWorkerException } from '../monitoring/sentry'; import type { CallbackEvent, RunTaskCallbacks, RunTaskContext } from './types'; import { fromRuntimeEnvelope } from './runtime-events/envelope'; import { + cancelPendingMissingChatCloseoutFallback, recordMissingChatCloseoutFallback, waitForMissingChatCloseoutFallbackDelivery, } from './missing-chat-closeout-fallback-settlement'; import { deliverShowWidgetFallback } from './show-widget-fallback-delivery'; +const NON_ACTIVITY_RUNTIME_EVENT_TYPES = new Set( + Object.values(ACP_LIVE_EVENT_TYPES), +); + interface PendingCompletionEvents { callbackTaskId: string; events: CallbackEvent[]; @@ -303,6 +309,10 @@ export function subscribeHarnessCallbacks({ ); const unsubscribeRuntimeOutput = harness.subscribeRuntimeOutput((event) => { + if (!NON_ACTIVITY_RUNTIME_EVENT_TYPES.has(event.eventType)) { + cancelPendingMissingChatCloseoutFallback(context); + } + if (event.role !== 'assistant') { return; } From 64f7db1eb8f059524b1493a0c21c16bbea82f93e Mon Sep 17 00:00:00 2001 From: "@daniel-lxs" <57051444+daniel-lxs@users.noreply.github.com> Date: Sat, 15 Aug 2026 19:16:26 +0000 Subject: [PATCH 2/2] fix: retain active tools across closeout settlement --- ...ng-chat-closeout-fallback-delivery.test.ts | 18 +++++++--- .../runtime-envelope-subscription.test.ts | 12 +++++++ ...ssing-chat-closeout-fallback-settlement.ts | 34 +++++++++++++++++-- .../run-task/subscribe-harness-callbacks.ts | 20 +++++++++++ 4 files changed, 76 insertions(+), 8 deletions(-) diff --git a/apps/worker/src/run-task/__tests__/missing-chat-closeout-fallback-delivery.test.ts b/apps/worker/src/run-task/__tests__/missing-chat-closeout-fallback-delivery.test.ts index 73c108d42..d9552d6ac 100644 --- a/apps/worker/src/run-task/__tests__/missing-chat-closeout-fallback-delivery.test.ts +++ b/apps/worker/src/run-task/__tests__/missing-chat-closeout-fallback-delivery.test.ts @@ -148,7 +148,7 @@ describe('deliverMissingChatCloseoutFallback', () => { vi.useRealTimers(); }); - it('cancels a settled fallback through the runtime subscription path', async () => { + it('does not deliver a fallback while a tool started before completion is active', async () => { vi.useFakeTimers(); const context = {}; let outputListener: ((event: unknown) => void) | undefined; @@ -173,6 +173,12 @@ describe('deliverMissingChatCloseoutFallback', () => { logger, mcpTaskEnv, }); + outputListener?.({ + eventType: 'roomote_runtime.tool_call_update', + metadata: { toolCallId: 'tool-before-completion', status: 'in_progress' }, + payload: { toolCallId: 'tool-before-completion', status: 'in_progress' }, + role: 'tool', + }); taskListener?.({ eventName: TaskEventName.TaskCompleted, payload: [ @@ -190,15 +196,17 @@ describe('deliverMissingChatCloseoutFallback', () => { context, 'completion-with-late-tool-activity', ); + await vi.runAllTimersAsync(); + + expect(replyToChatThread).not.toHaveBeenCalled(); outputListener?.({ eventType: 'roomote_runtime.tool_call_update', - role: 'assistant', + metadata: { toolCallId: 'tool-before-completion', status: 'completed' }, + payload: { toolCallId: 'tool-before-completion', status: 'completed' }, + role: 'tool', }); - await vi.runAllTimersAsync(); await unsubscribe(); - - expect(replyToChatThread).not.toHaveBeenCalled(); vi.useRealTimers(); }); diff --git a/apps/worker/src/run-task/__tests__/runtime-envelope-subscription.test.ts b/apps/worker/src/run-task/__tests__/runtime-envelope-subscription.test.ts index f34b142d3..0f25fdda0 100644 --- a/apps/worker/src/run-task/__tests__/runtime-envelope-subscription.test.ts +++ b/apps/worker/src/run-task/__tests__/runtime-envelope-subscription.test.ts @@ -3,11 +3,13 @@ import { TaskEventName, type TaskEvent } from '@roomote/types'; const { mockCancelPendingMissingChatCloseoutFallback, mockRecordMissingChatCloseoutFallback, + mockRecordMissingChatCloseoutToolActivity, mockDeliverShowWidgetFallback, mockWaitForMissingChatCloseoutFallbackDelivery, } = vi.hoisted(() => ({ mockCancelPendingMissingChatCloseoutFallback: vi.fn(), mockRecordMissingChatCloseoutFallback: vi.fn(), + mockRecordMissingChatCloseoutToolActivity: vi.fn(), mockDeliverShowWidgetFallback: vi.fn().mockResolvedValue(undefined), mockWaitForMissingChatCloseoutFallbackDelivery: vi .fn() @@ -32,6 +34,8 @@ vi.mock('../missing-chat-closeout-fallback-settlement', () => ({ cancelPendingMissingChatCloseoutFallback: mockCancelPendingMissingChatCloseoutFallback, recordMissingChatCloseoutFallback: mockRecordMissingChatCloseoutFallback, + recordMissingChatCloseoutToolActivity: + mockRecordMissingChatCloseoutToolActivity, waitForMissingChatCloseoutFallbackDelivery: mockWaitForMissingChatCloseoutFallbackDelivery, })); @@ -140,6 +144,7 @@ describe('subscribeHarnessCallbacks', () => { captureWorkerExceptionMock.mockClear(); mockCancelPendingMissingChatCloseoutFallback.mockClear(); mockRecordMissingChatCloseoutFallback.mockClear(); + mockRecordMissingChatCloseoutToolActivity.mockClear(); mockWaitForMissingChatCloseoutFallbackDelivery.mockClear(); mockDeliverShowWidgetFallback.mockClear(); }); @@ -712,6 +717,13 @@ describe('subscribeHarnessCallbacks', () => { expect(mockCancelPendingMissingChatCloseoutFallback).toHaveBeenCalledWith( context, ); + expect(mockRecordMissingChatCloseoutToolActivity).toHaveBeenCalledWith( + context, + { + toolCallId: 'call-raw-output', + status: 'completed', + }, + ); expect(recordMessageEnvelopeMock).not.toHaveBeenCalled(); await unsubscribe(); diff --git a/apps/worker/src/run-task/missing-chat-closeout-fallback-settlement.ts b/apps/worker/src/run-task/missing-chat-closeout-fallback-settlement.ts index a9246ec5a..e2fb8c51b 100644 --- a/apps/worker/src/run-task/missing-chat-closeout-fallback-settlement.ts +++ b/apps/worker/src/run-task/missing-chat-closeout-fallback-settlement.ts @@ -13,6 +13,7 @@ interface PendingMissingChatCloseoutFallback { interface MissingChatCloseoutFallbackState { pending: PendingMissingChatCloseoutFallback | null; settledCompletionIds: Set; + activeToolCallIds: Set; deliveryTimer: ReturnType | null; deliveryWork: Set>; } @@ -20,6 +21,13 @@ interface MissingChatCloseoutFallbackState { // OpenCode can emit a stale idle before trailing tool activity reaches Roomote. // Give that activity a chance to cancel the otherwise-terminal fallback. const MISSING_CHAT_CLOSEOUT_FALLBACK_GRACE_MS = 5_000; +const TERMINAL_TOOL_STATUSES = new Set([ + 'canceled', + 'cancelled', + 'completed', + 'error', + 'failed', +]); const stateByContext = new WeakMap< RunTaskContext, @@ -35,6 +43,7 @@ function getState(context: RunTaskContext): MissingChatCloseoutFallbackState { const state: MissingChatCloseoutFallbackState = { pending: null, settledCompletionIds: new Set(), + activeToolCallIds: new Set(), deliveryTimer: null, deliveryWork: new Set(), }; @@ -53,9 +62,14 @@ function clearDeliveryTimer(state: MissingChatCloseoutFallbackState): void { function startDeliveryNowIfSettled( state: MissingChatCloseoutFallbackState, + options?: { allowActiveTools?: boolean }, ): Promise | null { const pending = state.pending; - if (!pending || !state.settledCompletionIds.has(pending.completionId)) { + if ( + !pending || + !state.settledCompletionIds.has(pending.completionId) || + (!options?.allowActiveTools && state.activeToolCallIds.size > 0) + ) { return null; } @@ -76,7 +90,8 @@ function scheduleDeliveryIfSettled( if ( state.deliveryTimer || !pending || - !state.settledCompletionIds.has(pending.completionId) + !state.settledCompletionIds.has(pending.completionId) || + state.activeToolCallIds.size > 0 ) { return; } @@ -115,6 +130,19 @@ export function cancelPendingMissingChatCloseoutFallback( state.settledCompletionIds.clear(); } +export function recordMissingChatCloseoutToolActivity( + context: RunTaskContext, + input: { toolCallId: string; status: string | null }, +): void { + const state = getState(context); + if (input.status && TERMINAL_TOOL_STATUSES.has(input.status)) { + state.activeToolCallIds.delete(input.toolCallId); + return; + } + + state.activeToolCallIds.add(input.toolCallId); +} + export async function settleMissingChatCloseoutFallback( context: RunTaskContext, completionId: string, @@ -133,7 +161,7 @@ export async function waitForMissingChatCloseoutFallbackDelivery( } clearDeliveryTimer(state); - await startDeliveryNowIfSettled(state); + await startDeliveryNowIfSettled(state, { allowActiveTools: true }); while (state.deliveryWork.size > 0) { await Promise.allSettled([...state.deliveryWork]); diff --git a/apps/worker/src/run-task/subscribe-harness-callbacks.ts b/apps/worker/src/run-task/subscribe-harness-callbacks.ts index ed244cb2f..e33425382 100644 --- a/apps/worker/src/run-task/subscribe-harness-callbacks.ts +++ b/apps/worker/src/run-task/subscribe-harness-callbacks.ts @@ -1,4 +1,5 @@ import { + ACP_ENVELOPE_EVENT_TYPES, ACP_LIVE_EVENT_TYPES, type AcpPersistedEnvelope, type AcpTurnCompletedEvent, @@ -22,6 +23,7 @@ import { fromRuntimeEnvelope } from './runtime-events/envelope'; import { cancelPendingMissingChatCloseoutFallback, recordMissingChatCloseoutFallback, + recordMissingChatCloseoutToolActivity, waitForMissingChatCloseoutFallbackDelivery, } from './missing-chat-closeout-fallback-settlement'; import { deliverShowWidgetFallback } from './show-widget-fallback-delivery'; @@ -29,6 +31,10 @@ import { deliverShowWidgetFallback } from './show-widget-fallback-delivery'; const NON_ACTIVITY_RUNTIME_EVENT_TYPES = new Set( Object.values(ACP_LIVE_EVENT_TYPES), ); +const TOOL_RUNTIME_EVENT_TYPES = new Set([ + ACP_ENVELOPE_EVENT_TYPES.ToolCall, + ACP_ENVELOPE_EVENT_TYPES.ToolCallUpdate, +]); interface PendingCompletionEvents { callbackTaskId: string; @@ -309,6 +315,20 @@ export function subscribeHarnessCallbacks({ ); const unsubscribeRuntimeOutput = harness.subscribeRuntimeOutput((event) => { + if (TOOL_RUNTIME_EVENT_TYPES.has(event.eventType)) { + const metadata = asRecord(event.metadata); + const payload = asRecord(event.payload); + const toolCallId = + asString(metadata?.toolCallId) ?? asString(payload?.toolCallId); + if (toolCallId) { + recordMissingChatCloseoutToolActivity(context, { + toolCallId, + status: + asString(metadata?.status) ?? asString(payload?.status) ?? null, + }); + } + } + if (!NON_ACTIVITY_RUNTIME_EVENT_TYPES.has(event.eventType)) { cancelPendingMissingChatCloseoutFallback(context); }