Skip to content

Producer reads one chunk per Redis round trip while a client is rejoined #52

Description

@nikolailehbrink

Version: resumable-stream 2.2.13 (redis 6.3.0, Upstash)

Problem

In createNewResumableStream, the producer's read loop awaits a publish to every rejoined listener before it reads the next chunk:

for (const listenerId of listenerChannels) {
  promises.push(ctx.publisher.publish(`${ctx.keyPrefix}:chunk:${listenerId}`, value));
}
await Promise.all(promises);
read();

Once a client rejoins with resumeExistingStream, the producer reads at most one chunk per Redis round trip. The source stream is usually a tee of an AI SDK UI message stream, so backpressure slows the whole pipeline: generation, onFinish persistence, and reacting to an abort signal. With about 55 ms round trip time (a laptop talking to a regional Upstash database), a rejoined reply streamed about 350 characters in 3 seconds instead of about 2000, and a stop requested meanwhile only took effect after the buffered backlog was drained.

Stale listeners (#48, #49) make this worse, because every chunk then waits on publishes to channels nobody reads anymore.

Expected

Reading the source does not wait for publishes to listeners. One publisher connection already keeps publishes in order, so awaiting each publish is not needed for ordering.

Suggested fix

Publish without awaiting and handle rejections:

for (const listenerId of listenerChannels) {
  ctx.publisher.publish(`${ctx.keyPrefix}:chunk:${listenerId}`, value).catch((e) => {
    if (isDebug()) console.error(e);
  });
}
read();

We apply this as a patch on 2.2.13 and are happy to open a PR.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions