Skip to main content

RedisStreamsPubSub

RedisStreamsPubSub is a PubSub implementation backed by Redis Streams. It delivers events across processes and hosts, with persistence, consumer groups, and redelivery on failure. It also implements LeaseProvider. The signals layer can elect a single owner per resource across instances, which is what lets signals coordinate runs in distributed and serverless deployments.

Use it for distributed deployments where several services share an event stream. For single-process delivery, use EventEmitterPubSub. For Google Cloud, use GoogleCloudPubSub.

Each topic maps to a Redis stream key. Subscriptions with a group use a Redis consumer group. Members share work round-robin. Subscriptions without a group create a private consumer group. Every subscriber receives every event.

RedisStreamsPubSub is a pull transport: consumers read events with XREADGROUP, so Mastra runs an orchestration worker to read on its behalf.

Installation
Direct link to Installation

npm install @mastra/redis-streams

Requires Redis 7.0 or later. The reclaim loop relies on XCLAIM removing trimmed entries from the pending list, which earlier versions don't do.

Usage example
Direct link to Usage example

Provide a Redis connection URL.

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { RedisStreamsPubSub } from '@mastra/redis-streams'

export const mastra = new Mastra({
pubsub: new RedisStreamsPubSub({
url: 'redis://localhost:6379',
}),
})

Redis Cluster
Direct link to Redis Cluster

Pass cluster instead of url/redisOptions to connect to a Redis Cluster (for example, AWS ElastiCache with cluster mode enabled). The options are forwarded to createCluster() from redis.

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { RedisStreamsPubSub } from '@mastra/redis-streams'

export const mastra = new Mastra({
pubsub: new RedisStreamsPubSub({
cluster: {
rootNodes: [{ url: 'redis://node-1:6379' }, { url: 'redis://node-2:6379' }],
},
}),
})

Bring your own client
Direct link to Bring your own client

To control client construction yourself (TLS, credential providers, and so on), pass an unconnected redis client as client. The pubsub uses it as its writer and calls client.duplicate() for each subscription's blocking reader. Every connection stays distinct. The pubsub owns the client's lifecycle (it connects on first use and quits it on close()), so don't share it with the rest of your app.

src/mastra/index.ts
import { RedisStreamsPubSub } from '@mastra/redis-streams'
import { createClient } from 'redis'

const pubsub = new RedisStreamsPubSub({
client: createClient({ url: process.env.REDIS_URL, socket: { tls: true } }),
})

Constructor parameters
Direct link to Constructor parameters

url?:

string
= redis://localhost:6379
Redis connection URL. Falls back to redisOptions.url.

keyPrefix?:

string
= mastra:topic
Prefix for stream keys. Each topic maps to <keyPrefix>:<topic>.

blockMs?:

number
= 1000
How long, in milliseconds, each read blocks while waiting for new events.

redisOptions?:

RedisClientOptions
Options passed to the underlying redis client for advanced configuration.

cluster?:

RedisClusterOptions
Connect to a Redis Cluster. Options are passed to createCluster() from redis. Mutually exclusive with url, redisOptions, and client.

client?:

RedisClientType | RedisClusterType
A pre-configured, unconnected redis client (standalone or cluster) used as the writer; readers are created with client.duplicate(). The pubsub owns its lifecycle (connects lazily, quits on close()). Mutually exclusive with url, redisOptions, and cluster.

maxStreamLength?:

number
= 10000
Approximate maximum number of entries kept per stream. Set to 0 to disable trimming.

streamIdleTtlMs?:

number
= 0
Idle expiry in milliseconds: a sliding TTL refreshed on every write to a stream (publish, subscribe creating its consumer group, nack retry). Each write resets it, so an actively-written stream never expires mid-flight; a stream left idle for the full duration is deleted by Redis automatically, together with its entries and consumer groups. Only writes refresh the TTL โ€” a consumer slowly draining a backlog does not โ€” so keep it well above the longest expected gap between writes on a live topic, and above the longest a grouped worker may be offline while still expecting to replay its backlog on restart. For topics with a clear end of life (workflow and agent runs), clearTopic deletes the stream eagerly and the TTL only covers streams that never reach that call, such as a crashed run. For open-ended topics that are never cleared (for example, per-conversation streams), the TTL is the only mechanism that reclaims their memory โ€” with it disabled they stay in Redis forever, so production deployments should set it (for example, 7 days). Must be a non-negative integer. Defaults to 0 (disabled).

reclaimIntervalMs?:

number
= 30000
How often, in milliseconds, a subscription reclaims events that an earlier consumer read but never acknowledged. Set to 0 to disable.

reclaimIdleMs?:

number
= 60000
Minimum idle time, in milliseconds, before a pending event is eligible for reclaim. Keep this well above typical processing time to avoid double delivery.

maxDeliveryAttempts?:

number
= 5
Maximum times an event is redelivered through nack before it is dropped. Pass Infinity to disable the cap.

inFlightTimeoutMs?:

number
= 0
How long a handler may hold an event without acking or nacking before the reclaim loop nacks it on its behalf. The event is republished with an incremented deliveryAttempt, so maxDeliveryAttempts still applies. Set this to recover hung handlers in single-consumer groups. 0 disables it.

logger?:

{ debug?: Function; warn?: Function }
Optional logger for diagnostics. When omitted, suppressed errors are silent.

Properties
Direct link to Properties

supportedModes:

ReadonlyArray<"pull" | "push">
Returns ["pull"].

supportsOffsets:

boolean
Returns false. Redis stream anchors support start positions, but not the numeric indexes used by subscribeFromOffset().

Methods
Direct link to Methods

RedisStreamsPubSub implements the PubSub contract. The methods below have behavior specific to this implementation.

subscribe(topic, cb, options?)
Direct link to subscribetopic-cb-options

Subscribes to a topic. With options.group, members of the group share events through a Redis consumer group. Without a group, the subscriber receives every event through a private consumer group.

Set options.startFrom to "latest" to skip retained entries when Redis creates the consumer group. The default, "earliest", reads retained entries first. This option doesn't reset an existing group's checkpoint.

await pubsub.subscribe(
'workflow.events',
(event, ack, nack) => {
console.log(event)
},
{ group: 'live-workers', startFrom: 'latest' },
)

flush()
Direct link to flush

Waits for in-flight publishes to complete.

await pubsub.flush()

clearTopic(topic)
Direct link to cleartopictopic

Deletes a topic's stream and every consumer group on it, freeing the memory a finished topic would otherwise hold. Mastra's run lifecycles (durable agents and the evented workflow engine) call this automatically when a run reaches a terminal state. Call it yourself only once nothing will read the topic again. It's best-effort and never throws. Failures are logged at warn level. A subscriber still attached when the stream is deleted recovers on its own but misses the deleted entries.

Automatic cleanup requires @mastra/core and @mastra/redis-streams versions that both support clearTopic: the runtime routes the call through its caching layer, so upgrade the two packages together to get end-of-run stream deletion.

Not every topic reaches a clearTopic call. Topics without a defined end of life (per-conversation streams, application feeds) are never cleared, and a stream that misses its cleanup (for example, a crashed run) is never retried. Configure the streamIdleTtlMs sliding TTL to reclaim those streams once they go idle.

await pubsub.clearTopic('workflow.events.run-123')

trimTopic(topic, { runId, producedBefore? })
Direct link to trimtopictopic--runid-producedbefore-

Pages through the stream with XRANGE and runs XDEL for every entry published with runId. With producedBefore, only entries produced at or before that time are deleted, skipping approval and suspension prompts. Mastra uses this to trim a long run each time its messages are saved. Mastra calls this once a thread run's messages are persisted, so an idle thread's stream empties instead of growing to maxLen. Because entries are matched by runId, a run resumed after a restart also removes the entries it published before the restart. Entries from other runs, including those from other processes, are never touched. Best-effort. Failures are logged at warn level.

await pubsub.trimTopic('agent.thread.thread-abc', { runId: 'run-123' })

close()
Direct link to close

Closes the Redis connections and stops all subscriptions. Call this during graceful shutdown.

await pubsub.close()

Redelivery and reclaim
Direct link to Redelivery and reclaim

When a subscriber calls nack, the event is republished with an incremented deliveryAttempt and the original is acknowledged. Once an event reaches maxDeliveryAttempts, it's dropped instead of redelivered. Separately, each subscription periodically reclaims events that an earlier consumer in the group read but never acknowledged, controlled by reclaimIntervalMs and reclaimIdleMs.

A subscription never reclaims an event its own handler is still processing, so a slow handler isn't invoked twice for the same event. If a handler hangs, a different consumer in the group reclaims the event once it has been idle for reclaimIdleMs. In a single-consumer group there's no sibling to do that, so set inFlightTimeoutMs to have the subscription nack the event itself after that long.

Distributed leasing
Direct link to Distributed leasing

RedisStreamsPubSub implements the LeaseProvider contract on top of the same Redis connection. The signals runtime uses it to elect a single owner (usually per thread key) so that across instances only one process wakes and runs the agent, and others route follow-up work to the holder. This is what makes signals work on serverless and multi-instance deployments; without a shared lease, each instance would start its own competing run.

Lease keys are namespaced under the same keyPrefix as topics, as <keyPrefix>:lease:<key>. All operations are atomic: acquireLease uses SET NX PX and refreshes its own TTL idempotently, while releaseLease, renewLease, and transferLease use Lua scripts that check ownership before mutating, so a concurrent renewal from another owner is never clobbered.

You don't call these methods directly. Configuring RedisStreamsPubSub as the pubsub backend is enough for the runtime to detect and use the capability. See LeaseProvider for the full method contract.