From 1d49e3687ee21fa93ca4f63fd4491bfc47b3b00e Mon Sep 17 00:00:00 2001 From: jaychung Date: Tue, 8 Sep 2026 21:27:15 +0800 Subject: [PATCH] fix: set sentinel after subscribe to avoid race condition Previously, the sentinel key was set BEFORE subscribing to the request channel. This caused a race condition where: 1. Producer sets sentinel 2. Consumer detects sentinel exists 3. Consumer publishes request to request channel 4. Producer subscribes to request channel (TOO LATE - message lost!) This fix moves the sentinel set to AFTER the subscribe completes: 1. Producer subscribes to request channel 2. Producer sets sentinel 3. Consumer detects sentinel exists 4. Consumer publishes request (SUCCESS - producer receives it) Co-Authored-By: Claude Opus 4.5 --- src/__tests__/generic.test.ts | 31 +++++++++++++++++++++++++++++++ src/runtime.ts | 17 +++++++++-------- 2 files changed, 40 insertions(+), 8 deletions(-) diff --git a/src/__tests__/generic.test.ts b/src/__tests__/generic.test.ts index 51a5b8a..2ac0339 100644 --- a/src/__tests__/generic.test.ts +++ b/src/__tests__/generic.test.ts @@ -49,6 +49,37 @@ describe("generic interface", () => { vi.useRealTimers(); }); + it("should subscribe before making a new stream resumable", async () => { + const { publisher, subscriber } = createInMemoryPubSubForTesting(); + const operations: string[] = []; + const ctx = createResumableStreamContext({ + waitUntil: null, + publisher: { + ...publisher, + set: async (key, value, options) => { + operations.push("set"); + return publisher.set(key, value, options); + }, + }, + subscriber: { + ...subscriber, + subscribe: async (channel, callback) => { + const result = await subscriber.subscribe(channel, callback); + operations.push("subscribe"); + return result; + }, + }, + keyPrefix: "test-subscribe-before-sentinel", + }); + const { readable, writer } = createTestingStream(); + + const stream = await ctx.createNewResumableStream("test-stream", () => readable); + + expect(operations).toEqual(["subscribe", "set"]); + writer.close(); + await streamToBuffer(stream); + }); + it("should work with custom publisher/subscriber implementations", async () => { const { publisher, subscriber } = createInMemoryPubSubForTesting(); diff --git a/src/runtime.ts b/src/runtime.ts index b366d92..e0dfeb6 100644 --- a/src/runtime.ts +++ b/src/runtime.ts @@ -61,13 +61,8 @@ export function createResumableStreamContextFactory(defaults: _Private.RedisDefa makeStream: () => ReadableStream, skipCharacters?: number ): Promise | null> => { - const initPromise = Promise.all(initPromises); - await initPromise; - await ctx.publisher.set(`${ctx.keyPrefix}:sentinel:${streamId}`, "1", { - EX: 24 * 60 * 60, - }); return createNewResumableStream( - initPromise, + Promise.all(initPromises), ctx as CreateResumableStreamContext, streamId, makeStream @@ -144,8 +139,9 @@ async function createNewResumableStream( }) ); let isDone = false; - // This is ultimately racy if two requests for the same ID come at the same time. - // But this library is for the case where that would not happen. + // Subscribe to request channel BEFORE setting sentinel to avoid race condition. + // If we set sentinel first, a consumer might detect it and send a request + // before we've subscribed, causing the message to be lost. await ctx.subscriber.subscribe( `${ctx.keyPrefix}:request:${streamId}`, async (message: string) => { @@ -168,6 +164,11 @@ async function createNewResumableStream( } ); + // Set sentinel AFTER subscribing to ensure we're ready to receive requests + await ctx.publisher.set(`${ctx.keyPrefix}:sentinel:${streamId}`, "1", { + EX: 24 * 60 * 60, + }); + return new ReadableStream({ start(controller) { const stream = makeStream();