From 1f7101844b53c9f8be709b2b462f2bca4ca66fcc Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Tue, 25 Aug 2026 23:53:00 -0700 Subject: [PATCH] fix(memory): recover embedding batches exceeding provider item limits (#129927) Co-authored-by: shin4141 <128954611+shin4141@users.noreply.github.com> --- .../manager-embedding-ops.retry.test.ts | 63 +++++++++++++++++++ .../src/memory/manager-embedding-ops.ts | 6 +- .../memory/manager-embedding-policy.test.ts | 50 ++++++++++----- .../src/memory/manager-embedding-policy.ts | 19 +++--- 4 files changed, 109 insertions(+), 29 deletions(-) diff --git a/extensions/memory-core/src/memory/manager-embedding-ops.retry.test.ts b/extensions/memory-core/src/memory/manager-embedding-ops.retry.test.ts index 0eddd2cd4883..28c922666eab 100644 --- a/extensions/memory-core/src/memory/manager-embedding-ops.retry.test.ts +++ b/extensions/memory-core/src/memory/manager-embedding-ops.retry.test.ts @@ -31,6 +31,20 @@ function createEmbeddingQueryRetryHarness( }) as EmbeddingQueryRetryHarness; } +function createEmbeddingBatchRetryHarness(embedBatch: EmbeddingProvider["embedBatch"]) { + const manager = Object.assign( + createEmbeddingQueryRetryHarness(async () => []), + { + waitForEmbeddingRetry: vi.fn(async () => {}), + }, + ) as EmbeddingQueryRetryHarness & { + embedBatchWithRetry: (texts: string[]) => Promise; + waitForEmbeddingRetry: ReturnType; + }; + manager.provider.embedBatch = embedBatch; + return manager; +} + describe("memory embedding query retry cancellation", () => { afterEach(() => { vi.useRealTimers(); @@ -109,3 +123,52 @@ describe("memory embedding query retry cancellation", () => { expect(vi.getTimerCount()).toBe(0); }); }); + +describe("memory embedding batch retry boundary", () => { + it.each([ + [ + "explicit maximum and actual input counts", + (count: number) => + `Embeddings API input limit exceeded: max 10, got ${count}. Request id: fixture-000597000`, + ], + ["an explicit maximum input length", () => "embeddings max input length is 10"], + ])( + "splits provider errors with %s without retrying oversized requests", + async (_label, error) => { + const items = Array.from({ length: 33 }, (_, index) => `item-${index}`); + const embedBatch = vi.fn(async (texts: string[]) => { + if (texts.length > 10) { + throw new Error(`openai-compatible embeddings failed: HTTP 400: ${error(texts.length)}`); + } + return texts.map((text) => [Number.parseInt(text.slice(5), 10)]); + }); + const manager = createEmbeddingBatchRetryHarness(embedBatch); + + await expect(manager.embedBatchWithRetry(items)).resolves.toEqual( + items.map((_, index) => [index]), + ); + expect(embedBatch.mock.calls.map(([texts]) => texts.length)).toEqual([ + 33, 17, 9, 8, 16, 8, 8, + ]); + expect(manager.waitForEmbeddingRetry).not.toHaveBeenCalled(); + expect(manager.markLocalEmbeddingProviderDegraded).not.toHaveBeenCalled(); + }, + ); + + it("does not retry or split generic input validation errors containing request-id digits", async () => { + const embedBatch = vi.fn(async () => { + throw new Error( + 'openai-compatible embeddings failed: HTTP 400: {"error":{"code":"InvalidParameter","message":"The parameter input specified in the request is not valid. Request id: fixture-000597000","param":"input"}}', + ); + }); + const manager = createEmbeddingBatchRetryHarness(embedBatch); + + await expect(manager.embedBatchWithRetry(["one", "two"])).rejects.toMatchObject({ + code: "MEMORY_EMBEDDING_OPERATION_FAILED", + operation: "batch", + }); + expect(embedBatch).toHaveBeenCalledOnce(); + expect(manager.waitForEmbeddingRetry).not.toHaveBeenCalled(); + expect(manager.markLocalEmbeddingProviderDegraded).toHaveBeenCalledOnce(); + }); +}); diff --git a/extensions/memory-core/src/memory/manager-embedding-ops.ts b/extensions/memory-core/src/memory/manager-embedding-ops.ts index 3d556aa82c9b..d2dd3bf41a8e 100644 --- a/extensions/memory-core/src/memory/manager-embedding-ops.ts +++ b/extensions/memory-core/src/memory/manager-embedding-ops.ts @@ -51,7 +51,7 @@ import { buildTextEmbeddingInputs, filterNonEmptyMemoryChunks, isRetryableMemoryEmbeddingError, - isSplittableMemoryEmbeddingTransportError, + isSplittableMemoryEmbeddingBatchError, resolveMemoryEmbeddingRetryDelay, runMemoryEmbeddingBatchRetryWithSplit, runMemoryEmbeddingRetryLoop, @@ -679,7 +679,7 @@ export abstract class MemoryManagerEmbeddingOps extends MemoryManagerSyncOps { return result; }, isRetryable: isRetryableMemoryEmbeddingError, - isSplittable: isSplittableMemoryEmbeddingTransportError, + isSplittable: isSplittableMemoryEmbeddingBatchError, waitForRetry: async (delayMs) => { await this.waitForEmbeddingRetry( delayMs, @@ -690,7 +690,7 @@ export abstract class MemoryManagerEmbeddingOps extends MemoryManagerSyncOps { baseDelayMs: EMBEDDING_RETRY_BASE_DELAY_MS, onSplit: ({ itemCount, splitAt }) => { log.warn( - `memory embeddings transport failed after retries; splitting ${label} of ${itemCount} into ${splitAt} + ${itemCount - splitAt}`, + `memory embeddings ${label} failed; splitting ${itemCount} inputs into ${splitAt} + ${itemCount - splitAt}`, ); }, }), diff --git a/extensions/memory-core/src/memory/manager-embedding-policy.test.ts b/extensions/memory-core/src/memory/manager-embedding-policy.test.ts index 2a5f59b97b5a..470bef5a2d21 100644 --- a/extensions/memory-core/src/memory/manager-embedding-policy.test.ts +++ b/extensions/memory-core/src/memory/manager-embedding-policy.test.ts @@ -5,7 +5,7 @@ import { buildMemoryEmbeddingBatches, filterNonEmptyMemoryChunks, isRetryableMemoryEmbeddingError, - isSplittableMemoryEmbeddingTransportError, + isSplittableMemoryEmbeddingBatchError, resolveMemoryEmbeddingRetryDelay, runMemoryEmbeddingBatchRetryWithSplit, runMemoryEmbeddingRetryLoop, @@ -168,20 +168,40 @@ describe("memory embedding policy", () => { for (const message of splittableMessages) { expect(isRetryableMemoryEmbeddingError(message)).toBe(true); - expect(isSplittableMemoryEmbeddingTransportError(message)).toBe(true); + expect(isSplittableMemoryEmbeddingBatchError(message)).toBe(true); } expect(isRetryableMemoryEmbeddingError("ECONNREFUSED")).toBe(true); - expect(isSplittableMemoryEmbeddingTransportError("ECONNREFUSED")).toBe(false); + expect(isSplittableMemoryEmbeddingBatchError("ECONNREFUSED")).toBe(false); expect(isRetryableMemoryEmbeddingError("EHOSTUNREACH")).toBe(true); - expect(isSplittableMemoryEmbeddingTransportError("EHOSTUNREACH")).toBe(false); + expect(isSplittableMemoryEmbeddingBatchError("EHOSTUNREACH")).toBe(false); expect(isRetryableMemoryEmbeddingError("memory embeddings batch timed out")).toBe(true); - expect(isSplittableMemoryEmbeddingTransportError("memory embeddings batch timed out")).toBe( - false, - ); + expect(isSplittableMemoryEmbeddingBatchError("memory embeddings batch timed out")).toBe(false); expect(isRetryableMemoryEmbeddingError("worker terminated by user")).toBe(false); expect(isRetryableMemoryEmbeddingError("embedding validation failed")).toBe(false); }); + it("recognizes only provider errors with an explicit numeric embedding item limit", () => { + for (const message of [ + "Embeddings API input limit exceeded: max 10, got 33. Request id: fixture-000597000", + "embeddings max input length is 16", + ]) { + expect(isSplittableMemoryEmbeddingBatchError(message)).toBe(true); + expect(isRetryableMemoryEmbeddingError(message)).toBe(false); + } + + for (const message of [ + "embedding input exceeds maximum token length 4096", + "embeddings max input length is unknown", + "Embeddings API input limit exceeded", + 'HTTP 400: {"code":"InvalidParameter","param":"input","message":"input must be a string"}', + ]) { + expect(isSplittableMemoryEmbeddingBatchError(message)).toBe(false); + } + expect(isRetryableMemoryEmbeddingError("HTTP 400: request id fixture-000597000")).toBe(false); + expect(isRetryableMemoryEmbeddingError("HTTP 429: rate limit")).toBe(true); + expect(isRetryableMemoryEmbeddingError("HTTP 503: service unavailable")).toBe(true); + }); + it("splits OpenAI 431 oversized embedding batches without retrying the same request", async () => { const run = vi.fn(async (items: string[]) => { if (items.length > 1) { @@ -196,7 +216,7 @@ describe("memory embedding policy", () => { items: ["a", "b", "c", "d"], run, isRetryable: isRetryableMemoryEmbeddingError, - isSplittable: isSplittableMemoryEmbeddingTransportError, + isSplittable: isSplittableMemoryEmbeddingBatchError, waitForRetry: async () => {}, maxAttempts: 3, baseDelayMs: 500, @@ -205,10 +225,10 @@ describe("memory embedding policy", () => { expect(result).toEqual([[97], [98], [99], [100]]); expect(run.mock.calls.map(([items]) => items.length)).toEqual([4, 2, 1, 1, 2, 1, 1]); expect(isRetryableMemoryEmbeddingError("431 request_headers_too_large")).toBe(false); - expect(isSplittableMemoryEmbeddingTransportError("431 request_headers_too_large")).toBe(true); - expect( - isSplittableMemoryEmbeddingTransportError("embedding validation failed at item 4312"), - ).toBe(false); + expect(isSplittableMemoryEmbeddingBatchError("431 request_headers_too_large")).toBe(true); + expect(isSplittableMemoryEmbeddingBatchError("embedding validation failed at item 4312")).toBe( + false, + ); }); it("retries too-many-tokens-per-day errors", async () => { @@ -272,7 +292,7 @@ describe("memory embedding policy", () => { items: ["a", "b", "c", "d"], run, isRetryable: isRetryableMemoryEmbeddingError, - isSplittable: isSplittableMemoryEmbeddingTransportError, + isSplittable: isSplittableMemoryEmbeddingBatchError, waitForRetry: async (delayMs) => { waits.push(delayMs); }, @@ -299,7 +319,7 @@ describe("memory embedding policy", () => { items: ["a", "b"], run, isRetryable: isRetryableMemoryEmbeddingError, - isSplittable: isSplittableMemoryEmbeddingTransportError, + isSplittable: isSplittableMemoryEmbeddingBatchError, waitForRetry: async () => {}, maxAttempts: 1, baseDelayMs: 500, @@ -318,7 +338,7 @@ describe("memory embedding policy", () => { items: ["a", "b"], run, isRetryable: isRetryableMemoryEmbeddingError, - isSplittable: isSplittableMemoryEmbeddingTransportError, + isSplittable: isSplittableMemoryEmbeddingBatchError, waitForRetry: async () => {}, maxAttempts: 2, baseDelayMs: 500, diff --git a/extensions/memory-core/src/memory/manager-embedding-policy.ts b/extensions/memory-core/src/memory/manager-embedding-policy.ts index 4530b4f6636a..c6b7cf3c1c4b 100644 --- a/extensions/memory-core/src/memory/manager-embedding-policy.ts +++ b/extensions/memory-core/src/memory/manager-embedding-policy.ts @@ -83,26 +83,23 @@ export function buildMemoryEmbeddingBatches( } const RETRYABLE_MEMORY_EMBEDDING_SERVICE_ERROR_RE = - /(rate[_ ]limit|too many requests|429|resource has been exhausted|5\d\d|cloudflare|tokens per day)/i; + /(rate[_ ]limit|too many requests|\b(?:429|5\d\d)\b|resource has been exhausted|cloudflare|tokens per day)/i; const RETRYABLE_MEMORY_EMBEDDING_TRANSPORT_ERROR_RE = /(fetch failed|other side closed|ECONNRESET|ECONNREFUSED|ETIMEDOUT|EPIPE|UND_ERR_|socket hang up|socket terminated|network error|read ECONN|timed out|connection (?:reset|refused|aborted|timed out)|EHOSTUNREACH|ENETUNREACH|ECONNABORTED|EAI_AGAIN)/i; -const SPLITTABLE_MEMORY_EMBEDDING_TRANSPORT_ERROR_RE = - /(request_headers_too_large|request header fields too large|other side closed|ECONNRESET|EPIPE|UND_ERR_SOCKET|socket hang up|socket terminated|read ECONN|connection (?:reset|aborted))/i; +const SPLITTABLE_MEMORY_EMBEDDING_BATCH_ERROR_RE = + /(request_headers_too_large|request header fields too large|other side closed|ECONNRESET|EPIPE|UND_ERR_SOCKET|socket hang up|socket terminated|read ECONN|connection (?:reset|aborted)|\bembeddings (?:api input limit exceeded:\s*max\s+\d+\s*,\s*got\s+\d+|max input length is\s+\d+)\b)/i; -function isRetryableMemoryEmbeddingTransportError(message: string): boolean { - return RETRYABLE_MEMORY_EMBEDDING_TRANSPORT_ERROR_RE.test(message); -} - -export function isSplittableMemoryEmbeddingTransportError(message: string): boolean { - return SPLITTABLE_MEMORY_EMBEDDING_TRANSPORT_ERROR_RE.test(message); +export function isSplittableMemoryEmbeddingBatchError(message: string): boolean { + return SPLITTABLE_MEMORY_EMBEDDING_BATCH_ERROR_RE.test(message); } export function isRetryableMemoryEmbeddingError(message: string): boolean { return ( - RETRYABLE_MEMORY_EMBEDDING_SERVICE_ERROR_RE.test(message) || - isRetryableMemoryEmbeddingTransportError(message) + RETRYABLE_MEMORY_EMBEDDING_TRANSPORT_ERROR_RE.test(message) || + (!isSplittableMemoryEmbeddingBatchError(message) && + RETRYABLE_MEMORY_EMBEDDING_SERVICE_ERROR_RE.test(message)) ); }