From e67b2343f3223a17045bcdd97b5818adf1d4ae29 Mon Sep 17 00:00:00 2001 From: Justin Date: Sat, 14 Feb 2026 01:53:00 -0500 Subject: [PATCH] fix: replace dedup resolver closure chain with array to prevent hangs MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The previous implementation chained waiters by patching entry.resolve via closures. This was fragile and hard to reason about. More critically, removeInflight() (called on client disconnect or request error) deleted the inflight entry without resolving pending waiters — causing them to hang forever with no response. Changes: - Replace closure-chain pattern with a simple resolvers array - On complete(): iterate all resolvers and resolve each one - On removeInflight(): resolve waiters with 503 error instead of leaving them hanging, so clients can retry independently Co-Authored-By: Claude Opus 4.6 --- src/dedup.ts | 46 ++++++++++++++++++++++++++-------------------- 1 file changed, 26 insertions(+), 20 deletions(-) diff --git a/src/dedup.ts b/src/dedup.ts index 1abf82c0..1d6a723e 100644 --- a/src/dedup.ts +++ b/src/dedup.ts @@ -15,8 +15,7 @@ export type CachedResponse = { }; type InflightEntry = { - resolve: (result: CachedResponse) => void; - waiters: Promise[]; + resolvers: Array<(result: CachedResponse) => void>; }; const DEFAULT_TTL_MS = 30_000; // 30 seconds @@ -109,27 +108,15 @@ export class RequestDeduplicator { getInflight(key: string): Promise | undefined { const entry = this.inflight.get(key); if (!entry) return undefined; - const promise = new Promise((resolve) => { - // Will be resolved when the original request completes - entry.waiters.push( - new Promise((r) => { - const orig = entry.resolve; - entry.resolve = (result) => { - orig(result); - resolve(result); - r(result); - }; - }), - ); + return new Promise((resolve) => { + entry.resolvers.push(resolve); }); - return promise; } /** Mark a request as in-flight. */ markInflight(key: string): void { this.inflight.set(key, { - resolve: () => {}, - waiters: [], + resolvers: [], }); } @@ -142,16 +129,35 @@ export class RequestDeduplicator { const entry = this.inflight.get(key); if (entry) { - entry.resolve(result); + for (const resolve of entry.resolvers) { + resolve(result); + } this.inflight.delete(key); } this.prune(); } - /** Remove an in-flight entry on error (don't cache failures). */ + /** Remove an in-flight entry on error (don't cache failures). + * Also rejects any waiters so they can retry independently. */ removeInflight(key: string): void { - this.inflight.delete(key); + const entry = this.inflight.get(key); + if (entry) { + // Resolve waiters with a sentinel error response so they don't hang forever. + // Waiters will see a 503 and can retry on their own. + const errorBody = Buffer.from(JSON.stringify({ + error: { message: "Original request failed, please retry", type: "dedup_origin_failed" }, + })); + for (const resolve of entry.resolvers) { + resolve({ + status: 503, + headers: { "content-type": "application/json" }, + body: errorBody, + completedAt: Date.now(), + }); + } + this.inflight.delete(key); + } } /** Prune expired completed entries. */