From d49acffa5e1a696a145eb71534fd7eee70a71412 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sun, 12 Jul 2026 19:58:55 -0700 Subject: [PATCH] refactor(infra): consolidate retry scheduling (#105789) * refactor(infra): consolidate retry scheduling * fix(build): wire shared retry package * fix(build): resolve retry package in source graphs * fix(infra): keep retry adapter dependency-free --- .../.generated/plugin-sdk-api-baseline.sha256 | 4 +- extensions/qqbot/src/engine/api/retry.test.ts | 86 ++++++ extensions/qqbot/src/engine/api/retry.ts | 72 +++-- .../tsconfig.package-boundary.paths.json | 3 + extensions/xai/tsconfig.json | 3 + packages/memory-host-sdk/package.json | 3 +- .../src/host/batch-http.test.ts | 4 +- .../memory-host-sdk/src/host/batch-http.ts | 2 +- .../memory-host-sdk/src/host/embeddings.ts | 4 +- .../memory-host-sdk/src/host/read-retry.ts | 2 +- .../src/host/retry-utils.test.ts | 77 ----- .../memory-host-sdk/src/host/retry-utils.ts | 175 ----------- packages/plugin-sdk/tsconfig.json | 1 + packages/retry/package.json | 21 ++ packages/retry/src/index.test.ts | 75 +++++ packages/retry/src/index.ts | 274 +++++++++++++++++ pnpm-lock.yaml | 5 + scripts/lib/extension-package-boundary.ts | 1 + scripts/lib/tsdown-output-roots.mjs | 1 + ...e-extension-package-boundary-artifacts.mjs | 3 + scripts/run-node-watch-paths.mjs | 1 + scripts/write-plugin-sdk-entry-dts.ts | 1 + src/cron/isolated-agent/delivery-dispatch.ts | 43 +-- src/infra/clawhub-retry.test.ts | 13 + src/infra/clawhub-retry.ts | 67 +++-- src/infra/retry.ts | 284 ++---------------- src/media/store.ts | 36 ++- src/plugin-sdk/root-alias.cjs | 6 + ...tension-package-project-boundaries.test.ts | 1 + .../contracts/plugin-sdk-root-alias.test.ts | 3 + src/plugins/sdk-alias.test.ts | 20 ++ src/plugins/sdk-alias.ts | 8 + test/scripts/release-wrapper-scripts.test.ts | 2 +- test/vitest-scoped-config.test.ts | 4 + test/vitest/vitest.shared.config.ts | 1 + tsconfig.json | 1 + tsconfig.plugin-sdk.dts.json | 1 + tsdown.config.ts | 16 + 38 files changed, 706 insertions(+), 618 deletions(-) create mode 100644 extensions/qqbot/src/engine/api/retry.test.ts delete mode 100644 packages/memory-host-sdk/src/host/retry-utils.test.ts delete mode 100644 packages/memory-host-sdk/src/host/retry-utils.ts create mode 100644 packages/retry/package.json create mode 100644 packages/retry/src/index.test.ts create mode 100644 packages/retry/src/index.ts diff --git a/docs/.generated/plugin-sdk-api-baseline.sha256 b/docs/.generated/plugin-sdk-api-baseline.sha256 index ed79faaf060..35cb9c8d635 100644 --- a/docs/.generated/plugin-sdk-api-baseline.sha256 +++ b/docs/.generated/plugin-sdk-api-baseline.sha256 @@ -1,2 +1,2 @@ -75edbe496d57c19d3ea53e51908c6c7aa37ed8a5c7311234c7ce17eea0775db8 plugin-sdk-api-baseline.json -f21d330f28d93248bd6a7d96780452d6a969f5d1fcca1b9246ceb5a68e589b9f plugin-sdk-api-baseline.jsonl +e07c029f6f47cc853bbdcb4cdb8e042e81fb6a3ba82ecbacf9ce94b508038204 plugin-sdk-api-baseline.json +05a9890c8ee43b0c65ec24e820b16f89f2355e24372280198114ad0d30fcca24 plugin-sdk-api-baseline.jsonl diff --git a/extensions/qqbot/src/engine/api/retry.test.ts b/extensions/qqbot/src/engine/api/retry.test.ts new file mode 100644 index 00000000000..ea418c8fb6d --- /dev/null +++ b/extensions/qqbot/src/engine/api/retry.test.ts @@ -0,0 +1,86 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; +import type { EngineLogger } from "../types.js"; +import { withRetry } from "./retry.js"; + +const mocks = vi.hoisted(() => ({ + sleep: vi.fn(async () => {}), +})); + +vi.mock("openclaw/plugin-sdk/runtime-env", () => ({ sleep: mocks.sleep })); + +function createLogger(): EngineLogger { + return { + info: vi.fn(), + error: vi.fn(), + warn: vi.fn(), + debug: vi.fn(), + }; +} + +beforeEach(() => { + mocks.sleep.mockClear(); +}); + +describe("withRetry", () => { + it("uses the shared runner without changing exponential schedules", async () => { + const operation = vi + .fn<() => Promise>() + .mockRejectedValueOnce(new Error("first")) + .mockRejectedValueOnce(new Error("second")) + .mockResolvedValueOnce("ok"); + + await expect( + withRetry( + operation, + { maxRetries: 2, baseDelayMs: 100, backoff: "exponential" }, + undefined, + createLogger(), + ), + ).resolves.toBe("ok"); + expect(mocks.sleep).toHaveBeenNthCalledWith(1, 100); + expect(mocks.sleep).toHaveBeenNthCalledWith(2, 200); + }); + + it("preserves the policy's zero-based attempt index", async () => { + const shouldRetry = vi.fn(() => false); + await expect( + withRetry( + async () => { + throw new Error("stop"); + }, + { + maxRetries: 2, + baseDelayMs: 100, + backoff: "fixed", + shouldRetry, + }, + ), + ).rejects.toThrow("stop"); + expect(shouldRetry).toHaveBeenCalledWith(expect.any(Error), 0); + expect(mocks.sleep).not.toHaveBeenCalled(); + }); + + it("does not restart a persistent loop after its terminal failure", async () => { + const persistentTrigger = Object.assign(new Error("processing"), { bizCode: 42 }); + const terminal = new Error("permission denied"); + const operation = vi + .fn<() => Promise>() + .mockRejectedValueOnce(persistentTrigger) + .mockRejectedValueOnce(terminal); + + await expect( + withRetry( + operation, + { maxRetries: 2, baseDelayMs: 100, backoff: "fixed" }, + { + timeoutMs: 1_000, + intervalMs: 10, + shouldPersistRetry: (error) => + "bizCode" in error && (error as { bizCode?: number }).bizCode === 42, + }, + ), + ).rejects.toBe(terminal); + expect(operation).toHaveBeenCalledTimes(2); + expect(mocks.sleep).not.toHaveBeenCalled(); + }); +}); diff --git a/extensions/qqbot/src/engine/api/retry.ts b/extensions/qqbot/src/engine/api/retry.ts index e2ed50a5d73..d8a51ee805a 100644 --- a/extensions/qqbot/src/engine/api/retry.ts +++ b/extensions/qqbot/src/engine/api/retry.ts @@ -10,6 +10,7 @@ * parameterized by `RetryPolicy` and optional `PersistentRetryPolicy`. */ +import { retryAsync } from "openclaw/plugin-sdk/retry-runtime"; import { sleep } from "openclaw/plugin-sdk/runtime-env"; import { truncateUtf16Safe } from "openclaw/plugin-sdk/text-utility-runtime"; import type { EngineLogger } from "../types.js"; @@ -62,41 +63,54 @@ export async function withRetry( persistentPolicy?: PersistentRetryPolicy, logger?: EngineLogger, ): Promise { - let lastError: Error | null = null; - - for (let attempt = 0; attempt <= policy.maxRetries; attempt++) { - try { - return await fn(); - } catch (err) { - lastError = err instanceof Error ? err : new Error(formatErrorMessage(err)); - - // Check for persistent-retry trigger before standard retry logic. - if (persistentPolicy?.shouldPersistRetry(lastError)) { + // A persistent loop owns its terminal failure. Mark that Error so the outer + // bounded runner does not accidentally restart the completed deadline loop. + const persistentFailures = new WeakSet(); + return await retryAsync( + async () => { + try { + return await fn(); + } catch (err) { + const error = err instanceof Error ? err : new Error(formatErrorMessage(err)); + if (!persistentPolicy?.shouldPersistRetry(error)) { + throw error; + } (logger?.warn ?? logger?.error)?.( `[qqbot:retry] Hit persistent-retry trigger, entering persistent loop (timeout=${persistentPolicy.timeoutMs / 1000}s)`, ); - return await persistentRetryLoop(fn, persistentPolicy, logger); + try { + return await persistentRetryLoop(fn, persistentPolicy, logger); + } catch (persistentError) { + const terminal = + persistentError instanceof Error + ? persistentError + : new Error(formatErrorMessage(persistentError)); + persistentFailures.add(terminal); + throw terminal; + } } - - // Check whether this error is retryable under the standard policy. - if (policy.shouldRetry?.(lastError, attempt) === false) { - throw lastError; - } - - // Schedule the next retry with the configured backoff. - if (attempt < policy.maxRetries) { - const delay = - policy.backoff === "exponential" ? policy.baseDelayMs * 2 ** attempt : policy.baseDelayMs; - + }, + { + attempts: policy.maxRetries + 1, + minDelayMs: 0, + maxDelayMs: 2_147_000_000, + delayMs: ({ attempt }) => + policy.backoff === "exponential" + ? policy.baseDelayMs * 2 ** (attempt - 1) + : policy.baseDelayMs, + shouldRetry: (err, attempt) => { + const error = err instanceof Error ? err : new Error(formatErrorMessage(err)); + return !persistentFailures.has(error) && policy.shouldRetry?.(error, attempt - 1) !== false; + }, + onRetry: ({ attempt, delayMs, err }) => { + const error = err instanceof Error ? err : new Error(formatErrorMessage(err)); logger?.debug?.( - `[qqbot:retry] Attempt ${attempt + 1} failed, retrying in ${delay}ms: ${truncateUtf16Safe(lastError.message, 100)}`, + `[qqbot:retry] Attempt ${attempt} failed, retrying in ${delayMs}ms: ${truncateUtf16Safe(error.message, 100)}`, ); - await sleep(delay); - } - } - } - - throw lastError!; + }, + sleep, + }, + ); } /** diff --git a/extensions/tsconfig.package-boundary.paths.json b/extensions/tsconfig.package-boundary.paths.json index c0a2fd79942..02c4865a6c2 100644 --- a/extensions/tsconfig.package-boundary.paths.json +++ b/extensions/tsconfig.package-boundary.paths.json @@ -275,6 +275,9 @@ "@openclaw/normalization-core/*": [ "../dist/plugin-sdk/packages/normalization-core/src/*.d.ts" ], + "@openclaw/retry": [ + "../dist/plugin-sdk/packages/retry/src/index.d.ts" + ], "@openclaw/acp-core": [ "../dist/plugin-sdk/packages/acp-core/src/index.d.ts" ], diff --git a/extensions/xai/tsconfig.json b/extensions/xai/tsconfig.json index 88b08aae229..cb71f3123b6 100644 --- a/extensions/xai/tsconfig.json +++ b/extensions/xai/tsconfig.json @@ -258,6 +258,9 @@ "@openclaw/normalization-core/*": [ "../../dist/plugin-sdk/packages/normalization-core/src/*.d.ts" ], + "@openclaw/retry": [ + "../../dist/plugin-sdk/packages/retry/src/index.d.ts" + ], "@openclaw/acp-core": [ "../../dist/plugin-sdk/packages/acp-core/src/index.d.ts" ], diff --git a/packages/memory-host-sdk/package.json b/packages/memory-host-sdk/package.json index a24edd15cf5..5435a532627 100644 --- a/packages/memory-host-sdk/package.json +++ b/packages/memory-host-sdk/package.json @@ -4,7 +4,8 @@ "private": true, "type": "module", "dependencies": { - "@openclaw/normalization-core": "workspace:*" + "@openclaw/normalization-core": "workspace:*", + "@openclaw/retry": "workspace:*" }, "exports": { "./runtime-core": "./src/runtime-core.ts", diff --git a/packages/memory-host-sdk/src/host/batch-http.test.ts b/packages/memory-host-sdk/src/host/batch-http.test.ts index 5650f49b6aa..15ee3168fbe 100644 --- a/packages/memory-host-sdk/src/host/batch-http.test.ts +++ b/packages/memory-host-sdk/src/host/batch-http.test.ts @@ -69,7 +69,7 @@ describe("postJsonWithRetry", () => { headers: { Authorization: "Bearer test" }, body: { chunks: ["a", "b"] }, errorPrefix: "memory batch failed", - retryImpl: retryAsyncMock as typeof import("./retry-utils.js").retryAsync, + retryImpl: retryAsyncMock as typeof import("@openclaw/retry").retryAsync, }); expect(result).toEqual({ ok: true, ids: [1, 2] }); @@ -101,7 +101,7 @@ describe("postJsonWithRetry", () => { headers: {}, body: { chunks: [] }, errorPrefix: "memory batch failed", - retryImpl: retryAsyncMock as typeof import("./retry-utils.js").retryAsync, + retryImpl: retryAsyncMock as typeof import("@openclaw/retry").retryAsync, }); } catch (caught) { error = caught; diff --git a/packages/memory-host-sdk/src/host/batch-http.ts b/packages/memory-host-sdk/src/host/batch-http.ts index dbca28843db..101c43704c8 100644 --- a/packages/memory-host-sdk/src/host/batch-http.ts +++ b/packages/memory-host-sdk/src/host/batch-http.ts @@ -1,6 +1,6 @@ // Memory Host SDK module implements batch http behavior. +import { retryAsync } from "@openclaw/retry"; import { postJson } from "./post-json.js"; -import { retryAsync } from "./retry-utils.js"; import type { SsrFPolicy } from "./ssrf-policy.js"; // JSON POST helper for batch APIs with provider-style transient retry. diff --git a/packages/memory-host-sdk/src/host/embeddings.ts b/packages/memory-host-sdk/src/host/embeddings.ts index c80838a3439..5ed9782eaea 100644 --- a/packages/memory-host-sdk/src/host/embeddings.ts +++ b/packages/memory-host-sdk/src/host/embeddings.ts @@ -1,4 +1,5 @@ import { expectDefined } from "@openclaw/normalization-core"; +import { toRetryError } from "@openclaw/retry"; import { DEFAULT_LOCAL_MODEL } from "./embedding-defaults.js"; import { sanitizeAndNormalizeEmbedding } from "./embedding-vectors.js"; import { createLocalEmbeddingWorkerProvider } from "./embeddings-worker.js"; @@ -14,7 +15,6 @@ import { type LlamaModel, } from "./node-llama.js"; // Memory Host SDK module implements embeddings behavior. -import { toLintErrorObject } from "./retry-utils.js"; import { normalizeOptionalString } from "./string-utils.js"; type DisposableResource = { @@ -51,7 +51,7 @@ async function disposeResources( } } if (firstError) { - throw toLintErrorObject(firstError, "Non-Error thrown"); + throw toRetryError(firstError); } } diff --git a/packages/memory-host-sdk/src/host/read-retry.ts b/packages/memory-host-sdk/src/host/read-retry.ts index 5ce78bdeecf..d8bb59d589b 100644 --- a/packages/memory-host-sdk/src/host/read-retry.ts +++ b/packages/memory-host-sdk/src/host/read-retry.ts @@ -1,5 +1,5 @@ // Memory Host SDK module implements read retry behavior. -import { retryAsync } from "./retry-utils.js"; +import { retryAsync } from "@openclaw/retry"; // Retry helper for transient filesystem reads observed on memory stores. diff --git a/packages/memory-host-sdk/src/host/retry-utils.test.ts b/packages/memory-host-sdk/src/host/retry-utils.test.ts deleted file mode 100644 index 0a884fe1a3e..00000000000 --- a/packages/memory-host-sdk/src/host/retry-utils.test.ts +++ /dev/null @@ -1,77 +0,0 @@ -// Memory Host SDK tests cover retry utils behavior. -import { afterEach, describe, expect, it, vi } from "vitest"; -import { MAX_SAFE_TIMEOUT_DELAY_MS } from "../../../gateway-client/src/timeouts.js"; -import { retryAsync } from "./retry-utils.js"; - -afterEach(() => { - vi.restoreAllMocks(); -}); - -describe("retryAsync", () => { - it("falls back to the default attempt count for malformed numeric counts", async () => { - const run = vi.fn(async () => { - throw new Error("boom"); - }); - - await expect(retryAsync(run, Number.NaN, 0)).rejects.toThrow("boom"); - - expect(run).toHaveBeenCalledTimes(3); - }); - - it("falls back to the default attempt count for malformed option counts", async () => { - const run = vi.fn(async () => { - throw new Error("boom"); - }); - - await expect( - retryAsync(run, { - attempts: 1.5, - minDelayMs: 0, - maxDelayMs: 0, - }), - ).rejects.toThrow("boom"); - - expect(run).toHaveBeenCalledTimes(3); - }); - - it("caps legacy numeric retry sleeps at the timer-safe ceiling", async () => { - const run = vi - .fn<() => Promise>() - .mockRejectedValueOnce(new Error("boom")) - .mockResolvedValueOnce("ok"); - const timeoutSpy = vi.spyOn(globalThis, "setTimeout").mockImplementation((callback) => { - if (typeof callback === "function") { - callback(); - } - return 0 as unknown as ReturnType; - }); - - await expect(retryAsync(run, 2, Number.MAX_SAFE_INTEGER)).resolves.toBe("ok"); - - expect(timeoutSpy).toHaveBeenCalledWith(expect.any(Function), MAX_SAFE_TIMEOUT_DELAY_MS); - }); - - it("caps retryAfterMs sleeps at the timer-safe ceiling", async () => { - const run = vi - .fn<() => Promise>() - .mockRejectedValueOnce(new Error("boom")) - .mockResolvedValueOnce("ok"); - const timeoutSpy = vi.spyOn(globalThis, "setTimeout").mockImplementation((callback) => { - if (typeof callback === "function") { - callback(); - } - return 0 as unknown as ReturnType; - }); - - await expect( - retryAsync(run, { - attempts: 2, - minDelayMs: 0, - maxDelayMs: Number.MAX_SAFE_INTEGER, - retryAfterMs: () => Number.MAX_SAFE_INTEGER, - }), - ).resolves.toBe("ok"); - - expect(timeoutSpy).toHaveBeenCalledWith(expect.any(Function), MAX_SAFE_TIMEOUT_DELAY_MS); - }); -}); diff --git a/packages/memory-host-sdk/src/host/retry-utils.ts b/packages/memory-host-sdk/src/host/retry-utils.ts deleted file mode 100644 index e6b9888b23c..00000000000 --- a/packages/memory-host-sdk/src/host/retry-utils.ts +++ /dev/null @@ -1,175 +0,0 @@ -// Memory Host SDK helper module supports retry utils behavior. -import { resolveSafeTimeoutDelayMs } from "../../../gateway-client/src/timeouts.js"; - -/** Retry timing configuration with optional jitter. */ -type RetryConfig = { - attempts?: number; - minDelayMs?: number; - maxDelayMs?: number; - jitter?: number; -}; - -/** Retry callback payload. */ -type RetryInfo = { - attempt: number; - maxAttempts: number; - delayMs: number; - err: unknown; - label?: string; -}; - -/** Retry options for retryAsync. */ -type RetryOptions = RetryConfig & { - label?: string; - shouldRetry?: (err: unknown, attempt: number) => boolean; - retryAfterMs?: (err: unknown) => number | undefined; - onRetry?: (info: RetryInfo) => void; -}; - -const DEFAULT_RETRY_CONFIG = { - attempts: 3, - minDelayMs: 300, - maxDelayMs: 30_000, - jitter: 0, -}; - -function sleep(ms: number): Promise { - return new Promise((resolve) => { - setTimeout(resolve, ms); - }); -} - -function asFiniteNumber(value: unknown): number | undefined { - if (typeof value !== "number" || !Number.isFinite(value)) { - return undefined; - } - return value; -} - -function clampNumber(value: unknown, fallback: number, min?: number, max?: number): number { - const next = asFiniteNumber(value); - if (next === undefined) { - return fallback; - } - const floor = typeof min === "number" ? min : Number.NEGATIVE_INFINITY; - const ceiling = typeof max === "number" ? max : Number.POSITIVE_INFINITY; - return Math.min(Math.max(next, floor), ceiling); -} - -function resolveAttempts(value: unknown, fallback: number): number { - if (typeof value !== "number" || !Number.isSafeInteger(value)) { - return fallback; - } - return Math.max(1, value); -} - -/** Resolve retry settings with clamped positive timeout values. */ -function resolveRetryConfig( - defaults: Required = DEFAULT_RETRY_CONFIG, - overrides?: RetryConfig, -): Required { - const attempts = resolveAttempts(overrides?.attempts, defaults.attempts); - const minDelayMs = resolveSafeTimeoutDelayMs( - Math.round(clampNumber(overrides?.minDelayMs, defaults.minDelayMs, 0)), - { minMs: 0 }, - ); - const maxDelayMs = Math.max( - minDelayMs, - resolveSafeTimeoutDelayMs( - Math.round(clampNumber(overrides?.maxDelayMs, defaults.maxDelayMs, 0)), - { minMs: 0 }, - ), - ); - const jitter = clampNumber(overrides?.jitter, defaults.jitter, 0, 1); - return { attempts, minDelayMs, maxDelayMs, jitter }; -} - -function applyJitter(delayMs: number, jitter: number): number { - if (jitter <= 0) { - return delayMs; - } - const offset = (Math.random() * 2 - 1) * jitter; - return Math.max(0, Math.round(delayMs * (1 + offset))); -} - -/** Run an async operation with exponential backoff retry handling. */ -export async function retryAsync( - fn: () => Promise, - attemptsOrOptions: number | RetryOptions = 3, - initialDelayMs = 300, -): Promise { - if (typeof attemptsOrOptions === "number") { - const attempts = resolveAttempts(attemptsOrOptions, DEFAULT_RETRY_CONFIG.attempts); - let lastErr: unknown; - for (let i = 0; i < attempts; i += 1) { - try { - return await fn(); - } catch (err) { - lastErr = err; - if (i === attempts - 1) { - break; - } - await sleep(resolveSafeTimeoutDelayMs(initialDelayMs * 2 ** i, { minMs: 0 })); - } - } - throw toLintErrorObject(lastErr ?? new Error("Retry failed"), "Non-Error thrown"); - } - - const options = attemptsOrOptions; - const resolved = resolveRetryConfig(DEFAULT_RETRY_CONFIG, options); - const maxAttempts = resolved.attempts; - const minDelayMs = resolved.minDelayMs; - const maxDelayMs = - Number.isFinite(resolved.maxDelayMs) && resolved.maxDelayMs > 0 - ? resolved.maxDelayMs - : Number.POSITIVE_INFINITY; - const shouldRetry = options.shouldRetry ?? (() => true); - let lastErr: unknown; - - for (let attempt = 1; attempt <= maxAttempts; attempt += 1) { - try { - return await fn(); - } catch (err) { - lastErr = err; - if (attempt >= maxAttempts || !shouldRetry(err, attempt)) { - break; - } - - const retryAfterMs = options.retryAfterMs?.(err); - const hasRetryAfter = typeof retryAfterMs === "number" && Number.isFinite(retryAfterMs); - const baseDelay = hasRetryAfter - ? Math.max(resolveSafeTimeoutDelayMs(retryAfterMs, { minMs: 0 }), minDelayMs) - : resolveSafeTimeoutDelayMs(minDelayMs * 2 ** (attempt - 1), { minMs: 0 }); - let delay = Math.min(baseDelay, maxDelayMs); - delay = applyJitter(delay, resolved.jitter); - delay = Math.min(Math.max(delay, minDelayMs), maxDelayMs); - - options.onRetry?.({ - attempt, - maxAttempts, - delayMs: delay, - err, - label: options.label, - }); - if (delay > 0) { - await sleep(delay); - } - } - } - - throw toLintErrorObject(lastErr ?? new Error("Retry failed"), "Non-Error thrown"); -} - -export function toLintErrorObject(value: unknown, fallbackMessage: string): Error { - if (value instanceof Error) { - return value; - } - if (typeof value === "string") { - return new Error(value); - } - const error = new Error(fallbackMessage, { cause: value }); - if ((typeof value === "object" && value !== null) || typeof value === "function") { - Object.assign(error, value); - } - return error; -} diff --git a/packages/plugin-sdk/tsconfig.json b/packages/plugin-sdk/tsconfig.json index e5d9fa8122a..ab76f585f8d 100644 --- a/packages/plugin-sdk/tsconfig.json +++ b/packages/plugin-sdk/tsconfig.json @@ -19,6 +19,7 @@ "../../packages/media-generation-core/src/**/*.ts", "../../packages/model-catalog-core/src/**/*.ts", "../../packages/normalization-core/src/**/*.ts", + "../../packages/retry/src/**/*.ts", "../../packages/acp-core/src/**/*.ts", "../../packages/terminal-core/src/**/*.ts", "../../src/plugin-sdk/**/*.ts", diff --git a/packages/retry/package.json b/packages/retry/package.json new file mode 100644 index 00000000000..a184531ea60 --- /dev/null +++ b/packages/retry/package.json @@ -0,0 +1,21 @@ +{ + "name": "@openclaw/retry", + "version": "0.0.0-private", + "private": true, + "files": [ + "dist" + ], + "type": "module", + "main": "./dist/index.mjs", + "types": "./dist/index.d.mts", + "exports": { + ".": { + "types": "./dist/index.d.mts", + "import": "./dist/index.mjs", + "default": "./dist/index.mjs" + } + }, + "scripts": { + "build": "tsdown src/index.ts --no-config --platform node --format esm --dts --out-dir dist --clean" + } +} diff --git a/packages/retry/src/index.test.ts b/packages/retry/src/index.test.ts new file mode 100644 index 00000000000..1054efea3e3 --- /dev/null +++ b/packages/retry/src/index.test.ts @@ -0,0 +1,75 @@ +import { describe, expect, it, vi } from "vitest"; +import { createRetryRunner, retryAsync } from "./index.js"; + +describe("retryAsync", () => { + it.each([0, 0.5])( + "never rounds an honorable Retry-After below its floor with jitter=%s", + async (jitter) => { + const sleeps: number[] = []; + const run = createRetryRunner({ sleep: async (ms) => void sleeps.push(ms) }); + const operation = vi + .fn<() => Promise>() + .mockRejectedValueOnce(new Error("rate limited")) + .mockResolvedValueOnce("ok"); + + await expect( + run(operation, { + attempts: 2, + minDelayMs: 0, + maxDelayMs: 10, + jitter, + random: () => 0, + retryAfterMs: () => 1.4, + }), + ).resolves.toBe("ok"); + expect(sleeps).toEqual([2]); + }, + ); + + it("supports custom schedules, abortable sleeps, and async retry hooks", async () => { + const events: string[] = []; + const operation = vi + .fn<() => Promise>() + .mockRejectedValueOnce(new Error("first")) + .mockRejectedValueOnce(new Error("second")) + .mockResolvedValueOnce("ok"); + + await expect( + retryAsync(operation, { + attempts: 3, + minDelayMs: 0, + maxDelayMs: 100, + delayMs: ({ attempt }) => [10, 30][attempt - 1] ?? 0, + onRetry: async ({ attempt }) => void events.push(`retry:${attempt}`), + sleep: async (ms) => void events.push(`sleep:${ms}`), + }), + ).resolves.toBe("ok"); + expect(events).toEqual(["retry:1", "sleep:10", "retry:2", "sleep:30"]); + }); + + it("preserves terminal Error identity", async () => { + const terminal = new Error("terminal"); + await expect( + retryAsync( + async () => { + throw terminal; + }, + { + attempts: 1, + }, + ), + ).rejects.toBe(terminal); + }); + + it("clamps numeric overload delays to the Node timer ceiling", async () => { + const sleeps: number[] = []; + const run = createRetryRunner({ sleep: async (ms) => void sleeps.push(ms) }); + const operation = vi + .fn<() => Promise>() + .mockRejectedValueOnce(new Error("first")) + .mockResolvedValueOnce("ok"); + + await run(operation, 2, Number.POSITIVE_INFINITY); + expect(sleeps).toEqual([2_147_000_000]); + }); +}); diff --git a/packages/retry/src/index.ts b/packages/retry/src/index.ts new file mode 100644 index 00000000000..e81f6030969 --- /dev/null +++ b/packages/retry/src/index.ts @@ -0,0 +1,274 @@ +// Dependency-free retry scheduling shared across core and leaf workspace packages. + +// Keep a small margin below Node's signed 32-bit timeout ceiling. +const MAX_TIMER_TIMEOUT_MS = 2_147_000_000; + +/** Retry timing knobs shared by generic retry runners and channel retry policies. */ +export type RetryConfig = { + attempts?: number; + minDelayMs?: number; + maxDelayMs?: number; + /** + * Delay spread strategy. A fraction (0-1) spreads proportionally around the + * backoff delay. `"full"` draws uniformly from [delay, 2*delay). + */ + jitter?: number | "full"; +}; + +/** Metadata available while selecting the delay before the next retry. */ +type RetryDelayContext = { + attempt: number; + maxAttempts: number; + err: unknown; + label?: string; +}; + +/** Metadata emitted before a retry attempt sleeps and reruns the operation. */ +export type RetryInfo = RetryDelayContext & { + delayMs: number; +}; + +/** Retry execution options, including predicates, delay hooks, and callbacks. */ +export type RetryOptions = RetryConfig & { + label?: string; + shouldRetry?: (err: unknown, attempt: number) => boolean; + retryAfterMs?: (err: unknown) => number | undefined; + retryAfterMaxDelayMs?: number; + /** Overrides exponential backoff while retaining timer clamping and jitter. */ + delayMs?: number | ((context: RetryDelayContext) => number); + /** Runs before sleeping; returned promises are awaited. */ + onRetry?: (info: RetryInfo) => unknown; + /** Random fraction source in [0, 1); injectable for deterministic tests. */ + random?: () => number; + /** Sleep implementation; useful for abortable waits and deterministic tests. */ + sleep?: (ms: number) => Promise; +}; + +/** Runtime dependencies used to adapt the leaf scheduler to its host. */ +export type RetryRuntime = { + sleep?: (ms: number) => Promise; + random?: () => number; + createFailure?: (attemptErrors: readonly unknown[]) => Error; +}; + +const DEFAULT_RETRY_CONFIG: Required = { + attempts: 3, + minDelayMs: 300, + maxDelayMs: 30_000, + jitter: 0, +}; + +function defaultSleep(ms: number): Promise { + return new Promise((resolve) => { + setTimeout(resolve, ms); + }); +} + +function asFiniteNumber(value: unknown): number | undefined { + return typeof value === "number" && Number.isFinite(value) ? value : undefined; +} + +function clampNumber(value: unknown, fallback: number, min?: number, max?: number): number { + const next = asFiniteNumber(value); + if (next === undefined) { + return fallback; + } + const floor = min ?? Number.NEGATIVE_INFINITY; + const ceiling = max ?? Number.POSITIVE_INFINITY; + return Math.min(Math.max(next, floor), ceiling); +} + +function resolveAttemptCount(value: unknown, fallback: number): number { + const candidate = asFiniteNumber(value) ?? fallback; + return Math.max(1, Math.round(candidate)); +} + +function resolveRetryDelayMs(value: number): number { + if (value === Number.POSITIVE_INFINITY) { + return MAX_TIMER_TIMEOUT_MS; + } + const finite = asFiniteNumber(value) ?? 0; + return Math.min(Math.max(Math.round(finite), 0), MAX_TIMER_TIMEOUT_MS); +} + +function resolveJitterConfig(value: unknown, fallback: number | "full"): number | "full" { + if (value === "full") { + return "full"; + } + const fraction = asFiniteNumber(value); + return fraction === undefined ? fallback : Math.min(Math.max(fraction, 0), 1); +} + +/** Resolves retry overrides into clamped timer-safe settings. */ +export function resolveRetryConfig( + defaults: Required = DEFAULT_RETRY_CONFIG, + overrides?: RetryConfig, +): Required { + const attempts = resolveAttemptCount( + clampNumber(overrides?.attempts, defaults.attempts, 1), + defaults.attempts, + ); + const minDelayMs = resolveRetryDelayMs( + Math.round(clampNumber(overrides?.minDelayMs, defaults.minDelayMs, 0)), + ); + const maxDelayMs = Math.max( + minDelayMs, + resolveRetryDelayMs(Math.round(clampNumber(overrides?.maxDelayMs, defaults.maxDelayMs, 0))), + ); + return { + attempts, + minDelayMs, + maxDelayMs, + jitter: resolveJitterConfig(overrides?.jitter, defaults.jitter), + }; +} + +type JitterMode = "symmetric" | "positive"; + +function applyJitter( + delayMs: number, + jitter: number | "full", + mode: JitterMode, + random: () => number, +): number { + if (jitter === "full") { + if (mode === "symmetric") { + // Over-cap Retry-After cannot be honored. Spread downward instead of + // letting the final cap collapse every client onto the same instant. + return Math.max(0, Math.round(delayMs * (0.5 + random() * 0.5))); + } + return Math.max(0, Math.ceil(delayMs * (1 + random()))); + } + if (jitter <= 0) { + return mode === "positive" ? Math.ceil(delayMs) : delayMs; + } + const fraction = random(); + const offset = mode === "positive" ? fraction * jitter : (fraction * 2 - 1) * jitter; + const raw = delayMs * (1 + offset); + // Retry-After is a lower bound. Positive jitter must round upward or a + // fractional server hint can be undercut even with a zero random draw. + return Math.max(0, mode === "positive" ? Math.ceil(raw) : Math.round(raw)); +} + +/** Normalizes an arbitrary thrown value while preserving Error identity. */ +export function toRetryError(value: unknown, fallbackMessage = "Non-Error thrown"): Error { + if (value instanceof Error) { + return value; + } + if (typeof value === "string") { + return new Error(value); + } + const error = new Error(fallbackMessage, { cause: value }); + if ((typeof value === "object" && value !== null) || typeof value === "function") { + Object.assign(error, value); + } + return error; +} + +function defaultCreateFailure(attemptErrors: readonly unknown[]): Error { + return toRetryError(attemptErrors.at(-1) ?? new Error("Retry failed")); +} + +/** Creates a retry runner bound to host-specific sleep, randomness, and diagnostics. */ +export function createRetryRunner(runtime: RetryRuntime = {}) { + const runtimeSleep = runtime.sleep ?? defaultSleep; + const runtimeRandom = runtime.random ?? Math.random; + const createFailure = runtime.createFailure ?? defaultCreateFailure; + + return async function retryAsync( + fn: () => Promise, + attemptsOrOptions: number | RetryOptions = 3, + initialDelayMs = 300, + ): Promise { + const attemptErrors: unknown[] = []; + if (typeof attemptsOrOptions === "number") { + const attempts = resolveAttemptCount(attemptsOrOptions, DEFAULT_RETRY_CONFIG.attempts); + for (let index = 0; index < attempts; index += 1) { + try { + return await fn(); + } catch (err) { + attemptErrors.push(err); + if (index === attempts - 1) { + break; + } + await runtimeSleep(resolveRetryDelayMs(initialDelayMs * 2 ** index)); + } + } + throw createFailure(attemptErrors); + } + + const options = attemptsOrOptions; + const resolved = resolveRetryConfig(DEFAULT_RETRY_CONFIG, options); + const maxAttempts = resolved.attempts; + const minDelayMs = resolved.minDelayMs; + const maxDelayMs = resolved.maxDelayMs > 0 ? resolved.maxDelayMs : Number.POSITIVE_INFINITY; + const retryAfterMaxDelayMs = + options.retryAfterMaxDelayMs === undefined + ? maxDelayMs + : Math.max( + minDelayMs, + resolveRetryDelayMs( + Math.round(clampNumber(options.retryAfterMaxDelayMs, maxDelayMs, 0)), + ), + ); + const random = options.random ?? runtimeRandom; + const sleep = options.sleep ?? runtimeSleep; + const shouldRetry = options.shouldRetry ?? (() => true); + + for (let attempt = 1; attempt <= maxAttempts; attempt += 1) { + try { + return await fn(); + } catch (err) { + attemptErrors.push(err); + if (attempt >= maxAttempts || !shouldRetry(err, attempt)) { + break; + } + + const context: RetryDelayContext = { + attempt, + maxAttempts, + err, + label: options.label, + }; + const retryAfterMs = options.retryAfterMs?.(err); + const hasRetryAfter = typeof retryAfterMs === "number" && Number.isFinite(retryAfterMs); + const configuredDelay = + typeof options.delayMs === "function" ? options.delayMs(context) : options.delayMs; + const resolvedConfiguredDelay = + configuredDelay === undefined ? undefined : resolveRetryDelayMs(configuredDelay); + const baseDelay = hasRetryAfter + ? Math.max(retryAfterMs, minDelayMs) + : resolvedConfiguredDelay === undefined + ? minDelayMs * 2 ** (attempt - 1) + : Math.max(resolvedConfiguredDelay, minDelayMs); + const delayCap = hasRetryAfter ? retryAfterMaxDelayMs : maxDelayMs; + let delay = Math.min(baseDelay, delayCap); + + // Honorable Retry-After hints use positive jitter. Only an over-cap, + // already-unsatisfiable hint may spread downward to avoid lockstep. + const canHonorRetryAfter = + hasRetryAfter && typeof retryAfterMs === "number" && retryAfterMs <= delayCap; + const overCapRetryAfter = hasRetryAfter && !canHonorRetryAfter; + const wantsPositiveDraw = + resolved.jitter === "full" ? !overCapRetryAfter : canHonorRetryAfter; + delay = applyJitter( + delay, + resolved.jitter, + wantsPositiveDraw ? "positive" : "symmetric", + random, + ); + delay = Math.min(Math.max(delay, minDelayMs), delayCap); + + await options.onRetry?.({ ...context, delayMs: delay }); + if (delay > 0) { + await sleep(delay); + } + } + } + + throw createFailure(attemptErrors); + }; +} + +/** Default retry runner for dependency-leaf consumers. */ +export const retryAsync = createRetryRunner(); diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 45e6d79162a..c048026b85a 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -1961,6 +1961,9 @@ importers: '@openclaw/normalization-core': specifier: workspace:* version: link:../normalization-core + '@openclaw/retry': + specifier: workspace:* + version: link:../retry packages/model-catalog-core: {} @@ -1972,6 +1975,8 @@ importers: packages/normalization-core: {} + packages/retry: {} + packages/plugin-package-contract: {} packages/plugin-sdk: {} diff --git a/scripts/lib/extension-package-boundary.ts b/scripts/lib/extension-package-boundary.ts index 4a133525b82..3bfb70a5d27 100644 --- a/scripts/lib/extension-package-boundary.ts +++ b/scripts/lib/extension-package-boundary.ts @@ -188,6 +188,7 @@ export const EXTENSION_PACKAGE_BOUNDARY_BASE_PATHS = { "../dist/plugin-sdk/packages/normalization-core/src/string-coerce.d.ts", ], "@openclaw/normalization-core/*": ["../dist/plugin-sdk/packages/normalization-core/src/*.d.ts"], + "@openclaw/retry": ["../dist/plugin-sdk/packages/retry/src/index.d.ts"], ...buildPackageBoundaryDtsPaths({ packageName: "@openclaw/acp-core", packageDir: "acp-core", diff --git a/scripts/lib/tsdown-output-roots.mjs b/scripts/lib/tsdown-output-roots.mjs index 61be55152ec..b1ed9a7981d 100644 --- a/scripts/lib/tsdown-output-roots.mjs +++ b/scripts/lib/tsdown-output-roots.mjs @@ -12,6 +12,7 @@ const TSDOWN_PACKAGE_NAMES = [ "model-catalog-core", "net-policy", "normalization-core", + "retry", "speech-core", "terminal-core", "acp-core", diff --git a/scripts/prepare-extension-package-boundary-artifacts.mjs b/scripts/prepare-extension-package-boundary-artifacts.mjs index 0bb7efd9258..2e524a32637 100644 --- a/scripts/prepare-extension-package-boundary-artifacts.mjs +++ b/scripts/prepare-extension-package-boundary-artifacts.mjs @@ -92,6 +92,7 @@ const PLUGIN_SDK_TYPE_INPUTS = [ "packages/media-generation-core/src", "packages/media-understanding-common/src", "packages/normalization-core/src", + "packages/retry/src", "packages/acp-core/src", "packages/terminal-core/src", "src/video-generation/dashscope-compatible.ts", @@ -167,6 +168,7 @@ const ROOT_DTS_REQUIRED_OUTPUTS = [ "dist/plugin-sdk/packages/model-catalog-core/src/provider-id.d.ts", "dist/plugin-sdk/packages/model-catalog-core/src/provider-model-id-normalization.d.ts", "dist/plugin-sdk/packages/model-catalog-core/src/provider-model-id-normalize.d.ts", + "dist/plugin-sdk/packages/retry/src/index.d.ts", "dist/plugin-sdk/error-runtime.d.ts", "dist/plugin-sdk/plugin-entry.d.ts", "dist/plugin-sdk/provider-auth.d.ts", @@ -221,6 +223,7 @@ const PACKAGE_DTS_REQUIRED_OUTPUTS = [ "packages/plugin-sdk/dist/packages/normalization-core/src/record-coerce.d.ts", "packages/plugin-sdk/dist/packages/normalization-core/src/string-coerce.d.ts", "packages/plugin-sdk/dist/packages/normalization-core/src/string-normalization.d.ts", + "packages/plugin-sdk/dist/packages/retry/src/index.d.ts", "packages/plugin-sdk/dist/packages/terminal-core/src/ansi.d.ts", "packages/plugin-sdk/dist/packages/terminal-core/src/decorative-emoji.d.ts", "packages/plugin-sdk/dist/packages/terminal-core/src/health-style.d.ts", diff --git a/scripts/run-node-watch-paths.mjs b/scripts/run-node-watch-paths.mjs index b4050327126..6f493c1a60b 100644 --- a/scripts/run-node-watch-paths.mjs +++ b/scripts/run-node-watch-paths.mjs @@ -16,6 +16,7 @@ const RUN_NODE_PACKAGE_SOURCE_ROOTS = [ "packages/media-generation-core/src", "packages/media-understanding-common/src", "packages/normalization-core/src", + "packages/retry/src", "packages/acp-core/src", "packages/terminal-core/src", "packages/web-content-core/src", diff --git a/scripts/write-plugin-sdk-entry-dts.ts b/scripts/write-plugin-sdk-entry-dts.ts index 2a1624afac9..b35c11b1d2e 100644 --- a/scripts/write-plugin-sdk-entry-dts.ts +++ b/scripts/write-plugin-sdk-entry-dts.ts @@ -51,6 +51,7 @@ function isBareImportSpecifier(id: string): boolean { id === "@openclaw/llm-core" || id.startsWith("@openclaw/llm-core/") || id === "@openclaw/model-catalog-core/model-catalog-types" || + id === "@openclaw/retry" || id.startsWith("@openclaw/normalization-core/") || id.startsWith("@openclaw/media-core/") || id.startsWith("@openclaw/acp-core/") diff --git a/src/cron/isolated-agent/delivery-dispatch.ts b/src/cron/isolated-agent/delivery-dispatch.ts index 2cfd3dad448..606ce0ad979 100644 --- a/src/cron/isolated-agent/delivery-dispatch.ts +++ b/src/cron/isolated-agent/delivery-dispatch.ts @@ -36,6 +36,7 @@ import type { SourceDeliveryVisibleDelivery, } from "../../infra/outbound/source-delivery-plan.js"; import { normalizeTargetForProvider } from "../../infra/outbound/target-normalization.js"; +import { retryAsync } from "../../infra/retry.js"; import { hasReplyPayloadContent } from "../../interactive/payload.js"; import { stringifyRouteThreadId } from "../../plugin-sdk/channel-route.js"; import { @@ -970,28 +971,32 @@ async function retryTransientDirectCronDelivery(params: { run: () => Promise; }): Promise { const retryDelaysMs = resolveDirectCronRetryDelaysMs(); - for (const [retryIndex, delayMs] of retryDelaysMs.entries()) { - if (params.signal?.aborted) { - throw new Error("cron delivery aborted"); - } - try { - return await params.run(); - } catch (err) { - if (!isTransientDirectCronDeliveryError(err) || params.signal?.aborted) { - throw err; - } - const nextAttempt = retryIndex + 2; - const maxAttempts = retryDelaysMs.length + 1; - await logCronDeliveryWarn( - `[cron:${params.jobId}] transient direct announce delivery failure, retrying ${nextAttempt}/${maxAttempts} in ${Math.round(delayMs / 1000)}s: ${summarizeDirectCronDeliveryError(err)}`, - ); - await sleepWithAbort(delayMs, params.signal); - } - } if (params.signal?.aborted) { throw new Error("cron delivery aborted"); } - return await params.run(); + const runWithAbortCheck = async () => { + if (params.signal?.aborted) { + throw new Error("cron delivery aborted"); + } + return await params.run(); + }; + return await retryAsync(runWithAbortCheck, { + attempts: retryDelaysMs.length + 1, + minDelayMs: 0, + maxDelayMs: Math.max(...retryDelaysMs), + delayMs: ({ attempt }) => retryDelaysMs[attempt - 1] ?? 0, + shouldRetry: (err) => + params.signal?.aborted !== true && isTransientDirectCronDeliveryError(err), + onRetry: async ({ attempt, maxAttempts, delayMs, err }) => { + await logCronDeliveryWarn( + `[cron:${params.jobId}] transient direct announce delivery failure, retrying ${attempt + 1}/${maxAttempts} in ${Math.round(delayMs / 1000)}s: ${summarizeDirectCronDeliveryError(err)}`, + ); + if (delayMs === 0) { + await sleepWithAbort(0, params.signal); + } + }, + sleep: async (delayMs) => await sleepWithAbort(delayMs, params.signal), + }); } /** Dispatches cron run output through verified message-tool or direct delivery paths. */ diff --git a/src/infra/clawhub-retry.test.ts b/src/infra/clawhub-retry.test.ts index 8d086dea430..8f2fe78cff2 100644 --- a/src/infra/clawhub-retry.test.ts +++ b/src/infra/clawhub-retry.test.ts @@ -171,4 +171,17 @@ describe("retryClawHubRead", () => { expect(await optedInResult.response.text()).toBe("ok"); expect(optedInAttempts).toBe(2); }); + + it("returns the final retryable response for caller-owned HTTP handling", async () => { + const disposeRetry = vi.fn(async ({ response }: { response: Response }) => { + await response.body?.cancel(); + }); + const result = await retryClawHubRead( + async () => ({ response: new Response("unavailable", { status: 503 }) }), + { disposeRetry, sleep: async () => {} }, + ); + + expect(result.response.status).toBe(503); + expect(disposeRetry).toHaveBeenCalledTimes(3); + }); }); diff --git a/src/infra/clawhub-retry.ts b/src/infra/clawhub-retry.ts index f5c55a53ccc..75c3d5ac217 100644 --- a/src/infra/clawhub-retry.ts +++ b/src/infra/clawhub-retry.ts @@ -1,5 +1,6 @@ // Defines the bounded retry contract shared by ClawHub runtime and release reads. import { parseRetryAfterHttpDateMs } from "../../packages/ai/src/internal/retry-after.js"; +import { retryAsync } from "./retry.js"; const CLAWHUB_RETRY_DELAYS_MS = [1_000, 3_000, 10_000] as const; const CLAWHUB_MAX_RETRY_AFTER_MS = 60_000; @@ -14,6 +15,12 @@ type ClawHubRetryOptions = { sleep?: (ms: number) => Promise; }; +class RetryableClawHubResponse extends Error { + constructor(readonly result: T) { + super(`ClawHub request returned retryable status ${result.response.status}`); + } +} + function isRetryableClawHubStatus(status: number, retryRateLimit: boolean): boolean { return (retryRateLimit && status === 429) || status === 502 || status === 503 || status === 504; } @@ -36,14 +43,6 @@ function parseRetryAfterMs(headers: Headers): number | undefined { return delayMs <= CLAWHUB_MAX_RETRY_AFTER_MS ? delayMs : undefined; } -function retryDelayMs(response: Response | undefined, attempt: number): number { - return ( - (response ? parseRetryAfterMs(response.headers) : undefined) ?? - CLAWHUB_RETRY_DELAYS_MS[attempt] ?? - 0 - ); -} - async function defaultSleep(ms: number): Promise { await new Promise((resolve) => { setTimeout(resolve, ms); @@ -58,27 +57,37 @@ export async function retryClawHubRead( request: () => Promise, options: ClawHubRetryOptions, ): Promise { - for (let attempt = 0; ; attempt += 1) { - let result: T; - try { - result = await request(); - } catch (error) { - if (attempt >= CLAWHUB_RETRY_DELAYS_MS.length) { - throw error; - } - await (options.sleep ?? defaultSleep)(retryDelayMs(undefined, attempt)); - continue; + try { + return await retryAsync( + async () => { + const result = await request(); + if (isRetryableClawHubStatus(result.response.status, options.retryRateLimit === true)) { + throw new RetryableClawHubResponse(result); + } + return result; + }, + { + attempts: CLAWHUB_RETRY_DELAYS_MS.length + 1, + minDelayMs: 0, + maxDelayMs: CLAWHUB_MAX_RETRY_AFTER_MS, + delayMs: ({ attempt }) => CLAWHUB_RETRY_DELAYS_MS[attempt - 1] ?? 0, + retryAfterMs: (error) => + error instanceof RetryableClawHubResponse + ? parseRetryAfterMs(error.result.response.headers) + : undefined, + onRetry: async ({ err }) => { + if (err instanceof RetryableClawHubResponse) { + await options.disposeRetry(err.result); + } + }, + sleep: options.sleep ?? defaultSleep, + }, + ); + } catch (error) { + // Callers own final HTTP error handling and therefore need the response. + if (error instanceof RetryableClawHubResponse) { + return error.result; } - - if ( - !isRetryableClawHubStatus(result.response.status, options.retryRateLimit === true) || - attempt >= CLAWHUB_RETRY_DELAYS_MS.length - ) { - return result; - } - - const delayMs = retryDelayMs(result.response, attempt); - await options.disposeRetry(result); - await (options.sleep ?? defaultSleep)(delayMs); + throw error; } } diff --git a/src/infra/retry.ts b/src/infra/retry.ts index 5b912c138cb..aeaaeb00475 100644 --- a/src/infra/retry.ts +++ b/src/infra/retry.ts @@ -1,277 +1,33 @@ -// Provides generic retry timing and sleep helpers. -import { asFiniteNumber } from "@openclaw/normalization-core/number-coercion"; -import { MAX_TIMER_TIMEOUT_MS, resolveTimerTimeoutMs } from "../shared/number-coercion.js"; -import { sleep } from "../utils.js"; -import { toErrorObject } from "./errors.js"; +// Adapts the dependency-free retry scheduler to core runtime facilities. +import { + createRetryRunner, + resolveRetryConfig, + toRetryError, + type RetryConfig, + type RetryInfo, + type RetryOptions, +} from "../../packages/retry/src/index.js"; import { getRetryAttemptErrors, recordRetryAttemptErrors } from "./retry-attempt-errors.js"; import { generateSecureFraction } from "./secure-random.js"; -/** Retry timing knobs shared by generic retry runners and channel retry policies. */ -export type RetryConfig = { - attempts?: number; - minDelayMs?: number; - maxDelayMs?: number; - /** - * Delay spread strategy. A fraction (0-1) spreads proportionally around the - * backoff delay (existing behavior). `"full"` draws uniformly from - * [delay, 2*delay): the backoff delay stays a hard floor and `maxDelayMs` - * clamps after the draw, so capped attempts land exactly on the cap. - */ - jitter?: number | "full"; -}; +export { resolveRetryConfig, type RetryConfig, type RetryInfo, type RetryOptions }; -/** Metadata emitted before a retry attempt sleeps and reruns the operation. */ -export type RetryInfo = { - attempt: number; - maxAttempts: number; - delayMs: number; - err: unknown; - label?: string; -}; - -/** Retry execution options, including predicates, Retry-After hooks, and retry callbacks. */ -export type RetryOptions = RetryConfig & { - label?: string; - shouldRetry?: (err: unknown, attempt: number) => boolean; - retryAfterMs?: (err: unknown) => number | undefined; - retryAfterMaxDelayMs?: number; - onRetry?: (info: RetryInfo) => void; - /** Random fraction source in [0, 1); injectable for deterministic tests. */ - random?: () => number; -}; - -const DEFAULT_RETRY_CONFIG = { - attempts: 3, - minDelayMs: 300, - maxDelayMs: 30_000, - jitter: 0, -}; - -function appendRetryAttemptError(attemptErrors: unknown[], err: unknown): void { - const nestedAttempts = getRetryAttemptErrors(err); - attemptErrors.push(...(nestedAttempts ?? [err])); -} - -function createRetryFailure(attemptErrors: readonly unknown[]): Error { - const failure = toErrorObject( +function createRetryFailure(rawAttemptErrors: readonly unknown[]): Error { + const attemptErrors = rawAttemptErrors.flatMap((err) => getRetryAttemptErrors(err) ?? [err]); + const failure = toRetryError( attemptErrors.at(-1) ?? new Error("Retry failed"), "Non-Error thrown", ); if (attemptErrors.length > 1) { - // Preserve the public terminal-error identity while carrying every internal - // attempt into duplicate-send decisions made outside the channel adapter. + // Preserve terminal-error identity while carrying all attempts into + // duplicate-send decisions outside the channel adapter. recordRetryAttemptErrors(failure, attemptErrors); } return failure; } -const clampNumber = (value: unknown, fallback: number, min?: number, max?: number) => { - const next = asFiniteNumber(value); - if (next === undefined) { - return fallback; - } - const floor = typeof min === "number" ? min : Number.NEGATIVE_INFINITY; - const ceiling = typeof max === "number" ? max : Number.POSITIVE_INFINITY; - return Math.min(Math.max(next, floor), ceiling); -}; - -function resolveAttemptCount(value: unknown, fallback: number): number { - const candidate = typeof value === "number" && Number.isFinite(value) ? value : fallback; - return Math.max(1, Math.round(candidate)); -} - -function resolveRetryDelayMs(value: number): number { - if (value === Number.POSITIVE_INFINITY) { - return MAX_TIMER_TIMEOUT_MS; - } - return resolveTimerTimeoutMs(value, 0, 0); -} - -function resolveJitterConfig(value: unknown, fallback: number | "full"): number | "full" { - if (value === "full") { - return "full"; - } - const fraction = asFiniteNumber(value); - if (fraction === undefined) { - return fallback; - } - return Math.min(Math.max(fraction, 0), 1); -} - -/** Resolves retry config overrides into clamped timer-safe settings. */ -export function resolveRetryConfig( - defaults: Required = DEFAULT_RETRY_CONFIG, - overrides?: RetryConfig, -): Required { - const attempts = resolveAttemptCount( - clampNumber(overrides?.attempts, defaults.attempts, 1), - defaults.attempts, - ); - const minDelayMs = resolveRetryDelayMs( - Math.round(clampNumber(overrides?.minDelayMs, defaults.minDelayMs, 0)), - ); - const maxDelayMs = Math.max( - minDelayMs, - resolveRetryDelayMs(Math.round(clampNumber(overrides?.maxDelayMs, defaults.maxDelayMs, 0))), - ); - const jitter = resolveJitterConfig(overrides?.jitter, defaults.jitter); - return { attempts, minDelayMs, maxDelayMs, jitter }; -} - -type JitterMode = "symmetric" | "positive"; - -function applyJitter( - delayMs: number, - jitter: number | "full", - mode: JitterMode, - random: () => number, -): number { - if (jitter === "full") { - if (mode === "symmetric") { - // Unsatisfiable over-cap Retry-After: an upward draw would be erased by - // the caller's cap clamp and every client would land in lockstep at the - // cap, so draw downward across one half-period instead. That preserves - // spread (the invariant the numeric symmetric fallback below protects) - // while staying as close to the server's hint as the cap allows. - return Math.max(0, Math.round(delayMs * (0.5 + random() * 0.5))); - } - // Full jitter draws uniformly from [delay, 2*delay): the backoff delay is - // a hard floor (never fire early), which also keeps honorable Retry-After - // lower bounds. Callers clamp `maxDelayMs` after this, so capped attempts - // land exactly on the cap (same boundary trade-off as `positive` mode - // below). Ceil preserves the floor contract for fractional bases. - return Math.max(0, Math.ceil(delayMs * (1 + random()))); - } - if (jitter <= 0) { - return delayMs; - } - // `symmetric` spreads within ±jitter around the base delay; correct for pure - // exponential backoff where going slightly early is harmless. `positive` - // only adds to the base delay; use it when the base delay is already a - // lower bound the caller must respect (for example a server-supplied - // Retry-After) so concurrent clients still spread without ever dipping - // below the caller's floor. - const fraction = random(); - const offset = mode === "positive" ? fraction * jitter : (fraction * 2 - 1) * jitter; - const raw = delayMs * (1 + offset); - // Rounding choice preserves the mode's contract. `positive` guarantees - // `delay >= delayMs`, so a non-integer `delayMs` (e.g. retryAfterMs=1.4) - // must round *up* — plain `Math.round(1.4)=1` would drop the delay below - // the caller's lower bound and violate the Retry-After invariant the - // positive branch exists to enforce. Symmetric has no floor contract so - // it stays on `Math.round`. - return Math.max(0, mode === "positive" ? Math.ceil(raw) : Math.round(raw)); -} - -/** Runs an async operation until it succeeds, retry policy stops, or attempts are exhausted. */ -export async function retryAsync( - fn: () => Promise, - attemptsOrOptions: number | RetryOptions = 3, - initialDelayMs = 300, -): Promise { - if (typeof attemptsOrOptions === "number") { - const attempts = resolveAttemptCount(attemptsOrOptions, DEFAULT_RETRY_CONFIG.attempts); - const attemptErrors: unknown[] = []; - for (let i = 0; i < attempts; i += 1) { - try { - return await fn(); - } catch (err) { - appendRetryAttemptError(attemptErrors, err); - if (i === attempts - 1) { - break; - } - const delay = resolveRetryDelayMs(initialDelayMs * 2 ** i); - await sleep(delay); - } - } - throw createRetryFailure(attemptErrors); - } - - const options = attemptsOrOptions; - - const resolved = resolveRetryConfig(DEFAULT_RETRY_CONFIG, options); - const maxAttempts = resolved.attempts; - const minDelayMs = resolved.minDelayMs; - const maxDelayMs = - Number.isFinite(resolved.maxDelayMs) && resolved.maxDelayMs > 0 - ? resolved.maxDelayMs - : Number.POSITIVE_INFINITY; - const retryAfterMaxDelayMs = - options.retryAfterMaxDelayMs === undefined - ? maxDelayMs - : Math.max( - minDelayMs, - resolveRetryDelayMs(Math.round(clampNumber(options.retryAfterMaxDelayMs, maxDelayMs, 0))), - ); - const jitter = resolved.jitter; - const random = options.random ?? generateSecureFraction; - const shouldRetry = options.shouldRetry ?? (() => true); - const attemptErrors: unknown[] = []; - - for (let attempt = 1; attempt <= maxAttempts; attempt += 1) { - try { - return await fn(); - } catch (err) { - appendRetryAttemptError(attemptErrors, err); - if (attempt >= maxAttempts || !shouldRetry(err, attempt)) { - break; - } - - const retryAfterMs = options.retryAfterMs?.(err); - const hasRetryAfter = typeof retryAfterMs === "number" && Number.isFinite(retryAfterMs); - const baseDelay = hasRetryAfter - ? Math.max(retryAfterMs, minDelayMs) - : minDelayMs * 2 ** (attempt - 1); - const delayCap = hasRetryAfter ? retryAfterMaxDelayMs : maxDelayMs; - let delay = Math.min(baseDelay, delayCap); - // Server-supplied Retry-After is a lower-bound contract with the - // upstream rate limiter; symmetric jitter would let roughly half the - // retries land before the requested time and invite escalation. Use - // positive-only jitter in that case so clients still spread but never - // dip below the server's hint. - // - // Exception: when retryAfterMs > maxDelayMs the base is already capped - // to maxDelayMs, so positive jitter would be erased by the final clamp - // below and every retry would land at exactly maxDelayMs — reintroducing - // the thundering herd we are trying to avoid. In that case the server - // contract is already unsatisfiable, so fall back to symmetric jitter - // to preserve spread. - // Use `<=` so the `retryAfterMs === maxDelayMs` boundary keeps the - // positive-jitter contract. At the boundary, positive jitter followed by - // the final clamp collapses every retry to exactly maxDelayMs — clients - // do land in lockstep at that instant, which is thundering-herd-shaped - // locally. The trade-off is deliberate: symmetric jitter at the boundary - // would schedule roughly half the retries below maxDelayMs (= - // retryAfterMs), which is a *Retry-After contract violation* and invites - // upstream escalation (429 → extended cooldown / bans on Telegram, - // Discord, etc.). A synchronized retry at the exact server-cleared - // instant is strictly preferable to a spread that undercuts the server's - // hint. Only switch to symmetric when the hint exceeds our local cap - // (`retryAfterMs > maxDelayMs`), where the contract is already - // unsatisfiable and we gain spread without adding a violation. - const canHonorRetryAfter = - hasRetryAfter && typeof retryAfterMs === "number" && retryAfterMs <= delayCap; - // Full jitter's upward draw is inherently positive, so it serves both - // plain backoff and honorable Retry-After floors; only the unsatisfiable - // over-cap hint must switch to the symmetric downward spread. Numeric - // jitter keeps the original positive/symmetric split. - const overCapRetryAfter = hasRetryAfter && !canHonorRetryAfter; - const wantsPositiveDraw = jitter === "full" ? !overCapRetryAfter : canHonorRetryAfter; - delay = applyJitter(delay, jitter, wantsPositiveDraw ? "positive" : "symmetric", random); - delay = Math.min(Math.max(delay, minDelayMs), delayCap); - - options.onRetry?.({ - attempt, - maxAttempts, - delayMs: delay, - err, - label: options.label, - }); - if (delay > 0) { - await sleep(delay); - } - } - } - - throw createRetryFailure(attemptErrors); -} +/** Runs an async operation until it succeeds, policy stops, or attempts are exhausted. */ +export const retryAsync = createRetryRunner({ + random: generateSecureFraction, + createFailure: createRetryFailure, +}); diff --git a/src/media/store.ts b/src/media/store.ts index ea3977ac4f7..af9c87bec04 100644 --- a/src/media/store.ts +++ b/src/media/store.ts @@ -22,6 +22,7 @@ import { sanitizeUntrustedFileName } from "../infra/fs-safe-advanced.js"; import { isPathInside } from "../infra/fs-safe.js"; import { retainSafeHeadersForCrossOriginRedirect } from "../infra/net/redirect-headers.js"; import { resolvePinnedHostname } from "../infra/net/ssrf.js"; +import { retryAsync } from "../infra/retry.js"; import { writeSiblingTempFile } from "../infra/sibling-temp-file.js"; import { resolveConfigDir } from "../utils.js"; import { isFsSafeError, readLocalFileSafely, type FsSafeLikeError } from "./store.runtime.js"; @@ -175,21 +176,26 @@ function isMissingPathError(err: unknown): boolean { } async function retryAfterRecreatingDir(dir: string, run: () => Promise): Promise { - try { - return await run(); - } catch (err) { - const noSpaceError = findErrorWithCode(err, "ENOSPC"); - if (noSpaceError) { - throw noSpaceError; - } - if (!isMissingPathError(err)) { - throw err; - } - // Recursive cleanup can prune an empty directory between mkdir and the later - // file open/write. Recreate once and retry the media write path. - await fs.mkdir(dir, { recursive: true, mode: 0o700 }); - return await run(); - } + return await retryAsync( + async () => { + try { + return await run(); + } catch (err) { + throw findErrorWithCode(err, "ENOSPC") ?? err; + } + }, + { + attempts: 2, + minDelayMs: 0, + maxDelayMs: 0, + shouldRetry: isMissingPathError, + onRetry: async () => { + // Cleanup can prune the directory between mkdir and file open. Recreate + // it once; further failures remain terminal instead of looping. + await fs.mkdir(dir, { recursive: true, mode: 0o700 }); + }, + }, + ); } // Maps the cleanup mode onto the prune sweep depth. The fs-safe prune walker keys descent off diff --git a/src/plugin-sdk/root-alias.cjs b/src/plugin-sdk/root-alias.cjs index a0c5b87cd8c..bb10eb6339a 100644 --- a/src/plugin-sdk/root-alias.cjs +++ b/src/plugin-sdk/root-alias.cjs @@ -110,6 +110,12 @@ const workspacePackageAliasEntries = { }, }, }, + "@openclaw/retry": { + dir: "retry", + subpaths: { + "": { srcFile: "src/index.ts", distFile: "dist/index.mjs" }, + }, + }, }; const workspacePackageAliases = Object.entries(workspacePackageAliasEntries).flatMap( ([name, pkg]) => diff --git a/src/plugins/contracts/extension-package-project-boundaries.test.ts b/src/plugins/contracts/extension-package-project-boundaries.test.ts index 363bdba6d8b..43ca302fabc 100644 --- a/src/plugins/contracts/extension-package-project-boundaries.test.ts +++ b/src/plugins/contracts/extension-package-project-boundaries.test.ts @@ -202,6 +202,7 @@ describe("opt-in extension package boundaries", () => { "../../packages/media-generation-core/src/**/*.ts", "../../packages/model-catalog-core/src/**/*.ts", "../../packages/normalization-core/src/**/*.ts", + "../../packages/retry/src/**/*.ts", "../../packages/acp-core/src/**/*.ts", "../../packages/terminal-core/src/**/*.ts", "../../src/plugin-sdk/**/*.ts", diff --git a/src/plugins/contracts/plugin-sdk-root-alias.test.ts b/src/plugins/contracts/plugin-sdk-root-alias.test.ts index 50fde771ca0..5067381d17a 100644 --- a/src/plugins/contracts/plugin-sdk-root-alias.test.ts +++ b/src/plugins/contracts/plugin-sdk-root-alias.test.ts @@ -519,6 +519,7 @@ describe("plugin-sdk root alias", () => { "src", "number-coercion.ts", ), + retry: path.join(packageRoot, "packages", "retry", "src", "index.ts"), }; const lazyModule = loadRootAliasWithStubs({ existingPaths: Object.values(sourcePaths), @@ -534,6 +535,7 @@ describe("plugin-sdk root alias", () => { expect(aliasMap["@openclaw/normalization-core/number-coercion"]).toBe( sourcePaths.numberCoercion, ); + expect(aliasMap["@openclaw/retry"]).toBe(sourcePaths.retry); }); it("keeps bootstrap plugin-sdk aliases deterministic and ignores unsafe subpaths", () => { @@ -593,6 +595,7 @@ describe("plugin-sdk root alias", () => { "@openclaw/normalization-core/string-coerce", "@openclaw/normalization-core/string-normalization", "@openclaw/normalization-core/utf16-slice", + "@openclaw/retry", "openclaw/plugin-sdk", "@openclaw/plugin-sdk", ]); diff --git a/src/plugins/sdk-alias.test.ts b/src/plugins/sdk-alias.test.ts index 86030650166..dc94d99a198 100644 --- a/src/plugins/sdk-alias.test.ts +++ b/src/plugins/sdk-alias.test.ts @@ -1570,6 +1570,12 @@ describe("plugin sdk alias helpers", () => { srcFile: "string-coerce.ts", distFile: "string-coerce.mjs", }); + const retry = writeWorkspacePackageEntry({ + root: fixture.root, + packageDir: "retry", + srcFile: "index.ts", + distFile: "index.mjs", + }); const markdownCore = writeWorkspacePackageEntry({ root: fixture.root, packageDir: "markdown-core", @@ -1634,6 +1640,7 @@ describe("plugin sdk alias helpers", () => { fs.rmSync(normalizationCore.distFile); fs.rmSync(normalizationBooleanCoercion.distFile); fs.rmSync(normalizationStringCoerce.distFile); + fs.rmSync(retry.distFile); fs.rmSync(terminalCore.distFile); fs.rmSync(terminalCoreTheme.distFile); fs.rmSync(netPolicy.distFile); @@ -1697,6 +1704,7 @@ describe("plugin sdk alias helpers", () => { expect(fs.realpathSync(aliases["@openclaw/normalization-core/string-coerce"] ?? "")).toBe( fs.realpathSync(normalizationStringCoerce.srcFile), ); + expect(fs.realpathSync(aliases["@openclaw/retry"] ?? "")).toBe(fs.realpathSync(retry.srcFile)); expect(fs.realpathSync(aliases["@openclaw/terminal-core"] ?? "")).toBe( fs.realpathSync(terminalCore.srcFile), ); @@ -1766,6 +1774,15 @@ describe("plugin sdk alias helpers", () => { ); mkdirSafeDir(path.dirname(normalizationCoreRootDistFile)); fs.writeFileSync(normalizationCoreRootDistFile, "export {};\n", "utf-8"); + writeWorkspacePackageEntry({ + root: fixture.root, + packageDir: "retry", + srcFile: "index.ts", + distFile: "index.mjs", + }); + const retryRootDistFile = path.join(fixture.root, "dist", "retry", "index.js"); + mkdirSafeDir(path.dirname(retryRootDistFile)); + fs.writeFileSync(retryRootDistFile, "export {};\n", "utf-8"); const markdownCore = writeWorkspacePackageEntry({ root: fixture.root, packageDir: "markdown-core", @@ -1824,6 +1841,9 @@ describe("plugin sdk alias helpers", () => { expect(fs.realpathSync(aliases["@openclaw/normalization-core/record-coerce"] ?? "")).toBe( fs.realpathSync(normalizationCoreRootDistFile), ); + expect(fs.realpathSync(aliases["@openclaw/retry"] ?? "")).toBe( + fs.realpathSync(retryRootDistFile), + ); expect(fs.realpathSync(aliases["@openclaw/terminal-core/links"] ?? "")).toBe( fs.realpathSync(terminalCoreRootDistFile), ); diff --git a/src/plugins/sdk-alias.ts b/src/plugins/sdk-alias.ts index bcf0812c2e3..eb9abd20dfc 100644 --- a/src/plugins/sdk-alias.ts +++ b/src/plugins/sdk-alias.ts @@ -824,6 +824,13 @@ const WORKSPACE_PACKAGE_ALIAS_ENTRIES: WorkspacePackageAliasEntry[] = [ srcFile: "utf16-slice.ts", distFile: "utf16-slice.mjs", }, + { + packageName: "@openclaw/retry", + packageDir: "retry", + subpath: "", + srcFile: "index.ts", + distFile: "index.mjs", + }, { packageName: "@openclaw/terminal-core", packageDir: "terminal-core", @@ -1053,6 +1060,7 @@ const ROOT_PACKAGED_WORKSPACE_PACKAGE_DIRS = new Set([ "acp-core", "media-core", "normalization-core", + "retry", "terminal-core", ]); diff --git a/test/scripts/release-wrapper-scripts.test.ts b/test/scripts/release-wrapper-scripts.test.ts index 061a6b38fe9..cc2facb931b 100644 --- a/test/scripts/release-wrapper-scripts.test.ts +++ b/test/scripts/release-wrapper-scripts.test.ts @@ -103,7 +103,7 @@ describe("release wrapper scripts", () => { ], oldTarget, ); - expect(plan.status).toBe(0); + expect(plan.status, plan.stderr).toBe(0); expect(JSON.parse(plan.stdout)).toMatchObject({ bootstrapWorkflowSha: "b".repeat(40), bootstrap: { ref: "main", shouldDispatch: false }, diff --git a/test/vitest-scoped-config.test.ts b/test/vitest-scoped-config.test.ts index 59670b863bb..8b7ff91d71b 100644 --- a/test/vitest-scoped-config.test.ts +++ b/test/vitest-scoped-config.test.ts @@ -176,6 +176,10 @@ describe("resolveVitestIsolation", () => { replacement: path.join(process.cwd(), "packages", "acp-core", "src", "runtime", "types.ts"), }, ); + expect(findAlias(sharedVitestConfig.resolve.alias, "@openclaw/retry")).toEqual({ + find: "@openclaw/retry", + replacement: path.join(process.cwd(), "packages", "retry", "src", "index.ts"), + }); }); it("defaults shared scoped configs to the non-isolated runner", () => { diff --git a/test/vitest/vitest.shared.config.ts b/test/vitest/vitest.shared.config.ts index 7569455ca15..ec6ffd67350 100644 --- a/test/vitest/vitest.shared.config.ts +++ b/test/vitest/vitest.shared.config.ts @@ -471,6 +471,7 @@ export const sharedVitestConfig = { sourcePackageAlias("media-core", "mime"), sourcePackageAlias("media-core", "read-byte-stream-with-limit"), sourcePackageAlias("media-core"), + sourcePackageAlias("retry"), ...sourcePackageAliasesFromExports("acp-core", acpCorePackageJson.exports), ...sourcePluginSdkSubpaths.map((subpath) => ({ find: `openclaw/plugin-sdk/${subpath}`, diff --git a/tsconfig.json b/tsconfig.json index 09a522a694e..0359f0f44db 100644 --- a/tsconfig.json +++ b/tsconfig.json @@ -155,6 +155,7 @@ "./packages/normalization-core/src/utf16-slice.ts" ], "@openclaw/normalization-core/*": ["./packages/normalization-core/src/*"], + "@openclaw/retry": ["./packages/retry/src/index.ts"], "@openclaw/acp-core": ["./packages/acp-core/src/index.ts"], "@openclaw/acp-core/meta": ["./packages/acp-core/src/meta.ts"], "@openclaw/acp-core/numeric-options": ["./packages/acp-core/src/numeric-options.ts"], diff --git a/tsconfig.plugin-sdk.dts.json b/tsconfig.plugin-sdk.dts.json index b4373a7b4c8..5431cd6c1d5 100644 --- a/tsconfig.plugin-sdk.dts.json +++ b/tsconfig.plugin-sdk.dts.json @@ -21,6 +21,7 @@ "packages/model-catalog-core/src/**/*.ts", "packages/memory-host-sdk/src/**/*.ts", "packages/normalization-core/src/**/*.ts", + "packages/retry/src/**/*.ts", "packages/acp-core/src/**/*.ts", "packages/terminal-core/src/**/*.ts", "src/video-generation/dashscope-compatible.ts", diff --git a/tsdown.config.ts b/tsdown.config.ts index c83b6197d7c..8002185b91c 100644 --- a/tsdown.config.ts +++ b/tsdown.config.ts @@ -224,6 +224,7 @@ function shouldAlwaysBundleDependency(id: string): boolean { id.startsWith("@openclaw/fs-safe/") || id === "@openclaw/normalization-core" || id.startsWith("@openclaw/normalization-core/") || + id === "@openclaw/retry" || id === "@openclaw/media-core" || id.startsWith("@openclaw/media-core/") || id === "@openclaw/acp-core" || @@ -452,6 +453,12 @@ function buildNormalizationCoreDistEntries(): Record { }; } +function buildRetryDistEntries(): Record { + return { + index: "packages/retry/src/index.ts", + }; +} + function buildMediaCoreDistEntries(): Record { return { index: "packages/media-core/src/index.ts", @@ -627,6 +634,9 @@ function buildUnifiedDistEntries(): Record { source, ]), ), + ...Object.fromEntries( + Object.entries(buildRetryDistEntries()).map(([entry, source]) => [`retry/${entry}`, source]), + ), ...Object.fromEntries( Object.entries(buildMediaCoreDistEntries()).map(([entry, source]) => [ `media-core/${entry}`, @@ -731,6 +741,12 @@ const configs = [ entry: buildNormalizationCoreDistEntries(), outDir: tsdownPackageOutputRoot("normalization-core"), }), + nodeWorkspacePackageBuildConfig({ + clean: true, + dts: TSDOWN_DECLARATIONS, + entry: buildRetryDistEntries(), + outDir: tsdownPackageOutputRoot("retry"), + }), nodeWorkspacePackageBuildConfig({ clean: true, dts: TSDOWN_DECLARATIONS,