Pub/sub
Respire models subscriptions as IAsyncEnumerable<RespireMessage>. Leaving the loop and disposing the subscription performs cleanup—no delegate bookkeeping required.
For typed keyspace, keyevent, and Redis 8.8 subkey events, see Keyspace notifications.
Subscribe
SubscribeAsync returns once the server has acknowledged the SUBSCRIBE, so the subscription is live before the first message is published—no polling on PublishAsync's receiver count.
await using var subscription =
await redis.SubscribeAsync(["orders", "payments"], stoppingToken);
await foreach (RespireMessage message in
subscription.WithCancellation(stoppingToken))
{
if (message.Kind == RespireMessageKind.Gap)
{
Console.Error.WriteLine($"Delivery gap: {message.Gap}. Reload authoritative state before applying more messages.");
continue;
}
Console.WriteLine($"{message.Channel}: {message.Text}");
}
Pattern and Redis 7 sharded subscriptions use dedicated entry points:
await using var patterns = await redis.SubscribePatternAsync("events:*");
await using var shard = await redis.SubscribeShardedAsync("events:eu-west");
Messages are buffered from the moment the subscription is acknowledged. The configured capacity and overflow policy apply even before enumeration starts.
A subscription is a single-consumer stream: only one enumerator may be active at a time. Dispose
it before starting another. Kind and immutable Targets describe what the subscription covers,
while IsDisposed reports whether it has ended. Await Completion to distinguish explicit
disposal, disposal of the owning client, and ReconnectExhausted when a configured
reconnect attempt limit ends the subscription.
Publish
long receivers = await redis.PublishAsync("orders", orderJson);
long shardReceivers = await redis.PublishShardedAsync("events:eu-west", payload);
Binary channels and explicit subscription kinds
RespireChannel owns the exact channel bytes. Strings convert implicitly after UTF-16 validation;
byte arrays and ReadOnlyMemory<byte> convert implicitly by copying their contents. Mutating the
original buffer after construction, subscription, or publication cannot change a channel's identity.
Reuse a constructed channel to avoid copying it for each publication.
RespireChannel binary = new byte[] { 0xff, 0x00, 0x3a, 0x41 };
await using var subscription = await redis.SubscribeAsync(binary, stoppingToken);
await redis.PublishAsync(binary, "payload", stoppingToken);
var pattern = RespireChannel.Pattern(new byte[] { 0xff, 0x00, 0x3a, (byte)'*' });
await using var patterns = await redis.SubscribeAsync(pattern, stoppingToken);
var shard = RespireChannel.Sharded("events:{eu-west}");
await using var sharded = await redis.SubscribeAsync(shard, stoppingToken);
await redis.PublishAsync(shard, "payload", stoppingToken); // SPUBLISH
RespireChannel[] targets = ["orders", binary];
await using var multiple = await redis.SubscribeAsync(targets, stoppingToken);
The metadata-aware SubscribeAsync uses Kind to select SUBSCRIBE, PSUBSCRIBE, or SSUBSCRIBE.
Multi-target subscriptions require one kind; use separate subscriptions for mixed kinds. Named
SubscribePatternAsync and SubscribeShardedAsync also accept binary targets and explicitly select
their command family. Patterns cannot be published. Existing string subscription and publication
overloads remain available, including sharded subscriptions across Redis Cluster primaries.
Equality and hashing compare only bytes, independently of kind. Equivalent text and UTF-8 byte
targets deduplicate within a subscription. ClusterSlot uses raw bytes and Redis hash-tag rules.
Empty channels, including default(RespireChannel), are valid. Unpaired UTF-16 surrogates throw
ArgumentException before subscription work begins; arbitrary binary values are accepted.
Names such as __keyspace@0__:key remain ordinary channels, with no inferred notification routing.
Pre-release API change: RespireMessage.Channel, nullable Pattern, and subscription Targets
now contain RespireChannel values. Use .Bytes for lossless identity and .ToString() for UTF-8
display. Display replaces invalid UTF-8 and can make distinct channels look identical. For exact
diagnostics use Convert.ToHexString(channel.Bytes.Span); channel bytes are not telemetry tags.
Sharded subscriptions in Redis Cluster
With UseCluster = true, SubscribeShardedAsync groups channels by their hash-slot owner and
uses one dedicated subscription connection per primary. Channels on different slots can share
one subscription; channels on the same primary share its connection. Duplicate consumers share
one server-side subscription until the last consumer disposes. SPUBLISH uses the command
connection for the channel's slot. Channel names are never affected by a client's key prefix.
MOVED and ASK replies, unsolicited SUNSUBSCRIBE frames during resharding, and refreshed topology
all trigger routing to the current owner. Socket failures restore the affected channels without
resubscribing healthy primaries. The existing subscription and its buffer survive these changes.
A reconnect gap marker precedes messages from the replacement subscription; Redis pub/sub
cannot replay messages lost while a channel changes owners.
Sharded Cluster subscriptions share one recovery episode and configured attempt budget,
with health retained for every affected primary until the episode completes. An attempt on
another primary cannot clear an earlier primary's failure. Successful recovery clears all
affected endpoints; exhaustion marks them disconnected. Topology-driven recovery reports
subscription owners rather than an unrelated configured seed. New sharded subscriptions
fail while an episode is active, with or without a configured ReconnectPolicy, even for
channels whose primary is healthy, because the episode owns every sharded route; subscribe
again after recovery completes. Existing subscriptions keep their buffers and routes.
With a null policy, recovery retries indefinitely: the first attempt is immediate, then
exponential waits begin at 250 ms and stop growing at five seconds. New subscriptions do
not wait for that episode; during a prolonged outage they continue to fail promptly.
Configure ReconnectPolicy.MaxAttempts to bound recovery.
The sharded recovery budget is
separate from regular channel and pattern subscriptions. Exhausting that budget completes all
sharded subscriptions with ReconnectExhausted; regular subscriptions remain usable. Recreate
the client to create sharded subscriptions after exhaustion. Notifications across Cluster
primaries are a separate feature and remain unsupported.
Read message data
RespireMessage exposes text, bytes, channel and pattern metadata. Deserialize application messages with the client's configured serializer:
OrderCreated order = message.As<OrderCreated>();
Channel, pattern, and payload bytes remain valid after enumeration advances. Exact-channel
messages share immutable registered channel storage; pattern messages own a copy of the concrete
channel name from the incoming frame. Accessing .Bytes does not allocate. Text conversion is
explicit and can allocate; do not mutate the exposed read-only storage through unsafe APIs.
Reconnection and pressure
RespireOptions.ReconnectPolicy can bound and delay automatic reconnection and
resubscription. Its budget resets only after all live routes are acknowledged. Exhaustion
ends live subscriptions; recreate the client to subscribe again. New subscriptions during
configured recovery fail rather than bypassing its delay. Null preserves the existing
immediate attempt followed by 250 ms exponential waits capped at five seconds. See
connection recovery for
attempt semantics, cancellation, lifecycle events, and telemetry.
Subscriptions resubscribe after reconnection. The subscription buffer is bounded; configure
SubscriptionOverflow in RespireOptions to drop either the oldest buffered message or the
newest incoming message when a consumer falls behind. Override those defaults for one
subscription with RespireSubscriptionOptions:
var options = new RespireSubscriptionOptions(
BufferSize: 128,
Overflow: SubscriptionOverflow.DropNewest);
await using var telemetry = await redis.SubscribeAsync(
"telemetry", options, stoppingToken);
Blocking and throwing policies are intentionally unavailable because they would stop the shared
pub/sub reader and affect unrelated subscriptions. DroppedMessages reports the number discarded
for that subscription, and the respire.pubsub.messages.dropped counter exposes the same event to
metrics collectors.
Pub/sub is transient: Redis does not retain messages for disconnected subscribers. Use streams when delivery tracking and replay matter.
Detecting delivery gaps
Pre-release behavior change: subscription streams now include RespireMessageKind.Gap items.
Check Kind before reading or deserializing a published payload. A gap has empty channel/payload
fields, and As<T>() throws because there is no published value.
The following example logs each item. Replace the gap log with your application's state reload, and replace the message log with its normal update handler.
await using var subscription = await redis.SubscribeAsync("orders", stoppingToken);
await foreach (var message in subscription.WithCancellation(stoppingToken))
{
switch (message.Kind)
{
case RespireMessageKind.Gap:
Console.Error.WriteLine($"Delivery gap: {message.Gap}. Reload authoritative state before applying more messages.");
break;
case RespireMessageKind.Message:
Console.WriteLine($"{message.Channel}: {message.Text}");
break;
default:
throw new InvalidOperationException($"Unknown subscription item: {message.Kind}");
}
}
message.Gap reports Reason (Reconnect, BufferOverflow, or both), StartedAt, EndedAt,
Duration, and the known local DroppedMessages count. Connection loss counts are unknown.
The reconnect interval begins when the client observes a failed connection, which can be later
than the actual interruption, and ends at the target's resubscription acknowledgement.
Each affected target produces a reconnect gap before any messages received after its acknowledgement,
including when both arrive in one socket read. Multi-target subscriptions can report more than one
gap as targets resume. Messages already buffered before reconnect remain before the marker.
Overflow markers sit at the loss position: before retained data for DropOldest, after previously
buffered data for DropNewest. Adjacent markers coalesce by combining reasons, observed intervals,
and discard counts. Markers do not consume data capacity and cannot themselves be dropped; their
storage remains bounded by the configured message capacity plus one pending marker.
subscription.DeliveryGap and the respire.pubsub.delivery.gaps counter report each detected target
interruption or local discard, even without enumeration. Their count can exceed the number of
coalesced stream markers. Counter tags are respire.subscription.kind and
respire.subscription.gap.reason; channel names and payloads are not tags. Event handlers run
synchronously on the receive path: keep them short, signal background work, and never block on
Redis or subscription operations. Handler exceptions are logged without stopping message delivery.
Reloading is application-specific: use versions or idempotent updates when reconciling buffered messages with a fresh snapshot. Redis pub/sub cannot replay lost messages; use Streams when replay or acknowledged delivery is required.