1"""
2Redis Keyspace Notifications support for redis-py.
3
4This module provides utilities for subscribing to and parsing Redis keyspace
5notifications. Keyspace notifications allow clients to receive events when
6keys are modified in Redis.
7
8Note: Keyspace notifications must be enabled on the Redis server via the
9``notify-keyspace-events`` configuration option. This is a server-side
10configuration that should be done by your infrastructure/operations team.
11See the Redis documentation for details:
12https://redis.io/docs/latest/develop/pubsub/keyspace-notifications/
13
14Standalone Redis Example:
15 >>> from redis import Redis
16 >>> from redis.keyspace_notifications import (
17 ... KeyspaceNotifications,
18 ... KeyspaceChannel,
19 ... EventType,
20 ... )
21 >>>
22 >>> r = Redis()
23 >>> # Server must have notify-keyspace-events configured (e.g., "KEA")
24 >>> ksn = KeyspaceNotifications(r)
25 >>>
26 >>> # Subscribe using Channel class (patterns auto-detected)
27 >>> channel = KeyspaceChannel("user:*")
28 >>> ksn.subscribe(channel)
29 >>>
30 >>> # Or use convenience methods for specific event types
31 >>> ksn.subscribe_keyevent(EventType.SET)
32 >>>
33 >>> for notification in ksn.listen():
34 ... print(f"Key: {notification.key}, Event: {notification.event_type}")
35
36Redis Cluster Example:
37 >>> from redis.cluster import RedisCluster
38 >>> from redis.keyspace_notifications import (
39 ... ClusterKeyspaceNotifications,
40 ... KeyspaceChannel,
41 ... EventType,
42 ... )
43 >>>
44 >>> rc = RedisCluster(host="localhost", port=7000)
45 >>> # Server must have notify-keyspace-events configured (e.g., "KEA")
46 >>> ksn = ClusterKeyspaceNotifications(rc)
47 >>>
48 >>> # Subscribe using Channel class (patterns auto-detected)
49 >>> channel = KeyspaceChannel("user:*")
50 >>> ksn.subscribe(channel)
51 >>>
52 >>> # Or use convenience methods for specific event types
53 >>> ksn.subscribe_keyevent(EventType.SET)
54 >>>
55 >>> for notification in ksn.listen():
56 ... print(f"Key: {notification.key}, Event: {notification.event_type}")
57"""
58
59from __future__ import annotations
60
61import logging
62import re
63import threading
64import time
65from abc import ABC, abstractmethod
66from collections.abc import Callable
67from dataclasses import dataclass
68from enum import Enum
69from typing import TYPE_CHECKING, Any, ClassVar, Union
70
71from redis.client import Redis
72from redis.cluster import RedisCluster
73from redis.exceptions import (
74 ConnectionError,
75 RedisError,
76 TimeoutError,
77)
78from redis.utils import safe_str
79
80logger = logging.getLogger(__name__)
81
82if TYPE_CHECKING:
83 from typing import TypeAlias
84
85# Type alias for channel arguments - can be a string, bytes, or Channel object
86# This is defined here and the actual types are added after class definitions
87ChannelT: TypeAlias = Union[
88 str,
89 bytes,
90 "KeyspaceChannel",
91 "KeyeventChannel",
92 "SubkeyspaceChannel",
93 "SubkeyeventChannel",
94 "SubkeyspaceitemChannel",
95 "SubkeyspaceeventChannel",
96]
97
98
99# Type alias for sync handlers
100SyncHandlerT = Callable[["KeyNotification"], None]
101
102
103# =============================================================================
104# Event Type Constants
105# =============================================================================
106# These are common Redis keyspace notification event types provided for
107# convenience. You can use any string as an event type - these constants
108# are not exhaustive and Redis may add new events in future versions.
109
110
111class EventType:
112 """
113 Common Redis keyspace notification event type constants.
114
115 These are provided for convenience and IDE autocomplete. You can use
116 any string as an event type - new Redis events will work without
117 needing library updates.
118 """
119
120 # String commands
121 SET = "set"
122 SETEX = "setex"
123 SETNX = "setnx"
124 SETRANGE = "setrange"
125 INCR = "incr"
126 INCRBY = "incrby"
127 INCRBYFLOAT = "incrbyfloat"
128 DECR = "decr"
129 DECRBY = "decrby"
130 APPEND = "append"
131
132 # Generic commands
133 DEL = "del"
134 UNLINK = "unlink"
135 RENAME = "rename"
136 RENAME_FROM = "rename_from"
137 RENAME_TO = "rename_to"
138 COPY_TO = "copy_to"
139 MOVE = "move"
140 RESTORE = "restore"
141
142 # Expiration events
143 EXPIRE = "expire"
144 EXPIREAT = "expireat"
145 PEXPIRE = "pexpire"
146 PEXPIREAT = "pexpireat"
147 EXPIRED = "expired"
148 PERSIST = "persist"
149
150 # Eviction events
151 EVICTED = "evicted"
152
153 # List commands
154 LPUSH = "lpush"
155 RPUSH = "rpush"
156 LPOP = "lpop"
157 RPOP = "rpop"
158 LINSERT = "linsert"
159 LSET = "lset"
160 LTRIM = "ltrim"
161 LMOVE = "lmove"
162 BLPOP = "blpop"
163 BRPOP = "brpop"
164 BLMOVE = "blmove"
165
166 # Set commands
167 SADD = "sadd"
168 SREM = "srem"
169 SPOP = "spop"
170 SMOVE = "smove"
171 SINTERSTORE = "sinterstore"
172 SUNIONSTORE = "sunionstore"
173 SDIFFSTORE = "sdiffstore"
174
175 # Sorted set commands
176 ZADD = "zadd"
177 ZINCRBY = "zincrby"
178 ZREM = "zrem"
179 ZREMRANGEBYRANK = "zremrangebyrank"
180 ZREMRANGEBYSCORE = "zremrangebyscore"
181 ZREMRANGEBYLEX = "zremrangebylex"
182 ZPOPMIN = "zpopmin"
183 ZPOPMAX = "zpopmax"
184 BZPOPMIN = "bzpopmin"
185 BZPOPMAX = "bzpopmax"
186 ZINTERSTORE = "zinterstore"
187 ZUNIONSTORE = "zunionstore"
188 ZDIFFSTORE = "zdiffstore"
189 ZRANGESTORE = "zrangestore"
190
191 # Hash commands
192 HSET = "hset"
193 HSETNX = "hsetnx"
194 HDEL = "hdel"
195 HINCRBY = "hincrby"
196 HINCRBYFLOAT = "hincrbyfloat"
197
198 # Stream commands
199 XADD = "xadd"
200 XTRIM = "xtrim"
201 XDEL = "xdel"
202 XGROUP_CREATE = "xgroup-create"
203 XGROUP_CREATECONSUMER = "xgroup-createconsumer"
204 XGROUP_DELCONSUMER = "xgroup-delconsumer"
205 XGROUP_DESTROY = "xgroup-destroy"
206 XGROUP_SETID = "xgroup-setid"
207 XSETID = "xsetid"
208 XCLAIM = "xclaim"
209 XAUTOCLAIM = "xautoclaim"
210 XREADGROUP = "xreadgroup"
211
212 # Other
213 NEW = "new" # Key created (when tracking new keys)
214 SORTSTORE = "sortstore"
215 GETEX = "getex"
216 GETDEL = "getdel"
217 SETIFGT = "setifgt"
218 SETIFLT = "setiflt"
219 SETIFEQ = "setifeq"
220 SETIFNE = "setifne"
221
222
223def _parse_length_prefixed_subkeys(s: bytes | str) -> list[str]:
224 """Parse a length-prefixed subkey list.
225
226 The wire format is ``<len>:<subkey>[,<len>:<subkey>...]``.
227
228 The server writes ``<len>`` as the subkey size in **bytes** (it uses
229 ``sdslen`` when formatting the notification), so parsing must happen in
230 the byte domain; splitting the decoded string by character count
231 corrupts every multi-byte subkey and desynchronizes the whole list.
232
233 Returns:
234 A list of subkey strings.
235 """
236 raw = s.encode("utf-8") if isinstance(s, str) else bytes(s)
237 subkeys: list[str] = []
238 pos = 0
239 while pos < len(raw):
240 colon = raw.index(b":", pos)
241 length = int(raw[pos:colon])
242 start = colon + 1
243 subkeys.append(raw[start : start + length].decode("utf-8", "replace"))
244 pos = start + length
245 if pos < len(raw) and raw[pos] == ord(","):
246 pos += 1 # skip comma separator
247 return subkeys
248
249
250@dataclass
251class KeyNotification:
252 """
253 Represents a parsed Redis keyspace, keyevent, or subkey notification.
254
255 This class provides convenient access to the notification details
256 like key, event type, database number, and affected subkeys.
257
258 Attributes:
259 key: The Redis key that was affected (for keyspace notifications)
260 or the key name from the message data (for keyevent notifications)
261 event_type: The type of operation that occurred (e.g., "set", "del").
262 This is a plain string, so new Redis events work automatically.
263 Compare against EventType constants or any string.
264 database: The database number where the event occurred
265 channel: The original channel name
266 is_keyspace: True if this is a keyspace notification, False for keyevent
267 data: The raw data payload from the notification message.
268 subkeys: List of affected subkeys (fields) for subkey notifications.
269 Empty list for regular keyspace/keyevent notifications.
270 """
271
272 # Regex patterns for parsing keyspace/keyevent channels
273 # Pattern: __keyspace@<db>__:<key> or __keyevent@<db>__:<event>
274 _KEYSPACE_PATTERN: ClassVar[re.Pattern] = re.compile(
275 r"^__keyspace@(\d+|\*)__:(.+)$"
276 )
277 _KEYEVENT_PATTERN: ClassVar[re.Pattern] = re.compile(
278 r"^__keyevent@(\d+|\*)__:(.+)$"
279 )
280 _SUBKEYSPACE_PATTERN: ClassVar[re.Pattern] = re.compile(
281 r"^__subkeyspace@(\d+|\*)__:(.+)$"
282 )
283 _SUBKEYEVENT_PATTERN: ClassVar[re.Pattern] = re.compile(
284 r"^__subkeyevent@(\d+|\*)__:(.+)$"
285 )
286 _SUBKEYSPACEITEM_PATTERN: ClassVar[re.Pattern] = re.compile(
287 r"^__subkeyspaceitem@(\d+|\*)__:(.+)$", re.DOTALL
288 )
289 _SUBKEYSPACEEVENT_PATTERN: ClassVar[re.Pattern] = re.compile(
290 r"^__subkeyspaceevent@(\d+|\*)__:(.+)$"
291 )
292
293 key: str
294 event_type: str
295 database: int
296 channel: str
297 is_keyspace: bool
298 data: str
299 subkeys: list[str] = None # type: ignore[assignment]
300
301 def __post_init__(self):
302 if self.subkeys is None:
303 self.subkeys = []
304
305 @classmethod
306 def from_message(
307 cls,
308 message: dict[str, Any] | None,
309 key_prefix: str | bytes | None = None,
310 ) -> KeyNotification | None:
311 """
312 Parse a pub/sub message into a KeyNotification.
313
314 Args:
315 message: A pub/sub message dict with 'channel', 'data', and 'type' keys
316 key_prefix: Optional prefix to filter and strip from keys.
317 If provided, only notifications for keys starting with
318 this prefix will be returned, and the prefix will be
319 stripped from the key.
320
321 Returns:
322 A KeyNotification if the message is a valid keyspace/keyevent
323 notification, None otherwise.
324
325 Example:
326 >>> message = {
327 ... 'type': 'pmessage',
328 ... 'pattern': '__keyspace@0__:user:*',
329 ... 'channel': '__keyspace@0__:user:123',
330 ... 'data': 'set'
331 ... }
332 >>> notification = KeyNotification.from_message(message)
333 >>> notification.key
334 'user:123'
335 >>> notification.event_type
336 'set'
337 """
338 if message is None:
339 return None
340
341 msg_type = message.get("type")
342 if msg_type not in ("message", "pmessage"):
343 return None
344
345 channel = message.get("channel")
346 data = message.get("data")
347
348 if channel is None or data is None:
349 return None
350
351 return cls.try_parse(channel, data, key_prefix)
352
353 @classmethod
354 def try_parse(
355 cls,
356 channel: str | bytes,
357 data: str | bytes,
358 key_prefix: str | bytes | None = None,
359 ) -> KeyNotification | None:
360 """
361 Try to parse a channel and data into a KeyNotification.
362
363 This is a lower-level method that takes the channel and data directly,
364 useful when working with callback-based subscription handlers.
365
366 Args:
367 channel: The channel name (e.g., "__keyspace@0__:mykey")
368 data: The message data (event type for keyspace, key for keyevent)
369 key_prefix: Optional prefix to filter and strip from keys
370
371 Returns:
372 A KeyNotification if valid, None otherwise.
373 """
374 raw_data = bytes(data) if isinstance(data, (bytes, bytearray)) else None
375 channel = safe_str(channel)
376 data = safe_str(data)
377
378 return cls._parse(channel, data, key_prefix, raw_data=raw_data)
379
380 @classmethod
381 def _parse(
382 cls,
383 channel: str,
384 data: str,
385 key_prefix: str | bytes | None = None,
386 raw_data: bytes | None = None,
387 ) -> KeyNotification | None:
388 """Internal parsing logic."""
389 # Normalize key_prefix
390 key_prefix = safe_str(key_prefix) if key_prefix else None
391 # Byte-domain view of the payload: the server prefixes every element
392 # with its size in bytes (sdslen), so length-prefixed fields must be
393 # split before decoding. When the caller only has the decoded string
394 # (decode_responses=True) the exact bytes are recovered by re-encoding.
395 raw = raw_data if raw_data is not None else data.encode("utf-8")
396
397 # Try keyspace pattern first: __keyspace@<db>__:<key>
398 match = cls._KEYSPACE_PATTERN.match(channel)
399 if match:
400 db_str, key = match.groups()
401 database = int(db_str) if db_str != "*" else -1
402 event_type = data # For keyspace, the data is the event type
403
404 # Apply key prefix filter
405 if key_prefix:
406 if not key.startswith(key_prefix):
407 return None
408 key = key[len(key_prefix) :]
409
410 return cls(
411 key=key,
412 event_type=event_type,
413 database=database,
414 channel=channel,
415 is_keyspace=True,
416 data=data,
417 )
418
419 # Try keyevent pattern: __keyevent@<db>__:<event>
420 match = cls._KEYEVENT_PATTERN.match(channel)
421 if match:
422 db_str, event_type = match.groups()
423 database = int(db_str) if db_str != "*" else -1
424 key = data # For keyevent, the data is the key
425
426 # Apply key prefix filter
427 if key_prefix:
428 if not key.startswith(key_prefix):
429 return None
430 key = key[len(key_prefix) :]
431
432 return cls(
433 key=key,
434 event_type=event_type,
435 database=database,
436 channel=channel,
437 is_keyspace=False,
438 data=data,
439 )
440
441 # Try subkeyspace: channel=__subkeyspace@<db>__:<key>
442 # data=<event>|<subkey_len>:<subkey>[,<subkey_len>:<subkey>...]
443 match = cls._SUBKEYSPACE_PATTERN.match(channel)
444 if match:
445 db_str, key = match.groups()
446 database = int(db_str) if db_str != "*" else -1
447 pipe_idx = raw.index(b"|")
448 event_type = raw[:pipe_idx].decode("utf-8", "replace")
449 subkeys = _parse_length_prefixed_subkeys(raw[pipe_idx + 1 :])
450
451 if key_prefix:
452 if not key.startswith(key_prefix):
453 return None
454 key = key[len(key_prefix) :]
455
456 return cls(
457 key=key,
458 event_type=event_type,
459 database=database,
460 channel=channel,
461 is_keyspace=True,
462 data=data,
463 subkeys=subkeys,
464 )
465
466 # Try subkeyevent: channel=__subkeyevent@<db>__:<event>
467 # data=<key_len>:<key>|<subkey_len>:<subkey>[,...]
468 match = cls._SUBKEYEVENT_PATTERN.match(channel)
469 if match:
470 db_str, event_type = match.groups()
471 database = int(db_str) if db_str != "*" else -1
472 # Parse key by length prefix
473 colon_idx = raw.index(b":")
474 key_len = int(raw[:colon_idx])
475 key_start = colon_idx + 1
476 key = raw[key_start : key_start + key_len].decode("utf-8", "replace")
477 # After key, expect '|' then subkeys
478 subkeys_start = key_start + key_len + 1 # +1 for '|'
479 subkeys = _parse_length_prefixed_subkeys(raw[subkeys_start:])
480
481 if key_prefix:
482 if not key.startswith(key_prefix):
483 return None
484 key = key[len(key_prefix) :]
485
486 return cls(
487 key=key,
488 event_type=event_type,
489 database=database,
490 channel=channel,
491 is_keyspace=False,
492 data=data,
493 subkeys=subkeys,
494 )
495
496 # Try subkeyspaceitem: channel=__subkeyspaceitem@<db>__:<key>\n<subkey>
497 # data=<event>
498 match = cls._SUBKEYSPACEITEM_PATTERN.match(channel)
499 if match:
500 db_str, key_and_subkey = match.groups()
501 database = int(db_str) if db_str != "*" else -1
502 newline_idx = key_and_subkey.index("\n")
503 key = key_and_subkey[:newline_idx]
504 subkey = key_and_subkey[newline_idx + 1 :]
505 event_type = data
506
507 if key_prefix:
508 if not key.startswith(key_prefix):
509 return None
510 key = key[len(key_prefix) :]
511
512 return cls(
513 key=key,
514 event_type=event_type,
515 database=database,
516 channel=channel,
517 is_keyspace=True,
518 data=data,
519 subkeys=[subkey],
520 )
521
522 # Try subkeyspaceevent: channel=__subkeyspaceevent@<db>__:<event>|<key>
523 # data=<subkey_len>:<subkey>[,...]
524 match = cls._SUBKEYSPACEEVENT_PATTERN.match(channel)
525 if match:
526 db_str, event_and_key = match.groups()
527 database = int(db_str) if db_str != "*" else -1
528 pipe_idx = event_and_key.index("|")
529 event_type = event_and_key[:pipe_idx]
530 key = event_and_key[pipe_idx + 1 :]
531 subkeys = _parse_length_prefixed_subkeys(raw)
532
533 if key_prefix:
534 if not key.startswith(key_prefix):
535 return None
536 key = key[len(key_prefix) :]
537
538 return cls(
539 key=key,
540 event_type=event_type,
541 database=database,
542 channel=channel,
543 is_keyspace=False,
544 data=data,
545 subkeys=subkeys,
546 )
547
548 return None
549
550 def key_starts_with(self, prefix: str | bytes) -> bool:
551 """Check if the key starts with the given prefix."""
552 prefix = safe_str(prefix)
553 return self.key.startswith(prefix)
554
555
556# =============================================================================
557# Channel Classes
558# =============================================================================
559
560
561class KeyspaceChannel:
562 """
563 Represents a keyspace notification channel for subscribing to events on keys.
564
565 Keyspace notifications publish the event type (e.g., "set", "del") as the message
566 when a key matching the pattern is modified.
567
568 This class can be used directly with subscribe()/psubscribe() as it implements
569 __str__ to return the channel string.
570
571 Attributes:
572 key_or_pattern: The key or pattern to monitor (use '*' for wildcards)
573 db: The database number (defaults to 0, the only database in Redis Cluster)
574 is_pattern: Whether this channel contains wildcards
575
576 Examples:
577 >>> channel = KeyspaceChannel("user:123", db=0)
578 >>> str(channel)
579 '__keyspace@0__:user:123'
580
581 >>> # Pattern subscription (wildcards are auto-detected)
582 >>> channel = KeyspaceChannel("user:*", db=0)
583 >>> str(channel)
584 '__keyspace@0__:user:*'
585
586 >>> # Use with KeyspaceNotifications
587 >>> notifications = KeyspaceNotifications(redis_client)
588 >>> notifications.subscribe(channel)
589 """
590
591 PREFIX: ClassVar[str] = "__keyspace@"
592
593 def __init__(self, key_or_pattern: str, db: int = 0):
594 """
595 Create a keyspace notification channel.
596
597 Args:
598 key_or_pattern: The key or pattern to monitor. Use '*' for wildcards.
599 db: The database number. Defaults to 0 (the only database in Redis Cluster).
600 """
601 self.key_or_pattern = key_or_pattern
602 self.db = db
603 self._channel_str = self._build_channel_string()
604
605 def _build_channel_string(self) -> str:
606 return f"{self.PREFIX}{self.db}__:{self.key_or_pattern}"
607
608 @property
609 def is_pattern(self) -> bool:
610 """Check if this channel contains wildcards and should use psubscribe."""
611 return _is_pattern(self.key_or_pattern)
612
613 def __str__(self) -> str:
614 return self._channel_str
615
616 def __repr__(self) -> str:
617 return f"KeyspaceChannel({self.key_or_pattern!r}, db={self.db})"
618
619 def __eq__(self, other: object) -> bool:
620 if isinstance(other, KeyspaceChannel):
621 return self._channel_str == other._channel_str
622 if isinstance(other, str):
623 return self._channel_str == other
624 return NotImplemented
625
626 def __hash__(self) -> int:
627 return hash(self._channel_str)
628
629
630class KeyeventChannel:
631 """
632 Represents a keyevent notification channel for subscribing to event types.
633
634 Keyevent notifications publish the key name as the message when the specified
635 event type occurs on any key.
636
637 This class can be used directly with subscribe()/psubscribe() as it implements
638 __str__ to return the channel string.
639
640 Attributes:
641 event: The event type to monitor
642 db: The database number (defaults to 0, the only database in Redis Cluster)
643 is_pattern: Whether this channel contains wildcards
644
645 Examples:
646 >>> channel = KeyeventChannel(EventType.SET, db=0)
647 >>> str(channel)
648 '__keyevent@0__:set'
649
650 >>> channel = KeyeventChannel.all_events(db=0)
651 >>> str(channel)
652 '__keyevent@0__:*'
653
654 >>> # Use with KeyspaceNotifications
655 >>> notifications = KeyspaceNotifications(redis_client)
656 >>> notifications.subscribe(channel)
657 """
658
659 PREFIX: ClassVar[str] = "__keyevent@"
660
661 def __init__(self, event: str, db: int = 0):
662 """
663 Create a keyevent notification channel.
664
665 Args:
666 event: The event type to monitor (e.g., EventType.SET or "set")
667 db: The database number. Defaults to 0 (the only database in Redis Cluster).
668 """
669 self.event = event
670 self.db = db
671 self._channel_str = self._build_channel_string()
672
673 def _build_channel_string(self) -> str:
674 return f"{self.PREFIX}{self.db}__:{self.event}"
675
676 @property
677 def is_pattern(self) -> bool:
678 """Check if this channel contains wildcards and should use psubscribe."""
679 return _is_pattern(self.event)
680
681 @classmethod
682 def all_events(cls, db: int = 0) -> "KeyeventChannel":
683 """
684 Create a keyevent pattern for subscribing to all event types.
685
686 This is equivalent to KeyeventChannel("*").
687
688 Args:
689 db: The database number. Defaults to 0 (the only database in Redis Cluster).
690
691 Returns:
692 A KeyeventChannel configured to receive all events.
693
694 Examples:
695 >>> channel = KeyeventChannel.all_events()
696 >>> str(channel)
697 '__keyevent@0__:*'
698 """
699 return cls("*", db=db)
700
701 def __str__(self) -> str:
702 return self._channel_str
703
704 def __repr__(self) -> str:
705 return f"KeyeventChannel({self.event!r}, db={self.db})"
706
707 def __eq__(self, other: object) -> bool:
708 if isinstance(other, KeyeventChannel):
709 return self._channel_str == other._channel_str
710 if isinstance(other, str):
711 return self._channel_str == other
712 return NotImplemented
713
714 def __hash__(self) -> int:
715 return hash(self._channel_str)
716
717
718class SubkeyspaceChannel:
719 """
720 Represents a subkeyspace notification channel for subscribing to
721 subkey-level events on keys (e.g., hash field changes).
722
723 The channel format is ``__subkeyspace@<db>__:<key>``.
724 The message payload is ``<event>|<subkey_len>:<subkey>[,...]``.
725
726 Examples:
727 >>> channel = SubkeyspaceChannel("myhash", db=0)
728 >>> str(channel)
729 '__subkeyspace@0__:myhash'
730 """
731
732 PREFIX: ClassVar[str] = "__subkeyspace@"
733
734 def __init__(self, key_or_pattern: str, db: int = 0):
735 self.key_or_pattern = key_or_pattern
736 self.db = db
737 self._channel_str = self._build_channel_string()
738
739 def _build_channel_string(self) -> str:
740 return f"{self.PREFIX}{self.db}__:{self.key_or_pattern}"
741
742 @property
743 def is_pattern(self) -> bool:
744 return _is_pattern(self.key_or_pattern)
745
746 def __str__(self) -> str:
747 return self._channel_str
748
749 def __repr__(self) -> str:
750 return f"SubkeyspaceChannel({self.key_or_pattern!r}, db={self.db})"
751
752 def __eq__(self, other: object) -> bool:
753 if isinstance(other, SubkeyspaceChannel):
754 return self._channel_str == other._channel_str
755 if isinstance(other, str):
756 return self._channel_str == other
757 return NotImplemented
758
759 def __hash__(self) -> int:
760 return hash(self._channel_str)
761
762
763class SubkeyeventChannel:
764 """
765 Represents a subkeyevent notification channel for subscribing to
766 specific event types with subkey-level detail.
767
768 The channel format is ``__subkeyevent@<db>__:<event>``.
769 The message payload is ``<key_len>:<key>|<subkey_len>:<subkey>[,...]``.
770
771 Examples:
772 >>> channel = SubkeyeventChannel("hdel", db=0)
773 >>> str(channel)
774 '__subkeyevent@0__:hdel'
775 """
776
777 PREFIX: ClassVar[str] = "__subkeyevent@"
778
779 def __init__(self, event: str, db: int = 0):
780 self.event = event
781 self.db = db
782 self._channel_str = self._build_channel_string()
783
784 def _build_channel_string(self) -> str:
785 return f"{self.PREFIX}{self.db}__:{self.event}"
786
787 @property
788 def is_pattern(self) -> bool:
789 return _is_pattern(self.event)
790
791 @classmethod
792 def all_events(cls, db: int = 0) -> SubkeyeventChannel:
793 """Create a channel for all subkeyevent types."""
794 return cls("*", db=db)
795
796 def __str__(self) -> str:
797 return self._channel_str
798
799 def __repr__(self) -> str:
800 return f"SubkeyeventChannel({self.event!r}, db={self.db})"
801
802 def __eq__(self, other: object) -> bool:
803 if isinstance(other, SubkeyeventChannel):
804 return self._channel_str == other._channel_str
805 if isinstance(other, str):
806 return self._channel_str == other
807 return NotImplemented
808
809 def __hash__(self) -> int:
810 return hash(self._channel_str)
811
812
813class SubkeyspaceitemChannel:
814 """
815 Represents a subkeyspaceitem notification channel for subscribing to
816 events on a specific subkey (field) of a specific key.
817
818 The channel format is ``__subkeyspaceitem@<db>__:<key>\\n<subkey>``.
819 The message payload is the event type (e.g., ``"hset"``).
820
821 Note:
822 The server only emits this notification when the key does not
823 contain a newline character.
824
825 Examples:
826 >>> channel = SubkeyspaceitemChannel("myhash", "myfield", db=0)
827 >>> str(channel)
828 '__subkeyspaceitem@0__:myhash\\nmyfield'
829 """
830
831 PREFIX: ClassVar[str] = "__subkeyspaceitem@"
832
833 def __init__(self, key_or_pattern: str, subkey_or_pattern: str, db: int = 0):
834 self.key_or_pattern = key_or_pattern
835 self.subkey_or_pattern = subkey_or_pattern
836 self.db = db
837 self._channel_str = self._build_channel_string()
838
839 def _build_channel_string(self) -> str:
840 return (
841 f"{self.PREFIX}{self.db}__:{self.key_or_pattern}\n{self.subkey_or_pattern}"
842 )
843
844 @property
845 def is_pattern(self) -> bool:
846 return _is_pattern(self.key_or_pattern) or _is_pattern(self.subkey_or_pattern)
847
848 def __str__(self) -> str:
849 return self._channel_str
850
851 def __repr__(self) -> str:
852 return (
853 f"SubkeyspaceitemChannel({self.key_or_pattern!r}, "
854 f"{self.subkey_or_pattern!r}, db={self.db})"
855 )
856
857 def __eq__(self, other: object) -> bool:
858 if isinstance(other, SubkeyspaceitemChannel):
859 return self._channel_str == other._channel_str
860 if isinstance(other, str):
861 return self._channel_str == other
862 return NotImplemented
863
864 def __hash__(self) -> int:
865 return hash(self._channel_str)
866
867
868class SubkeyspaceeventChannel:
869 """
870 Represents a subkeyspaceevent notification channel for subscribing to
871 a specific event on a specific key, receiving affected subkeys.
872
873 The channel format is ``__subkeyspaceevent@<db>__:<event>|<key>``.
874 The message payload is a length-prefixed subkey list.
875
876 Examples:
877 >>> channel = SubkeyspaceeventChannel("hset", "myhash", db=0)
878 >>> str(channel)
879 '__subkeyspaceevent@0__:hset|myhash'
880 """
881
882 PREFIX: ClassVar[str] = "__subkeyspaceevent@"
883
884 def __init__(self, event: str, key_or_pattern: str, db: int = 0):
885 self.event = event
886 self.key_or_pattern = key_or_pattern
887 self.db = db
888 self._channel_str = self._build_channel_string()
889
890 def _build_channel_string(self) -> str:
891 return f"{self.PREFIX}{self.db}__:{self.event}|{self.key_or_pattern}"
892
893 @property
894 def is_pattern(self) -> bool:
895 return _is_pattern(self.event) or _is_pattern(self.key_or_pattern)
896
897 def __str__(self) -> str:
898 return self._channel_str
899
900 def __repr__(self) -> str:
901 return (
902 f"SubkeyspaceeventChannel({self.event!r}, "
903 f"{self.key_or_pattern!r}, db={self.db})"
904 )
905
906 def __eq__(self, other: object) -> bool:
907 if isinstance(other, SubkeyspaceeventChannel):
908 return self._channel_str == other._channel_str
909 if isinstance(other, str):
910 return self._channel_str == other
911 return NotImplemented
912
913 def __hash__(self) -> int:
914 return hash(self._channel_str)
915
916
917class ChannelType(Enum):
918 """
919 Enum representing the type of a Redis keyspace notification channel.
920
921 Redis provides two types of keyspace notifications and four subkey
922 notification types:
923
924 - KEYSPACE: ``__keyspace@{db}__:{key}`` — data is the event type.
925 - KEYEVENT: ``__keyevent@{db}__:{event}`` — data is the key name.
926 - SUBKEYSPACE: ``__subkeyspace@{db}__:{key}`` — data is event + subkeys.
927 - SUBKEYEVENT: ``__subkeyevent@{db}__:{event}`` — data is key + subkeys.
928 - SUBKEYSPACEITEM: ``__subkeyspaceitem@{db}__:{key}\\n{subkey}`` — data
929 is the event type.
930 - SUBKEYSPACEEVENT: ``__subkeyspaceevent@{db}__:{event}|{key}`` — data
931 is a subkey list.
932
933 Examples:
934 >>> get_channel_type("__keyspace@0__:mykey")
935 ChannelType.KEYSPACE
936 >>> get_channel_type("__subkeyspace@0__:myhash")
937 ChannelType.SUBKEYSPACE
938 """
939
940 KEYSPACE = "keyspace"
941 KEYEVENT = "keyevent"
942 SUBKEYSPACE = "subkeyspace"
943 SUBKEYEVENT = "subkeyevent"
944 SUBKEYSPACEITEM = "subkeyspaceitem"
945 SUBKEYSPACEEVENT = "subkeyspaceevent"
946
947
948def get_channel_type(channel: str | bytes) -> ChannelType | None:
949 """
950 Determine the type of a Redis keyspace notification channel.
951
952 Args:
953 channel: The channel name to check (string or bytes).
954
955 Returns:
956 ChannelType.KEYSPACE if it's a keyspace notification channel,
957 ChannelType.KEYEVENT if it's a keyevent notification channel,
958 None if it's not a keyspace notification channel.
959
960 Examples:
961 >>> get_channel_type("__keyspace@0__:mykey")
962 ChannelType.KEYSPACE
963 >>> get_channel_type("__keyevent@0__:set")
964 ChannelType.KEYEVENT
965 >>> get_channel_type("regular_channel") is None
966 True
967 >>> get_channel_type(b"__keyspace@0__:mykey")
968 ChannelType.KEYSPACE
969 """
970 channel_str = safe_str(channel)
971 # Check subkey prefixes first (they are longer and more specific)
972 if channel_str.startswith(SubkeyspaceitemChannel.PREFIX):
973 return ChannelType.SUBKEYSPACEITEM
974 if channel_str.startswith(SubkeyspaceeventChannel.PREFIX):
975 return ChannelType.SUBKEYSPACEEVENT
976 if channel_str.startswith(SubkeyspaceChannel.PREFIX):
977 return ChannelType.SUBKEYSPACE
978 if channel_str.startswith(SubkeyeventChannel.PREFIX):
979 return ChannelType.SUBKEYEVENT
980 if channel_str.startswith(KeyspaceChannel.PREFIX):
981 return ChannelType.KEYSPACE
982 if channel_str.startswith(KeyeventChannel.PREFIX):
983 return ChannelType.KEYEVENT
984 return None
985
986
987def _is_pattern(
988 channel: str | bytes | KeyspaceChannel | KeyeventChannel,
989) -> bool:
990 """
991 Check if a channel string contains glob-style pattern characters.
992
993 Redis uses glob-style patterns for psubscribe:
994 - * matches any sequence of characters
995 - ? matches any single character
996 - [...] matches any character in the brackets
997
998 Args:
999 channel: The channel string to check. Can be a string, bytes,
1000 or a KeyspaceChannel/KeyeventChannel object.
1001
1002 Returns:
1003 True if the channel contains pattern characters, False otherwise.
1004 """
1005 # Handle Channel objects that have _channel_str attribute
1006 # (KeyspaceChannel, KeyeventChannel)
1007 if hasattr(channel, "_channel_str"):
1008 channel = channel._channel_str
1009 channel = safe_str(channel)
1010 # Check for unescaped glob pattern characters.
1011 # * and ? are always pattern characters.
1012 # [ is only a pattern character when followed by a matching unescaped ],
1013 # forming a bracket expression like [abc] or [a-z]. A lone [ (e.g. in
1014 # a key named "my[key") is treated as a literal by Redis.
1015 i = 0
1016 while i < len(channel):
1017 char = channel[i]
1018 if char == "\\":
1019 # Skip escaped character
1020 i += 2
1021 continue
1022 if char in ("*", "?"):
1023 return True
1024 if char == "[":
1025 # Look for a matching unescaped ]
1026 j = i + 1
1027 while j < len(channel):
1028 if channel[j] == "\\":
1029 j += 2
1030 continue
1031 if channel[j] == "]":
1032 return True
1033 j += 1
1034 # No matching ] found — literal [
1035 i += 1
1036 return False
1037
1038
1039# =============================================================================
1040# Abstract Base Class for Keyspace Notifications
1041# =============================================================================
1042
1043
1044class KeyspaceNotificationsInterface(ABC):
1045 """
1046 Interface for keyspace notification managers.
1047
1048 This interface provides a consistent API for both standalone (KeyspaceNotifications)
1049 and cluster (ClusterKeyspaceNotifications) implementations, allowing the same
1050 code patterns to work with both standalone and cluster Redis deployments.
1051 """
1052
1053 @abstractmethod
1054 def subscribe(
1055 self,
1056 *channels: ChannelT,
1057 handler: SyncHandlerT | None = None,
1058 ):
1059 """Subscribe to keyspace notification channels."""
1060 pass
1061
1062 @abstractmethod
1063 def unsubscribe(self, *channels: ChannelT):
1064 """Unsubscribe from keyspace notification channels."""
1065 pass
1066
1067 @abstractmethod
1068 def subscribe_keyspace(
1069 self,
1070 key_or_pattern: str,
1071 db: int = 0,
1072 handler: SyncHandlerT | None = None,
1073 ):
1074 """Subscribe to keyspace notifications for specific keys."""
1075 pass
1076
1077 @abstractmethod
1078 def subscribe_keyevent(
1079 self,
1080 event: str,
1081 db: int = 0,
1082 handler: SyncHandlerT | None = None,
1083 ):
1084 """Subscribe to keyevent notifications for specific event types."""
1085 pass
1086
1087 @abstractmethod
1088 def subscribe_subkeyspace(
1089 self,
1090 key_or_pattern: str,
1091 db: int = 0,
1092 handler: SyncHandlerT | None = None,
1093 ):
1094 """Subscribe to subkeyspace notifications for specific keys."""
1095 pass
1096
1097 @abstractmethod
1098 def subscribe_subkeyevent(
1099 self,
1100 event: str,
1101 db: int = 0,
1102 handler: SyncHandlerT | None = None,
1103 ):
1104 """Subscribe to subkeyevent notifications for specific event types."""
1105 pass
1106
1107 @abstractmethod
1108 def subscribe_subkeyspaceitem(
1109 self,
1110 key_or_pattern: str,
1111 subkey_or_pattern: str,
1112 db: int = 0,
1113 handler: SyncHandlerT | None = None,
1114 ):
1115 """Subscribe to subkeyspaceitem notifications for a specific subkey."""
1116 pass
1117
1118 @abstractmethod
1119 def subscribe_subkeyspaceevent(
1120 self,
1121 event: str,
1122 key_or_pattern: str,
1123 db: int = 0,
1124 handler: SyncHandlerT | None = None,
1125 ):
1126 """Subscribe to subkeyspaceevent notifications for an event on a key."""
1127 pass
1128
1129 @abstractmethod
1130 def get_message(
1131 self,
1132 ignore_subscribe_messages: bool | None = None,
1133 timeout: float = 0.0,
1134 ) -> KeyNotification | None:
1135 """Get the next keyspace notification if one is available."""
1136 pass
1137
1138 @abstractmethod
1139 def listen(self):
1140 """Listen for keyspace notifications."""
1141 pass
1142
1143 @abstractmethod
1144 def close(self):
1145 """Close the notification manager and clean up resources."""
1146 pass
1147
1148 @abstractmethod
1149 def __enter__(self):
1150 pass
1151
1152 @abstractmethod
1153 def __exit__(self, _exc_type, _exc_val, _exc_tb):
1154 pass
1155
1156 @property
1157 @abstractmethod
1158 def subscribed(self) -> bool:
1159 """Check if there are any active subscriptions and not closed."""
1160 pass
1161
1162 @abstractmethod
1163 def run_in_thread(
1164 self,
1165 poll_timeout: float = 0.0,
1166 daemon: bool = False,
1167 exception_handler: Callable[
1168 [
1169 BaseException,
1170 KeyspaceNotificationsInterface,
1171 KeyspaceWorkerThread,
1172 ],
1173 None,
1174 ]
1175 | None = None,
1176 ) -> KeyspaceWorkerThread:
1177 """Start a background thread that polls for notifications."""
1178 pass
1179
1180
1181class AbstractKeyspaceNotifications(KeyspaceNotificationsInterface):
1182 """
1183 Abstract base class for keyspace notification managers.
1184
1185 Provides shared implementation for subscribe/unsubscribe logic.
1186 Subclasses must implement:
1187 - _execute_subscribe: Execute the subscribe operation
1188 - _execute_unsubscribe: Execute the unsubscribe operation
1189 - get_message: Get the next notification
1190 - listen: Generator for notifications
1191 - close: Clean up resources
1192 """
1193
1194 def __init__(
1195 self,
1196 key_prefix: str | bytes | None = None,
1197 ignore_subscribe_messages: bool = True,
1198 ):
1199 """
1200 Initialize the base keyspace notification manager.
1201
1202 Args:
1203 key_prefix: Optional prefix to filter and strip from keys in notifications
1204 ignore_subscribe_messages: If True, subscribe/unsubscribe confirmations
1205 are not returned by get_message/listen
1206 """
1207 self.key_prefix = key_prefix
1208 self.ignore_subscribe_messages = ignore_subscribe_messages
1209 self._closed = False
1210
1211 def subscribe(
1212 self,
1213 *channels: ChannelT,
1214 handler: SyncHandlerT | None = None,
1215 ):
1216 """
1217 Subscribe to keyspace notification channels.
1218
1219 Automatically detects whether each channel is a pattern (contains
1220 wildcards like *, ?, [) or an exact channel name and uses the
1221 appropriate Redis subscribe command internally.
1222
1223 Args:
1224 *channels: Channels to subscribe to. Can be strings, KeyspaceChannel,
1225 or KeyeventChannel objects. Patterns are auto-detected.
1226 handler: Optional callback function that receives KeyNotification
1227 objects. If provided, notifications are passed to the handler
1228 instead of being returned by get_message()/listen().
1229 """
1230 # Wrap the handler to convert raw messages to KeyNotification objects
1231 wrapped_handler: Callable | None = None
1232 if handler is not None:
1233 # Capture key_prefix in closure for consistent filtering/stripping
1234 key_prefix = self.key_prefix
1235
1236 def _wrap_handler(message):
1237 notification = KeyNotification.from_message(
1238 message, key_prefix=key_prefix
1239 )
1240 if notification is not None:
1241 handler(notification)
1242
1243 wrapped_handler = _wrap_handler
1244
1245 patterns = {}
1246 exact_channels = {}
1247
1248 for channel in channels:
1249 if hasattr(channel, "_channel_str"):
1250 channel_str = str(channel)
1251 else:
1252 channel_str = safe_str(channel)
1253 if _is_pattern(channel):
1254 patterns[channel_str] = wrapped_handler
1255 else:
1256 exact_channels[channel_str] = wrapped_handler
1257
1258 # Delegate to subclass implementation first. For standalone Redis
1259 # this raises on failure, keeping tracking state clean. For cluster
1260 # implementations the operation is best-effort (partial failures are
1261 # logged, not raised) so tracking state is always updated afterwards.
1262 self._execute_subscribe(patterns, exact_channels)
1263 self._track_subscribe(patterns, exact_channels)
1264
1265 @abstractmethod
1266 def _execute_subscribe(
1267 self, patterns: dict[str, Any], exact_channels: dict[str, Any]
1268 ) -> None:
1269 """
1270 Execute the subscribe operation.
1271
1272 Args:
1273 patterns: Dict mapping pattern strings to handlers (for psubscribe)
1274 exact_channels: Dict mapping channel strings to handlers (for subscribe)
1275 """
1276 pass
1277
1278 def unsubscribe(self, *channels: ChannelT):
1279 """
1280 Unsubscribe from keyspace notification channels.
1281
1282 Automatically detects whether each channel is a pattern or exact
1283 channel and uses the appropriate Redis unsubscribe command.
1284
1285 Args:
1286 *channels: Channels to unsubscribe from.
1287 """
1288 patterns = []
1289 exact_channels = []
1290
1291 for channel in channels:
1292 if hasattr(channel, "_channel_str"):
1293 channel_str = str(channel)
1294 else:
1295 channel_str = safe_str(channel)
1296 if _is_pattern(channel):
1297 patterns.append(channel_str)
1298 else:
1299 exact_channels.append(channel_str)
1300
1301 # Delegate to subclass implementation first. For standalone Redis
1302 # this raises on failure, keeping tracking state intact. For cluster
1303 # implementations the operation is best-effort (partial failures are
1304 # logged, not raised) so tracking state is always removed afterwards
1305 # — this is intentional: the user asked to unsubscribe, so
1306 # refresh_subscriptions should not re-subscribe these channels.
1307 self._execute_unsubscribe(patterns, exact_channels)
1308 self._untrack_subscribe(patterns, exact_channels)
1309
1310 @abstractmethod
1311 def _execute_unsubscribe(
1312 self, patterns: list[str], exact_channels: list[str]
1313 ) -> None:
1314 """
1315 Execute the unsubscribe operation.
1316
1317 Args:
1318 patterns: List of pattern strings to punsubscribe from
1319 exact_channels: List of channel strings to unsubscribe from
1320 """
1321 pass
1322
1323 def _track_subscribe(
1324 self, patterns: dict[str, Any], exact_channels: dict[str, Any]
1325 ) -> None:
1326 """Track newly subscribed patterns/channels.
1327
1328 Override in subclasses that need to maintain their own subscription
1329 registry (e.g. cluster implementations that must re-subscribe
1330 new/failed-over nodes). The default is a no-op because standalone
1331 implementations delegate tracking to the underlying PubSub object.
1332 """
1333
1334 def _untrack_subscribe(
1335 self, patterns: list[str], exact_channels: list[str]
1336 ) -> None:
1337 """Remove patterns/channels from the subscription registry.
1338
1339 Override in subclasses that maintain their own subscription registry.
1340 The default is a no-op.
1341 """
1342
1343 def subscribe_keyspace(
1344 self,
1345 key_or_pattern: str,
1346 db: int = 0,
1347 handler: SyncHandlerT | None = None,
1348 ):
1349 """
1350 Subscribe to keyspace notifications for specific keys.
1351
1352 Args:
1353 key_or_pattern: The key or pattern to monitor. Use '*' for wildcards.
1354 db: The database number (default 0).
1355 handler: Optional callback for notifications.
1356
1357 Example:
1358 >>> ksn.subscribe_keyspace("user:123", db=0)
1359 >>> ksn.subscribe_keyspace("user:*", db=0)
1360 """
1361 channel = KeyspaceChannel(key_or_pattern, db=db)
1362 self.subscribe(channel, handler=handler)
1363
1364 def subscribe_keyevent(
1365 self,
1366 event: str,
1367 db: int = 0,
1368 handler: SyncHandlerT | None = None,
1369 ):
1370 """
1371 Subscribe to keyevent notifications for specific event types.
1372
1373 Args:
1374 event: The event type to monitor (e.g., EventType.SET or "set")
1375 db: The database number (default 0).
1376 handler: Optional callback for notifications.
1377
1378 Example:
1379 >>> ksn.subscribe_keyevent(EventType.SET)
1380 >>> ksn.subscribe_keyevent(EventType.EXPIRED, handler=my_handler)
1381 """
1382 channel = KeyeventChannel(event, db=db)
1383 self.subscribe(channel, handler=handler)
1384
1385 def subscribe_subkeyspace(
1386 self,
1387 key_or_pattern: str,
1388 db: int = 0,
1389 handler: SyncHandlerT | None = None,
1390 ):
1391 """
1392 Subscribe to subkeyspace notifications for specific keys.
1393
1394 Receives events with affected subkeys (fields) for the given key.
1395
1396 Args:
1397 key_or_pattern: The key or pattern to monitor.
1398 db: The database number (default 0).
1399 handler: Optional callback for notifications.
1400 """
1401 channel = SubkeyspaceChannel(key_or_pattern, db=db)
1402 self.subscribe(channel, handler=handler)
1403
1404 def subscribe_subkeyevent(
1405 self,
1406 event: str,
1407 db: int = 0,
1408 handler: SyncHandlerT | None = None,
1409 ):
1410 """
1411 Subscribe to subkeyevent notifications for specific event types.
1412
1413 Receives the affected key and subkeys when the given event occurs.
1414
1415 Args:
1416 event: The event type to monitor (e.g., "hset", "hdel").
1417 db: The database number (default 0).
1418 handler: Optional callback for notifications.
1419 """
1420 channel = SubkeyeventChannel(event, db=db)
1421 self.subscribe(channel, handler=handler)
1422
1423 def subscribe_subkeyspaceitem(
1424 self,
1425 key_or_pattern: str,
1426 subkey_or_pattern: str,
1427 db: int = 0,
1428 handler: SyncHandlerT | None = None,
1429 ):
1430 """
1431 Subscribe to subkeyspaceitem notifications for a specific subkey.
1432
1433 Receives the event type when the given subkey of the given key is
1434 modified.
1435
1436 Args:
1437 key_or_pattern: The key or pattern to monitor.
1438 subkey_or_pattern: The subkey (field) or pattern to monitor.
1439 db: The database number (default 0).
1440 handler: Optional callback for notifications.
1441 """
1442 channel = SubkeyspaceitemChannel(key_or_pattern, subkey_or_pattern, db=db)
1443 self.subscribe(channel, handler=handler)
1444
1445 def subscribe_subkeyspaceevent(
1446 self,
1447 event: str,
1448 key_or_pattern: str,
1449 db: int = 0,
1450 handler: SyncHandlerT | None = None,
1451 ):
1452 """
1453 Subscribe to subkeyspaceevent notifications for an event on a key.
1454
1455 Receives the affected subkeys when the given event occurs on the
1456 given key.
1457
1458 Args:
1459 event: The event type to monitor.
1460 key_or_pattern: The key or pattern to monitor.
1461 db: The database number (default 0).
1462 handler: Optional callback for notifications.
1463 """
1464 channel = SubkeyspaceeventChannel(event, key_or_pattern, db=db)
1465 self.subscribe(channel, handler=handler)
1466
1467 def __enter__(self):
1468 return self
1469
1470 def __exit__(self, _exc_type, _exc_val, _exc_tb):
1471 self.close()
1472 return False
1473
1474 def run_in_thread(
1475 self,
1476 poll_timeout: float = 0.0,
1477 daemon: bool = False,
1478 exception_handler: Callable[
1479 [
1480 BaseException,
1481 KeyspaceNotificationsInterface,
1482 KeyspaceWorkerThread,
1483 ],
1484 None,
1485 ]
1486 | None = None,
1487 ) -> KeyspaceWorkerThread:
1488 """
1489 Start a background thread that polls for notifications and triggers handlers.
1490
1491 This method spawns a thread that continuously calls get_message() to
1492 process incoming notifications. When a notification arrives, any
1493 registered handler for that channel/pattern is invoked automatically.
1494
1495 All subscriptions must have handlers registered before calling this method.
1496
1497 Args:
1498 poll_timeout: Timeout in seconds for get_message() calls. When no message
1499 is available, the thread waits up to this long before checking
1500 again. Default 0.0 (non-blocking). WARNING: the default
1501 causes a CPU spin-loop. It is preferred to pass a positive
1502 value (e.g. 0.1 or 1.0).
1503 daemon: If True, the thread will be a daemon thread and will be
1504 terminated when the main program exits. Default False.
1505 exception_handler: Optional callback invoked when an exception occurs
1506 in the worker thread. Receives (exception, notifications,
1507 thread) as arguments. If None, exceptions are raised.
1508
1509 Returns:
1510 KeyspaceWorkerThread: The started worker thread. Call stop() on it
1511 to stop the thread and close the notifications.
1512
1513 Raises:
1514 RedisError: If any subscription doesn't have a handler registered.
1515
1516 Example:
1517 >>> def my_handler(notification):
1518 ... print(f"Got: {notification.key} - {notification.event_type}")
1519 >>>
1520 >>> notifications.subscribe(KeyspaceChannel("user:*"), handler=my_handler)
1521 >>> thread = notifications.run_in_thread(poll_timeout=0.1, daemon=True)
1522 >>> # ... handlers are called automatically ...
1523 >>> thread.stop()
1524 """
1525 self._validate_all_handlers()
1526
1527 thread = KeyspaceWorkerThread(
1528 self,
1529 poll_timeout,
1530 daemon=daemon,
1531 exception_handler=exception_handler,
1532 )
1533 thread.start()
1534 return thread
1535
1536 @abstractmethod
1537 def _validate_all_handlers(self) -> None:
1538 """Raise :class:`~redis.RedisError` if any subscription lacks a handler.
1539
1540 Subclasses inspect their own subscription state to perform
1541 the validation.
1542 """
1543 pass
1544
1545
1546class KeyspaceWorkerThread(threading.Thread):
1547 """
1548 Background thread for processing keyspace notifications.
1549
1550 This thread continuously polls for notifications and invokes registered
1551 handlers. It works with both KeyspaceNotifications (standalone) and
1552 ClusterKeyspaceNotifications.
1553
1554 Example:
1555 >>> thread = notifications.run_in_thread(poll_timeout=0.1)
1556 >>> # ... handlers are called automatically ...
1557 >>> thread.stop()
1558 """
1559
1560 def __init__(
1561 self,
1562 notifications: KeyspaceNotificationsInterface,
1563 poll_timeout: float,
1564 daemon: bool = False,
1565 exception_handler: Callable[
1566 [
1567 BaseException,
1568 KeyspaceNotificationsInterface,
1569 KeyspaceWorkerThread,
1570 ],
1571 None,
1572 ]
1573 | None = None,
1574 ):
1575 super().__init__()
1576 self.daemon = daemon
1577 self.notifications = notifications
1578 self.poll_timeout = poll_timeout
1579 self.exception_handler = exception_handler
1580 self._running = threading.Event()
1581
1582 def run(self) -> None:
1583 """Main loop that polls for notifications and triggers handlers."""
1584 if self._running.is_set():
1585 return
1586 self._running.set()
1587 notifications = self.notifications
1588 poll_timeout = self.poll_timeout
1589 while self._running.is_set():
1590 try:
1591 notifications.get_message(
1592 ignore_subscribe_messages=True, timeout=poll_timeout
1593 )
1594 except BaseException as e:
1595 if self.exception_handler is None:
1596 raise
1597 self.exception_handler(e, notifications, self)
1598 notifications.close()
1599
1600 def stop(self) -> None:
1601 """
1602 Stop the worker thread.
1603
1604 This signals the thread to exit its run loop. The thread will close
1605 the notifications object before terminating.
1606 """
1607 self._running.clear()
1608
1609
1610# =============================================================================
1611# Standalone Keyspace Notification Manager
1612# =============================================================================
1613
1614
1615class KeyspaceNotifications(AbstractKeyspaceNotifications):
1616 """
1617 Manages keyspace notification subscriptions for standalone Redis.
1618
1619 For standalone Redis, keyspace notifications work with a single PubSub
1620 connection. This class wraps that connection and provides:
1621 - Automatic pattern vs exact channel detection
1622 - KeyNotification parsing with optional key_prefix filtering
1623 - Convenience methods for keyspace and keyevent subscriptions
1624 - Context manager and run_in_thread support
1625 """
1626
1627 def __init__(
1628 self,
1629 redis_client: Redis,
1630 key_prefix: str | bytes | None = None,
1631 ignore_subscribe_messages: bool = True,
1632 ):
1633 """
1634 Initialize the standalone keyspace notification manager.
1635
1636 Note: Keyspace notifications must be enabled on the Redis server via
1637 the ``notify-keyspace-events`` configuration option. This is a server-side
1638 configuration that should be done by your infrastructure/operations team.
1639
1640 Args:
1641 redis_client: A Redis client instance
1642 key_prefix: Optional prefix to filter and strip from keys in notifications
1643 ignore_subscribe_messages: If True, subscribe/unsubscribe confirmations
1644 are not returned by get_message/listen
1645 """
1646 super().__init__(key_prefix, ignore_subscribe_messages)
1647 self.redis = redis_client
1648
1649 # Create the PubSub instance with ignore_subscribe_messages=False
1650 # so that the per-call argument in get_message() can control behavior
1651 self._pubsub = redis_client.pubsub(ignore_subscribe_messages=False)
1652
1653 def _execute_subscribe(
1654 self, patterns: dict[str, Any], exact_channels: dict[str, Any]
1655 ) -> None:
1656 """Execute subscribe on the single pubsub connection."""
1657 if patterns:
1658 self._pubsub.psubscribe(**patterns)
1659 if exact_channels:
1660 self._pubsub.subscribe(**exact_channels)
1661
1662 def _execute_unsubscribe(
1663 self, patterns: list[str], exact_channels: list[str]
1664 ) -> None:
1665 """Execute unsubscribe on the single pubsub connection."""
1666 if patterns:
1667 self._pubsub.punsubscribe(*patterns)
1668 if exact_channels:
1669 self._pubsub.unsubscribe(*exact_channels)
1670
1671 def get_message(
1672 self,
1673 ignore_subscribe_messages: bool | None = None,
1674 timeout: float = 0.0,
1675 ) -> KeyNotification | None:
1676 """
1677 Get the next keyspace notification if one is available.
1678
1679 Note: If a handler was registered for the channel, pubsub will call
1680 the handler directly and this method returns None for that message.
1681
1682 Args:
1683 ignore_subscribe_messages: If True, skip subscribe/unsubscribe messages.
1684 Defaults to the value set in __init__ (True).
1685 timeout: Time to wait for a message.
1686
1687 Returns:
1688 A KeyNotification if a notification is available and no handler
1689 was registered for the channel, None otherwise.
1690 """
1691 if ignore_subscribe_messages is None:
1692 ignore_subscribe_messages = self.ignore_subscribe_messages
1693
1694 if self._closed:
1695 return None
1696
1697 # Pubsub's get_message will call wrapped handlers directly for channels
1698 # with registered handlers and return None. For channels without handlers,
1699 # it returns the raw message which we parse to KeyNotification.
1700 message = self._pubsub.get_message(
1701 ignore_subscribe_messages=ignore_subscribe_messages,
1702 timeout=timeout,
1703 )
1704
1705 if message is not None:
1706 return KeyNotification.from_message(message, key_prefix=self.key_prefix)
1707
1708 return None
1709
1710 def listen(self):
1711 """
1712 Listen for keyspace notifications.
1713
1714 This is a generator that yields KeyNotification objects as they arrive.
1715 It blocks until a notification is received.
1716
1717 Yields:
1718 KeyNotification objects for each keyspace/keyevent notification.
1719
1720 Example:
1721 >>> for notification in ksn.listen():
1722 ... print(f"{notification.key}: {notification.event_type}")
1723 """
1724 while self.subscribed:
1725 notification = self.get_message(timeout=1.0)
1726 if notification is not None:
1727 yield notification
1728
1729 @property
1730 def subscribed(self) -> bool:
1731 """Check if there are any active subscriptions and not closed."""
1732 return not self._closed and self._pubsub.subscribed
1733
1734 def _validate_all_handlers(self) -> None:
1735 """Raise if any subscription in the underlying PubSub lacks a handler."""
1736 for channel, handler in self._pubsub.channels.items():
1737 if handler is None:
1738 raise RedisError(f"Channel '{channel}' has no handler registered")
1739 for pattern, handler in self._pubsub.patterns.items():
1740 if handler is None:
1741 raise RedisError(f"Pattern '{pattern}' has no handler registered")
1742
1743 def close(self):
1744 """Close the pubsub connection and clean up resources."""
1745 self._closed = True
1746 try:
1747 self._pubsub.close()
1748 except Exception:
1749 pass
1750
1751
1752# =============================================================================
1753# Cluster-Aware Keyspace Notification Manager
1754# =============================================================================
1755
1756
1757class ClusterKeyspaceNotifications(AbstractKeyspaceNotifications):
1758 """
1759 Manages keyspace notification subscriptions across all nodes in a Redis Cluster.
1760
1761 In Redis Cluster, keyspace notifications are NOT broadcast between nodes.
1762 Each node only emits notifications for keys it owns. This class automatically
1763 subscribes to all primary nodes in the cluster and handles topology changes.
1764 """
1765
1766 def __init__(
1767 self,
1768 redis_cluster: RedisCluster,
1769 key_prefix: str | bytes | None = None,
1770 ignore_subscribe_messages: bool = True,
1771 ):
1772 """
1773 Initialize the cluster keyspace notification manager.
1774
1775 Note: Keyspace notifications must be enabled on all Redis cluster nodes via
1776 the ``notify-keyspace-events`` configuration option. This is a server-side
1777 configuration that should be done by your infrastructure/operations team.
1778
1779 Args:
1780 redis_cluster: A RedisCluster instance
1781 key_prefix: Optional prefix to filter and strip from keys in notifications
1782 ignore_subscribe_messages: If True, subscribe/unsubscribe confirmations
1783 are not returned by get_message/listen
1784 """
1785 super().__init__(key_prefix, ignore_subscribe_messages)
1786 self.cluster = redis_cluster
1787
1788 # Canonical subscription registry: pattern/channel -> wrapped handler.
1789 # In cluster mode there are multiple PubSub objects (one per node), so
1790 # this is the single source of truth used to (re-)subscribe new or
1791 # failed-over nodes.
1792 self._subscribed_patterns: dict[str, Any] = {}
1793 self._subscribed_channels: dict[str, Any] = {}
1794
1795 # Track subscriptions per node
1796 self._node_pubsubs: dict[str, Any] = {}
1797
1798 # Lock for topology refresh operations
1799 self._refresh_lock = threading.Lock()
1800
1801 # Current pubsub index for round-robin polling
1802 self._poll_index = 0
1803
1804 @property
1805 def subscribed(self) -> bool:
1806 """Check if there are any active subscriptions and not closed."""
1807 return not self._closed and bool(
1808 self._subscribed_patterns or self._subscribed_channels
1809 )
1810
1811 def _track_subscribe(
1812 self, patterns: dict[str, Any], exact_channels: dict[str, Any]
1813 ) -> None:
1814 """Track newly subscribed patterns/channels in the cluster registry."""
1815 if patterns:
1816 self._subscribed_patterns.update(patterns)
1817 if exact_channels:
1818 self._subscribed_channels.update(exact_channels)
1819
1820 def _untrack_subscribe(
1821 self, patterns: list[str], exact_channels: list[str]
1822 ) -> None:
1823 """Remove patterns/channels from the cluster registry."""
1824 for p in patterns:
1825 self._subscribed_patterns.pop(p, None)
1826 for c in exact_channels:
1827 self._subscribed_channels.pop(c, None)
1828
1829 def _validate_all_handlers(self) -> None:
1830 """Raise if any subscription in the cluster registry lacks a handler."""
1831 for channel, handler in self._subscribed_channels.items():
1832 if handler is None:
1833 raise RedisError(f"Channel '{channel}' has no handler registered")
1834 for pattern, handler in self._subscribed_patterns.items():
1835 if handler is None:
1836 raise RedisError(f"Pattern '{pattern}' has no handler registered")
1837
1838 def _get_all_primary_nodes(self):
1839 """Get all primary nodes in the cluster."""
1840 return self.cluster.get_primaries()
1841
1842 def _cleanup_node(self, node_name: str) -> None:
1843 """Remove and close a node's PubSub.
1844
1845 Closing the ``PubSub`` disconnects its connection so it is not
1846 left in a subscribed state inside the connection pool.
1847 """
1848 pubsub = self._node_pubsubs.pop(node_name, None)
1849 if pubsub:
1850 try:
1851 pubsub.close()
1852 except Exception:
1853 pass
1854
1855 def _ensure_node_pubsub(self, node) -> Any:
1856 """Get or create a PubSub instance for a node."""
1857 if node.name not in self._node_pubsubs:
1858 redis_conn = self.cluster.get_redis_connection(node)
1859 # Always create PubSub with ignore_subscribe_messages=False
1860 # so that the per-call argument in get_message() can control
1861 # the behavior reliably
1862 pubsub = redis_conn.pubsub(ignore_subscribe_messages=False)
1863 self._node_pubsubs[node.name] = pubsub
1864 return self._node_pubsubs[node.name]
1865
1866 def _execute_subscribe(
1867 self, patterns: dict[str, Any], exact_channels: dict[str, Any]
1868 ) -> None:
1869 """Execute subscribe on all cluster nodes.
1870
1871 Patterns and exact channels are subscribed in a single pass over
1872 nodes so that a mid-batch node failure cannot create a
1873 partially-caught-up replacement. If a node fails during this
1874 call it is removed from ``_node_pubsubs`` and will be fully
1875 re-subscribed on the next ``refresh_subscriptions`` cycle.
1876
1877 If a newly discovered node is encountered (not yet in
1878 ``_node_pubsubs``), it is also subscribed to all *previously*
1879 tracked patterns/channels so it doesn't miss notifications for
1880 subscriptions that were established before this node joined.
1881 """
1882 if not patterns and not exact_channels:
1883 return
1884
1885 failed_nodes: list[str] = []
1886 for node in self._get_all_primary_nodes():
1887 is_new_node = node.name not in self._node_pubsubs
1888 pubsub = self._ensure_node_pubsub(node)
1889 try:
1890 # If this is a brand-new node, catch it up on existing
1891 # subscriptions before adding the new channels.
1892 if is_new_node:
1893 if self._subscribed_patterns:
1894 pubsub.psubscribe(**self._subscribed_patterns)
1895 if self._subscribed_channels:
1896 pubsub.subscribe(**self._subscribed_channels)
1897
1898 if patterns:
1899 pubsub.psubscribe(**patterns)
1900 if exact_channels:
1901 pubsub.subscribe(**exact_channels)
1902 except Exception:
1903 # Remove the broken pubsub so refresh_subscriptions can
1904 # re-create it later.
1905 self._cleanup_node(node.name)
1906 failed_nodes.append(node.name)
1907
1908 if failed_nodes:
1909 logger.warning(
1910 "Failed to subscribe on cluster nodes: %s. "
1911 "These nodes will be retried on the next refresh cycle.",
1912 ", ".join(failed_nodes),
1913 )
1914
1915 def _execute_unsubscribe(
1916 self, patterns: list[str], exact_channels: list[str]
1917 ) -> None:
1918 """Execute unsubscribe on all cluster nodes."""
1919 if patterns:
1920 self._unsubscribe_from_all_nodes(patterns, use_punsubscribe=True)
1921 if exact_channels:
1922 self._unsubscribe_from_all_nodes(exact_channels, use_punsubscribe=False)
1923
1924 def _unsubscribe_from_all_nodes(self, channels: list[str], use_punsubscribe: bool):
1925 """Unsubscribe from patterns/channels on all nodes.
1926
1927 Best-effort: tries every node so that a single broken connection
1928 does not prevent the remaining nodes from being unsubscribed.
1929 Broken pubsubs are cleaned up; the tracking state is still removed
1930 by the caller, so ``refresh_subscriptions`` will *not* re-subscribe
1931 these channels on replacement nodes.
1932 """
1933 failed_nodes: list[str] = []
1934 for node_name, pubsub in list(self._node_pubsubs.items()):
1935 try:
1936 if use_punsubscribe:
1937 pubsub.punsubscribe(*channels)
1938 else:
1939 pubsub.unsubscribe(*channels)
1940 except Exception:
1941 self._cleanup_node(node_name)
1942 failed_nodes.append(node_name)
1943
1944 if failed_nodes:
1945 logger.warning(
1946 "Failed to unsubscribe on cluster nodes: %s. "
1947 "These nodes will be re-created on the next refresh cycle.",
1948 ", ".join(failed_nodes),
1949 )
1950
1951 def get_message(
1952 self,
1953 ignore_subscribe_messages: bool | None = None,
1954 timeout: float = 0.0,
1955 ) -> KeyNotification | None:
1956 """
1957 Get the next keyspace notification if one is available.
1958
1959 This method polls all node pubsubs in round-robin fashion until
1960 a message is received or the timeout expires.
1961 If a connection error occurs, subscriptions are automatically refreshed.
1962
1963 Args:
1964 ignore_subscribe_messages: If True, skip subscribe/unsubscribe messages.
1965 Defaults to the value set in __init__ (True).
1966 timeout: Total time to wait for a message (distributed across all nodes)
1967
1968 Returns:
1969 A KeyNotification if a notification is available, None otherwise.
1970 """
1971 if self._closed:
1972 return None
1973
1974 total_nodes = len(self._node_pubsubs)
1975 if total_nodes == 0:
1976 # Sleep for the requested timeout so callers that loop
1977 # (run_in_thread, listen) don't spin the CPU when all node
1978 # connections have been cleaned up.
1979 if timeout > 0:
1980 time.sleep(timeout)
1981 return None
1982
1983 # Use instance default if not specified
1984 if ignore_subscribe_messages is None:
1985 ignore_subscribe_messages = self.ignore_subscribe_messages
1986
1987 # Handle timeout=0 as a single non-blocking poll over all pubsubs
1988 # This matches the expected semantics of PubSub.get_message(timeout=0)
1989 if timeout == 0.0:
1990 return self._poll_all_nodes_once(ignore_subscribe_messages)
1991
1992 # Calculate per-node timeout for each poll
1993 # Use a small timeout per node to allow round-robin polling
1994 per_node_timeout = min(0.1, timeout / max(total_nodes, 1))
1995
1996 start_time = time.monotonic()
1997 end_time = start_time + timeout
1998
1999 while True:
2000 # Check if we've exceeded the total timeout
2001 if time.monotonic() >= end_time:
2002 return None
2003
2004 pubsubs = list(self._node_pubsubs.values())
2005 if not pubsubs:
2006 return None
2007
2008 # Round-robin polling
2009 self._poll_index = self._poll_index % len(pubsubs)
2010 pubsub = pubsubs[self._poll_index]
2011 self._poll_index += 1
2012
2013 try:
2014 message = pubsub.get_message(
2015 ignore_subscribe_messages=ignore_subscribe_messages,
2016 timeout=per_node_timeout,
2017 )
2018 except (ConnectionError, TimeoutError, RedisError):
2019 # Connection error - refresh subscriptions and continue
2020 self._refresh_subscriptions_on_error()
2021 continue
2022
2023 if message is not None:
2024 # Note: If a handler was registered, PubSub already invoked it
2025 # and returned None, so we only reach here for handler-less subscriptions
2026 notification = KeyNotification.from_message(
2027 message, key_prefix=self.key_prefix
2028 )
2029 if notification is not None:
2030 return notification
2031 # If not a keyspace notification, continue checking other nodes
2032
2033 def _poll_all_nodes_once(
2034 self, ignore_subscribe_messages: bool
2035 ) -> KeyNotification | None:
2036 """
2037 Perform a single non-blocking poll over all node pubsubs.
2038
2039 This is used when timeout=0 to match the expected semantics of
2040 PubSub.get_message(timeout=0) - a non-blocking check for messages.
2041
2042 Returns:
2043 A KeyNotification if one is available, None otherwise.
2044 """
2045 had_error = False
2046 for pubsub in list(self._node_pubsubs.values()):
2047 try:
2048 message = pubsub.get_message(
2049 ignore_subscribe_messages=ignore_subscribe_messages,
2050 timeout=0.0,
2051 )
2052 except (ConnectionError, TimeoutError, RedisError):
2053 # Record the error but continue polling remaining healthy
2054 # nodes so that already-buffered notifications are not lost.
2055 had_error = True
2056 continue
2057
2058 if message is not None:
2059 # Note: If a handler was registered, PubSub already invoked it
2060 # and returned None, so we only reach here for handler-less subscriptions
2061 notification = KeyNotification.from_message(
2062 message, key_prefix=self.key_prefix
2063 )
2064 if notification is not None:
2065 # Refresh before returning if any node had an error,
2066 # so the next poll cycle has fresh state.
2067 if had_error:
2068 self._refresh_subscriptions_on_error()
2069 return notification
2070
2071 # Refresh after polling all nodes if any had errors
2072 if had_error:
2073 self._refresh_subscriptions_on_error()
2074 return None
2075
2076 def listen(self):
2077 """
2078 Listen for keyspace notifications from all cluster nodes.
2079
2080 This is a generator that yields KeyNotification objects as they arrive.
2081 It blocks until a notification is received.
2082
2083 Yields:
2084 KeyNotification objects for each keyspace/keyevent notification.
2085
2086 Example:
2087 >>> for notification in ksn.listen():
2088 ... print(f"{notification.key}: {notification.event_type}")
2089 """
2090 while self.subscribed:
2091 notification = self.get_message(timeout=1.0)
2092 if notification is not None:
2093 yield notification
2094
2095 def _refresh_subscriptions_on_error(self):
2096 """
2097 Refresh subscriptions after a connection error.
2098
2099 This is called automatically when a connection error occurs during
2100 get_message(). It checks if nodes changed before refreshing.
2101 """
2102 self._poll_index = 0 # Reset round-robin index
2103
2104 try:
2105 self.refresh_subscriptions()
2106 except Exception:
2107 logger.warning(
2108 "Failed to refresh cluster subscriptions, will retry on next error",
2109 exc_info=True,
2110 )
2111
2112 def _is_pubsub_connected(self, pubsub) -> bool:
2113 """Check if a pubsub connection is still alive."""
2114 try:
2115 conn = pubsub.connection
2116 if conn is None:
2117 return False
2118 return conn.is_connected
2119 except Exception:
2120 return False
2121
2122 def refresh_subscriptions(self):
2123 """
2124 Refresh subscriptions after a topology change.
2125
2126 This method is called automatically when topology changes are detected
2127 or when connection errors occur. You can also call it manually if needed.
2128
2129 This method:
2130 1. Discovers any new primary nodes and subscribes them
2131 2. Removes pubsubs for nodes that are no longer primaries
2132 3. Re-creates broken pubsub connections for existing nodes
2133 """
2134 with self._refresh_lock:
2135 current_primaries = {
2136 node.name: node for node in self._get_all_primary_nodes()
2137 }
2138
2139 # Remove pubsubs for nodes that are no longer primaries
2140 removed_nodes = set(self._node_pubsubs.keys()) - set(
2141 current_primaries.keys()
2142 )
2143 for node_name in removed_nodes:
2144 self._cleanup_node(node_name)
2145
2146 # Detect broken connections for existing nodes and remove them
2147 # so they get re-created below
2148 existing_nodes = set(self._node_pubsubs.keys()) & set(
2149 current_primaries.keys()
2150 )
2151 for node_name in existing_nodes:
2152 pubsub = self._node_pubsubs.get(node_name)
2153 if pubsub and not self._is_pubsub_connected(pubsub):
2154 # Connection is broken, remove it so it gets re-created
2155 self._cleanup_node(node_name)
2156
2157 # Subscribe new nodes (and nodes with broken connections) to existing
2158 # patterns/channels
2159 new_nodes = set(current_primaries.keys()) - set(self._node_pubsubs.keys())
2160 failed_nodes: list[str] = []
2161 for node_name in new_nodes:
2162 node = current_primaries[node_name]
2163 pubsub = self._ensure_node_pubsub(node)
2164
2165 try:
2166 if self._subscribed_patterns:
2167 pubsub.psubscribe(**self._subscribed_patterns)
2168 if self._subscribed_channels:
2169 pubsub.subscribe(**self._subscribed_channels)
2170 except Exception:
2171 # Subscription failed - remove from dict so retry is possible
2172 self._cleanup_node(node_name)
2173 failed_nodes.append(node_name)
2174
2175 # Raise after attempting all nodes so we don't skip any
2176 if failed_nodes:
2177 raise ConnectionError(
2178 f"Failed to subscribe to cluster nodes: {', '.join(failed_nodes)}"
2179 )
2180
2181 def close(self):
2182 """Close all pubsub connections and clean up resources."""
2183 self._closed = True
2184 for node_name in list(self._node_pubsubs.keys()):
2185 self._cleanup_node(node_name)
2186 self._subscribed_patterns.clear()
2187 self._subscribed_channels.clear()