Skip to content

Commit ab09620

Browse files
committed
[AI] fix(outbound): clear recoveryState on connect-phase errors so drain can retry (#100979)
When a send fails with a connect-phase network error (UND_ERR_CONNECT_TIMEOUT, ECONNREFUSED, ENOTFOUND, EAI_AGAIN, ENETUNREACH), the request body never reached the platform. However, failDelivery() does not clear recoveryState, so the entry remains in send_attempt_started state. On gateway reconnect, the drain refuses blind replay and permanently drops the message. Add clearDeliveryRecoveryState() to the storage layer and a connect-phase error detection + cleanup branch in deliver.ts, so these entries become eligible for normal retry/replay. Non-connect-phase errors (ECONNRESET, ETIMEDOUT) keep the existing guard. Co-Authored-By: Claude Opus 4.8 <[email protected]> Related to #100979
1 parent 96e676c commit ab09620

5 files changed

Lines changed: 143 additions & 0 deletions

File tree

src/infra/outbound/deliver.test.ts

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,7 @@ const internalHookMocks = vi.hoisted(() => ({
5555
const queueMocks = vi.hoisted(() => ({
5656
enqueueDelivery: vi.fn(async () => "mock-queue-id"),
5757
ackDelivery: vi.fn(async () => {}),
58+
clearDeliveryRecoveryState: vi.fn(async () => {}),
5859
failDelivery: vi.fn(async () => {}),
5960
failDeliveryAfterPlatformSend: vi.fn(async () => {}),
6061
markDeliveryPlatformOutcomeUnknown: vi.fn(async () => {}),
@@ -99,6 +100,7 @@ vi.mock("../../hooks/internal-hooks.js", () => ({
99100
vi.mock("./delivery-queue.js", () => ({
100101
enqueueDelivery: queueMocks.enqueueDelivery,
101102
ackDelivery: queueMocks.ackDelivery,
103+
clearDeliveryRecoveryState: queueMocks.clearDeliveryRecoveryState,
102104
failDelivery: queueMocks.failDelivery,
103105
failDeliveryAfterPlatformSend: queueMocks.failDeliveryAfterPlatformSend,
104106
markDeliveryPlatformOutcomeUnknown: queueMocks.markDeliveryPlatformOutcomeUnknown,
@@ -332,6 +334,8 @@ describe("deliverOutboundPayloads", () => {
332334
queueMocks.markDeliveryPlatformSendAttemptStarted.mockResolvedValue(undefined);
333335
queueMocks.markDeliveryPlatformSendDispatched.mockClear();
334336
queueMocks.markDeliveryPlatformSendDispatched.mockResolvedValue(undefined);
337+
queueMocks.clearDeliveryRecoveryState.mockClear();
338+
queueMocks.clearDeliveryRecoveryState.mockResolvedValue(undefined);
335339
queueMocks.withActiveDeliveryClaim.mockClear();
336340
queueMocks.withActiveDeliveryClaim.mockImplementation(async (_entryId, fn) => ({
337341
status: "claimed",
@@ -1497,6 +1501,50 @@ describe("deliverOutboundPayloads", () => {
14971501
expect(queueMocks.ackDelivery).not.toHaveBeenCalled();
14981502
});
14991503

1504+
it("clears recoveryState on connect-phase network error so drain can retry (#100979)", async () => {
1505+
const connectErr = Object.assign(new Error("connect refused"), { code: "ECONNREFUSED" });
1506+
const sendMatrix = vi.fn().mockRejectedValueOnce(connectErr);
1507+
1508+
await expect(
1509+
deliverOutboundPayloads({
1510+
cfg: {},
1511+
channel: "matrix",
1512+
to: "!room:example",
1513+
payloads: [{ text: "hello" }],
1514+
deps: { matrix: sendMatrix },
1515+
queuePolicy: "required",
1516+
}),
1517+
).rejects.toThrow("connect refused");
1518+
1519+
expect(queueMocks.clearDeliveryRecoveryState).toHaveBeenCalledWith("mock-queue-id", undefined);
1520+
expect(queueMocks.failDelivery).toHaveBeenCalledWith(
1521+
"mock-queue-id",
1522+
expect.stringContaining("connect refused"),
1523+
);
1524+
});
1525+
1526+
it("does NOT clear recoveryState on non-connect-phase error (ECONNRESET) (#100979)", async () => {
1527+
const resetErr = Object.assign(new Error("connection reset"), { code: "ECONNRESET" });
1528+
const sendMatrix = vi.fn().mockRejectedValueOnce(resetErr);
1529+
1530+
await expect(
1531+
deliverOutboundPayloads({
1532+
cfg: {},
1533+
channel: "matrix",
1534+
to: "!room:example",
1535+
payloads: [{ text: "hello" }],
1536+
deps: { matrix: sendMatrix },
1537+
queuePolicy: "required",
1538+
}),
1539+
).rejects.toThrow("connection reset");
1540+
1541+
expect(queueMocks.clearDeliveryRecoveryState).not.toHaveBeenCalled();
1542+
expect(queueMocks.failDelivery).toHaveBeenCalledWith(
1543+
"mock-queue-id",
1544+
expect.stringContaining("connection reset"),
1545+
);
1546+
});
1547+
15001548
it("directly acks a sent delivery when the post-send unknown marker cannot be written", async () => {
15011549
queueMocks.markDeliveryPlatformOutcomeUnknown.mockRejectedValueOnce(
15021550
new Error("unknown marker offline"),

src/infra/outbound/deliver.ts

Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,7 @@ import {
6565
} from "./delivery-commit-hooks.js";
6666
import {
6767
ackDelivery,
68+
clearDeliveryRecoveryState,
6869
enqueueDelivery,
6970
failDelivery,
7071
failDeliveryAfterPlatformSend,
@@ -656,6 +657,48 @@ const isDeliveryAbortError = (err: unknown): boolean =>
656657
(err instanceof OutboundDeliveryError &&
657658
isAbortError((err as Error & { cause?: unknown }).cause));
658659

660+
/**
661+
* Error codes from the connect phase of HTTP/TCP that prove the request body
662+
* was never transmitted to the platform. These are distinct from errors that
663+
* can occur after partial data transfer (ECONNRESET, ETIMEDOUT), where the
664+
* platform may have committed the message.
665+
*
666+
* - UND_ERR_CONNECT_TIMEOUT: connection-level timeout, no bytes sent
667+
* - ECONNREFUSED: TCP handshake rejected
668+
* - ENOTFOUND / EAI_AGAIN: DNS resolution failure
669+
* - ENETUNREACH: network unreachable
670+
*/
671+
const CONNECT_PHASE_ERROR_CODES = new Set([
672+
"UND_ERR_CONNECT_TIMEOUT",
673+
"ECONNREFUSED",
674+
"ENOTFOUND",
675+
"EAI_AGAIN",
676+
"ENETUNREACH",
677+
]);
678+
679+
function getErrorCode(err: unknown): string | undefined {
680+
if (err && typeof err === "object" && "code" in err) {
681+
return (err as { code?: unknown }).code as string | undefined;
682+
}
683+
return undefined;
684+
}
685+
686+
/** Walks the error cause chain to find a connect-phase error code. */
687+
function isConnectPhaseError(err: unknown): boolean {
688+
let current: unknown = err;
689+
const seen = new Set<unknown>();
690+
while (current && typeof current === "object" && !seen.has(current)) {
691+
seen.add(current);
692+
const code = getErrorCode(current);
693+
if (code && CONNECT_PHASE_ERROR_CODES.has(code)) {
694+
return true;
695+
}
696+
// Walk cause chain (undici wraps errors with a cause property).
697+
current = (current as { cause?: unknown }).cause;
698+
}
699+
return false;
700+
}
701+
659702
type QueuedPostSendState = "marked" | "acked" | "failed";
660703

661704
type QueuedPreSendState = "marked" | "acked";
@@ -1610,6 +1653,23 @@ async function deliverOutboundPayloadsWithQueueCleanup(
16101653
);
16111654
}
16121655
await runCommitHooksAfterAck();
1656+
} else if (isConnectPhaseError(err)) {
1657+
// Connect-phase error: the HTTP request body never reached the
1658+
// platform. Clear recoveryState so the entry is eligible for normal
1659+
// retry/replay without requiring an adapter reconcileUnknownSend.
1660+
// (#100979)
1661+
await clearDeliveryRecoveryState(queueId, platformQueueStateDir).catch(
1662+
(clearErr: unknown) => {
1663+
log.warn(
1664+
`failed to clear delivery recovery state for ${queueId}: ${formatErrorMessage(clearErr)}`,
1665+
);
1666+
},
1667+
);
1668+
await failDelivery(queueId, formatErrorMessage(err)).catch((failErr: unknown) => {
1669+
log.warn(
1670+
`failed to mark queued delivery ${queueId} as failed: ${formatErrorMessage(failErr)}`,
1671+
);
1672+
});
16131673
} else {
16141674
await failDelivery(queueId, formatErrorMessage(err)).catch((failErr: unknown) => {
16151675
log.warn(

src/infra/outbound/delivery-queue-storage.ts

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -215,6 +215,20 @@ export async function markDeliveryPlatformOutcomeUnknown(
215215
}));
216216
}
217217

218+
/**
219+
* Clear the recovery state of a queued delivery entry so it becomes eligible
220+
* for normal retry / replay. Used when a send failed before any platform I/O
221+
* completed (connect-phase network error) — the request never reached the
222+
* channel, so the entry can safely be re-attempted without reconciliation.
223+
* Must NOT be called when send evidence exists (results or sentBeforeError).
224+
*/
225+
export async function clearDeliveryRecoveryState(id: string, stateDir?: string): Promise<void> {
226+
updateQueuedDelivery(id, stateDir, (entry) => ({
227+
...entry,
228+
recoveryState: undefined,
229+
}));
230+
}
231+
218232
/** Load a single pending delivery entry by ID from the queue directory. */
219233
export async function loadPendingDelivery(
220234
id: string,

src/infra/outbound/delivery-queue.storage.test.ts

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import { describe, expect, it, vi } from "vitest";
55
import { openOpenClawStateDatabase } from "../../state/openclaw-state-db.js";
66
import {
77
ackDelivery,
8+
clearDeliveryRecoveryState,
89
enqueueDelivery,
910
failDelivery,
1011
failDeliveryAfterPlatformSend,
@@ -228,6 +229,25 @@ describe("delivery-queue storage", () => {
228229
});
229230
});
230231

232+
describe("clearDeliveryRecoveryState", () => {
233+
it("clears recoveryState so the entry becomes eligible for normal retry", async () => {
234+
const id = await enqueueTextDelivery({
235+
channel: "forum",
236+
to: "123",
237+
payloads: [{ text: "hello" }],
238+
});
239+
240+
await markDeliveryPlatformSendAttemptStarted(id, tmpDir());
241+
let entry = readQueuedEntry(tmpDir(), id);
242+
expect(entry.recoveryState).toBe("send_attempt_started");
243+
244+
await clearDeliveryRecoveryState(id, tmpDir());
245+
246+
entry = readQueuedEntry(tmpDir(), id);
247+
expect(entry.recoveryState).toBeUndefined();
248+
});
249+
});
250+
231251
describe("moveToFailed", () => {
232252
it("moves entry to failed/ subdirectory", async () => {
233253
const id = await enqueueTextDelivery(

src/infra/outbound/delivery-queue.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
// Public outbound delivery queue facade for storage and recovery operations.
22
export {
33
ackDelivery,
4+
clearDeliveryRecoveryState,
45
enqueueDelivery,
56
failDelivery,
67
failDeliveryAfterPlatformSend,

0 commit comments

Comments
 (0)