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

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

800 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.index(b":", pos) 

241 length = int(raw[pos:colon]) 

242 start = colon + 1 

243 subkeys.append(raw[start : start + length].decode("utf-8", "replace")) 

244 pos = start + length 

245 if pos < len(raw) and raw[pos] == ord(","): 

246 pos += 1 # skip comma separator 

247 return subkeys 

248 

249 

250@dataclass 

251class KeyNotification: 

252 """ 

253 Represents a parsed Redis keyspace, keyevent, or subkey notification. 

254 

255 This class provides convenient access to the notification details 

256 like key, event type, database number, and affected subkeys. 

257 

258 Attributes: 

259 key: The Redis key that was affected (for keyspace notifications) 

260 or the key name from the message data (for keyevent notifications) 

261 event_type: The type of operation that occurred (e.g., "set", "del"). 

262 This is a plain string, so new Redis events work automatically. 

263 Compare against EventType constants or any string. 

264 database: The database number where the event occurred 

265 channel: The original channel name 

266 is_keyspace: True if this is a keyspace notification, False for keyevent 

267 data: The raw data payload from the notification message. 

268 subkeys: List of affected subkeys (fields) for subkey notifications. 

269 Empty list for regular keyspace/keyevent notifications. 

270 """ 

271 

272 # Regex patterns for parsing keyspace/keyevent channels 

273 # Pattern: __keyspace@<db>__:<key> or __keyevent@<db>__:<event> 

274 _KEYSPACE_PATTERN: ClassVar[re.Pattern] = re.compile( 

275 r"^__keyspace@(\d+|\*)__:(.+)$" 

276 ) 

277 _KEYEVENT_PATTERN: ClassVar[re.Pattern] = re.compile( 

278 r"^__keyevent@(\d+|\*)__:(.+)$" 

279 ) 

280 _SUBKEYSPACE_PATTERN: ClassVar[re.Pattern] = re.compile( 

281 r"^__subkeyspace@(\d+|\*)__:(.+)$" 

282 ) 

283 _SUBKEYEVENT_PATTERN: ClassVar[re.Pattern] = re.compile( 

284 r"^__subkeyevent@(\d+|\*)__:(.+)$" 

285 ) 

286 _SUBKEYSPACEITEM_PATTERN: ClassVar[re.Pattern] = re.compile( 

287 r"^__subkeyspaceitem@(\d+|\*)__:(.+)$", re.DOTALL 

288 ) 

289 _SUBKEYSPACEEVENT_PATTERN: ClassVar[re.Pattern] = re.compile( 

290 r"^__subkeyspaceevent@(\d+|\*)__:(.+)$" 

291 ) 

292 

293 key: str 

294 event_type: str 

295 database: int 

296 channel: str 

297 is_keyspace: bool 

298 data: str 

299 subkeys: list[str] = None # type: ignore[assignment] 

300 

301 def __post_init__(self): 

302 if self.subkeys is None: 

303 self.subkeys = [] 

304 

305 @classmethod 

306 def from_message( 

307 cls, 

308 message: dict[str, Any] | None, 

309 key_prefix: str | bytes | None = None, 

310 ) -> KeyNotification | None: 

311 """ 

312 Parse a pub/sub message into a KeyNotification. 

313 

314 Args: 

315 message: A pub/sub message dict with 'channel', 'data', and 'type' keys 

316 key_prefix: Optional prefix to filter and strip from keys. 

317 If provided, only notifications for keys starting with 

318 this prefix will be returned, and the prefix will be 

319 stripped from the key. 

320 

321 Returns: 

322 A KeyNotification if the message is a valid keyspace/keyevent 

323 notification, None otherwise. 

324 

325 Example: 

326 >>> message = { 

327 ... 'type': 'pmessage', 

328 ... 'pattern': '__keyspace@0__:user:*', 

329 ... 'channel': '__keyspace@0__:user:123', 

330 ... 'data': 'set' 

331 ... } 

332 >>> notification = KeyNotification.from_message(message) 

333 >>> notification.key 

334 'user:123' 

335 >>> notification.event_type 

336 'set' 

337 """ 

338 if message is None: 

339 return None 

340 

341 msg_type = message.get("type") 

342 if msg_type not in ("message", "pmessage"): 

343 return None 

344 

345 channel = message.get("channel") 

346 data = message.get("data") 

347 

348 if channel is None or data is None: 

349 return None 

350 

351 return cls.try_parse(channel, data, key_prefix) 

352 

353 @classmethod 

354 def try_parse( 

355 cls, 

356 channel: str | bytes, 

357 data: str | bytes, 

358 key_prefix: str | bytes | None = None, 

359 ) -> KeyNotification | None: 

360 """ 

361 Try to parse a channel and data into a KeyNotification. 

362 

363 This is a lower-level method that takes the channel and data directly, 

364 useful when working with callback-based subscription handlers. 

365 

366 Args: 

367 channel: The channel name (e.g., "__keyspace@0__:mykey") 

368 data: The message data (event type for keyspace, key for keyevent) 

369 key_prefix: Optional prefix to filter and strip from keys 

370 

371 Returns: 

372 A KeyNotification if valid, None otherwise. 

373 """ 

374 raw_data = bytes(data) if isinstance(data, (bytes, bytearray)) else None 

375 channel = safe_str(channel) 

376 data = safe_str(data) 

377 

378 return cls._parse(channel, data, key_prefix, raw_data=raw_data) 

379 

380 @classmethod 

381 def _parse( 

382 cls, 

383 channel: str, 

384 data: str, 

385 key_prefix: str | bytes | None = None, 

386 raw_data: bytes | None = None, 

387 ) -> KeyNotification | None: 

388 """Internal parsing logic.""" 

389 # Normalize key_prefix 

390 key_prefix = safe_str(key_prefix) if key_prefix else None 

391 # Byte-domain view of the payload: the server prefixes every element 

392 # with its size in bytes (sdslen), so length-prefixed fields must be 

393 # split before decoding. When the caller only has the decoded string 

394 # (decode_responses=True) the exact bytes are recovered by re-encoding. 

395 raw = raw_data if raw_data is not None else data.encode("utf-8") 

396 

397 # Try keyspace pattern first: __keyspace@<db>__:<key> 

398 match = cls._KEYSPACE_PATTERN.match(channel) 

399 if match: 

400 db_str, key = match.groups() 

401 database = int(db_str) if db_str != "*" else -1 

402 event_type = data # For keyspace, the data is the event type 

403 

404 # Apply key prefix filter 

405 if key_prefix: 

406 if not key.startswith(key_prefix): 

407 return None 

408 key = key[len(key_prefix) :] 

409 

410 return cls( 

411 key=key, 

412 event_type=event_type, 

413 database=database, 

414 channel=channel, 

415 is_keyspace=True, 

416 data=data, 

417 ) 

418 

419 # Try keyevent pattern: __keyevent@<db>__:<event> 

420 match = cls._KEYEVENT_PATTERN.match(channel) 

421 if match: 

422 db_str, event_type = match.groups() 

423 database = int(db_str) if db_str != "*" else -1 

424 key = data # For keyevent, the data is the key 

425 

426 # Apply key prefix filter 

427 if key_prefix: 

428 if not key.startswith(key_prefix): 

429 return None 

430 key = key[len(key_prefix) :] 

431 

432 return cls( 

433 key=key, 

434 event_type=event_type, 

435 database=database, 

436 channel=channel, 

437 is_keyspace=False, 

438 data=data, 

439 ) 

440 

441 # Try subkeyspace: channel=__subkeyspace@<db>__:<key> 

442 # data=<event>|<subkey_len>:<subkey>[,<subkey_len>:<subkey>...] 

443 match = cls._SUBKEYSPACE_PATTERN.match(channel) 

444 if match: 

445 db_str, key = match.groups() 

446 database = int(db_str) if db_str != "*" else -1 

447 pipe_idx = raw.index(b"|") 

448 event_type = raw[:pipe_idx].decode("utf-8", "replace") 

449 subkeys = _parse_length_prefixed_subkeys(raw[pipe_idx + 1 :]) 

450 

451 if key_prefix: 

452 if not key.startswith(key_prefix): 

453 return None 

454 key = key[len(key_prefix) :] 

455 

456 return cls( 

457 key=key, 

458 event_type=event_type, 

459 database=database, 

460 channel=channel, 

461 is_keyspace=True, 

462 data=data, 

463 subkeys=subkeys, 

464 ) 

465 

466 # Try subkeyevent: channel=__subkeyevent@<db>__:<event> 

467 # data=<key_len>:<key>|<subkey_len>:<subkey>[,...] 

468 match = cls._SUBKEYEVENT_PATTERN.match(channel) 

469 if match: 

470 db_str, event_type = match.groups() 

471 database = int(db_str) if db_str != "*" else -1 

472 # Parse key by length prefix 

473 colon_idx = raw.index(b":") 

474 key_len = int(raw[:colon_idx]) 

475 key_start = colon_idx + 1 

476 key = raw[key_start : key_start + key_len].decode("utf-8", "replace") 

477 # After key, expect '|' then subkeys 

478 subkeys_start = key_start + key_len + 1 # +1 for '|' 

479 subkeys = _parse_length_prefixed_subkeys(raw[subkeys_start:]) 

480 

481 if key_prefix: 

482 if not key.startswith(key_prefix): 

483 return None 

484 key = key[len(key_prefix) :] 

485 

486 return cls( 

487 key=key, 

488 event_type=event_type, 

489 database=database, 

490 channel=channel, 

491 is_keyspace=False, 

492 data=data, 

493 subkeys=subkeys, 

494 ) 

495 

496 # Try subkeyspaceitem: channel=__subkeyspaceitem@<db>__:<key>\n<subkey> 

497 # data=<event> 

498 match = cls._SUBKEYSPACEITEM_PATTERN.match(channel) 

499 if match: 

500 db_str, key_and_subkey = match.groups() 

501 database = int(db_str) if db_str != "*" else -1 

502 newline_idx = key_and_subkey.index("\n") 

503 key = key_and_subkey[:newline_idx] 

504 subkey = key_and_subkey[newline_idx + 1 :] 

505 event_type = data 

506 

507 if key_prefix: 

508 if not key.startswith(key_prefix): 

509 return None 

510 key = key[len(key_prefix) :] 

511 

512 return cls( 

513 key=key, 

514 event_type=event_type, 

515 database=database, 

516 channel=channel, 

517 is_keyspace=True, 

518 data=data, 

519 subkeys=[subkey], 

520 ) 

521 

522 # Try subkeyspaceevent: channel=__subkeyspaceevent@<db>__:<event>|<key> 

523 # data=<subkey_len>:<subkey>[,...] 

524 match = cls._SUBKEYSPACEEVENT_PATTERN.match(channel) 

525 if match: 

526 db_str, event_and_key = match.groups() 

527 database = int(db_str) if db_str != "*" else -1 

528 pipe_idx = event_and_key.index("|") 

529 event_type = event_and_key[:pipe_idx] 

530 key = event_and_key[pipe_idx + 1 :] 

531 subkeys = _parse_length_prefixed_subkeys(raw) 

532 

533 if key_prefix: 

534 if not key.startswith(key_prefix): 

535 return None 

536 key = key[len(key_prefix) :] 

537 

538 return cls( 

539 key=key, 

540 event_type=event_type, 

541 database=database, 

542 channel=channel, 

543 is_keyspace=False, 

544 data=data, 

545 subkeys=subkeys, 

546 ) 

547 

548 return None 

549 

550 def key_starts_with(self, prefix: str | bytes) -> bool: 

551 """Check if the key starts with the given prefix.""" 

552 prefix = safe_str(prefix) 

553 return self.key.startswith(prefix) 

554 

555 

556# ============================================================================= 

557# Channel Classes 

558# ============================================================================= 

559 

560 

561class KeyspaceChannel: 

562 """ 

563 Represents a keyspace notification channel for subscribing to events on keys. 

564 

565 Keyspace notifications publish the event type (e.g., "set", "del") as the message 

566 when a key matching the pattern is modified. 

567 

568 This class can be used directly with subscribe()/psubscribe() as it implements 

569 __str__ to return the channel string. 

570 

571 Attributes: 

572 key_or_pattern: The key or pattern to monitor (use '*' for wildcards) 

573 db: The database number (defaults to 0, the only database in Redis Cluster) 

574 is_pattern: Whether this channel contains wildcards 

575 

576 Examples: 

577 >>> channel = KeyspaceChannel("user:123", db=0) 

578 >>> str(channel) 

579 '__keyspace@0__:user:123' 

580 

581 >>> # Pattern subscription (wildcards are auto-detected) 

582 >>> channel = KeyspaceChannel("user:*", db=0) 

583 >>> str(channel) 

584 '__keyspace@0__:user:*' 

585 

586 >>> # Use with KeyspaceNotifications 

587 >>> notifications = KeyspaceNotifications(redis_client) 

588 >>> notifications.subscribe(channel) 

589 """ 

590 

591 PREFIX: ClassVar[str] = "__keyspace@" 

592 

593 def __init__(self, key_or_pattern: str, db: int = 0): 

594 """ 

595 Create a keyspace notification channel. 

596 

597 Args: 

598 key_or_pattern: The key or pattern to monitor. Use '*' for wildcards. 

599 db: The database number. Defaults to 0 (the only database in Redis Cluster). 

600 """ 

601 self.key_or_pattern = key_or_pattern 

602 self.db = db 

603 self._channel_str = self._build_channel_string() 

604 

605 def _build_channel_string(self) -> str: 

606 return f"{self.PREFIX}{self.db}__:{self.key_or_pattern}" 

607 

608 @property 

609 def is_pattern(self) -> bool: 

610 """Check if this channel contains wildcards and should use psubscribe.""" 

611 return _is_pattern(self.key_or_pattern) 

612 

613 def __str__(self) -> str: 

614 return self._channel_str 

615 

616 def __repr__(self) -> str: 

617 return f"KeyspaceChannel({self.key_or_pattern!r}, db={self.db})" 

618 

619 def __eq__(self, other: object) -> bool: 

620 if isinstance(other, KeyspaceChannel): 

621 return self._channel_str == other._channel_str 

622 if isinstance(other, str): 

623 return self._channel_str == other 

624 return NotImplemented 

625 

626 def __hash__(self) -> int: 

627 return hash(self._channel_str) 

628 

629 

630class KeyeventChannel: 

631 """ 

632 Represents a keyevent notification channel for subscribing to event types. 

633 

634 Keyevent notifications publish the key name as the message when the specified 

635 event type occurs on any key. 

636 

637 This class can be used directly with subscribe()/psubscribe() as it implements 

638 __str__ to return the channel string. 

639 

640 Attributes: 

641 event: The event type to monitor 

642 db: The database number (defaults to 0, the only database in Redis Cluster) 

643 is_pattern: Whether this channel contains wildcards 

644 

645 Examples: 

646 >>> channel = KeyeventChannel(EventType.SET, db=0) 

647 >>> str(channel) 

648 '__keyevent@0__:set' 

649 

650 >>> channel = KeyeventChannel.all_events(db=0) 

651 >>> str(channel) 

652 '__keyevent@0__:*' 

653 

654 >>> # Use with KeyspaceNotifications 

655 >>> notifications = KeyspaceNotifications(redis_client) 

656 >>> notifications.subscribe(channel) 

657 """ 

658 

659 PREFIX: ClassVar[str] = "__keyevent@" 

660 

661 def __init__(self, event: str, db: int = 0): 

662 """ 

663 Create a keyevent notification channel. 

664 

665 Args: 

666 event: The event type to monitor (e.g., EventType.SET or "set") 

667 db: The database number. Defaults to 0 (the only database in Redis Cluster). 

668 """ 

669 self.event = event 

670 self.db = db 

671 self._channel_str = self._build_channel_string() 

672 

673 def _build_channel_string(self) -> str: 

674 return f"{self.PREFIX}{self.db}__:{self.event}" 

675 

676 @property 

677 def is_pattern(self) -> bool: 

678 """Check if this channel contains wildcards and should use psubscribe.""" 

679 return _is_pattern(self.event) 

680 

681 @classmethod 

682 def all_events(cls, db: int = 0) -> "KeyeventChannel": 

683 """ 

684 Create a keyevent pattern for subscribing to all event types. 

685 

686 This is equivalent to KeyeventChannel("*"). 

687 

688 Args: 

689 db: The database number. Defaults to 0 (the only database in Redis Cluster). 

690 

691 Returns: 

692 A KeyeventChannel configured to receive all events. 

693 

694 Examples: 

695 >>> channel = KeyeventChannel.all_events() 

696 >>> str(channel) 

697 '__keyevent@0__:*' 

698 """ 

699 return cls("*", db=db) 

700 

701 def __str__(self) -> str: 

702 return self._channel_str 

703 

704 def __repr__(self) -> str: 

705 return f"KeyeventChannel({self.event!r}, db={self.db})" 

706 

707 def __eq__(self, other: object) -> bool: 

708 if isinstance(other, KeyeventChannel): 

709 return self._channel_str == other._channel_str 

710 if isinstance(other, str): 

711 return self._channel_str == other 

712 return NotImplemented 

713 

714 def __hash__(self) -> int: 

715 return hash(self._channel_str) 

716 

717 

718class SubkeyspaceChannel: 

719 """ 

720 Represents a subkeyspace notification channel for subscribing to 

721 subkey-level events on keys (e.g., hash field changes). 

722 

723 The channel format is ``__subkeyspace@<db>__:<key>``. 

724 The message payload is ``<event>|<subkey_len>:<subkey>[,...]``. 

725 

726 Examples: 

727 >>> channel = SubkeyspaceChannel("myhash", db=0) 

728 >>> str(channel) 

729 '__subkeyspace@0__:myhash' 

730 """ 

731 

732 PREFIX: ClassVar[str] = "__subkeyspace@" 

733 

734 def __init__(self, key_or_pattern: str, db: int = 0): 

735 self.key_or_pattern = key_or_pattern 

736 self.db = db 

737 self._channel_str = self._build_channel_string() 

738 

739 def _build_channel_string(self) -> str: 

740 return f"{self.PREFIX}{self.db}__:{self.key_or_pattern}" 

741 

742 @property 

743 def is_pattern(self) -> bool: 

744 return _is_pattern(self.key_or_pattern) 

745 

746 def __str__(self) -> str: 

747 return self._channel_str 

748 

749 def __repr__(self) -> str: 

750 return f"SubkeyspaceChannel({self.key_or_pattern!r}, db={self.db})" 

751 

752 def __eq__(self, other: object) -> bool: 

753 if isinstance(other, SubkeyspaceChannel): 

754 return self._channel_str == other._channel_str 

755 if isinstance(other, str): 

756 return self._channel_str == other 

757 return NotImplemented 

758 

759 def __hash__(self) -> int: 

760 return hash(self._channel_str) 

761 

762 

763class SubkeyeventChannel: 

764 """ 

765 Represents a subkeyevent notification channel for subscribing to 

766 specific event types with subkey-level detail. 

767 

768 The channel format is ``__subkeyevent@<db>__:<event>``. 

769 The message payload is ``<key_len>:<key>|<subkey_len>:<subkey>[,...]``. 

770 

771 Examples: 

772 >>> channel = SubkeyeventChannel("hdel", db=0) 

773 >>> str(channel) 

774 '__subkeyevent@0__:hdel' 

775 """ 

776 

777 PREFIX: ClassVar[str] = "__subkeyevent@" 

778 

779 def __init__(self, event: str, db: int = 0): 

780 self.event = event 

781 self.db = db 

782 self._channel_str = self._build_channel_string() 

783 

784 def _build_channel_string(self) -> str: 

785 return f"{self.PREFIX}{self.db}__:{self.event}" 

786 

787 @property 

788 def is_pattern(self) -> bool: 

789 return _is_pattern(self.event) 

790 

791 @classmethod 

792 def all_events(cls, db: int = 0) -> SubkeyeventChannel: 

793 """Create a channel for all subkeyevent types.""" 

794 return cls("*", db=db) 

795 

796 def __str__(self) -> str: 

797 return self._channel_str 

798 

799 def __repr__(self) -> str: 

800 return f"SubkeyeventChannel({self.event!r}, db={self.db})" 

801 

802 def __eq__(self, other: object) -> bool: 

803 if isinstance(other, SubkeyeventChannel): 

804 return self._channel_str == other._channel_str 

805 if isinstance(other, str): 

806 return self._channel_str == other 

807 return NotImplemented 

808 

809 def __hash__(self) -> int: 

810 return hash(self._channel_str) 

811 

812 

813class SubkeyspaceitemChannel: 

814 """ 

815 Represents a subkeyspaceitem notification channel for subscribing to 

816 events on a specific subkey (field) of a specific key. 

817 

818 The channel format is ``__subkeyspaceitem@<db>__:<key>\\n<subkey>``. 

819 The message payload is the event type (e.g., ``"hset"``). 

820 

821 Note: 

822 The server only emits this notification when the key does not 

823 contain a newline character. 

824 

825 Examples: 

826 >>> channel = SubkeyspaceitemChannel("myhash", "myfield", db=0) 

827 >>> str(channel) 

828 '__subkeyspaceitem@0__:myhash\\nmyfield' 

829 """ 

830 

831 PREFIX: ClassVar[str] = "__subkeyspaceitem@" 

832 

833 def __init__(self, key_or_pattern: str, subkey_or_pattern: str, db: int = 0): 

834 self.key_or_pattern = key_or_pattern 

835 self.subkey_or_pattern = subkey_or_pattern 

836 self.db = db 

837 self._channel_str = self._build_channel_string() 

838 

839 def _build_channel_string(self) -> str: 

840 return ( 

841 f"{self.PREFIX}{self.db}__:{self.key_or_pattern}\n{self.subkey_or_pattern}" 

842 ) 

843 

844 @property 

845 def is_pattern(self) -> bool: 

846 return _is_pattern(self.key_or_pattern) or _is_pattern(self.subkey_or_pattern) 

847 

848 def __str__(self) -> str: 

849 return self._channel_str 

850 

851 def __repr__(self) -> str: 

852 return ( 

853 f"SubkeyspaceitemChannel({self.key_or_pattern!r}, " 

854 f"{self.subkey_or_pattern!r}, db={self.db})" 

855 ) 

856 

857 def __eq__(self, other: object) -> bool: 

858 if isinstance(other, SubkeyspaceitemChannel): 

859 return self._channel_str == other._channel_str 

860 if isinstance(other, str): 

861 return self._channel_str == other 

862 return NotImplemented 

863 

864 def __hash__(self) -> int: 

865 return hash(self._channel_str) 

866 

867 

868class SubkeyspaceeventChannel: 

869 """ 

870 Represents a subkeyspaceevent notification channel for subscribing to 

871 a specific event on a specific key, receiving affected subkeys. 

872 

873 The channel format is ``__subkeyspaceevent@<db>__:<event>|<key>``. 

874 The message payload is a length-prefixed subkey list. 

875 

876 Examples: 

877 >>> channel = SubkeyspaceeventChannel("hset", "myhash", db=0) 

878 >>> str(channel) 

879 '__subkeyspaceevent@0__:hset|myhash' 

880 """ 

881 

882 PREFIX: ClassVar[str] = "__subkeyspaceevent@" 

883 

884 def __init__(self, event: str, key_or_pattern: str, db: int = 0): 

885 self.event = event 

886 self.key_or_pattern = key_or_pattern 

887 self.db = db 

888 self._channel_str = self._build_channel_string() 

889 

890 def _build_channel_string(self) -> str: 

891 return f"{self.PREFIX}{self.db}__:{self.event}|{self.key_or_pattern}" 

892 

893 @property 

894 def is_pattern(self) -> bool: 

895 return _is_pattern(self.event) or _is_pattern(self.key_or_pattern) 

896 

897 def __str__(self) -> str: 

898 return self._channel_str 

899 

900 def __repr__(self) -> str: 

901 return ( 

902 f"SubkeyspaceeventChannel({self.event!r}, " 

903 f"{self.key_or_pattern!r}, db={self.db})" 

904 ) 

905 

906 def __eq__(self, other: object) -> bool: 

907 if isinstance(other, SubkeyspaceeventChannel): 

908 return self._channel_str == other._channel_str 

909 if isinstance(other, str): 

910 return self._channel_str == other 

911 return NotImplemented 

912 

913 def __hash__(self) -> int: 

914 return hash(self._channel_str) 

915 

916 

917class ChannelType(Enum): 

918 """ 

919 Enum representing the type of a Redis keyspace notification channel. 

920 

921 Redis provides two types of keyspace notifications and four subkey 

922 notification types: 

923 

924 - KEYSPACE: ``__keyspace@{db}__:{key}`` — data is the event type. 

925 - KEYEVENT: ``__keyevent@{db}__:{event}`` — data is the key name. 

926 - SUBKEYSPACE: ``__subkeyspace@{db}__:{key}`` — data is event + subkeys. 

927 - SUBKEYEVENT: ``__subkeyevent@{db}__:{event}`` — data is key + subkeys. 

928 - SUBKEYSPACEITEM: ``__subkeyspaceitem@{db}__:{key}\\n{subkey}`` — data 

929 is the event type. 

930 - SUBKEYSPACEEVENT: ``__subkeyspaceevent@{db}__:{event}|{key}`` — data 

931 is a subkey list. 

932 

933 Examples: 

934 >>> get_channel_type("__keyspace@0__:mykey") 

935 ChannelType.KEYSPACE 

936 >>> get_channel_type("__subkeyspace@0__:myhash") 

937 ChannelType.SUBKEYSPACE 

938 """ 

939 

940 KEYSPACE = "keyspace" 

941 KEYEVENT = "keyevent" 

942 SUBKEYSPACE = "subkeyspace" 

943 SUBKEYEVENT = "subkeyevent" 

944 SUBKEYSPACEITEM = "subkeyspaceitem" 

945 SUBKEYSPACEEVENT = "subkeyspaceevent" 

946 

947 

948def get_channel_type(channel: str | bytes) -> ChannelType | None: 

949 """ 

950 Determine the type of a Redis keyspace notification channel. 

951 

952 Args: 

953 channel: The channel name to check (string or bytes). 

954 

955 Returns: 

956 ChannelType.KEYSPACE if it's a keyspace notification channel, 

957 ChannelType.KEYEVENT if it's a keyevent notification channel, 

958 None if it's not a keyspace notification channel. 

959 

960 Examples: 

961 >>> get_channel_type("__keyspace@0__:mykey") 

962 ChannelType.KEYSPACE 

963 >>> get_channel_type("__keyevent@0__:set") 

964 ChannelType.KEYEVENT 

965 >>> get_channel_type("regular_channel") is None 

966 True 

967 >>> get_channel_type(b"__keyspace@0__:mykey") 

968 ChannelType.KEYSPACE 

969 """ 

970 channel_str = safe_str(channel) 

971 # Check subkey prefixes first (they are longer and more specific) 

972 if channel_str.startswith(SubkeyspaceitemChannel.PREFIX): 

973 return ChannelType.SUBKEYSPACEITEM 

974 if channel_str.startswith(SubkeyspaceeventChannel.PREFIX): 

975 return ChannelType.SUBKEYSPACEEVENT 

976 if channel_str.startswith(SubkeyspaceChannel.PREFIX): 

977 return ChannelType.SUBKEYSPACE 

978 if channel_str.startswith(SubkeyeventChannel.PREFIX): 

979 return ChannelType.SUBKEYEVENT 

980 if channel_str.startswith(KeyspaceChannel.PREFIX): 

981 return ChannelType.KEYSPACE 

982 if channel_str.startswith(KeyeventChannel.PREFIX): 

983 return ChannelType.KEYEVENT 

984 return None 

985 

986 

987def _is_pattern( 

988 channel: str | bytes | KeyspaceChannel | KeyeventChannel, 

989) -> bool: 

990 """ 

991 Check if a channel string contains glob-style pattern characters. 

992 

993 Redis uses glob-style patterns for psubscribe: 

994 - * matches any sequence of characters 

995 - ? matches any single character 

996 - [...] matches any character in the brackets 

997 

998 Args: 

999 channel: The channel string to check. Can be a string, bytes, 

1000 or a KeyspaceChannel/KeyeventChannel object. 

1001 

1002 Returns: 

1003 True if the channel contains pattern characters, False otherwise. 

1004 """ 

1005 # Handle Channel objects that have _channel_str attribute 

1006 # (KeyspaceChannel, KeyeventChannel) 

1007 if hasattr(channel, "_channel_str"): 

1008 channel = channel._channel_str 

1009 channel = safe_str(channel) 

1010 # Check for unescaped glob pattern characters. 

1011 # * and ? are always pattern characters. 

1012 # [ is only a pattern character when followed by a matching unescaped ], 

1013 # forming a bracket expression like [abc] or [a-z]. A lone [ (e.g. in 

1014 # a key named "my[key") is treated as a literal by Redis. 

1015 i = 0 

1016 while i < len(channel): 

1017 char = channel[i] 

1018 if char == "\\": 

1019 # Skip escaped character 

1020 i += 2 

1021 continue 

1022 if char in ("*", "?"): 

1023 return True 

1024 if char == "[": 

1025 # Look for a matching unescaped ] 

1026 j = i + 1 

1027 while j < len(channel): 

1028 if channel[j] == "\\": 

1029 j += 2 

1030 continue 

1031 if channel[j] == "]": 

1032 return True 

1033 j += 1 

1034 # No matching ] found — literal [ 

1035 i += 1 

1036 return False 

1037 

1038 

1039# ============================================================================= 

1040# Abstract Base Class for Keyspace Notifications 

1041# ============================================================================= 

1042 

1043 

1044class KeyspaceNotificationsInterface(ABC): 

1045 """ 

1046 Interface for keyspace notification managers. 

1047 

1048 This interface provides a consistent API for both standalone (KeyspaceNotifications) 

1049 and cluster (ClusterKeyspaceNotifications) implementations, allowing the same 

1050 code patterns to work with both standalone and cluster Redis deployments. 

1051 """ 

1052 

1053 @abstractmethod 

1054 def subscribe( 

1055 self, 

1056 *channels: ChannelT, 

1057 handler: SyncHandlerT | None = None, 

1058 ): 

1059 """Subscribe to keyspace notification channels.""" 

1060 pass 

1061 

1062 @abstractmethod 

1063 def unsubscribe(self, *channels: ChannelT): 

1064 """Unsubscribe from keyspace notification channels.""" 

1065 pass 

1066 

1067 @abstractmethod 

1068 def subscribe_keyspace( 

1069 self, 

1070 key_or_pattern: str, 

1071 db: int = 0, 

1072 handler: SyncHandlerT | None = None, 

1073 ): 

1074 """Subscribe to keyspace notifications for specific keys.""" 

1075 pass 

1076 

1077 @abstractmethod 

1078 def subscribe_keyevent( 

1079 self, 

1080 event: str, 

1081 db: int = 0, 

1082 handler: SyncHandlerT | None = None, 

1083 ): 

1084 """Subscribe to keyevent notifications for specific event types.""" 

1085 pass 

1086 

1087 @abstractmethod 

1088 def subscribe_subkeyspace( 

1089 self, 

1090 key_or_pattern: str, 

1091 db: int = 0, 

1092 handler: SyncHandlerT | None = None, 

1093 ): 

1094 """Subscribe to subkeyspace notifications for specific keys.""" 

1095 pass 

1096 

1097 @abstractmethod 

1098 def subscribe_subkeyevent( 

1099 self, 

1100 event: str, 

1101 db: int = 0, 

1102 handler: SyncHandlerT | None = None, 

1103 ): 

1104 """Subscribe to subkeyevent notifications for specific event types.""" 

1105 pass 

1106 

1107 @abstractmethod 

1108 def subscribe_subkeyspaceitem( 

1109 self, 

1110 key_or_pattern: str, 

1111 subkey_or_pattern: str, 

1112 db: int = 0, 

1113 handler: SyncHandlerT | None = None, 

1114 ): 

1115 """Subscribe to subkeyspaceitem notifications for a specific subkey.""" 

1116 pass 

1117 

1118 @abstractmethod 

1119 def subscribe_subkeyspaceevent( 

1120 self, 

1121 event: str, 

1122 key_or_pattern: str, 

1123 db: int = 0, 

1124 handler: SyncHandlerT | None = None, 

1125 ): 

1126 """Subscribe to subkeyspaceevent notifications for an event on a key.""" 

1127 pass 

1128 

1129 @abstractmethod 

1130 def get_message( 

1131 self, 

1132 ignore_subscribe_messages: bool | None = None, 

1133 timeout: float = 0.0, 

1134 ) -> KeyNotification | None: 

1135 """Get the next keyspace notification if one is available.""" 

1136 pass 

1137 

1138 @abstractmethod 

1139 def listen(self): 

1140 """Listen for keyspace notifications.""" 

1141 pass 

1142 

1143 @abstractmethod 

1144 def close(self): 

1145 """Close the notification manager and clean up resources.""" 

1146 pass 

1147 

1148 @abstractmethod 

1149 def __enter__(self): 

1150 pass 

1151 

1152 @abstractmethod 

1153 def __exit__(self, _exc_type, _exc_val, _exc_tb): 

1154 pass 

1155 

1156 @property 

1157 @abstractmethod 

1158 def subscribed(self) -> bool: 

1159 """Check if there are any active subscriptions and not closed.""" 

1160 pass 

1161 

1162 @abstractmethod 

1163 def run_in_thread( 

1164 self, 

1165 poll_timeout: float = 0.0, 

1166 daemon: bool = False, 

1167 exception_handler: Callable[ 

1168 [ 

1169 BaseException, 

1170 KeyspaceNotificationsInterface, 

1171 KeyspaceWorkerThread, 

1172 ], 

1173 None, 

1174 ] 

1175 | None = None, 

1176 ) -> KeyspaceWorkerThread: 

1177 """Start a background thread that polls for notifications.""" 

1178 pass 

1179 

1180 

1181class AbstractKeyspaceNotifications(KeyspaceNotificationsInterface): 

1182 """ 

1183 Abstract base class for keyspace notification managers. 

1184 

1185 Provides shared implementation for subscribe/unsubscribe logic. 

1186 Subclasses must implement: 

1187 - _execute_subscribe: Execute the subscribe operation 

1188 - _execute_unsubscribe: Execute the unsubscribe operation 

1189 - get_message: Get the next notification 

1190 - listen: Generator for notifications 

1191 - close: Clean up resources 

1192 """ 

1193 

1194 def __init__( 

1195 self, 

1196 key_prefix: str | bytes | None = None, 

1197 ignore_subscribe_messages: bool = True, 

1198 ): 

1199 """ 

1200 Initialize the base keyspace notification manager. 

1201 

1202 Args: 

1203 key_prefix: Optional prefix to filter and strip from keys in notifications 

1204 ignore_subscribe_messages: If True, subscribe/unsubscribe confirmations 

1205 are not returned by get_message/listen 

1206 """ 

1207 self.key_prefix = key_prefix 

1208 self.ignore_subscribe_messages = ignore_subscribe_messages 

1209 self._closed = False 

1210 

1211 def subscribe( 

1212 self, 

1213 *channels: ChannelT, 

1214 handler: SyncHandlerT | None = None, 

1215 ): 

1216 """ 

1217 Subscribe to keyspace notification channels. 

1218 

1219 Automatically detects whether each channel is a pattern (contains 

1220 wildcards like *, ?, [) or an exact channel name and uses the 

1221 appropriate Redis subscribe command internally. 

1222 

1223 Args: 

1224 *channels: Channels to subscribe to. Can be strings, KeyspaceChannel, 

1225 or KeyeventChannel objects. Patterns are auto-detected. 

1226 handler: Optional callback function that receives KeyNotification 

1227 objects. If provided, notifications are passed to the handler 

1228 instead of being returned by get_message()/listen(). 

1229 """ 

1230 # Wrap the handler to convert raw messages to KeyNotification objects 

1231 wrapped_handler: Callable | None = None 

1232 if handler is not None: 

1233 # Capture key_prefix in closure for consistent filtering/stripping 

1234 key_prefix = self.key_prefix 

1235 

1236 def _wrap_handler(message): 

1237 notification = KeyNotification.from_message( 

1238 message, key_prefix=key_prefix 

1239 ) 

1240 if notification is not None: 

1241 handler(notification) 

1242 

1243 wrapped_handler = _wrap_handler 

1244 

1245 patterns = {} 

1246 exact_channels = {} 

1247 

1248 for channel in channels: 

1249 if hasattr(channel, "_channel_str"): 

1250 channel_str = str(channel) 

1251 else: 

1252 channel_str = safe_str(channel) 

1253 if _is_pattern(channel): 

1254 patterns[channel_str] = wrapped_handler 

1255 else: 

1256 exact_channels[channel_str] = wrapped_handler 

1257 

1258 # Delegate to subclass implementation first. For standalone Redis 

1259 # this raises on failure, keeping tracking state clean. For cluster 

1260 # implementations the operation is best-effort (partial failures are 

1261 # logged, not raised) so tracking state is always updated afterwards. 

1262 self._execute_subscribe(patterns, exact_channels) 

1263 self._track_subscribe(patterns, exact_channels) 

1264 

1265 @abstractmethod 

1266 def _execute_subscribe( 

1267 self, patterns: dict[str, Any], exact_channels: dict[str, Any] 

1268 ) -> None: 

1269 """ 

1270 Execute the subscribe operation. 

1271 

1272 Args: 

1273 patterns: Dict mapping pattern strings to handlers (for psubscribe) 

1274 exact_channels: Dict mapping channel strings to handlers (for subscribe) 

1275 """ 

1276 pass 

1277 

1278 def unsubscribe(self, *channels: ChannelT): 

1279 """ 

1280 Unsubscribe from keyspace notification channels. 

1281 

1282 Automatically detects whether each channel is a pattern or exact 

1283 channel and uses the appropriate Redis unsubscribe command. 

1284 

1285 Args: 

1286 *channels: Channels to unsubscribe from. 

1287 """ 

1288 patterns = [] 

1289 exact_channels = [] 

1290 

1291 for channel in channels: 

1292 if hasattr(channel, "_channel_str"): 

1293 channel_str = str(channel) 

1294 else: 

1295 channel_str = safe_str(channel) 

1296 if _is_pattern(channel): 

1297 patterns.append(channel_str) 

1298 else: 

1299 exact_channels.append(channel_str) 

1300 

1301 # Delegate to subclass implementation first. For standalone Redis 

1302 # this raises on failure, keeping tracking state intact. For cluster 

1303 # implementations the operation is best-effort (partial failures are 

1304 # logged, not raised) so tracking state is always removed afterwards 

1305 # — this is intentional: the user asked to unsubscribe, so 

1306 # refresh_subscriptions should not re-subscribe these channels. 

1307 self._execute_unsubscribe(patterns, exact_channels) 

1308 self._untrack_subscribe(patterns, exact_channels) 

1309 

1310 @abstractmethod 

1311 def _execute_unsubscribe( 

1312 self, patterns: list[str], exact_channels: list[str] 

1313 ) -> None: 

1314 """ 

1315 Execute the unsubscribe operation. 

1316 

1317 Args: 

1318 patterns: List of pattern strings to punsubscribe from 

1319 exact_channels: List of channel strings to unsubscribe from 

1320 """ 

1321 pass 

1322 

1323 def _track_subscribe( 

1324 self, patterns: dict[str, Any], exact_channels: dict[str, Any] 

1325 ) -> None: 

1326 """Track newly subscribed patterns/channels. 

1327 

1328 Override in subclasses that need to maintain their own subscription 

1329 registry (e.g. cluster implementations that must re-subscribe 

1330 new/failed-over nodes). The default is a no-op because standalone 

1331 implementations delegate tracking to the underlying PubSub object. 

1332 """ 

1333 

1334 def _untrack_subscribe( 

1335 self, patterns: list[str], exact_channels: list[str] 

1336 ) -> None: 

1337 """Remove patterns/channels from the subscription registry. 

1338 

1339 Override in subclasses that maintain their own subscription registry. 

1340 The default is a no-op. 

1341 """ 

1342 

1343 def subscribe_keyspace( 

1344 self, 

1345 key_or_pattern: str, 

1346 db: int = 0, 

1347 handler: SyncHandlerT | None = None, 

1348 ): 

1349 """ 

1350 Subscribe to keyspace notifications for specific keys. 

1351 

1352 Args: 

1353 key_or_pattern: The key or pattern to monitor. Use '*' for wildcards. 

1354 db: The database number (default 0). 

1355 handler: Optional callback for notifications. 

1356 

1357 Example: 

1358 >>> ksn.subscribe_keyspace("user:123", db=0) 

1359 >>> ksn.subscribe_keyspace("user:*", db=0) 

1360 """ 

1361 channel = KeyspaceChannel(key_or_pattern, db=db) 

1362 self.subscribe(channel, handler=handler) 

1363 

1364 def subscribe_keyevent( 

1365 self, 

1366 event: str, 

1367 db: int = 0, 

1368 handler: SyncHandlerT | None = None, 

1369 ): 

1370 """ 

1371 Subscribe to keyevent notifications for specific event types. 

1372 

1373 Args: 

1374 event: The event type to monitor (e.g., EventType.SET or "set") 

1375 db: The database number (default 0). 

1376 handler: Optional callback for notifications. 

1377 

1378 Example: 

1379 >>> ksn.subscribe_keyevent(EventType.SET) 

1380 >>> ksn.subscribe_keyevent(EventType.EXPIRED, handler=my_handler) 

1381 """ 

1382 channel = KeyeventChannel(event, db=db) 

1383 self.subscribe(channel, handler=handler) 

1384 

1385 def subscribe_subkeyspace( 

1386 self, 

1387 key_or_pattern: str, 

1388 db: int = 0, 

1389 handler: SyncHandlerT | None = None, 

1390 ): 

1391 """ 

1392 Subscribe to subkeyspace notifications for specific keys. 

1393 

1394 Receives events with affected subkeys (fields) for the given key. 

1395 

1396 Args: 

1397 key_or_pattern: The key or pattern to monitor. 

1398 db: The database number (default 0). 

1399 handler: Optional callback for notifications. 

1400 """ 

1401 channel = SubkeyspaceChannel(key_or_pattern, db=db) 

1402 self.subscribe(channel, handler=handler) 

1403 

1404 def subscribe_subkeyevent( 

1405 self, 

1406 event: str, 

1407 db: int = 0, 

1408 handler: SyncHandlerT | None = None, 

1409 ): 

1410 """ 

1411 Subscribe to subkeyevent notifications for specific event types. 

1412 

1413 Receives the affected key and subkeys when the given event occurs. 

1414 

1415 Args: 

1416 event: The event type to monitor (e.g., "hset", "hdel"). 

1417 db: The database number (default 0). 

1418 handler: Optional callback for notifications. 

1419 """ 

1420 channel = SubkeyeventChannel(event, db=db) 

1421 self.subscribe(channel, handler=handler) 

1422 

1423 def subscribe_subkeyspaceitem( 

1424 self, 

1425 key_or_pattern: str, 

1426 subkey_or_pattern: str, 

1427 db: int = 0, 

1428 handler: SyncHandlerT | None = None, 

1429 ): 

1430 """ 

1431 Subscribe to subkeyspaceitem notifications for a specific subkey. 

1432 

1433 Receives the event type when the given subkey of the given key is 

1434 modified. 

1435 

1436 Args: 

1437 key_or_pattern: The key or pattern to monitor. 

1438 subkey_or_pattern: The subkey (field) or pattern to monitor. 

1439 db: The database number (default 0). 

1440 handler: Optional callback for notifications. 

1441 """ 

1442 channel = SubkeyspaceitemChannel(key_or_pattern, subkey_or_pattern, db=db) 

1443 self.subscribe(channel, handler=handler) 

1444 

1445 def subscribe_subkeyspaceevent( 

1446 self, 

1447 event: str, 

1448 key_or_pattern: str, 

1449 db: int = 0, 

1450 handler: SyncHandlerT | None = None, 

1451 ): 

1452 """ 

1453 Subscribe to subkeyspaceevent notifications for an event on a key. 

1454 

1455 Receives the affected subkeys when the given event occurs on the 

1456 given key. 

1457 

1458 Args: 

1459 event: The event type to monitor. 

1460 key_or_pattern: The key or pattern to monitor. 

1461 db: The database number (default 0). 

1462 handler: Optional callback for notifications. 

1463 """ 

1464 channel = SubkeyspaceeventChannel(event, key_or_pattern, db=db) 

1465 self.subscribe(channel, handler=handler) 

1466 

1467 def __enter__(self): 

1468 return self 

1469 

1470 def __exit__(self, _exc_type, _exc_val, _exc_tb): 

1471 self.close() 

1472 return False 

1473 

1474 def run_in_thread( 

1475 self, 

1476 poll_timeout: float = 0.0, 

1477 daemon: bool = False, 

1478 exception_handler: Callable[ 

1479 [ 

1480 BaseException, 

1481 KeyspaceNotificationsInterface, 

1482 KeyspaceWorkerThread, 

1483 ], 

1484 None, 

1485 ] 

1486 | None = None, 

1487 ) -> KeyspaceWorkerThread: 

1488 """ 

1489 Start a background thread that polls for notifications and triggers handlers. 

1490 

1491 This method spawns a thread that continuously calls get_message() to 

1492 process incoming notifications. When a notification arrives, any 

1493 registered handler for that channel/pattern is invoked automatically. 

1494 

1495 All subscriptions must have handlers registered before calling this method. 

1496 

1497 Args: 

1498 poll_timeout: Timeout in seconds for get_message() calls. When no message 

1499 is available, the thread waits up to this long before checking 

1500 again. Default 0.0 (non-blocking). WARNING: the default 

1501 causes a CPU spin-loop. It is preferred to pass a positive 

1502 value (e.g. 0.1 or 1.0). 

1503 daemon: If True, the thread will be a daemon thread and will be 

1504 terminated when the main program exits. Default False. 

1505 exception_handler: Optional callback invoked when an exception occurs 

1506 in the worker thread. Receives (exception, notifications, 

1507 thread) as arguments. If None, exceptions are raised. 

1508 

1509 Returns: 

1510 KeyspaceWorkerThread: The started worker thread. Call stop() on it 

1511 to stop the thread and close the notifications. 

1512 

1513 Raises: 

1514 RedisError: If any subscription doesn't have a handler registered. 

1515 

1516 Example: 

1517 >>> def my_handler(notification): 

1518 ... print(f"Got: {notification.key} - {notification.event_type}") 

1519 >>> 

1520 >>> notifications.subscribe(KeyspaceChannel("user:*"), handler=my_handler) 

1521 >>> thread = notifications.run_in_thread(poll_timeout=0.1, daemon=True) 

1522 >>> # ... handlers are called automatically ... 

1523 >>> thread.stop() 

1524 """ 

1525 self._validate_all_handlers() 

1526 

1527 thread = KeyspaceWorkerThread( 

1528 self, 

1529 poll_timeout, 

1530 daemon=daemon, 

1531 exception_handler=exception_handler, 

1532 ) 

1533 thread.start() 

1534 return thread 

1535 

1536 @abstractmethod 

1537 def _validate_all_handlers(self) -> None: 

1538 """Raise :class:`~redis.RedisError` if any subscription lacks a handler. 

1539 

1540 Subclasses inspect their own subscription state to perform 

1541 the validation. 

1542 """ 

1543 pass 

1544 

1545 

1546class KeyspaceWorkerThread(threading.Thread): 

1547 """ 

1548 Background thread for processing keyspace notifications. 

1549 

1550 This thread continuously polls for notifications and invokes registered 

1551 handlers. It works with both KeyspaceNotifications (standalone) and 

1552 ClusterKeyspaceNotifications. 

1553 

1554 Example: 

1555 >>> thread = notifications.run_in_thread(poll_timeout=0.1) 

1556 >>> # ... handlers are called automatically ... 

1557 >>> thread.stop() 

1558 """ 

1559 

1560 def __init__( 

1561 self, 

1562 notifications: KeyspaceNotificationsInterface, 

1563 poll_timeout: float, 

1564 daemon: bool = False, 

1565 exception_handler: Callable[ 

1566 [ 

1567 BaseException, 

1568 KeyspaceNotificationsInterface, 

1569 KeyspaceWorkerThread, 

1570 ], 

1571 None, 

1572 ] 

1573 | None = None, 

1574 ): 

1575 super().__init__() 

1576 self.daemon = daemon 

1577 self.notifications = notifications 

1578 self.poll_timeout = poll_timeout 

1579 self.exception_handler = exception_handler 

1580 self._running = threading.Event() 

1581 

1582 def run(self) -> None: 

1583 """Main loop that polls for notifications and triggers handlers.""" 

1584 if self._running.is_set(): 

1585 return 

1586 self._running.set() 

1587 notifications = self.notifications 

1588 poll_timeout = self.poll_timeout 

1589 while self._running.is_set(): 

1590 try: 

1591 notifications.get_message( 

1592 ignore_subscribe_messages=True, timeout=poll_timeout 

1593 ) 

1594 except BaseException as e: 

1595 if self.exception_handler is None: 

1596 raise 

1597 self.exception_handler(e, notifications, self) 

1598 notifications.close() 

1599 

1600 def stop(self) -> None: 

1601 """ 

1602 Stop the worker thread. 

1603 

1604 This signals the thread to exit its run loop. The thread will close 

1605 the notifications object before terminating. 

1606 """ 

1607 self._running.clear() 

1608 

1609 

1610# ============================================================================= 

1611# Standalone Keyspace Notification Manager 

1612# ============================================================================= 

1613 

1614 

1615class KeyspaceNotifications(AbstractKeyspaceNotifications): 

1616 """ 

1617 Manages keyspace notification subscriptions for standalone Redis. 

1618 

1619 For standalone Redis, keyspace notifications work with a single PubSub 

1620 connection. This class wraps that connection and provides: 

1621 - Automatic pattern vs exact channel detection 

1622 - KeyNotification parsing with optional key_prefix filtering 

1623 - Convenience methods for keyspace and keyevent subscriptions 

1624 - Context manager and run_in_thread support 

1625 """ 

1626 

1627 def __init__( 

1628 self, 

1629 redis_client: Redis, 

1630 key_prefix: str | bytes | None = None, 

1631 ignore_subscribe_messages: bool = True, 

1632 ): 

1633 """ 

1634 Initialize the standalone keyspace notification manager. 

1635 

1636 Note: Keyspace notifications must be enabled on the Redis server via 

1637 the ``notify-keyspace-events`` configuration option. This is a server-side 

1638 configuration that should be done by your infrastructure/operations team. 

1639 

1640 Args: 

1641 redis_client: A Redis client instance 

1642 key_prefix: Optional prefix to filter and strip from keys in notifications 

1643 ignore_subscribe_messages: If True, subscribe/unsubscribe confirmations 

1644 are not returned by get_message/listen 

1645 """ 

1646 super().__init__(key_prefix, ignore_subscribe_messages) 

1647 self.redis = redis_client 

1648 

1649 # Create the PubSub instance with ignore_subscribe_messages=False 

1650 # so that the per-call argument in get_message() can control behavior 

1651 self._pubsub = redis_client.pubsub(ignore_subscribe_messages=False) 

1652 

1653 def _execute_subscribe( 

1654 self, patterns: dict[str, Any], exact_channels: dict[str, Any] 

1655 ) -> None: 

1656 """Execute subscribe on the single pubsub connection.""" 

1657 if patterns: 

1658 self._pubsub.psubscribe(**patterns) 

1659 if exact_channels: 

1660 self._pubsub.subscribe(**exact_channels) 

1661 

1662 def _execute_unsubscribe( 

1663 self, patterns: list[str], exact_channels: list[str] 

1664 ) -> None: 

1665 """Execute unsubscribe on the single pubsub connection.""" 

1666 if patterns: 

1667 self._pubsub.punsubscribe(*patterns) 

1668 if exact_channels: 

1669 self._pubsub.unsubscribe(*exact_channels) 

1670 

1671 def get_message( 

1672 self, 

1673 ignore_subscribe_messages: bool | None = None, 

1674 timeout: float = 0.0, 

1675 ) -> KeyNotification | None: 

1676 """ 

1677 Get the next keyspace notification if one is available. 

1678 

1679 Note: If a handler was registered for the channel, pubsub will call 

1680 the handler directly and this method returns None for that message. 

1681 

1682 Args: 

1683 ignore_subscribe_messages: If True, skip subscribe/unsubscribe messages. 

1684 Defaults to the value set in __init__ (True). 

1685 timeout: Time to wait for a message. 

1686 

1687 Returns: 

1688 A KeyNotification if a notification is available and no handler 

1689 was registered for the channel, None otherwise. 

1690 """ 

1691 if ignore_subscribe_messages is None: 

1692 ignore_subscribe_messages = self.ignore_subscribe_messages 

1693 

1694 if self._closed: 

1695 return None 

1696 

1697 # Pubsub's get_message will call wrapped handlers directly for channels 

1698 # with registered handlers and return None. For channels without handlers, 

1699 # it returns the raw message which we parse to KeyNotification. 

1700 message = self._pubsub.get_message( 

1701 ignore_subscribe_messages=ignore_subscribe_messages, 

1702 timeout=timeout, 

1703 ) 

1704 

1705 if message is not None: 

1706 return KeyNotification.from_message(message, key_prefix=self.key_prefix) 

1707 

1708 return None 

1709 

1710 def listen(self): 

1711 """ 

1712 Listen for keyspace notifications. 

1713 

1714 This is a generator that yields KeyNotification objects as they arrive. 

1715 It blocks until a notification is received. 

1716 

1717 Yields: 

1718 KeyNotification objects for each keyspace/keyevent notification. 

1719 

1720 Example: 

1721 >>> for notification in ksn.listen(): 

1722 ... print(f"{notification.key}: {notification.event_type}") 

1723 """ 

1724 while self.subscribed: 

1725 notification = self.get_message(timeout=1.0) 

1726 if notification is not None: 

1727 yield notification 

1728 

1729 @property 

1730 def subscribed(self) -> bool: 

1731 """Check if there are any active subscriptions and not closed.""" 

1732 return not self._closed and self._pubsub.subscribed 

1733 

1734 def _validate_all_handlers(self) -> None: 

1735 """Raise if any subscription in the underlying PubSub lacks a handler.""" 

1736 for channel, handler in self._pubsub.channels.items(): 

1737 if handler is None: 

1738 raise RedisError(f"Channel '{channel}' has no handler registered") 

1739 for pattern, handler in self._pubsub.patterns.items(): 

1740 if handler is None: 

1741 raise RedisError(f"Pattern '{pattern}' has no handler registered") 

1742 

1743 def close(self): 

1744 """Close the pubsub connection and clean up resources.""" 

1745 self._closed = True 

1746 try: 

1747 self._pubsub.close() 

1748 except Exception: 

1749 pass 

1750 

1751 

1752# ============================================================================= 

1753# Cluster-Aware Keyspace Notification Manager 

1754# ============================================================================= 

1755 

1756 

1757class ClusterKeyspaceNotifications(AbstractKeyspaceNotifications): 

1758 """ 

1759 Manages keyspace notification subscriptions across all nodes in a Redis Cluster. 

1760 

1761 In Redis Cluster, keyspace notifications are NOT broadcast between nodes. 

1762 Each node only emits notifications for keys it owns. This class automatically 

1763 subscribes to all primary nodes in the cluster and handles topology changes. 

1764 """ 

1765 

1766 def __init__( 

1767 self, 

1768 redis_cluster: RedisCluster, 

1769 key_prefix: str | bytes | None = None, 

1770 ignore_subscribe_messages: bool = True, 

1771 ): 

1772 """ 

1773 Initialize the cluster keyspace notification manager. 

1774 

1775 Note: Keyspace notifications must be enabled on all Redis cluster nodes via 

1776 the ``notify-keyspace-events`` configuration option. This is a server-side 

1777 configuration that should be done by your infrastructure/operations team. 

1778 

1779 Args: 

1780 redis_cluster: A RedisCluster instance 

1781 key_prefix: Optional prefix to filter and strip from keys in notifications 

1782 ignore_subscribe_messages: If True, subscribe/unsubscribe confirmations 

1783 are not returned by get_message/listen 

1784 """ 

1785 super().__init__(key_prefix, ignore_subscribe_messages) 

1786 self.cluster = redis_cluster 

1787 

1788 # Canonical subscription registry: pattern/channel -> wrapped handler. 

1789 # In cluster mode there are multiple PubSub objects (one per node), so 

1790 # this is the single source of truth used to (re-)subscribe new or 

1791 # failed-over nodes. 

1792 self._subscribed_patterns: dict[str, Any] = {} 

1793 self._subscribed_channels: dict[str, Any] = {} 

1794 

1795 # Track subscriptions per node 

1796 self._node_pubsubs: dict[str, Any] = {} 

1797 

1798 # Lock for topology refresh operations 

1799 self._refresh_lock = threading.Lock() 

1800 

1801 # Current pubsub index for round-robin polling 

1802 self._poll_index = 0 

1803 

1804 @property 

1805 def subscribed(self) -> bool: 

1806 """Check if there are any active subscriptions and not closed.""" 

1807 return not self._closed and bool( 

1808 self._subscribed_patterns or self._subscribed_channels 

1809 ) 

1810 

1811 def _track_subscribe( 

1812 self, patterns: dict[str, Any], exact_channels: dict[str, Any] 

1813 ) -> None: 

1814 """Track newly subscribed patterns/channels in the cluster registry.""" 

1815 if patterns: 

1816 self._subscribed_patterns.update(patterns) 

1817 if exact_channels: 

1818 self._subscribed_channels.update(exact_channels) 

1819 

1820 def _untrack_subscribe( 

1821 self, patterns: list[str], exact_channels: list[str] 

1822 ) -> None: 

1823 """Remove patterns/channels from the cluster registry.""" 

1824 for p in patterns: 

1825 self._subscribed_patterns.pop(p, None) 

1826 for c in exact_channels: 

1827 self._subscribed_channels.pop(c, None) 

1828 

1829 def _validate_all_handlers(self) -> None: 

1830 """Raise if any subscription in the cluster registry lacks a handler.""" 

1831 for channel, handler in self._subscribed_channels.items(): 

1832 if handler is None: 

1833 raise RedisError(f"Channel '{channel}' has no handler registered") 

1834 for pattern, handler in self._subscribed_patterns.items(): 

1835 if handler is None: 

1836 raise RedisError(f"Pattern '{pattern}' has no handler registered") 

1837 

1838 def _get_all_primary_nodes(self): 

1839 """Get all primary nodes in the cluster.""" 

1840 return self.cluster.get_primaries() 

1841 

1842 def _cleanup_node(self, node_name: str) -> None: 

1843 """Remove and close a node's PubSub. 

1844 

1845 Closing the ``PubSub`` disconnects its connection so it is not 

1846 left in a subscribed state inside the connection pool. 

1847 """ 

1848 pubsub = self._node_pubsubs.pop(node_name, None) 

1849 if pubsub: 

1850 try: 

1851 pubsub.close() 

1852 except Exception: 

1853 pass 

1854 

1855 def _ensure_node_pubsub(self, node) -> Any: 

1856 """Get or create a PubSub instance for a node.""" 

1857 if node.name not in self._node_pubsubs: 

1858 redis_conn = self.cluster.get_redis_connection(node) 

1859 # Always create PubSub with ignore_subscribe_messages=False 

1860 # so that the per-call argument in get_message() can control 

1861 # the behavior reliably 

1862 pubsub = redis_conn.pubsub(ignore_subscribe_messages=False) 

1863 self._node_pubsubs[node.name] = pubsub 

1864 return self._node_pubsubs[node.name] 

1865 

1866 def _execute_subscribe( 

1867 self, patterns: dict[str, Any], exact_channels: dict[str, Any] 

1868 ) -> None: 

1869 """Execute subscribe on all cluster nodes. 

1870 

1871 Patterns and exact channels are subscribed in a single pass over 

1872 nodes so that a mid-batch node failure cannot create a 

1873 partially-caught-up replacement. If a node fails during this 

1874 call it is removed from ``_node_pubsubs`` and will be fully 

1875 re-subscribed on the next ``refresh_subscriptions`` cycle. 

1876 

1877 If a newly discovered node is encountered (not yet in 

1878 ``_node_pubsubs``), it is also subscribed to all *previously* 

1879 tracked patterns/channels so it doesn't miss notifications for 

1880 subscriptions that were established before this node joined. 

1881 """ 

1882 if not patterns and not exact_channels: 

1883 return 

1884 

1885 failed_nodes: list[str] = [] 

1886 for node in self._get_all_primary_nodes(): 

1887 is_new_node = node.name not in self._node_pubsubs 

1888 pubsub = self._ensure_node_pubsub(node) 

1889 try: 

1890 # If this is a brand-new node, catch it up on existing 

1891 # subscriptions before adding the new channels. 

1892 if is_new_node: 

1893 if self._subscribed_patterns: 

1894 pubsub.psubscribe(**self._subscribed_patterns) 

1895 if self._subscribed_channels: 

1896 pubsub.subscribe(**self._subscribed_channels) 

1897 

1898 if patterns: 

1899 pubsub.psubscribe(**patterns) 

1900 if exact_channels: 

1901 pubsub.subscribe(**exact_channels) 

1902 except Exception: 

1903 # Remove the broken pubsub so refresh_subscriptions can 

1904 # re-create it later. 

1905 self._cleanup_node(node.name) 

1906 failed_nodes.append(node.name) 

1907 

1908 if failed_nodes: 

1909 logger.warning( 

1910 "Failed to subscribe on cluster nodes: %s. " 

1911 "These nodes will be retried on the next refresh cycle.", 

1912 ", ".join(failed_nodes), 

1913 ) 

1914 

1915 def _execute_unsubscribe( 

1916 self, patterns: list[str], exact_channels: list[str] 

1917 ) -> None: 

1918 """Execute unsubscribe on all cluster nodes.""" 

1919 if patterns: 

1920 self._unsubscribe_from_all_nodes(patterns, use_punsubscribe=True) 

1921 if exact_channels: 

1922 self._unsubscribe_from_all_nodes(exact_channels, use_punsubscribe=False) 

1923 

1924 def _unsubscribe_from_all_nodes(self, channels: list[str], use_punsubscribe: bool): 

1925 """Unsubscribe from patterns/channels on all nodes. 

1926 

1927 Best-effort: tries every node so that a single broken connection 

1928 does not prevent the remaining nodes from being unsubscribed. 

1929 Broken pubsubs are cleaned up; the tracking state is still removed 

1930 by the caller, so ``refresh_subscriptions`` will *not* re-subscribe 

1931 these channels on replacement nodes. 

1932 """ 

1933 failed_nodes: list[str] = [] 

1934 for node_name, pubsub in list(self._node_pubsubs.items()): 

1935 try: 

1936 if use_punsubscribe: 

1937 pubsub.punsubscribe(*channels) 

1938 else: 

1939 pubsub.unsubscribe(*channels) 

1940 except Exception: 

1941 self._cleanup_node(node_name) 

1942 failed_nodes.append(node_name) 

1943 

1944 if failed_nodes: 

1945 logger.warning( 

1946 "Failed to unsubscribe on cluster nodes: %s. " 

1947 "These nodes will be re-created on the next refresh cycle.", 

1948 ", ".join(failed_nodes), 

1949 ) 

1950 

1951 def get_message( 

1952 self, 

1953 ignore_subscribe_messages: bool | None = None, 

1954 timeout: float = 0.0, 

1955 ) -> KeyNotification | None: 

1956 """ 

1957 Get the next keyspace notification if one is available. 

1958 

1959 This method polls all node pubsubs in round-robin fashion until 

1960 a message is received or the timeout expires. 

1961 If a connection error occurs, subscriptions are automatically refreshed. 

1962 

1963 Args: 

1964 ignore_subscribe_messages: If True, skip subscribe/unsubscribe messages. 

1965 Defaults to the value set in __init__ (True). 

1966 timeout: Total time to wait for a message (distributed across all nodes) 

1967 

1968 Returns: 

1969 A KeyNotification if a notification is available, None otherwise. 

1970 """ 

1971 if self._closed: 

1972 return None 

1973 

1974 total_nodes = len(self._node_pubsubs) 

1975 if total_nodes == 0: 

1976 # Sleep for the requested timeout so callers that loop 

1977 # (run_in_thread, listen) don't spin the CPU when all node 

1978 # connections have been cleaned up. 

1979 if timeout > 0: 

1980 time.sleep(timeout) 

1981 return None 

1982 

1983 # Use instance default if not specified 

1984 if ignore_subscribe_messages is None: 

1985 ignore_subscribe_messages = self.ignore_subscribe_messages 

1986 

1987 # Handle timeout=0 as a single non-blocking poll over all pubsubs 

1988 # This matches the expected semantics of PubSub.get_message(timeout=0) 

1989 if timeout == 0.0: 

1990 return self._poll_all_nodes_once(ignore_subscribe_messages) 

1991 

1992 # Calculate per-node timeout for each poll 

1993 # Use a small timeout per node to allow round-robin polling 

1994 per_node_timeout = min(0.1, timeout / max(total_nodes, 1)) 

1995 

1996 start_time = time.monotonic() 

1997 end_time = start_time + timeout 

1998 

1999 while True: 

2000 # Check if we've exceeded the total timeout 

2001 if time.monotonic() >= end_time: 

2002 return None 

2003 

2004 pubsubs = list(self._node_pubsubs.values()) 

2005 if not pubsubs: 

2006 return None 

2007 

2008 # Round-robin polling 

2009 self._poll_index = self._poll_index % len(pubsubs) 

2010 pubsub = pubsubs[self._poll_index] 

2011 self._poll_index += 1 

2012 

2013 try: 

2014 message = pubsub.get_message( 

2015 ignore_subscribe_messages=ignore_subscribe_messages, 

2016 timeout=per_node_timeout, 

2017 ) 

2018 except (ConnectionError, TimeoutError, RedisError): 

2019 # Connection error - refresh subscriptions and continue 

2020 self._refresh_subscriptions_on_error() 

2021 continue 

2022 

2023 if message is not None: 

2024 # Note: If a handler was registered, PubSub already invoked it 

2025 # and returned None, so we only reach here for handler-less subscriptions 

2026 notification = KeyNotification.from_message( 

2027 message, key_prefix=self.key_prefix 

2028 ) 

2029 if notification is not None: 

2030 return notification 

2031 # If not a keyspace notification, continue checking other nodes 

2032 

2033 def _poll_all_nodes_once( 

2034 self, ignore_subscribe_messages: bool 

2035 ) -> KeyNotification | None: 

2036 """ 

2037 Perform a single non-blocking poll over all node pubsubs. 

2038 

2039 This is used when timeout=0 to match the expected semantics of 

2040 PubSub.get_message(timeout=0) - a non-blocking check for messages. 

2041 

2042 Returns: 

2043 A KeyNotification if one is available, None otherwise. 

2044 """ 

2045 had_error = False 

2046 for pubsub in list(self._node_pubsubs.values()): 

2047 try: 

2048 message = pubsub.get_message( 

2049 ignore_subscribe_messages=ignore_subscribe_messages, 

2050 timeout=0.0, 

2051 ) 

2052 except (ConnectionError, TimeoutError, RedisError): 

2053 # Record the error but continue polling remaining healthy 

2054 # nodes so that already-buffered notifications are not lost. 

2055 had_error = True 

2056 continue 

2057 

2058 if message is not None: 

2059 # Note: If a handler was registered, PubSub already invoked it 

2060 # and returned None, so we only reach here for handler-less subscriptions 

2061 notification = KeyNotification.from_message( 

2062 message, key_prefix=self.key_prefix 

2063 ) 

2064 if notification is not None: 

2065 # Refresh before returning if any node had an error, 

2066 # so the next poll cycle has fresh state. 

2067 if had_error: 

2068 self._refresh_subscriptions_on_error() 

2069 return notification 

2070 

2071 # Refresh after polling all nodes if any had errors 

2072 if had_error: 

2073 self._refresh_subscriptions_on_error() 

2074 return None 

2075 

2076 def listen(self): 

2077 """ 

2078 Listen for keyspace notifications from all cluster nodes. 

2079 

2080 This is a generator that yields KeyNotification objects as they arrive. 

2081 It blocks until a notification is received. 

2082 

2083 Yields: 

2084 KeyNotification objects for each keyspace/keyevent notification. 

2085 

2086 Example: 

2087 >>> for notification in ksn.listen(): 

2088 ... print(f"{notification.key}: {notification.event_type}") 

2089 """ 

2090 while self.subscribed: 

2091 notification = self.get_message(timeout=1.0) 

2092 if notification is not None: 

2093 yield notification 

2094 

2095 def _refresh_subscriptions_on_error(self): 

2096 """ 

2097 Refresh subscriptions after a connection error. 

2098 

2099 This is called automatically when a connection error occurs during 

2100 get_message(). It checks if nodes changed before refreshing. 

2101 """ 

2102 self._poll_index = 0 # Reset round-robin index 

2103 

2104 try: 

2105 self.refresh_subscriptions() 

2106 except Exception: 

2107 logger.warning( 

2108 "Failed to refresh cluster subscriptions, will retry on next error", 

2109 exc_info=True, 

2110 ) 

2111 

2112 def _is_pubsub_connected(self, pubsub) -> bool: 

2113 """Check if a pubsub connection is still alive.""" 

2114 try: 

2115 conn = pubsub.connection 

2116 if conn is None: 

2117 return False 

2118 return conn.is_connected 

2119 except Exception: 

2120 return False 

2121 

2122 def refresh_subscriptions(self): 

2123 """ 

2124 Refresh subscriptions after a topology change. 

2125 

2126 This method is called automatically when topology changes are detected 

2127 or when connection errors occur. You can also call it manually if needed. 

2128 

2129 This method: 

2130 1. Discovers any new primary nodes and subscribes them 

2131 2. Removes pubsubs for nodes that are no longer primaries 

2132 3. Re-creates broken pubsub connections for existing nodes 

2133 """ 

2134 with self._refresh_lock: 

2135 current_primaries = { 

2136 node.name: node for node in self._get_all_primary_nodes() 

2137 } 

2138 

2139 # Remove pubsubs for nodes that are no longer primaries 

2140 removed_nodes = set(self._node_pubsubs.keys()) - set( 

2141 current_primaries.keys() 

2142 ) 

2143 for node_name in removed_nodes: 

2144 self._cleanup_node(node_name) 

2145 

2146 # Detect broken connections for existing nodes and remove them 

2147 # so they get re-created below 

2148 existing_nodes = set(self._node_pubsubs.keys()) & set( 

2149 current_primaries.keys() 

2150 ) 

2151 for node_name in existing_nodes: 

2152 pubsub = self._node_pubsubs.get(node_name) 

2153 if pubsub and not self._is_pubsub_connected(pubsub): 

2154 # Connection is broken, remove it so it gets re-created 

2155 self._cleanup_node(node_name) 

2156 

2157 # Subscribe new nodes (and nodes with broken connections) to existing 

2158 # patterns/channels 

2159 new_nodes = set(current_primaries.keys()) - set(self._node_pubsubs.keys()) 

2160 failed_nodes: list[str] = [] 

2161 for node_name in new_nodes: 

2162 node = current_primaries[node_name] 

2163 pubsub = self._ensure_node_pubsub(node) 

2164 

2165 try: 

2166 if self._subscribed_patterns: 

2167 pubsub.psubscribe(**self._subscribed_patterns) 

2168 if self._subscribed_channels: 

2169 pubsub.subscribe(**self._subscribed_channels) 

2170 except Exception: 

2171 # Subscription failed - remove from dict so retry is possible 

2172 self._cleanup_node(node_name) 

2173 failed_nodes.append(node_name) 

2174 

2175 # Raise after attempting all nodes so we don't skip any 

2176 if failed_nodes: 

2177 raise ConnectionError( 

2178 f"Failed to subscribe to cluster nodes: {', '.join(failed_nodes)}" 

2179 ) 

2180 

2181 def close(self): 

2182 """Close all pubsub connections and clean up resources.""" 

2183 self._closed = True 

2184 for node_name in list(self._node_pubsubs.keys()): 

2185 self._cleanup_node(node_name) 

2186 self._subscribed_patterns.clear() 

2187 self._subscribed_channels.clear()