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();