From 86cfcd3833f75fa2168776425b61bd562de88155 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sun, 16 Aug 2026 14:28:26 -0700 Subject: [PATCH] fix(delivery): unify terminal settlement ownership (#124825) * fix(delivery): unify terminal settlement ownership Treat identityless adapter returns as potentially visible across channel, queue, and cron paths. Let recovery own terminal completion so ambiguity persists as notice debt instead of being double-settled or silently suppressed. Co-authored-by: ruel225 * refactor(delivery): narrow terminal internals Remove now-unused internal exports after terminal-settlement ownership was consolidated. * test(tts): preserve message runtime exports Import and spread the actual message runtime so the focused mock retains every runtime binding while overriding only the durable send core. --------- Co-authored-by: ruel225 --- src/auto-reply/reply/route-reply.ts | 7 +- src/channels/message/runtime.ts | 1 + src/channels/message/send.ts | 19 +++++ src/channels/turn/durable-delivery.ts | 11 ++- .../turn/run-channel-turn.delivery.test.ts | 37 +++++++++ src/cron/isolated-agent/delivery-dispatch.ts | 10 +-- .../delivery-outbound.runtime.ts | 5 +- src/infra/outbound/deliver-queue-execute.ts | 44 ++++------- src/infra/outbound/deliver.test.ts | 12 +++ .../delivery-completion.integration.test.ts | 61 +++++++++++++++ src/infra/outbound/delivery-completion.ts | 19 ++++- src/infra/outbound/delivery-queue-recovery.ts | 59 +++++++------- src/infra/outbound/delivery-queue-storage.ts | 1 - .../outbound/delivery-queue.recovery.test.ts | 78 ++++++++++++++++++- src/tts/tts-runtime-fallbacks.test.ts | 16 ++-- 15 files changed, 297 insertions(+), 83 deletions(-) diff --git a/src/auto-reply/reply/route-reply.ts b/src/auto-reply/reply/route-reply.ts index c6294f7f6e2d..cea3d767d655 100644 --- a/src/auto-reply/reply/route-reply.ts +++ b/src/auto-reply/reply/route-reply.ts @@ -315,7 +315,8 @@ export async function routeReply(params: RouteReplyParams): Promise + outcome.status === "failed" + ? outcome.sentBeforeError + : outcome.status === "sent" || outcome.reason === "adapter_returned_no_identity", + ) === true + ); +} + export type SerializedDurableMessagePayloadOutcome = | { index: number; status: "sent"; resultCount: number } | { diff --git a/src/channels/turn/durable-delivery.ts b/src/channels/turn/durable-delivery.ts index aecfc1d1d505..34483db8b708 100644 --- a/src/channels/turn/durable-delivery.ts +++ b/src/channels/turn/durable-delivery.ts @@ -13,7 +13,10 @@ import { } from "../../infra/outbound/deliver.js"; import { buildOutboundSessionContext } from "../../infra/outbound/session-context.js"; import { deriveDurableFinalDeliveryRequirements } from "../message/capabilities.js"; -import { sendDurableMessageBatchCore } from "../message/send.js"; +import { + durableMessageBatchMayHaveReachedRecipient, + sendDurableMessageBatchCore, +} from "../message/send.js"; import { createChannelDeliveryResultFromReceipt } from "./delivery-result.js"; import type { ChannelDeliveryInfo, ChannelDeliveryResult } from "./types.js"; @@ -241,7 +244,7 @@ export async function deliverInboundReplyWithMessageSendContextCore( receipt: send.receipt, threadId: stringifyThreadId(threadId), ...(replyToId ? { replyToId } : {}), - visibleReplySent: send.status === "sent", + visibleReplySent: durableMessageBatchMayHaveReachedRecipient(send), ...(send.deliveryIntent ? { deliveryIntent: toDeliveryIntent(send.deliveryIntent) } : {}), }); const delivery: ChannelDeliveryResult = @@ -249,7 +252,9 @@ export async function deliverInboundReplyWithMessageSendContextCore( ? { ...receiptDelivery, suppression: resolveDurableSuppression(send) } : receiptDelivery; if (send.status === "suppressed") { - return { status: "handled_no_send", reason: "no_visible_result", delivery }; + return delivery.visibleReplySent === true + ? { status: "handled_visible", delivery } + : { status: "handled_no_send", reason: "no_visible_result", delivery }; } return { status: "handled_visible", delivery }; } diff --git a/src/channels/turn/run-channel-turn.delivery.test.ts b/src/channels/turn/run-channel-turn.delivery.test.ts index 46689587afef..bdd302090cc9 100644 --- a/src/channels/turn/run-channel-turn.delivery.test.ts +++ b/src/channels/turn/run-channel-turn.delivery.test.ts @@ -761,6 +761,43 @@ describe("channel turn delivery", () => { }); }); + it("keeps no-identity durable sends visible through lifecycle settlement", async () => { + sendDurableMessageBatch.mockResolvedValueOnce({ + status: "suppressed", + results: [], + receipt: { platformMessageIds: [], parts: [], sentAt: 1 }, + reason: "adapter_returned_no_identity", + }); + const onDelivered = vi.fn(); + + const result = await dispatchRoutedChannelTurn({ + cfg, + channel: "telegram", + route: { agentId: "main", sessionKey: "agent:main:telegram:peer" }, + ctxPayload: createCtx({ Surface: "telegram", To: "chat-1" }), + delivery: { + deliver: vi.fn(), + durable: { replyToMode: "first" }, + onDelivered, + }, + }); + + expect(onDelivered).toHaveBeenCalledWith( + { text: "reply" }, + { kind: "final" }, + expect.objectContaining({ + visibleReplySent: true, + suppression: { reason: "adapter_returned_no_identity" }, + }), + ); + expectDispatched(result); + expect(result.dispatchResult).toMatchObject({ + queuedFinal: true, + counts: { tool: 0, block: 0, final: 1 }, + }); + expect(hasVisibleChannelTurnDispatch(result.dispatchResult)).toBe(true); + }); + it("prepares payloads before durable enqueue and observes handled delivery", async () => { sendDurableMessageBatch.mockResolvedValueOnce(createDurableSendResult(["tlon-1"])); const onDelivered = vi.fn(); diff --git a/src/cron/isolated-agent/delivery-dispatch.ts b/src/cron/isolated-agent/delivery-dispatch.ts index cb604774e3b8..deede609dabe 100644 --- a/src/cron/isolated-agent/delivery-dispatch.ts +++ b/src/cron/isolated-agent/delivery-dispatch.ts @@ -163,6 +163,7 @@ export async function dispatchCronDelivery( const { buildOutboundSessionContext, createOutboundSendDeps, + durableMessageBatchMayHaveReachedRecipient, resolveAgentOutboundIdentity, resolveCronChannelReplyTransform, sendDurableMessageBatchCore, @@ -323,15 +324,8 @@ export async function dispatchCronDelivery( attemptedPayloadsForMirror.push(payload); }, }); - // No durable id is still ambiguous: the adapter was already invoked. payloadMayHaveReachedRecipientBeforeFailure ||= - send.payloadOutcomes?.some( - (outcome) => - outcome.status === "sent" || - (outcome.status === "failed" && outcome.sentBeforeError) || - (outcome.status === "suppressed" && - outcome.reason === "adapter_returned_no_identity"), - ) ?? false; + durableMessageBatchMayHaveReachedRecipient(send); if ( send.status === "failed" && (await waitForCompletedDirectCronDelivery({ diff --git a/src/cron/isolated-agent/delivery-outbound.runtime.ts b/src/cron/isolated-agent/delivery-outbound.runtime.ts index 28260e79b78f..8edc3e55ef3b 100644 --- a/src/cron/isolated-agent/delivery-outbound.runtime.ts +++ b/src/cron/isolated-agent/delivery-outbound.runtime.ts @@ -6,7 +6,10 @@ import { normalizeAnyChannelId } from "../../channels/registry-normalize.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; export { createOutboundSendDeps } from "../../cli/outbound-send-deps.js"; -export { sendDurableMessageBatchCore } from "../../channels/message/runtime.js"; +export { + durableMessageBatchMayHaveReachedRecipient, + sendDurableMessageBatchCore, +} from "../../channels/message/runtime.js"; export { type OutboundDeliveryResult } from "../../infra/outbound/deliver.js"; export { resolveAgentOutboundIdentity } from "../../infra/outbound/identity.js"; export { buildOutboundSessionContext } from "../../infra/outbound/session-context.js"; diff --git a/src/infra/outbound/deliver-queue-execute.ts b/src/infra/outbound/deliver-queue-execute.ts index 49767eb16b41..f4e3cef9ab23 100644 --- a/src/infra/outbound/deliver-queue-execute.ts +++ b/src/infra/outbound/deliver-queue-execute.ts @@ -24,11 +24,7 @@ import { type OutboundPayloadDeliveryOutcome, } from "./deliver-types.js"; import { runOutboundDeliveryCommitHooks } from "./delivery-commit-hooks.js"; -import { - completeDurableDelivery, - rejectDurableDelivery, - suppressDurableDelivery, -} from "./delivery-completion.js"; +import { rejectDurableDelivery, settleDurableDelivery } from "./delivery-completion.js"; import { failDelivery, failDeliveryAfterPlatformSend, @@ -87,9 +83,22 @@ export async function deliverOutboundPayloadsWithQueueCleanup( params.requireUnknownSendReconciliation === true && platformQueueId !== undefined; let queuedPreSendState: QueuedPreSendState | undefined; let queuedPostSendState: QueuedPostSendState | undefined; + let platformSendStarted = false; let platformSendRoute: PlatformSendRoute | undefined; let deliveredResults: OutboundDeliveryResult[] = []; let commitHooksRun = false; + const settleDeliveryCompletion = async ( + result: OutboundDeliveryResult | undefined, + ): Promise => { + if (!params.deliveryCompletion) { + return; + } + await settleDurableDelivery( + params.deliveryCompletion, + result ? { result } : { platformSendStarted }, + platformQueueStateDir, + ); + }; // Deliberately process-local: message_sent is best-effort after queue // settlement, not a durable plugin outbox or a reason to retry delivery. const messageSentEvents: MessageSentEvent[] = []; @@ -207,6 +216,7 @@ export async function deliverOutboundPayloadsWithQueueCleanup( params.abortSignal?.throwIfAborted(); await params.onPlatformSendStart?.(route); params.abortSignal?.throwIfAborted(); + platformSendStarted = true; }, onPlatformSendDispatch: async () => { params.abortSignal?.throwIfAborted(); @@ -301,17 +311,7 @@ export async function deliverOutboundPayloadsWithQueueCleanup( }); } if (!queueId) { - if (params.deliveryCompletion) { - if (results.length > 0) { - await completeDurableDelivery( - params.deliveryCompletion, - results.at(-1)!, - platformQueueStateDir, - ); - } else { - await suppressDurableDelivery(params.deliveryCompletion, platformQueueStateDir); - } - } + await settleDeliveryCompletion(results.at(-1)); if (!params.deferCommitHooks) { flushMessageSentEvents(); await runOutboundDeliveryCommitHooks(results); @@ -365,17 +365,6 @@ export async function deliverOutboundPayloadsWithQueueCleanup( ); } } else { - if (params.deliveryCompletion) { - if (results.length > 0) { - await completeDurableDelivery( - params.deliveryCompletion, - results.at(-1)!, - platformQueueStateDir, - ); - } else { - await suppressDurableDelivery(params.deliveryCompletion, platformQueueStateDir); - } - } const postSendState = queuedPostSendState ?? (results.length > 0 || queuedPreSendState === "marked" @@ -383,6 +372,7 @@ export async function deliverOutboundPayloadsWithQueueCleanup( : queuedPreSendState === "acked" ? "acked" : undefined); + await settleDeliveryCompletion(results.at(-1)); if (results.length === 0 && postSendState === "marked") { // The provider was invoked but returned no recipient-visible identity; // never convert that ambiguous platform outcome into a success receipt. diff --git a/src/infra/outbound/deliver.test.ts b/src/infra/outbound/deliver.test.ts index d23f029a8dfc..a77089c2c898 100644 --- a/src/infra/outbound/deliver.test.ts +++ b/src/infra/outbound/deliver.test.ts @@ -96,6 +96,7 @@ const queueMocks = vi.hoisted(() => ({ })); const completionMocks = vi.hoisted(() => ({ completeDurableDelivery: vi.fn(), + failDurableDelivery: vi.fn(), markDurableDeliveryQueued: vi.fn(async () => ({ state: "queued" as const })), rejectDurableDelivery: vi.fn(), suppressDurableDelivery: vi.fn(), @@ -196,6 +197,16 @@ vi.mock("./delivery-completion.js", () => ({ markDurableDeliveryQueued: completionMocks.markDurableDeliveryQueued, rejectDurableDelivery: completionMocks.rejectDurableDelivery, suppressDurableDelivery: completionMocks.suppressDurableDelivery, + settleDurableDelivery: ( + completion: unknown, + evidence: { result: unknown } | { platformSendStarted: boolean }, + stateDir?: string, + ) => + "result" in evidence + ? completionMocks.completeDurableDelivery(completion, evidence.result, stateDir) + : evidence.platformSendStarted + ? completionMocks.failDurableDelivery(completion, stateDir) + : completionMocks.suppressDurableDelivery(completion, stateDir), })); vi.mock("../../logging/subsystem.js", () => ({ createSubsystemLogger: () => { @@ -508,6 +519,7 @@ describe("deliverOutboundPayloads", () => { }, ); completionMocks.completeDurableDelivery.mockClear(); + completionMocks.failDurableDelivery.mockClear(); completionMocks.markDurableDeliveryQueued.mockClear(); completionMocks.rejectDurableDelivery.mockClear(); completionMocks.suppressDurableDelivery.mockClear(); diff --git a/src/infra/outbound/delivery-completion.integration.test.ts b/src/infra/outbound/delivery-completion.integration.test.ts index dd4325ca9a4f..30e5262abdec 100644 --- a/src/infra/outbound/delivery-completion.integration.test.ts +++ b/src/infra/outbound/delivery-completion.integration.test.ts @@ -88,4 +88,65 @@ describe("pending-final durable delivery completion", () => { expect(sendMatrix).toHaveBeenCalledOnce(); expect(await loadPendingDeliveries(tmpDir)).toEqual([]); }); + + it("keeps an uncertainty notice owed when a live send returns no delivery identity", async () => { + process.env.OPENCLAW_STATE_DIR = tmpDir; + const sessionKey = "agent:main:matrix:direct:unknown-live"; + const storePath = path.join(tmpDir, "sessions.json"); + const deliveryId = "pending-final-unknown-live"; + const completion = { + kind: "pending-final" as const, + deliveryId, + intentId: "pending-final-intent-unknown-live", + sessionId: "session-unknown-live", + sessionKey, + storePath, + }; + const context = { channel: "matrix", to: "!room:example" }; + await replaceSessionEntry( + { sessionKey, storePath }, + { + sessionId: completion.sessionId, + status: "running", + updatedAt: Date.now(), + pendingFinalDelivery: { + kind: "replayable", + text: "delivery identity may have been lost", + context, + createdAt: Date.now(), + intentId: completion.intentId, + deliveries: [{ id: deliveryId, state: "prepared" }], + }, + }, + ); + const sendMatrix = vi.fn().mockResolvedValue({}); + + await expect( + deliverOutboundPayloads({ + cfg: {} as OpenClawConfig, + channel: "matrix", + to: "!room:example", + payloads: [{ text: "delivery identity may have been lost" }], + deps: { matrix: sendMatrix }, + queuePolicy: "required", + deliveryIntentId: deliveryId, + deliveryCompletion: completion, + }), + ).resolves.toEqual([]); + + expect((await loadPendingDeliveries(tmpDir))[0]).toMatchObject({ + id: deliveryId, + recoveryState: "unknown_after_send", + }); + expect(loadSessionEntry({ sessionKey, storePath })).toMatchObject({ + pendingFinalDelivery: { + deliveries: [{ id: deliveryId, state: "unknown" }], + }, + pendingDeliveryNotice: { + intentId: completion.intentId, + state: "owed", + context, + }, + }); + }); }); diff --git a/src/infra/outbound/delivery-completion.ts b/src/infra/outbound/delivery-completion.ts index 78a13b0241db..9e3a4303b810 100644 --- a/src/infra/outbound/delivery-completion.ts +++ b/src/infra/outbound/delivery-completion.ts @@ -203,7 +203,7 @@ export async function completeDurableDelivery( } /** Finalizes a policy-suppressed send before its durable intent is acknowledged. */ -export async function suppressDurableDelivery( +async function suppressDurableDelivery( completion: DurableDeliveryCompletion, stateDir?: string, ): Promise { @@ -244,3 +244,20 @@ export async function failDurableDelivery( markConversationDeliveryUnknown(scopeForCompletion(completion), completion.operationId), ); } + +type DurableDeliveryTerminalEvidence = + | { result: OutboundDeliveryResult } + | { platformSendStarted: boolean }; + +/** Settles the completion owner from the final evidence held by its lifecycle owner. */ +export async function settleDurableDelivery( + completion: DurableDeliveryCompletion, + evidence: DurableDeliveryTerminalEvidence, + stateDir?: string, +): Promise { + return "result" in evidence + ? completeDurableDelivery(completion, evidence.result, stateDir) + : evidence.platformSendStarted + ? failDurableDelivery(completion, stateDir) + : suppressDurableDelivery(completion, stateDir); +} diff --git a/src/infra/outbound/delivery-queue-recovery.ts b/src/infra/outbound/delivery-queue-recovery.ts index f8caef65db17..b94a6638b660 100644 --- a/src/infra/outbound/delivery-queue-recovery.ts +++ b/src/infra/outbound/delivery-queue-recovery.ts @@ -21,6 +21,7 @@ import { import { formatErrorMessage } from "../errors.js"; import { resolveOutboundChannelMessageAdapter } from "./channel-resolution.js"; import { resolveDeferredDeliveryAdmission } from "./deferred-delivery-admission.js"; +import type { DeliverOutboundPayloadsParams } from "./deliver-contracts.js"; import { OUTBOUND_DELIVERY_LOG_SCOPE } from "./deliver-log.js"; import { buildPayloadSummary } from "./deliver-payload.js"; import { @@ -42,7 +43,7 @@ import { failDurableDelivery, markDurableDeliveryQueued, rejectDurableDelivery, - suppressDurableDelivery, + settleDurableDelivery, } from "./delivery-completion.js"; import { collectEntrySpoolPaths, releaseSpoolArtifacts } from "./delivery-queue-media-spool.js"; import { @@ -64,7 +65,6 @@ import { moveToFailed, reserveDeliveryAttempt, type QueuedDelivery, - type QueuedDeliveryPayload, } from "./delivery-queue-storage.js"; import { createMessageSentEmitter, type MessageSentEvent } from "./message-sent-hook.js"; import { @@ -75,23 +75,7 @@ import { } from "./outbound-audit.js"; import { acceptedPreparedOutboundEntries } from "./prepared-batch.js"; -export type DeliverFn = ( - params: { - cfg: OpenClawConfig; - } & QueuedDeliveryPayload & { - payloads: ReturnType; - deliveryQueueId?: string; - deliveryQueueStateDir?: string; - deliveryProducerClaimId?: string; - deliveryProducerLeaseRequired?: boolean; - skipQueue?: boolean; - deferredDeliveryAdmissionPassed?: true; - deferCommitHooks?: boolean; - onMessageSentEvent?: (event: MessageSentEvent, sourceIndex: number) => void; - onPayloadDeliveryOutcome?: (outcome: OutboundPayloadDeliveryOutcome) => void; - onDeliveryResult?: (result: OutboundDeliveryResult) => Promise | void; - }, -) => Promise; +export type DeliverFn = (params: DeliverOutboundPayloadsParams) => Promise; export interface RecoveryLogger { info(msg: string): void; @@ -326,7 +310,8 @@ function buildRecoveryDeliverParams( session: entry.session, gatewayClientScopes: entry.gatewayClientScopes, preparedMessageId: entry.preparedMessageId, - deliveryCompletion: entry.deliveryCompletion, + // Recovery owns terminal completion because nested delivery only reports + // process-local evidence that cannot survive another restart. deliveryQueueId: entry.id, deliveryQueueStateDir: stateDir, ...(producerClaimId ? { deliveryProducerClaimId: producerClaimId } : {}), @@ -884,6 +869,7 @@ async function drainQueuedEntry(opts: { // persisting plugin callbacks must never become part of delivery custody. const messageSentEvents: IndexedMessageSentEvent[] = []; let postSendState: QueuedPostSendState | undefined; + let platformSendStarted = false; let deliveredResults: OutboundDeliveryResult[] = []; let commitHooksRun = false; const collectResults = (results: readonly OutboundDeliveryResult[]): void => { @@ -973,6 +959,9 @@ async function drainQueuedEntry(opts: { ...buildRecoveryDeliverParams(entry, opts.cfg, opts.stateDir, producerClaimId), onPayloadDeliveryOutcome: collectPayloadOutcome, onMessageSentEvent: (event, sourceIndex) => messageSentEvents.push({ sourceIndex, event }), + onPlatformSendStart: async () => { + platformSendStarted = true; + }, onDeliveryResult: async (deliveryResult) => { collectResults([deliveryResult]); postSendState ??= await persistRecoveredPostSendState({ @@ -984,13 +973,11 @@ async function drainQueuedEntry(opts: { }, }); const results = isOutboundDeliveryResultArray(result) ? result : []; - if ( - producerClaimId !== undefined && - payloadOutcomes.some( - (outcome) => - outcome.status === "suppressed" && outcome.reason === "adapter_returned_no_identity", - ) - ) { + const adapterReturnedNoIdentity = payloadOutcomes.some( + (outcome) => + outcome.status === "suppressed" && outcome.reason === "adapter_returned_no_identity", + ); + if (adapterReturnedNoIdentity || (results.length === 0 && platformSendStarted)) { const error = "recovered platform send returned no delivery identity"; await recordRecoveredFailure( failDeliveryAfterPlatformSend, @@ -999,6 +986,13 @@ async function drainQueuedEntry(opts: { opts.stateDir, producerClaimId, ); + if (entry.deliveryCompletion) { + await settleDurableDelivery( + entry.deliveryCompletion, + { platformSendStarted: true }, + opts.stateDir, + ); + } opts.onFailed?.(entry, error); opts.log.warn(`Delivery entry ${entry.id} ${error}; preserving unknown_after_send`); emitQueuedAuditTerminals(entry, () => queuedUnknownAuditTerminals(entry)); @@ -1044,11 +1038,12 @@ async function drainQueuedEntry(opts: { return "failed"; } if (entry.deliveryCompletion) { - if (results.length > 0) { - await completeDurableDelivery(entry.deliveryCompletion, results.at(-1)!, opts.stateDir); - } else { - await suppressDurableDelivery(entry.deliveryCompletion, opts.stateDir); - } + const terminalResult = results.at(-1); + await settleDurableDelivery( + entry.deliveryCompletion, + terminalResult ? { result: terminalResult } : { platformSendStarted: false }, + opts.stateDir, + ); } postSendState ??= results.length > 0 diff --git a/src/infra/outbound/delivery-queue-storage.ts b/src/infra/outbound/delivery-queue-storage.ts index 0bf0fb0b6d49..b6316e21a140 100644 --- a/src/infra/outbound/delivery-queue-storage.ts +++ b/src/infra/outbound/delivery-queue-storage.ts @@ -55,7 +55,6 @@ export type { LegacyQueuedDelivery, LegacyQueuedDeliveryPreparation, QueuedDelivery, - QueuedDeliveryPayload, QueuedReplyPayloadSendingHook, QueuedRenderedMessageBatchPlan, } from "./delivery-queue-types.js"; diff --git a/src/infra/outbound/delivery-queue.recovery.test.ts b/src/infra/outbound/delivery-queue.recovery.test.ts index 3dc493e7b2c0..542a13388c12 100644 --- a/src/infra/outbound/delivery-queue.recovery.test.ts +++ b/src/infra/outbound/delivery-queue.recovery.test.ts @@ -11,7 +11,11 @@ import { markConversationDeliveryRejected, markConversationDeliverySuppressed, } from "../../config/sessions/conversation-delivery-store.js"; -import { upsertSessionEntryCore } from "../../config/sessions/session-accessor.js"; +import { + loadSessionEntry, + replaceSessionEntry, + upsertSessionEntryCore, +} from "../../config/sessions/session-accessor.js"; import { buildConversationRef } from "../../routing/conversation-ref.js"; import { createDeferredCore } from "../../shared/deferred.js"; import { closeOpenClawAgentDatabasesForTest } from "../../state/openclaw-agent-db.js"; @@ -354,6 +358,47 @@ describe("delivery-queue recovery", () => { ); return scope; } + async function createPendingFinalRecoveryFixture(deliveryId: string) { + const sessionKey = "agent:main:demo-channel-a:direct:pending-final"; + const storePath = path.join(tmpDir(), "pending-final-sessions.json"); + const completion = { + kind: "pending-final" as const, + deliveryId, + intentId: "pending-final-recovery-intent", + sessionId: "pending-final-recovery-session", + sessionKey, + storePath, + }; + const context = { channel: "demo-channel-a", to: "+1" }; + await replaceSessionEntry( + { sessionKey, storePath }, + { + sessionId: completion.sessionId, + status: "running", + updatedAt: Date.now(), + pendingFinalDelivery: { + kind: "replayable", + text: "recovered delivery identity may have been lost", + context, + createdAt: Date.now(), + intentId: completion.intentId, + deliveries: [{ id: deliveryId, state: "prepared" }], + }, + }, + ); + await enqueueDeliveryOnce( + { + channel: "demo-channel-a", + to: "+1", + queuePolicy: "required", + payloads: [{ text: "recovered delivery identity may have been lost" }], + deliveryCompletion: completion, + }, + deliveryId, + tmpDir(), + ); + return { completion, context }; + } it("recovers entries from a simulated crash", async () => { await enqueueCrashRecoveryEntries(); const deliver = vi.fn().mockResolvedValue([]); @@ -392,6 +437,37 @@ describe("delivery-queue recovery", () => { closeOpenClawAgentDatabasesForTest(); } }); + it("keeps an uncertainty notice owed when recovery returns no delivery identity", async () => { + const deliveryId = "pending-final-unknown-recovery"; + const { completion, context } = await createPendingFinalRecoveryFixture(deliveryId); + const deliver = vi.fn(async (params: Parameters[0]) => { + expect(params.deliveryCompletion).toBeUndefined(); + await markDeliveryPlatformSendAttemptStarted(deliveryId, tmpDir()); + await params.onPlatformSendStart?.({}); + return []; + }); + + const { result } = await runRecovery({ deliver }); + + expect( + loadSessionEntry({ sessionKey: completion.sessionKey, storePath: completion.storePath }), + ).toMatchObject({ + pendingFinalDelivery: { + deliveries: [{ id: deliveryId, state: "unknown" }], + }, + pendingDeliveryNotice: { + intentId: completion.intentId, + state: "owed", + context, + }, + }); + expect(result).toEqual(RECOVERY_SUMMARY.failed); + await expectPendingEntry({ + id: deliveryId, + recoveryState: "unknown_after_send", + retryCount: 1, + }); + }); it.each([ "acks a persisted suppressed conversation operation without replaying it", "acks a persisted rejected conversation operation without replaying it", diff --git a/src/tts/tts-runtime-fallbacks.test.ts b/src/tts/tts-runtime-fallbacks.test.ts index 6befb97e8b6f..728c08daa4da 100644 --- a/src/tts/tts-runtime-fallbacks.test.ts +++ b/src/tts/tts-runtime-fallbacks.test.ts @@ -37,12 +37,16 @@ import { const routedPayloads = vi.hoisted(() => [] as ReplyPayload[]); -vi.mock("../channels/message/runtime.js", () => ({ - sendDurableMessageBatchCore: async ({ payloads }: { payloads: ReplyPayload[] }) => { - routedPayloads.push(...payloads); - return { status: "sent", results: [{ messageId: "tts-route-1" }] }; - }, -})); +vi.mock("../channels/message/runtime.js", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + sendDurableMessageBatchCore: async ({ payloads }: { payloads: ReplyPayload[] }) => { + routedPayloads.push(...payloads); + return { status: "sent", results: [{ messageId: "tts-route-1" }] }; + }, + }; +}); function installStructuredReplyTestChannel(loaded: boolean): () => void { const previousRegistry = captureActivePluginRegistrySnapshot();