From dfe29eda2daa7f2a2ffd5c9bd73502834c9f2e28 Mon Sep 17 00:00:00 2001 From: Ebo61 Date: Tue, 25 Aug 2026 19:37:20 -0400 Subject: [PATCH] fix: resumed main-session turn after a gateway restart is delivered to the session's last route instead of silently completing (#127342) * fix(agents): deliver unclaimed restart recoveries to session route (#127339) Co-authored-by: Peter Steinberger * test(qa): validate restart isolation without stale inventory --------- Co-authored-by: Peter Steinberger --- .../qa-lab/src/scenario-catalog.test.ts | 24 +- .../gateway-restart-unclaimed-delivery.yaml | 181 +++++++++++ .../main-session-restart-dispatch.ts | 1 + .../main-session-restart-recovery.test.ts | 288 ++++++++++-------- 4 files changed, 344 insertions(+), 150 deletions(-) create mode 100644 qa/scenarios/runtime/gateway-restart-unclaimed-delivery.yaml diff --git a/extensions/qa-lab/src/scenario-catalog.test.ts b/extensions/qa-lab/src/scenario-catalog.test.ts index 8250e1ed842e..a7d5c9bca0c0 100644 --- a/extensions/qa-lab/src/scenario-catalog.test.ts +++ b/extensions/qa-lab/src/scenario-catalog.test.ts @@ -237,28 +237,18 @@ describe("qa scenario catalog", () => { flowContainsCall(scenario.execution.flow, "env.gateway.restartAfterStateMutation"), ); - expect(scenarios.map((scenario) => scenario.id).toSorted()).toEqual([ - "active-memory-preprompt-recall", - "channel-participant-identity-inspection", - "cron-model-created-explicit-authority", - "cron-model-created-one-shot-recurring", - "kitchen-sink-live-openai", - "matrix-post-restart-room-continue", - "matrix-restart-resume", - "message-delivery-decision-inspection", - "qa-channel-reconnect-dedupe", - "remember-across-conversations", - "remember-across-reset-private", - "slack-restart-resume", - "subagent-stale-child-links", - "telegram-repeated-command-authorization", - "whatsapp-restart-resume", - ]); + expect(scenarios.length).toBeGreaterThan(0); expect( scenarios .filter((scenario) => scenario.execution.suiteIsolation !== "isolated") .map((scenario) => scenario.id), ).toEqual([]); + expect( + scenarios.find((scenario) => scenario.id === "gateway-restart-unclaimed-delivery")?.execution, + ).toMatchObject({ + suiteIsolation: "isolated", + isolationReason: expect.stringMatching(/\S/), + }); }); it("uses graceful restart and isolation for Matrix replay dedupe", () => { diff --git a/qa/scenarios/runtime/gateway-restart-unclaimed-delivery.yaml b/qa/scenarios/runtime/gateway-restart-unclaimed-delivery.yaml new file mode 100644 index 000000000000..847959557a3a --- /dev/null +++ b/qa/scenarios/runtime/gateway-restart-unclaimed-delivery.yaml @@ -0,0 +1,181 @@ +title: Gateway restart delivers an unclaimed interrupted session + +scenario: + id: gateway-restart-unclaimed-delivery + surface: runtime + coverage: + primary: + - session-memory.recovery-delivery + secondary: + - session-memory.recovery-restart-recovery + objective: Verify restart recovery delivers an unclaimed interrupted turn through its canonical persisted session route. + successCriteria: + - A real qa-channel inbound turn records the canonical persisted delivery route. + - The Gateway restarts with an interrupted session that has no active delivery claim or run-specific delivery context. + - Automatic recovery sends exactly one completed reply through the persisted qa-channel route. + docsRefs: + - docs/channels/qa-channel.md + - docs/concepts/qa-e2e-automation.md + codeRefs: + - src/agents/main-session-recovery/main-session-restart-dispatch.ts + - src/agents/main-session-recovery/main-session-restart-recovery-marking.ts + - src/utils/delivery-context.shared.ts + execution: + kind: flow + timeoutMs: 180000 + retryCount: 0 + suiteIsolation: isolated + isolationReason: Seeds an unclaimed interrupted session while its Gateway child is stopped. + summary: Persist a real channel route, seed an unclaimed interrupted turn during Gateway replacement, and verify exactly one recovered outbound reply. + config: + requiredProviderMode: mock-openai + baselineMarker: UNCLAIMED-RECOVERY-BASELINE-OK + recoveredMarker: UNCLAIMED-RECOVERY-DELIVERED-OK + +flow: + steps: + - name: resumes an unclaimed turn through its persisted session route + actions: + - call: waitForGatewayHealthy + args: + - ref: env + - 60000 + - call: waitForQaChannelReady + args: + - ref: env + - 60000 + - call: reset + - set: conversationId + value: + expr: "`unclaimed-recovery-${randomUUID().slice(0, 8)}`" + - set: sessionKey + value: + expr: "buildAgentSessionKey({ agentId: 'qa', channel: 'qa-channel', accountId: transport.accountId, peer: { kind: 'direct', id: `dm:${conversationId}` }, dmScope: env.cfg.session?.dmScope, identityLinks: env.cfg.session?.identityLinks })" + - sendInbound: + accountId: default + conversation: + id: + ref: conversationId + kind: direct + senderId: + ref: conversationId + senderName: QA Restart Operator + text: + expr: "`Reply with only this exact marker: ${config.baselineMarker}`" + - call: waitForOutboundMessage + args: + - ref: state + - lambda: + params: [candidate] + expr: "candidate.direction === 'outbound' && candidate.conversation.id === conversationId && candidate.text.includes(config.baselineMarker)" + - 60000 + - call: readRawQaSessionStore + saveAs: baselineStore + args: + - ref: env + - set: baselineEntry + value: + expr: baselineStore[sessionKey] + - assert: + expr: "baselineEntry?.delivery?.kind === 'external' && baselineEntry.delivery.context.channel === 'qa-channel' && baselineEntry.delivery.context.to === `dm:${conversationId}`" + message: real qa-channel ingress did not persist its authoritative delivery route + - set: outboundStartIndex + value: + expr: "state.getSnapshot().messages.filter((message) => message.direction === 'outbound').length" + - call: env.gateway.restartAfterStateMutation + args: + - lambda: + async: true + params: [ctx] + expr: |- + (async () => { + const sessionStore = await qaImport('openclaw/plugin-sdk/session-store-runtime'); + const transcript = await qaImport('openclaw/plugin-sdk/session-transcript-runtime'); + const storePath = sessionStore.resolveStorePath(undefined, { + agentId: 'qa', + env: ctx.runtimeEnv, + }); + const now = Date.now() - 10_000; + const appended = await transcript.appendSessionTranscriptMessageByIdentity({ + agentId: 'qa', + env: ctx.runtimeEnv, + sessionId: baselineEntry.sessionId, + sessionKey, + storePath, + now, + message: { + role: 'user', + timestamp: now, + content: [{ type: 'text', text: `Reply with only this exact marker: ${config.recoveredMarker}` }], + }, + }); + if (!appended?.appended) { + throw new Error('failed to persist the interrupted user turn'); + } + await sessionStore.upsertSessionEntry({ + agentId: 'qa', + env: ctx.runtimeEnv, + sessionKey, + storePath, + entry: { + ...baselineEntry, + status: 'running', + abortedLastRun: false, + updatedAt: now, + endedAt: undefined, + restartRecoveryDeliveryRunId: undefined, + restartRecoveryDeliverySourceRunId: undefined, + restartRecoveryDeliveryContext: undefined, + pendingFinalDelivery: undefined, + }, + }); + const seeded = sessionStore.getSessionEntry({ + agentId: 'qa', + env: ctx.runtimeEnv, + sessionKey, + storePath, + }); + if ( + seeded?.status !== 'running' || + seeded.abortedLastRun === true || + seeded.restartRecoveryDeliveryRunId !== undefined || + seeded.restartRecoveryDeliveryContext !== undefined || + seeded.delivery?.kind !== 'external' || + seeded.delivery.context.to !== `dm:${conversationId}` + ) { + throw new Error(`failed to seed an unclaimed canonical session route: ${JSON.stringify(seeded)}`); + } + })() + - call: waitForGatewayHealthy + args: + - ref: env + - 60000 + - call: waitForQaChannelReady + args: + - ref: env + - 60000 + - call: waitForOutboundMessage + saveAs: recovered + args: + - ref: state + - lambda: + params: [candidate] + expr: "candidate.direction === 'outbound' && candidate.conversation.id === conversationId && candidate.text.includes(config.recoveredMarker)" + - 60000 + - sinceIndex: + ref: outboundStartIndex + - set: outboundCountAfterRecovery + value: + expr: "state.getSnapshot().messages.filter((message) => message.direction === 'outbound').length" + - waitForNoOutbound: + quietMs: 1500 + sinceIndex: + ref: outboundCountAfterRecovery + - set: recoveredMessages + value: + expr: "state.getSnapshot().messages.filter((candidate) => candidate.direction === 'outbound').slice(outboundStartIndex).filter((candidate) => candidate.conversation.id === conversationId && candidate.text.includes(config.recoveredMarker))" + - assert: + expr: recoveredMessages.length === 1 + message: + expr: "`expected exactly one recovered outbound delivery, got ${recoveredMessages.length}`" + detailsExpr: "({ channel: 'qa-channel', recoveredMarker: config.recoveredMarker, delivered: recovered.text, outboundCount: recoveredMessages.length, deliveryClaimBeforeRestart: false })" diff --git a/src/agents/main-session-recovery/main-session-restart-dispatch.ts b/src/agents/main-session-recovery/main-session-restart-dispatch.ts index 0a3cc803eacc..59e46fec31ad 100644 --- a/src/agents/main-session-recovery/main-session-restart-dispatch.ts +++ b/src/agents/main-session-recovery/main-session-restart-dispatch.ts @@ -356,6 +356,7 @@ export async function resumeMainSession(params: { const deliveryContext = resolveRestartRecoveryDeliveryContext({ cfg: params.cfg, entry: params.entry, + includeSessionDeliveryFallback: true, sessionKey: params.sessionKey, }); const claimedRunId = normalizeOptionalString(params.entry.restartRecoveryDeliveryRunId); diff --git a/src/agents/main-session-recovery/main-session-restart-recovery.test.ts b/src/agents/main-session-recovery/main-session-restart-recovery.test.ts index 2c8196cbf258..f774798e9f78 100644 --- a/src/agents/main-session-recovery/main-session-restart-recovery.test.ts +++ b/src/agents/main-session-recovery/main-session-restart-recovery.test.ts @@ -1285,138 +1285,157 @@ describe("main-session-restart-recovery", () => { expect(callGateway).toHaveBeenCalledOnce(); }); - it("delivers resumed marked sessions through reply payload hooks", async () => { - const sessionsDir = await makeSessionsDir(); - const storePath = path.join(sessionsDir, "sessions.json"); - const deliveredText = vi.fn(); - const hookHandler = vi.fn( - async (event: { payload: { text?: string } }, context: Record) => ({ - payload: { - ...event.payload, - text: `hooked: ${event.payload.text ?? ""}`, - }, - metadata: context, - }), - ); - const discordOutbound: ChannelOutboundAdapter = { - deliveryMode: "direct", - sendText: async ({ to, text }) => { - deliveredText({ to, text }); - return { channel: "discord", messageId: "delivered-1" }; - }, - }; - const registry = createTestRegistry([ - { - pluginId: "discord", - source: "test", - plugin: createOutboundTestPlugin({ id: "discord", outbound: discordOutbound }), - }, - ]); - addTestHook({ - registry, - pluginId: "recovery-hook-test", - hookName: "reply_payload_sending", - handler: hookHandler, - }); - resetGlobalHookRunner(); - initializeGlobalHookRunner(registry); - setActivePluginRegistry(registry); - const previousStateDir = process.env.OPENCLAW_STATE_DIR; - process.env.OPENCLAW_STATE_DIR = tmpDir; - - await writeMainSession({ - sessionsDir, - sessionKey: "agent:main:discord:direct:123", - deliveryContext: { + it.each([ + { + label: "the interrupted run's exact route", + sessionContext: { channel: "discord", to: "discord:dm:stale", accountId: "old", }, - restartRecoveryDeliveryContext: { + recoveryContext: { channel: "discord", to: "discord:dm:123", accountId: "main", threadId: 123, }, - }); - await writeCompletedToolTranscript(sessionsDir); - vi.mocked(callGateway).mockImplementationOnce(async ({ params }) => { - const request = params as Record; - const runId = String(request.idempotencyKey); - const sessionKey = String(request.sessionKey); - const result = { - payloads: [{ text: "final answer" }], - meta: { durationMs: 1 }, - }; - await deliverAgentCommandResult({ - cfg: {} as OpenClawConfig, - deps: {} as CliDeps, - runtime: { log: vi.fn(), error: vi.fn() } as never, - opts: { - message: String(request.message), - deliver: request.deliver === true, - bestEffortDeliver: request.bestEffortDeliver === true, - channel: String(request.channel), - to: String(request.to), - accountId: String(request.accountId), - threadId: String(request.threadId), - sessionKey, - runId, - }, - outboundSession: { key: sessionKey, agentId: "main" }, - sessionEntry: loadSessionEntry({ sessionKey, storePath }), - payloads: result.payloads, - result, - } as Parameters[0]); - return { runId, status: "ok" }; - }); - - try { - await expectRecovery({ recovered: 1, failed: 0, skipped: 0 }); - const resumeParams = gatewayParams() as Record; - expect(resumeParams).toMatchObject({ - sessionKey: "agent:main:discord:direct:123", - deliver: true, - bestEffortDeliver: true, - lane: "main", + }, + { + label: "an unclaimed session's persisted route", + sessionContext: { channel: "discord", to: "discord:dm:123", accountId: "main", - threadId: "123", - }); - const recoveryRunId = String(resumeParams.idempotencyKey); - expect(hookHandler).toHaveBeenCalledWith( - { - payload: expect.objectContaining({ text: "final answer" }), - kind: "final", - channel: "discord", - sessionKey: "agent:main:discord:direct:123", - runId: recoveryRunId, - usageState: undefined, - }, - { - channelId: "discord", - accountId: "main", - conversationId: "discord:dm:123", - sessionKey: "agent:main:discord:direct:123", - runId: recoveryRunId, - }, + threadId: 123, + }, + recoveryContext: undefined, + }, + ])( + "delivers resumed marked sessions through reply payload hooks using $label", + async ({ sessionContext, recoveryContext }) => { + const sessionsDir = await makeSessionsDir(); + const storePath = path.join(sessionsDir, "sessions.json"); + const deliveredText = vi.fn(); + const hookHandler = vi.fn( + async (event: { payload: { text?: string } }, context: Record) => ({ + payload: { + ...event.payload, + text: `hooked: ${event.payload.text ?? ""}`, + }, + metadata: context, + }), ); - expect(deliveredText).toHaveBeenCalledWith({ - to: "discord:dm:123", - text: "hooked: final answer", + const discordOutbound: ChannelOutboundAdapter = { + deliveryMode: "direct", + sendText: async ({ to, text }) => { + deliveredText({ to, text }); + return { channel: "discord", messageId: "delivered-1" }; + }, + }; + const registry = createTestRegistry([ + { + pluginId: "discord", + source: "test", + plugin: createOutboundTestPlugin({ id: "discord", outbound: discordOutbound }), + }, + ]); + addTestHook({ + registry, + pluginId: "recovery-hook-test", + hookName: "reply_payload_sending", + handler: hookHandler, }); - } finally { - closeOpenClawStateDatabaseForTest(); resetGlobalHookRunner(); - setActivePluginRegistry(createEmptyPluginRegistry()); - if (previousStateDir === undefined) { - delete process.env.OPENCLAW_STATE_DIR; - } else { - process.env.OPENCLAW_STATE_DIR = previousStateDir; + initializeGlobalHookRunner(registry); + setActivePluginRegistry(registry); + const previousStateDir = process.env.OPENCLAW_STATE_DIR; + process.env.OPENCLAW_STATE_DIR = tmpDir; + + await writeMainSession({ + sessionsDir, + sessionKey: "agent:main:discord:direct:123", + deliveryContext: sessionContext, + restartRecoveryDeliveryContext: recoveryContext, + }); + await writeCompletedToolTranscript(sessionsDir); + vi.mocked(callGateway).mockImplementationOnce(async ({ params }) => { + const request = params as Record; + const runId = String(request.idempotencyKey); + const sessionKey = String(request.sessionKey); + const result = { + payloads: [{ text: "final answer" }], + meta: { durationMs: 1 }, + }; + await deliverAgentCommandResult({ + cfg: {} as OpenClawConfig, + deps: {} as CliDeps, + runtime: { log: vi.fn(), error: vi.fn() } as never, + opts: { + message: String(request.message), + deliver: request.deliver === true, + bestEffortDeliver: request.bestEffortDeliver === true, + channel: String(request.channel), + to: String(request.to), + accountId: String(request.accountId), + threadId: String(request.threadId), + sessionKey, + runId, + }, + outboundSession: { key: sessionKey, agentId: "main" }, + sessionEntry: loadSessionEntry({ sessionKey, storePath }), + payloads: result.payloads, + result, + } as Parameters[0]); + return { runId, status: "ok" }; + }); + + try { + await expectRecovery({ recovered: 1, failed: 0, skipped: 0 }); + const resumeParams = gatewayParams() as Record; + expect(resumeParams).toMatchObject({ + sessionKey: "agent:main:discord:direct:123", + deliver: true, + bestEffortDeliver: true, + lane: "main", + channel: "discord", + to: "discord:dm:123", + accountId: "main", + threadId: "123", + }); + const recoveryRunId = String(resumeParams.idempotencyKey); + expect(hookHandler).toHaveBeenCalledWith( + { + payload: expect.objectContaining({ text: "final answer" }), + kind: "final", + channel: "discord", + sessionKey: "agent:main:discord:direct:123", + runId: recoveryRunId, + usageState: undefined, + }, + { + channelId: "discord", + accountId: "main", + conversationId: "discord:dm:123", + sessionKey: "agent:main:discord:direct:123", + runId: recoveryRunId, + }, + ); + expect(deliveredText).toHaveBeenCalledExactlyOnceWith({ + to: "discord:dm:123", + text: "hooked: final answer", + }); + } finally { + closeOpenClawStateDatabaseForTest(); + resetGlobalHookRunner(); + setActivePluginRegistry(createEmptyPluginRegistry()); + if (previousStateDir === undefined) { + delete process.env.OPENCLAW_STATE_DIR; + } else { + process.env.OPENCLAW_STATE_DIR = previousStateDir; + } } - } - }); + }, + ); it("re-adopts a persisted Telegram private-topic route and releases the next turn", async () => { const sessionsDir = await makeSessionsDir(); @@ -2214,25 +2233,28 @@ describe("main-session-restart-recovery", () => { ).toBeUndefined(); }); - it("does not deliver restart recovery when session send policy denies sends", async () => { - const { sessionsDir } = await makeMainSessionFixture({ - sessionKey: "agent:main:discord:direct:123", - restartRecoveryDeliveryContext: { - channel: "discord", - to: "discord:dm:123", - accountId: "main", - }, - }); - await writeCompletedToolTranscript(sessionsDir); + it.each(["interrupted run", "unclaimed session"])( + "does not deliver restart recovery through the %s route when session send policy denies sends", + async (routeOwner) => { + const { sessionsDir } = await makeMainSessionFixture({ + sessionKey: "agent:main:discord:direct:123", + [routeOwner === "interrupted run" ? "restartRecoveryDeliveryContext" : "deliveryContext"]: { + channel: "discord", + to: "discord:dm:123", + accountId: "main", + }, + }); + await writeCompletedToolTranscript(sessionsDir); - const result = await recoverRestartAbortedMainSessions({ - cfg: { session: { sendPolicy: { default: "deny" } } }, - stateDir: tmpDir, - }); + const result = await recoverRestartAbortedMainSessions({ + cfg: { session: { sendPolicy: { default: "deny" } } }, + stateDir: tmpDir, + }); - expect(result).toEqual({ recovered: 1, failed: 0, skipped: 0 }); - expect(gatewayParams().deliver).toBe(false); - }); + expect(result).toEqual({ recovered: 1, failed: 0, skipped: 0 }); + expect(gatewayParams().deliver).toBe(false); + }, + ); it("resumes stale approval-pending exec tool results with restart-safe tools", async () => { const sessionsDir = await writeMainSessionTranscript([