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