diff --git a/extensions/line/src/send-retry.test.ts b/extensions/line/src/send-retry.test.ts new file mode 100644 index 000000000000..0a9705aba2a5 --- /dev/null +++ b/extensions/line/src/send-retry.test.ts @@ -0,0 +1,202 @@ +// Line tests cover push retry and retry-key deduplication behavior. +import { HTTPFetchError } from "@line/bot-sdk"; +import { isChannelPartialDeliveryError } from "openclaw/plugin-sdk/channel-inbound"; +import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; +import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; + +const { + requireRuntimeConfigMock, + resolveLineAccountMock, + resolveLineChannelAccessTokenMock, + recordChannelActivityMock, + logVerboseMock, +} = vi.hoisted(() => ({ + requireRuntimeConfigMock: vi.fn((cfg: unknown) => cfg ?? {}), + resolveLineAccountMock: vi.fn(() => ({ accountId: "default" })), + resolveLineChannelAccessTokenMock: vi.fn(() => "test-token-placeholder"), + recordChannelActivityMock: vi.fn(), + logVerboseMock: vi.fn(), +})); + +vi.mock("openclaw/plugin-sdk/plugin-config-runtime", () => ({ + requireRuntimeConfig: requireRuntimeConfigMock, +})); + +vi.mock("./accounts.js", () => ({ + resolveLineAccount: resolveLineAccountMock, +})); + +vi.mock("./channel-access-token.js", () => ({ + resolveLineChannelAccessToken: resolveLineChannelAccessTokenMock, +})); + +vi.mock("openclaw/plugin-sdk/channel-activity-runtime", () => ({ + recordChannelActivity: recordChannelActivityMock, +})); + +vi.mock("openclaw/plugin-sdk/runtime-env", async () => { + const actual = await vi.importActual( + "openclaw/plugin-sdk/runtime-env", + ); + return { ...actual, logVerbose: logVerboseMock }; +}); + +let sendModule: typeof import("./send.js"); + +const LINE_TEST_CFG = { + channels: { line: { accounts: { default: {} } } }, +} satisfies OpenClawConfig; +const LINE_TARGET = "line:user:U0123456789abcdef0123456789abcdef"; +const UUID_PATTERN = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/; + +function jsonResponse(body: unknown, status = 200): Response { + return new Response(JSON.stringify(body), { + status, + headers: { "content-type": "application/json" }, + }); +} + +function transportFailure(code: string): TypeError { + return Object.assign(new TypeError("fetch failed"), { + cause: Object.assign(new Error(`socket ${code}`), { code }), + }); +} + +function retryKeysOf(fetchMock: ReturnType>): (string | null)[] { + return fetchMock.mock.calls.map(([, init]) => new Headers(init?.headers).get("X-Line-Retry-Key")); +} + +describe("LINE push retries", () => { + const fetchMock = vi.fn(); + + beforeAll(async () => { + sendModule = await import("./send.js"); + }); + + afterAll(() => { + vi.doUnmock("openclaw/plugin-sdk/plugin-config-runtime"); + vi.doUnmock("./accounts.js"); + vi.doUnmock("./channel-access-token.js"); + vi.doUnmock("openclaw/plugin-sdk/channel-activity-runtime"); + vi.doUnmock("openclaw/plugin-sdk/runtime-env"); + vi.resetModules(); + }); + + beforeEach(() => { + vi.clearAllMocks(); + fetchMock.mockReset(); + requireRuntimeConfigMock.mockImplementation((cfg: unknown) => cfg ?? LINE_TEST_CFG); + resolveLineAccountMock.mockReturnValue({ accountId: "default" }); + resolveLineChannelAccessTokenMock.mockReturnValue("test-token-placeholder"); + vi.stubGlobal("fetch", fetchMock); + vi.useFakeTimers(); + }); + + afterEach(async () => { + await vi.runOnlyPendingTimersAsync(); + vi.useRealTimers(); + vi.unstubAllGlobals(); + }); + + async function resolveRetryRun(run: Promise): Promise { + run.catch(() => {}); + await vi.runAllTimersAsync(); + return await run; + } + + function pushText(text = "hello") { + return sendModule.pushMessagesLine(LINE_TARGET, [{ type: "text", text }], { + cfg: LINE_TEST_CFG, + }); + } + + it("retries a LINE server error under one retry key and delivers once", async () => { + fetchMock + .mockResolvedValueOnce(jsonResponse({ message: "Internal server error" }, 500)) + .mockResolvedValueOnce(jsonResponse({ sentMessages: [{ id: "delivered-1" }] })); + + const result = await resolveRetryRun(pushText()); + + expect(result.messageId).toBe("delivered-1"); + expect(fetchMock).toHaveBeenCalledTimes(2); + const retryKeys = retryKeysOf(fetchMock); + expect(retryKeys[0]).toMatch(UUID_PATTERN); + expect(retryKeys[1]).toBe(retryKeys[0]); + expect(recordChannelActivityMock).toHaveBeenCalledTimes(1); + }); + + it("keys each push separately so an unrelated send cannot be deduplicated away", async () => { + fetchMock.mockImplementation(async () => + jsonResponse({ sentMessages: [{ id: "delivered-1" }] }), + ); + + await resolveRetryRun(pushText("first")); + await resolveRetryRun(pushText("second")); + + const [firstKey, secondKey] = retryKeysOf(fetchMock); + expect(firstKey).toMatch(UUID_PATTERN); + expect(secondKey).toMatch(UUID_PATTERN); + expect(secondKey).not.toBe(firstKey); + }); + + it("retries a transport failure and keeps the accepted delivery when LINE reports a conflict", async () => { + fetchMock.mockRejectedValueOnce(transportFailure("ETIMEDOUT")).mockResolvedValueOnce( + jsonResponse( + { + message: "The retry key is already accepted", + sentMessages: [{ id: "accepted-earlier" }], + }, + 409, + ), + ); + + const result = await resolveRetryRun(pushText()); + + // The first attempt landed even though its outcome never reached us, so the + // accepted request's message id is the delivery — not a second send. + expect(result.messageId).toBe("accepted-earlier"); + expect(fetchMock).toHaveBeenCalledTimes(2); + expect(new Set(retryKeysOf(fetchMock)).size).toBe(1); + }); + + it("gives up after the configured attempts and surfaces the LINE failure", async () => { + fetchMock.mockImplementation(async () => + jsonResponse({ message: "Internal server error" }, 500), + ); + + await expect(resolveRetryRun(pushText())).rejects.toMatchObject({ status: 500 }); + expect(fetchMock).toHaveBeenCalledTimes(3); + expect(new Set(retryKeysOf(fetchMock)).size).toBe(1); + }); + + it.each([ + { label: "quota rejection", status: 429, message: "You have reached your monthly limit." }, + { label: "request rejection", status: 400, message: "The request body has 1 error(s)" }, + ])("does not retry a LINE $label", async ({ status, message }) => { + fetchMock.mockResolvedValue(jsonResponse({ message }, status)); + + await expect(resolveRetryRun(pushText())).rejects.toBeInstanceOf(HTTPFetchError); + expect(fetchMock).toHaveBeenCalledTimes(1); + }); + + it("does not retry once LINE accepted a request with an unreadable receipt", async () => { + fetchMock.mockResolvedValue(jsonResponse({ sentMessages: [{}] })); + + await expect(resolveRetryRun(pushText())).rejects.toSatisfy(isChannelPartialDeliveryError); + expect(fetchMock).toHaveBeenCalledTimes(1); + }); + + it("never retries a reply, which LINE cannot deduplicate", async () => { + fetchMock.mockResolvedValue(jsonResponse({ message: "Internal server error" }, 500)); + + await expect( + resolveRetryRun( + sendModule.replyMessageLine("reply-token", [{ type: "text", text: "hello" }], { + cfg: LINE_TEST_CFG, + }), + ), + ).rejects.toMatchObject({ status: 500 }); + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(retryKeysOf(fetchMock)).toEqual([null]); + }); +}); diff --git a/extensions/line/src/send-retry.ts b/extensions/line/src/send-retry.ts new file mode 100644 index 000000000000..e8b7300425f1 --- /dev/null +++ b/extensions/line/src/send-retry.ts @@ -0,0 +1,37 @@ +// Line plugin module implements push retry policy behavior. +import { HTTPFetchError } from "@line/bot-sdk"; +import { collectErrorGraphCandidates, extractErrorCode } from "openclaw/plugin-sdk/error-runtime"; +import { + classifyTransientNetworkErrorCode, + createChannelApiRetryRunner, +} from "openclaw/plugin-sdk/retry-runtime"; + +function isRetryableLinePushError(error: unknown): boolean { + const candidates = collectErrorGraphCandidates(error, (candidate) => [ + candidate.cause, + candidate.error, + ]); + const httpError = candidates.find( + (candidate): candidate is HTTPFetchError => candidate instanceof HTTPFetchError, + ); + if (httpError) { + // LINE documents server errors and transport failures as the retriable + // outcomes; every 4xx (429 included) answers "retries don't change the result". + return httpError.status >= 500; + } + // A transport failure never reached a LINE response, so the retry key decides + // whether the earlier attempt already landed. + return candidates.some( + (candidate) => classifyTransientNetworkErrorCode(extractErrorCode(candidate)) !== undefined, + ); +} + +/** + * Pushes are non-idempotent without a retry key, so the generic message-matching + * fallback stays off and only the classification above may replay a request. + */ +export const runLinePushWithRetries = createChannelApiRetryRunner({ + shouldRetry: isRetryableLinePushError, + strictShouldRetry: true, + verbose: true, +}); diff --git a/extensions/line/src/send.ts b/extensions/line/src/send.ts index 50ff68dc3360..a38806c9d300 100644 --- a/extensions/line/src/send.ts +++ b/extensions/line/src/send.ts @@ -1,4 +1,5 @@ // Line plugin module implements send behavior. +import { randomUUID } from "node:crypto"; import { HTTPFetchError, messagingApi } from "@line/bot-sdk"; import lineBotSdkPackage from "@line/bot-sdk/package.json" with { type: "json" }; import { recordChannelActivity } from "openclaw/plugin-sdk/channel-activity-runtime"; @@ -14,6 +15,7 @@ import { messageAction, normalizeLineMessageActions } from "./actions.js"; import { resolveLineChannelAccessToken } from "./channel-access-token.js"; import { validateLineMediaUrl } from "./outbound-media.js"; import { createLineSendReceipt } from "./send-receipt.js"; +import { runLinePushWithRetries } from "./send-retry.js"; import type { LineChannelData, LineOutboundMediaKind, LineSendResult } from "./types.js"; type Message = messagingApi.Message; @@ -172,6 +174,7 @@ async function sendLineProviderMessages( operation: "push" | "reply", token: string, request: messagingApi.PushMessageRequest | messagingApi.ReplyMessageRequest, + retryKey?: string, ): Promise { const response = await fetchWithRuntimeDispatcherOrMockedGlobal( `https://api.line.me/v2/bot/message/${operation}`, @@ -181,12 +184,18 @@ async function sendLineProviderMessages( "Content-Type": "application/json", Authorization: `Bearer ${token}`, "User-Agent": `@line/bot-sdk/${lineBotSdkPackage.version}`, + ...(retryKey ? { "X-Line-Retry-Key": retryKey } : {}), }, body: JSON.stringify(request), }, ); - if (!response.ok) { + // LINE answers a retried key with 409 and the accepted request's sent messages + // instead of delivering the batch a second time, so that conflict is the + // earlier attempt's success rather than a failure of this one. + const acceptedRetryConflict = retryKey !== undefined && response.status === 409; + + if (!response.ok && !acceptedRetryConflict) { throw new HTTPFetchError(`${response.status} - ${response.statusText}`, { status: response.status, statusText: response.statusText, @@ -325,17 +334,25 @@ async function pushLineMessages( const { account, token, chatId } = createLinePushContext(to, opts); const normalizedMessages = messages.map(normalizeLineMessageActions); - const pushRequest = sendLineProviderMessages("push", token, { - to: chatId, - messages: normalizedMessages, - }); + // One retry key per logical push: every attempt reuses it so LINE deduplicates + // an attempt that was accepted before its outcome reached us. + const retryKey = randomUUID(); - const response = behavior.errorContext - ? await pushRequest.catch((err: unknown) => { - logLineHttpError(err, behavior.errorContext!); - throw err; - }) - : await pushRequest; + const response = await runLinePushWithRetries(async () => { + try { + return await sendLineProviderMessages( + "push", + token, + { to: chatId, messages: normalizedMessages }, + retryKey, + ); + } catch (err) { + if (behavior.errorContext) { + logLineHttpError(err, behavior.errorContext); + } + throw err; + } + }, "line:push"); const { messageId, messageIds } = resolveLineProviderMessageIds(response, "push"); const result: LineSendResult = { messageId,