diff --git a/src/auto-reply/reply/agent-runner-cli-dispatch.test.ts b/src/auto-reply/reply/agent-runner-cli-dispatch.test.ts index aed0dd27bca9..66386666510d 100644 --- a/src/auto-reply/reply/agent-runner-cli-dispatch.test.ts +++ b/src/auto-reply/reply/agent-runner-cli-dispatch.test.ts @@ -214,8 +214,57 @@ describe("runCliAgentWithLifecycle", () => { }); // A run in a long pure-reasoning stretch must keep stamping activity, or - // stale-takeover reclaims it while it is visibly thinking. + // stale-takeover reclaims it while it is visibly thinking. Stamps are + // per-event (delivery-independent), so dedupe downstream does not matter. expect(onReasoningProgress).toHaveBeenCalledTimes(2); + expect(onActivity).toHaveBeenCalledTimes(3); + }); + + it("keeps stamping onActivity when bridges are suppressed for silent runs", async () => { + cliDispatchState.runCliAgentMock.mockImplementationOnce(async (params: { runId: string }) => { + emitAgentEvent({ + runId: params.runId, + stream: "thinking", + data: { progressTokens: 50 }, + }); + emitAgentEvent({ + runId: params.runId, + stream: "assistant", + data: { text: "Silent answer", delta: "Silent answer" }, + }); + return { payloads: [], meta: { durationMs: 1 } }; + }); + const onActivity = vi.fn(); + const onAssistantText = vi.fn<(text: string) => Promise>(async () => undefined); + const onReasoningProgress = vi.fn<(payload: ReasoningProgressPayload) => Promise>( + async () => undefined, + ); + + await runCliAgentWithLifecycle({ + runId: "run-activity-suppressed", + provider: "claude-cli", + suppressAssistantBridge: true, + onActivity, + onAssistantText, + onReasoningProgress, + runParams: { + sessionId: "session-1", + sessionFile: "/tmp/session.jsonl", + workspaceDir: "/tmp/workspace", + prompt: "hello", + provider: "claude-cli", + model: "claude", + thinkLevel: "high", + timeoutMs: 1_000, + runId: "run-activity-suppressed", + }, + }); + + // silentExpected runs suppress deliveries, but their events are still real + // liveness evidence — without these stamps a healthy silent run would be + // reclaimed as run_stalled at the takeover window. + expect(onAssistantText).not.toHaveBeenCalled(); + expect(onReasoningProgress).not.toHaveBeenCalled(); expect(onActivity).toHaveBeenCalledTimes(2); }); @@ -254,8 +303,7 @@ describe("runCliAgentWithLifecycle", () => { }, }); - // Reasoning text stamps even with no onReasoningText caller: the durable - // reasoning deliver always runs, and thinking is activity evidence. + // Every real event stamps, independent of which callbacks are registered. expect(onAssistantText).toHaveBeenCalledTimes(1); expect(onActivity).toHaveBeenCalledTimes(2); }); diff --git a/src/auto-reply/reply/agent-runner-cli-dispatch.ts b/src/auto-reply/reply/agent-runner-cli-dispatch.ts index f85ce88231d4..e35eab882a4e 100644 --- a/src/auto-reply/reply/agent-runner-cli-dispatch.ts +++ b/src/auto-reply/reply/agent-runner-cli-dispatch.ts @@ -470,13 +470,10 @@ async function runCliAgentWithLifecycleInternal( ...(params.runParams.sessionKey ? { sessionKey: params.runParams.sessionKey } : {}), }); try { - if (params.onFastModeAutoProgress) { - params.onActivity?.(); - await params.onFastModeAutoProgress({ - text: summary, - channelData: { openclawProgressKind: FAST_MODE_AUTO_PROGRESS_KIND }, - }); - } + await params.onFastModeAutoProgress?.({ + text: summary, + channelData: { openclawProgressKind: FAST_MODE_AUTO_PROGRESS_KIND }, + }); } catch { // Progress hints are best-effort; a channel failure must not fail the agent turn. } @@ -528,46 +525,46 @@ async function runCliAgentWithLifecycleInternal( }, }); } - // One activity seam for every delivered CLI event: stamping per-callback at - // call sites drifted (a queued-followup callback missed one), which lets - // stale-takeover reclaim a run that is visibly producing output. - const withActivity = ( - deliver: ((payload: T) => Promise) | undefined, - ): ((payload: T) => Promise) | undefined => - deliver && params.onActivity - ? async (payload: T) => { + // One delivery-independent activity seam for every CLI agent event. + // Suppressed (silentExpected) runs still emit real events and must keep + // stamping, or a healthy silent stream looks stale to the takeover window. + const activityBridge = params.onActivity + ? createAgentEventBridge>({ + runId: params.runId, + read: () => ({}), + deliver: async () => { params.onActivity?.(); - await deliver(payload); - } - : deliver; + }, + }) + : undefined; const assistantBridge = createAssistantTextBridge({ runId: params.runId, suppressed: params.suppressAssistantBridge, - deliver: withActivity(params.onAssistantText), + deliver: params.onAssistantText, }); let finalReasoningText: string | undefined; const reasoningBridge = createReasoningTextBridge({ runId: params.runId, suppressed: params.suppressAssistantBridge, - deliver: withActivity(async (payload: ReasoningTextPayload) => { + deliver: async (payload: ReasoningTextPayload) => { finalReasoningText = normalizeOptionalString(payload.text); await params.onReasoningText?.(payload); - }), + }, }); const reasoningProgressBridge = createReasoningProgressBridge({ runId: params.runId, suppressed: params.suppressAssistantBridge, - deliver: withActivity(params.onReasoningProgress), + deliver: params.onReasoningProgress, }); const toolBridge = createToolEventBridge({ runId: params.runId, suppressed: params.suppressAssistantBridge, - deliver: withActivity(params.onToolEvent), + deliver: params.onToolEvent, }); const commentaryBridge = createCommentaryEventBridge({ runId: params.runId, suppressed: params.suppressAssistantBridge, - deliver: withActivity(params.onCommentaryText), + deliver: params.onCommentaryText, }); const toolBoundaryBridge = createToolBoundaryBridge({ runId: params.runId, @@ -575,6 +572,7 @@ async function runCliAgentWithLifecycleInternal( deliver: maybeAnnounceFastModeAutoOff, }); const bridges = [ + activityBridge, assistantBridge, reasoningBridge, reasoningProgressBridge, diff --git a/src/logging/diagnostic.ts b/src/logging/diagnostic.ts index e036fcc77852..e301e363106b 100644 --- a/src/logging/diagnostic.ts +++ b/src/logging/diagnostic.ts @@ -40,6 +40,7 @@ import { } from "./diagnostic-session-context.js"; import { requestStuckSessionRecovery, + requestStuckSessionRecoveryOutcome, resetDiagnosticSessionRecoveryCoordinatorForTest, type RecoverStuckSession, } from "./diagnostic-session-recovery-coordinator.js"; @@ -174,6 +175,36 @@ async function recoverStuckSession( }); } +/** + * @deprecated Unused by core since the dispatch-side recovery loop was removed + * (#101910); reply admission owns stale-run reclaim now. Kept only because the + * plugin SDK re-exports this module; scheduled for removal in the next SDK major. + */ +export function isStuckSessionRecoveryEnabled(config?: OpenClawConfig): boolean { + return areDiagnosticsEnabledForProcess() && isDiagnosticsEnabled(config); +} + +/** + * @deprecated Unused by core since the dispatch-side recovery loop was removed + * (#101910); reply admission owns stale-run reclaim now. Kept only because the + * plugin SDK re-exports this module; scheduled for removal in the next SDK major. + */ +export async function requestStuckDiagnosticSessionRecovery( + params: StuckSessionRecoveryRequest, +): Promise { + return requestStuckSessionRecoveryOutcome({ + recover: recoverStuckSession, + classification: { + eventType: "session.stalled", + reason: "visible_reply_wait_timeout", + classification: "stalled_agent_run", + activeWorkKind: "embedded_run", + recoveryEligible: false, + }, + request: params, + }); +} + function formatDiagnosticWorkLabel( state: { sessionId?: string;