Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/redis/keyspace_notifications.py: 35%

Shortcuts on this page

r m x   toggle line displays

j k   next/prev highlighted chunk

0   (zero) top of page

1   (one) first highlighted chunk

813 statements  

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()