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.
InstallationDirect link to Installation
- npm
- pnpm
- Yarn
- Bun
npm install @mastra/redis-streams
pnpm add @mastra/redis-streams
yarn add @mastra/redis-streams
bun add @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 exampleDirect link to Usage example
Provide a Redis connection URL.
import { Mastra } from '@mastra/core'
import { RedisStreamsPubSub } from '@mastra/redis-streams'
export const mastra = new Mastra({
pubsub: new RedisStreamsPubSub({
url: 'redis://localhost:6379',
}),
})
Redis ClusterDirect 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.
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 clientDirect 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.
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 parametersDirect link to Constructor parameters
url?:
redisOptions.url.keyPrefix?:
<keyPrefix>:<topic>.blockMs?:
redisOptions?:
redis client for advanced configuration.cluster?:
createCluster() from redis. Mutually exclusive with url, redisOptions, and client.client?:
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?:
streamIdleTtlMs?:
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?:
reclaimIdleMs?:
maxDeliveryAttempts?:
nack before it is dropped. Pass Infinity to disable the cap.inFlightTimeoutMs?:
deliveryAttempt, so maxDeliveryAttempts still applies. Set this to recover hung handlers in single-consumer groups. 0 disables it.logger?:
PropertiesDirect link to Properties
supportedModes:
["pull"].supportsOffsets:
false. Redis stream anchors support start positions, but not the numeric indexes used by subscribeFromOffset().MethodsDirect 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 reclaimDirect 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 leasingDirect 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.