1from __future__ import annotations
2
3import asyncio
4import warnings
5from typing import (
6 TYPE_CHECKING,
7 Any,
8 AsyncIterator,
9 Awaitable,
10 Dict,
11 Iterable,
12 Iterator,
13 List,
14 Literal,
15 Mapping,
16 NoReturn,
17 Sequence,
18 overload,
19)
20
21from redis.crc import key_slot
22from redis.exceptions import RedisClusterException, RedisError
23from redis.typing import (
24 AnyKeyT,
25 AsyncClientProtocol,
26 ClusterCommandsProtocol,
27 ClusterLinksResponse,
28 ClusterNodeDetail,
29 ClusterShardsResponse,
30 EncodableT,
31 KeysT,
32 KeyT,
33 PatternT,
34 ResponseT,
35 StralgoResponse,
36 SyncClientProtocol,
37)
38from redis.utils import deprecated_function
39
40from .core import (
41 ACLCommands,
42 AsyncACLCommands,
43 AsyncBlessCommands,
44 AsyncDataAccessCommands,
45 AsyncFunctionCommands,
46 AsyncManagementCommands,
47 AsyncModuleCommands,
48 AsyncScriptCommands,
49 BlessCommands,
50 BlessFlag,
51 DataAccessCommands,
52 FunctionCommands,
53 HotkeysMetricsTypes,
54 ManagementCommands,
55 ModuleCommands,
56 PubSubCommands,
57 ScriptCommands,
58)
59from .helpers import list_or_args
60from .redismodules import AsyncRedisModuleCommands, RedisModuleCommands
61
62if TYPE_CHECKING:
63 from redis.asyncio.cluster import TargetNodesT
64 from redis.cluster import LoadBalancingStrategy
65
66# DEPRECATED - no longer consulted by the default metadata routing, and it will be removed in a future release.
67#
68# Replica safety is now decided from command metadata by the metadata resolver,
69# so this set no longer describes what commands are safe to execute on replicas.
70# It is kept as a public attribute only so an external caller reading it keeps working;
71# editing it changes nothing.
72#
73# To change replica safety, edit ``redis.commands.metadata._STATIC_COMMAND_METADATA`` or pass
74# a ``metadata_resolver`` to the client. Nothing here.
75READ_COMMANDS = frozenset(
76 [
77 # Bit Operations
78 "BITCOUNT",
79 "BITFIELD_RO",
80 "BITPOS",
81 # Scripting
82 "EVAL_RO",
83 "EVALSHA_RO",
84 "FCALL_RO",
85 # Key Operations
86 "DBSIZE",
87 "DIGEST",
88 "DUMP",
89 "EXISTS",
90 "EXPIRETIME",
91 "PEXPIRETIME",
92 "KEYS",
93 "SCAN",
94 "PTTL",
95 "RANDOMKEY",
96 "TTL",
97 "TYPE",
98 # String Operations
99 "GET",
100 "GETBIT",
101 "GETRANGE",
102 "MGET",
103 "STRLEN",
104 "LCS",
105 # Geo Operations
106 "GEODIST",
107 "GEOHASH",
108 "GEOPOS",
109 "GEOSEARCH",
110 # Hash Operations
111 "HEXISTS",
112 "HGET",
113 "HGETALL",
114 "HKEYS",
115 "HLEN",
116 "HMGET",
117 "HSTRLEN",
118 "HVALS",
119 "HRANDFIELD",
120 "HEXPIRETIME",
121 "HPEXPIRETIME",
122 "HTTL",
123 "HPTTL",
124 "HSCAN",
125 # List Operations
126 "LINDEX",
127 "LPOS",
128 "LLEN",
129 "LRANGE",
130 # Set Operations
131 "SCARD",
132 "SDIFF",
133 "SDIFFCARD",
134 "SINTER",
135 "SINTERCARD",
136 "SISMEMBER",
137 "SMISMEMBER",
138 "SMEMBERS",
139 "SRANDMEMBER",
140 "SUNION",
141 "SUNIONCARD",
142 "SSCAN",
143 # Sorted Set Operations
144 "ZCARD",
145 "ZCOUNT",
146 "ZDIFF",
147 "ZINTER",
148 "ZINTERCARD",
149 "ZLEXCOUNT",
150 "ZMSCORE",
151 "ZRANDMEMBER",
152 "ZRANGE",
153 "ZRANGEBYLEX",
154 "ZRANGEBYSCORE",
155 "ZRANK",
156 "ZREVRANGE",
157 "ZREVRANGEBYLEX",
158 "ZREVRANGEBYSCORE",
159 "ZREVRANK",
160 "ZSCAN",
161 "ZSCORE",
162 "ZUNION",
163 # Stream Operations
164 "XLEN",
165 "XPENDING",
166 "XRANGE",
167 "XREAD",
168 "XREVRANGE",
169 # JSON Module
170 "JSON.ARRINDEX",
171 "JSON.ARRLEN",
172 "JSON.GET",
173 "JSON.MGET",
174 "JSON.OBJKEYS",
175 "JSON.OBJLEN",
176 "JSON.RESP",
177 "JSON.STRLEN",
178 "JSON.TYPE",
179 # RediSearch Module
180 "FT.ALIASLIST",
181 "FT.EXPLAIN",
182 "FT.INFO",
183 "FT.PROFILE",
184 "FT.SEARCH",
185 ]
186)
187
188
189class ClusterMultiKeyCommands(ClusterCommandsProtocol):
190 """
191 A class containing commands that handle more than one key
192 """
193
194 # Read-routing configuration, which the cluster clients all set on the instance. The
195 # defaults are here for a host class that mixes these commands in without one - see
196 # ``_is_replica_safe`` below, which is the same fallback for the same audience - so
197 # that ``_execute_pipeline_by_slot`` reads a value rather than raising
198 # ``AttributeError``. Typed under TYPE_CHECKING because ``redis.cluster`` imports this
199 # module.
200 read_from_replicas: bool = False
201 load_balancing_strategy: LoadBalancingStrategy | None = None
202
203 def _partition_keys_by_slot(self, keys: Iterable[KeyT]) -> Dict[int, List[KeyT]]:
204 """Split keys into a dictionary that maps a slot to a list of keys."""
205
206 slots_to_keys = {}
207 for key in keys:
208 slot = key_slot(self.encoder.encode(key))
209 slots_to_keys.setdefault(slot, []).append(key)
210
211 return slots_to_keys
212
213 def _partition_pairs_by_slot(
214 self, mapping: Mapping[AnyKeyT, EncodableT]
215 ) -> Dict[int, List[EncodableT]]:
216 """Split pairs into a dictionary that maps a slot to a list of pairs."""
217
218 slots_to_pairs = {}
219 for pair in mapping.items():
220 slot = key_slot(self.encoder.encode(pair[0]))
221 slots_to_pairs.setdefault(slot, []).extend(pair)
222
223 return slots_to_pairs
224
225 def _is_replica_safe(self, command_name: str) -> bool:
226 warnings.warn(
227 "Using READ_COMMANDS for replica-safe checks is deprecated and will be removed in a future release. "
228 "Use a MetadataResolver instead.",
229 DeprecationWarning,
230 stacklevel=2,
231 )
232 return isinstance(command_name, str) and command_name.upper() in READ_COMMANDS
233
234 def _execute_pipeline_by_slot(
235 self, command: str, slots_to_args: Mapping[int, Iterable[EncodableT]]
236 ) -> List[Any]:
237 replica_safe = (
238 self.read_from_replicas or self.load_balancing_strategy is not None
239 ) and self._is_replica_safe(command)
240 pipe = self.pipeline()
241 [
242 pipe.execute_command(
243 command,
244 *slot_args,
245 target_nodes=[
246 self.nodes_manager.get_node_from_slot(
247 slot,
248 replica_safe,
249 self.load_balancing_strategy if replica_safe else None,
250 )
251 ],
252 )
253 for slot, slot_args in slots_to_args.items()
254 ]
255 return pipe.execute()
256
257 def _reorder_keys_by_command(
258 self,
259 keys: Iterable[KeyT],
260 slots_to_args: Mapping[int, Iterable[EncodableT]],
261 responses: Iterable[Any],
262 ) -> List[Any]:
263 results = {
264 k: v
265 for slot_values, response in zip(slots_to_args.values(), responses)
266 for k, v in zip(slot_values, response)
267 }
268 return [results[key] for key in keys]
269
270 def mget_nonatomic(self, keys: KeysT, *args: KeyT) -> List[Any | None]:
271 """
272 Splits the keys into different slots and then calls MGET
273 for the keys of every slot. This operation will not be atomic
274 if keys belong to more than one slot.
275
276 Returns a list of values ordered identically to ``keys``
277
278 For more information see https://redis.io/commands/mget
279 """
280
281 # Concatenate all keys into a list
282 keys = list_or_args(keys, args)
283
284 # Split keys into slots
285 slots_to_keys = self._partition_keys_by_slot(keys)
286
287 # Execute commands using a pipeline
288 res = self._execute_pipeline_by_slot("MGET", slots_to_keys)
289
290 # Reorder keys in the order the user provided & return
291 return self._reorder_keys_by_command(keys, slots_to_keys, res)
292
293 def mset_nonatomic(self, mapping: Mapping[AnyKeyT, EncodableT]) -> List[bool]:
294 """
295 Sets key/values based on a mapping. Mapping is a dictionary of
296 key/value pairs. Both keys and values should be strings or types that
297 can be cast to a string via str().
298
299 Splits the keys into different slots and then calls MSET
300 for the keys of every slot. This operation will not be atomic
301 if keys belong to more than one slot.
302
303 For more information see https://redis.io/commands/mset
304 """
305
306 # Partition the keys by slot
307 slots_to_pairs = self._partition_pairs_by_slot(mapping)
308
309 # Execute commands using a pipeline & return list of replies
310 return self._execute_pipeline_by_slot("MSET", slots_to_pairs)
311
312 def _split_command_across_slots(self, command: str, *keys: KeyT) -> int:
313 """
314 Runs the given command once for the keys
315 of each slot. Returns the sum of the return values.
316 """
317
318 # Partition the keys by slot
319 slots_to_keys = self._partition_keys_by_slot(keys)
320
321 # Sum up the reply from each command
322 return sum(self._execute_pipeline_by_slot(command, slots_to_keys))
323
324 @overload
325 def exists(self: SyncClientProtocol, *keys: KeyT) -> int: ...
326
327 @overload
328 def exists(self: AsyncClientProtocol, *keys: KeyT) -> Awaitable[int]: ...
329
330 def exists(self, *keys: KeyT) -> int | Awaitable[int]:
331 """
332 Returns the number of ``names`` that exist in the
333 whole cluster. The keys are first split up into slots
334 and then an EXISTS command is sent for every slot
335
336 For more information see https://redis.io/commands/exists
337 """
338 return self._split_command_across_slots("EXISTS", *keys)
339
340 @overload
341 def delete(self: SyncClientProtocol, *keys: KeyT) -> int: ...
342
343 @overload
344 def delete(self: AsyncClientProtocol, *keys: KeyT) -> Awaitable[int]: ...
345
346 def delete(self, *keys: KeyT) -> int | Awaitable[int]:
347 """
348 Deletes the given keys in the cluster.
349 The keys are first split up into slots
350 and then an DEL command is sent for every slot
351
352 Non-existent keys are ignored.
353 Returns the number of keys that were deleted.
354
355 For more information see https://redis.io/commands/del
356 """
357 return self._split_command_across_slots("DEL", *keys)
358
359 @overload
360 def touch(self: SyncClientProtocol, *keys: KeyT) -> int: ...
361
362 @overload
363 def touch(self: AsyncClientProtocol, *keys: KeyT) -> Awaitable[int]: ...
364
365 def touch(self, *keys: KeyT) -> int | Awaitable[int]:
366 """
367 Updates the last access time of given keys across the
368 cluster.
369
370 The keys are first split up into slots
371 and then an TOUCH command is sent for every slot
372
373 Non-existent keys are ignored.
374 Returns the number of keys that were touched.
375
376 For more information see https://redis.io/commands/touch
377 """
378 return self._split_command_across_slots("TOUCH", *keys)
379
380 @overload
381 def unlink(self: SyncClientProtocol, *keys: KeyT) -> int: ...
382
383 @overload
384 def unlink(self: AsyncClientProtocol, *keys: KeyT) -> Awaitable[int]: ...
385
386 def unlink(self, *keys: KeyT) -> int | Awaitable[int]:
387 """
388 Remove the specified keys in a different thread.
389
390 The keys are first split up into slots
391 and then an TOUCH command is sent for every slot
392
393 Non-existent keys are ignored.
394 Returns the number of keys that were unlinked.
395
396 For more information see https://redis.io/commands/unlink
397 """
398 return self._split_command_across_slots("UNLINK", *keys)
399
400
401class AsyncClusterMultiKeyCommands(ClusterMultiKeyCommands):
402 """
403 A class containing commands that handle more than one key
404 """
405
406 async def mget_nonatomic(self, keys: KeysT, *args: KeyT) -> List[Any | None]:
407 """
408 Splits the keys into different slots and then calls MGET
409 for the keys of every slot. This operation will not be atomic
410 if keys belong to more than one slot.
411
412 Returns a list of values ordered identically to ``keys``
413
414 For more information see https://redis.io/commands/mget
415 """
416
417 # Concatenate all keys into a list
418 keys = list_or_args(keys, args)
419
420 # Split keys into slots
421 slots_to_keys = self._partition_keys_by_slot(keys)
422
423 # Execute commands using a pipeline
424 res = await self._execute_pipeline_by_slot("MGET", slots_to_keys)
425
426 # Reorder keys in the order the user provided & return
427 return self._reorder_keys_by_command(keys, slots_to_keys, res)
428
429 async def mset_nonatomic(self, mapping: Mapping[AnyKeyT, EncodableT]) -> List[bool]:
430 """
431 Sets key/values based on a mapping. Mapping is a dictionary of
432 key/value pairs. Both keys and values should be strings or types that
433 can be cast to a string via str().
434
435 Splits the keys into different slots and then calls MSET
436 for the keys of every slot. This operation will not be atomic
437 if keys belong to more than one slot.
438
439 For more information see https://redis.io/commands/mset
440 """
441
442 # Partition the keys by slot
443 slots_to_pairs = self._partition_pairs_by_slot(mapping)
444
445 # Execute commands using a pipeline & return list of replies
446 return await self._execute_pipeline_by_slot("MSET", slots_to_pairs)
447
448 async def _split_command_across_slots(self, command: str, *keys: KeyT) -> int:
449 """
450 Runs the given command once for the keys
451 of each slot. Returns the sum of the return values.
452 """
453
454 # Partition the keys by slot
455 slots_to_keys = self._partition_keys_by_slot(keys)
456
457 # Sum up the reply from each command
458 return sum(await self._execute_pipeline_by_slot(command, slots_to_keys))
459
460 async def _is_replica_safe(self, command_name: str) -> bool:
461 warnings.warn(
462 "Using READ_COMMANDS for replica-safe checks is deprecated and will be removed in a future release. "
463 "Use an AsyncMetadataResolver instead.",
464 DeprecationWarning,
465 stacklevel=2,
466 )
467 return isinstance(command_name, str) and command_name.upper() in READ_COMMANDS
468
469 async def _execute_pipeline_by_slot(
470 self, command: str, slots_to_args: Mapping[int, Iterable[EncodableT]]
471 ) -> List[Any]:
472 if self._initialize:
473 await self.initialize()
474 replica_safe = (
475 self.read_from_replicas or self.load_balancing_strategy is not None
476 ) and await self._is_replica_safe(command)
477 pipe = self.pipeline()
478 [
479 pipe.execute_command(
480 command,
481 *slot_args,
482 target_nodes=[
483 self.nodes_manager.get_node_from_slot(
484 slot,
485 replica_safe,
486 self.load_balancing_strategy if replica_safe else None,
487 )
488 ],
489 )
490 for slot, slot_args in slots_to_args.items()
491 ]
492 return await pipe.execute()
493
494
495class ClusterManagementCommands(ManagementCommands):
496 """
497 A class for Redis Cluster management commands
498
499 The class inherits from Redis's core ManagementCommands class and do the
500 required adjustments to work with cluster mode
501 """
502
503 def slaveof(self, *args, **kwargs) -> NoReturn:
504 """
505 Make the server a replica of another instance, or promote it as master.
506
507 For more information see https://redis.io/commands/slaveof
508 """
509 raise RedisClusterException("SLAVEOF is not supported in cluster mode")
510
511 def replicaof(self, *args, **kwargs) -> NoReturn:
512 """
513 Make the server a replica of another instance, or promote it as master.
514
515 For more information see https://redis.io/commands/replicaof
516 """
517 raise RedisClusterException("REPLICAOF is not supported in cluster mode")
518
519 def swapdb(self, *args, **kwargs) -> NoReturn:
520 """
521 Swaps two Redis databases.
522
523 For more information see https://redis.io/commands/swapdb
524 """
525 raise RedisClusterException("SWAPDB is not supported in cluster mode")
526
527 @overload
528 def cluster_myid(
529 self: SyncClientProtocol, target_node: "TargetNodesT"
530 ) -> bytes | str: ...
531
532 @overload
533 def cluster_myid(
534 self: AsyncClientProtocol, target_node: "TargetNodesT"
535 ) -> Awaitable[bytes | str]: ...
536
537 def cluster_myid(self, target_node: "TargetNodesT") -> (bytes | str) | Awaitable[
538 bytes | str
539 ]:
540 """
541 Returns the node's id.
542
543 :target_node: 'ClusterNode'
544 The node to execute the command on
545
546 For more information check https://redis.io/commands/cluster-myid/
547 """
548 return self.execute_command("CLUSTER MYID", target_nodes=target_node)
549
550 @overload
551 def cluster_addslots(
552 self: SyncClientProtocol, target_node: "TargetNodesT", *slots: EncodableT
553 ) -> bool: ...
554
555 @overload
556 def cluster_addslots(
557 self: AsyncClientProtocol, target_node: "TargetNodesT", *slots: EncodableT
558 ) -> Awaitable[bool]: ...
559
560 def cluster_addslots(
561 self, target_node: "TargetNodesT", *slots: EncodableT
562 ) -> bool | Awaitable[bool]:
563 """
564 Assign new hash slots to receiving node. Sends to specified node.
565
566 :target_node: 'ClusterNode'
567 The node to execute the command on
568
569 For more information see https://redis.io/commands/cluster-addslots
570 """
571 return self.execute_command(
572 "CLUSTER ADDSLOTS", *slots, target_nodes=target_node
573 )
574
575 @overload
576 def cluster_addslotsrange(
577 self: SyncClientProtocol, target_node: "TargetNodesT", *slots: EncodableT
578 ) -> bool: ...
579
580 @overload
581 def cluster_addslotsrange(
582 self: AsyncClientProtocol, target_node: "TargetNodesT", *slots: EncodableT
583 ) -> Awaitable[bool]: ...
584
585 def cluster_addslotsrange(
586 self, target_node: "TargetNodesT", *slots: EncodableT
587 ) -> bool | Awaitable[bool]:
588 """
589 Similar to the CLUSTER ADDSLOTS command.
590 The difference between the two commands is that ADDSLOTS takes a list of slots
591 to assign to the node, while ADDSLOTSRANGE takes a list of slot ranges
592 (specified by start and end slots) to assign to the node.
593
594 :target_node: 'ClusterNode'
595 The node to execute the command on
596
597 For more information see https://redis.io/commands/cluster-addslotsrange
598 """
599 return self.execute_command(
600 "CLUSTER ADDSLOTSRANGE", *slots, target_nodes=target_node
601 )
602
603 @overload
604 def cluster_countkeysinslot(self: SyncClientProtocol, slot_id: int) -> int: ...
605
606 @overload
607 def cluster_countkeysinslot(
608 self: AsyncClientProtocol, slot_id: int
609 ) -> Awaitable[int]: ...
610
611 def cluster_countkeysinslot(self, slot_id: int) -> int | Awaitable[int]:
612 """
613 Return the number of local keys in the specified hash slot
614 Send to node based on specified slot_id
615
616 For more information see https://redis.io/commands/cluster-countkeysinslot
617 """
618 return self.execute_command("CLUSTER COUNTKEYSINSLOT", slot_id)
619
620 @overload
621 def cluster_count_failure_report(self: SyncClientProtocol, node_id: str) -> int: ...
622
623 @overload
624 def cluster_count_failure_report(
625 self: AsyncClientProtocol, node_id: str
626 ) -> Awaitable[int]: ...
627
628 def cluster_count_failure_report(self, node_id: str) -> int | Awaitable[int]:
629 """
630 Return the number of failure reports active for a given node
631 Sends to a random node
632
633 For more information see https://redis.io/commands/cluster-count-failure-reports
634 """
635 return self.execute_command("CLUSTER COUNT-FAILURE-REPORTS", node_id)
636
637 def cluster_delslots(self, *slots: EncodableT) -> List[bool]:
638 """
639 Set hash slots as unbound in the cluster.
640 It determines by it self what node the slot is in and sends it there
641
642 Returns a list of the results for each processed slot.
643
644 For more information see https://redis.io/commands/cluster-delslots
645 """
646 return [self.execute_command("CLUSTER DELSLOTS", slot) for slot in slots]
647
648 @overload
649 def cluster_delslotsrange(self: SyncClientProtocol, *slots: EncodableT) -> bool: ...
650
651 @overload
652 def cluster_delslotsrange(
653 self: AsyncClientProtocol, *slots: EncodableT
654 ) -> Awaitable[bool]: ...
655
656 def cluster_delslotsrange(self, *slots: EncodableT) -> bool | Awaitable[bool]:
657 """
658 Similar to the CLUSTER DELSLOTS command.
659 The difference is that CLUSTER DELSLOTS takes a list of hash slots to remove
660 from the node, while CLUSTER DELSLOTSRANGE takes a list of slot ranges to remove
661 from the node.
662
663 For more information see https://redis.io/commands/cluster-delslotsrange
664 """
665 return self.execute_command("CLUSTER DELSLOTSRANGE", *slots)
666
667 @overload
668 def cluster_failover(
669 self: SyncClientProtocol,
670 target_node: "TargetNodesT",
671 option: str | None = None,
672 ) -> bool: ...
673
674 @overload
675 def cluster_failover(
676 self: AsyncClientProtocol,
677 target_node: "TargetNodesT",
678 option: str | None = None,
679 ) -> Awaitable[bool]: ...
680
681 def cluster_failover(
682 self, target_node: "TargetNodesT", option: str | None = None
683 ) -> bool | Awaitable[bool]:
684 """
685 Forces a slave to perform a manual failover of its master
686 Sends to specified node
687
688 :target_node: 'ClusterNode'
689 The node to execute the command on
690
691 For more information see https://redis.io/commands/cluster-failover
692 """
693 if option:
694 if option.upper() not in ["FORCE", "TAKEOVER"]:
695 raise RedisError(
696 f"Invalid option for CLUSTER FAILOVER command: {option}"
697 )
698 else:
699 return self.execute_command(
700 "CLUSTER FAILOVER", option, target_nodes=target_node
701 )
702 else:
703 return self.execute_command("CLUSTER FAILOVER", target_nodes=target_node)
704
705 @overload
706 def cluster_info(
707 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None
708 ) -> dict[str, str]: ...
709
710 @overload
711 def cluster_info(
712 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None
713 ) -> Awaitable[dict[str, str]]: ...
714
715 def cluster_info(
716 self, target_nodes: "TargetNodesT" | None = None
717 ) -> dict[str, str] | Awaitable[dict[str, str]]:
718 """
719 Provides info about Redis Cluster node state.
720 The command will be sent to a random node in the cluster if no target
721 node is specified.
722
723 For more information see https://redis.io/commands/cluster-info
724 """
725 return self.execute_command("CLUSTER INFO", target_nodes=target_nodes)
726
727 @overload
728 def cluster_keyslot(self: SyncClientProtocol, key: str) -> int: ...
729
730 @overload
731 def cluster_keyslot(self: AsyncClientProtocol, key: str) -> Awaitable[int]: ...
732
733 def cluster_keyslot(self, key: str) -> int | Awaitable[int]:
734 """
735 Returns the hash slot of the specified key
736 Sends to random node in the cluster
737
738 For more information see https://redis.io/commands/cluster-keyslot
739 """
740 return self.execute_command("CLUSTER KEYSLOT", key)
741
742 @overload
743 def cluster_meet(
744 self: SyncClientProtocol,
745 host: str,
746 port: int,
747 target_nodes: "TargetNodesT" | None = None,
748 ) -> bool: ...
749
750 @overload
751 def cluster_meet(
752 self: AsyncClientProtocol,
753 host: str,
754 port: int,
755 target_nodes: "TargetNodesT" | None = None,
756 ) -> Awaitable[bool]: ...
757
758 def cluster_meet(
759 self, host: str, port: int, target_nodes: "TargetNodesT" | None = None
760 ) -> bool | Awaitable[bool]:
761 """
762 Force a node cluster to handshake with another node.
763 Sends to specified node.
764
765 For more information see https://redis.io/commands/cluster-meet
766 """
767 return self.execute_command(
768 "CLUSTER MEET", host, port, target_nodes=target_nodes
769 )
770
771 @overload
772 def cluster_nodes(self: SyncClientProtocol) -> dict[str, ClusterNodeDetail]: ...
773
774 @overload
775 def cluster_nodes(
776 self: AsyncClientProtocol,
777 ) -> Awaitable[dict[str, ClusterNodeDetail]]: ...
778
779 def cluster_nodes(
780 self,
781 ) -> dict[str, ClusterNodeDetail] | Awaitable[dict[str, ClusterNodeDetail]]:
782 """
783 Get Cluster config for the node.
784 Sends to random node in the cluster
785
786 For more information see https://redis.io/commands/cluster-nodes
787 """
788 return self.execute_command("CLUSTER NODES")
789
790 @overload
791 def cluster_replicate(
792 self: SyncClientProtocol, target_nodes: "TargetNodesT", node_id: str
793 ) -> bool: ...
794
795 @overload
796 def cluster_replicate(
797 self: AsyncClientProtocol, target_nodes: "TargetNodesT", node_id: str
798 ) -> Awaitable[bool]: ...
799
800 def cluster_replicate(
801 self, target_nodes: "TargetNodesT", node_id: str
802 ) -> bool | Awaitable[bool]:
803 """
804 Reconfigure a node as a slave of the specified master node
805
806 For more information see https://redis.io/commands/cluster-replicate
807 """
808 return self.execute_command(
809 "CLUSTER REPLICATE", node_id, target_nodes=target_nodes
810 )
811
812 @overload
813 def cluster_reset(
814 self: SyncClientProtocol,
815 soft: bool = True,
816 target_nodes: "TargetNodesT" | None = None,
817 ) -> bool: ...
818
819 @overload
820 def cluster_reset(
821 self: AsyncClientProtocol,
822 soft: bool = True,
823 target_nodes: "TargetNodesT" | None = None,
824 ) -> Awaitable[bool]: ...
825
826 def cluster_reset(
827 self, soft: bool = True, target_nodes: "TargetNodesT" | None = None
828 ) -> bool | Awaitable[bool]:
829 """
830 Reset a Redis Cluster node
831
832 If 'soft' is True then it will send 'SOFT' argument
833 If 'soft' is False then it will send 'HARD' argument
834
835 For more information see https://redis.io/commands/cluster-reset
836 """
837 return self.execute_command(
838 "CLUSTER RESET", b"SOFT" if soft else b"HARD", target_nodes=target_nodes
839 )
840
841 @overload
842 def cluster_save_config(
843 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None
844 ) -> bool: ...
845
846 @overload
847 def cluster_save_config(
848 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None
849 ) -> Awaitable[bool]: ...
850
851 def cluster_save_config(
852 self, target_nodes: "TargetNodesT" | None = None
853 ) -> bool | Awaitable[bool]:
854 """
855 Forces the node to save cluster state on disk
856
857 For more information see https://redis.io/commands/cluster-saveconfig
858 """
859 return self.execute_command("CLUSTER SAVECONFIG", target_nodes=target_nodes)
860
861 @overload
862 def cluster_get_keys_in_slot(
863 self: SyncClientProtocol, slot: int, num_keys: int
864 ) -> list[bytes | str]: ...
865
866 @overload
867 def cluster_get_keys_in_slot(
868 self: AsyncClientProtocol, slot: int, num_keys: int
869 ) -> Awaitable[list[bytes | str]]: ...
870
871 def cluster_get_keys_in_slot(
872 self, slot: int, num_keys: int
873 ) -> list[bytes | str] | Awaitable[list[bytes | str]]:
874 """
875 Returns the number of keys in the specified cluster slot
876
877 For more information see https://redis.io/commands/cluster-getkeysinslot
878 """
879 return self.execute_command("CLUSTER GETKEYSINSLOT", slot, num_keys)
880
881 @overload
882 def cluster_set_config_epoch(
883 self: SyncClientProtocol, epoch: int, target_nodes: "TargetNodesT" | None = None
884 ) -> bool: ...
885
886 @overload
887 def cluster_set_config_epoch(
888 self: AsyncClientProtocol,
889 epoch: int,
890 target_nodes: "TargetNodesT" | None = None,
891 ) -> Awaitable[bool]: ...
892
893 def cluster_set_config_epoch(
894 self, epoch: int, target_nodes: "TargetNodesT" | None = None
895 ) -> bool | Awaitable[bool]:
896 """
897 Set the configuration epoch in a new node
898
899 For more information see https://redis.io/commands/cluster-set-config-epoch
900 """
901 return self.execute_command(
902 "CLUSTER SET-CONFIG-EPOCH", epoch, target_nodes=target_nodes
903 )
904
905 @overload
906 def cluster_setslot(
907 self: SyncClientProtocol,
908 target_node: "TargetNodesT",
909 node_id: str,
910 slot_id: int,
911 state: str,
912 ) -> bool: ...
913
914 @overload
915 def cluster_setslot(
916 self: AsyncClientProtocol,
917 target_node: "TargetNodesT",
918 node_id: str,
919 slot_id: int,
920 state: str,
921 ) -> Awaitable[bool]: ...
922
923 def cluster_setslot(
924 self, target_node: "TargetNodesT", node_id: str, slot_id: int, state: str
925 ) -> bool | Awaitable[bool]:
926 """
927 Bind an hash slot to a specific node
928
929 :target_node: 'ClusterNode'
930 The node to execute the command on
931
932 For more information see https://redis.io/commands/cluster-setslot
933 """
934 if state.upper() in ("IMPORTING", "NODE", "MIGRATING"):
935 return self.execute_command(
936 "CLUSTER SETSLOT", slot_id, state, node_id, target_nodes=target_node
937 )
938 elif state.upper() == "STABLE":
939 raise RedisError('For "stable" state please use cluster_setslot_stable')
940 else:
941 raise RedisError(f"Invalid slot state: {state}")
942
943 @overload
944 def cluster_setslot_stable(self: SyncClientProtocol, slot_id: int) -> bool: ...
945
946 @overload
947 def cluster_setslot_stable(
948 self: AsyncClientProtocol, slot_id: int
949 ) -> Awaitable[bool]: ...
950
951 def cluster_setslot_stable(self, slot_id: int) -> bool | Awaitable[bool]:
952 """
953 Clears migrating / importing state from the slot.
954 It determines by it self what node the slot is in and sends it there.
955
956 For more information see https://redis.io/commands/cluster-setslot
957 """
958 return self.execute_command("CLUSTER SETSLOT", slot_id, "STABLE")
959
960 @overload
961 def cluster_replicas(
962 self: SyncClientProtocol,
963 node_id: str,
964 target_nodes: "TargetNodesT" | None = None,
965 ) -> dict[str, ClusterNodeDetail]: ...
966
967 @overload
968 def cluster_replicas(
969 self: AsyncClientProtocol,
970 node_id: str,
971 target_nodes: "TargetNodesT" | None = None,
972 ) -> Awaitable[dict[str, ClusterNodeDetail]]: ...
973
974 def cluster_replicas(
975 self, node_id: str, target_nodes: "TargetNodesT" | None = None
976 ) -> dict[str, ClusterNodeDetail] | Awaitable[dict[str, ClusterNodeDetail]]:
977 """
978 Provides a list of replica nodes replicating from the specified primary
979 target node.
980
981 For more information see https://redis.io/commands/cluster-replicas
982 """
983 return self.execute_command(
984 "CLUSTER REPLICAS", node_id, target_nodes=target_nodes
985 )
986
987 @overload
988 def cluster_slots(
989 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None
990 ) -> list[Any]: ...
991
992 @overload
993 def cluster_slots(
994 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None
995 ) -> Awaitable[list[Any]]: ...
996
997 def cluster_slots(
998 self, target_nodes: "TargetNodesT" | None = None
999 ) -> list[Any] | Awaitable[list[Any]]:
1000 """
1001 Get array of Cluster slot to node mappings
1002
1003 For more information see https://redis.io/commands/cluster-slots
1004 """
1005 return self.execute_command("CLUSTER SLOTS", target_nodes=target_nodes)
1006
1007 @overload
1008 def cluster_shards(
1009 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None
1010 ) -> ClusterShardsResponse: ...
1011
1012 @overload
1013 def cluster_shards(
1014 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None
1015 ) -> Awaitable[ClusterShardsResponse]: ...
1016
1017 def cluster_shards(
1018 self, target_nodes: "TargetNodesT" | None = None
1019 ) -> ClusterShardsResponse | Awaitable[ClusterShardsResponse]:
1020 """
1021 Returns details about the shards of the cluster.
1022
1023 For more information see https://redis.io/commands/cluster-shards
1024 """
1025 return self.execute_command("CLUSTER SHARDS", target_nodes=target_nodes)
1026
1027 @overload
1028 def cluster_myshardid(
1029 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None
1030 ) -> bytes | str: ...
1031
1032 @overload
1033 def cluster_myshardid(
1034 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None
1035 ) -> Awaitable[bytes | str]: ...
1036
1037 def cluster_myshardid(self, target_nodes: "TargetNodesT" | None = None) -> (
1038 bytes | str
1039 ) | Awaitable[bytes | str]:
1040 """
1041 Returns the shard ID of the node.
1042
1043 For more information see https://redis.io/commands/cluster-myshardid/
1044 """
1045 return self.execute_command("CLUSTER MYSHARDID", target_nodes=target_nodes)
1046
1047 @overload
1048 def cluster_links(
1049 self: SyncClientProtocol, target_node: "TargetNodesT"
1050 ) -> ClusterLinksResponse: ...
1051
1052 @overload
1053 def cluster_links(
1054 self: AsyncClientProtocol, target_node: "TargetNodesT"
1055 ) -> Awaitable[ClusterLinksResponse]: ...
1056
1057 def cluster_links(
1058 self, target_node: "TargetNodesT"
1059 ) -> ClusterLinksResponse | Awaitable[ClusterLinksResponse]:
1060 """
1061 Each node in a Redis Cluster maintains a pair of long-lived TCP link with each
1062 peer in the cluster: One for sending outbound messages towards the peer and one
1063 for receiving inbound messages from the peer.
1064
1065 This command outputs information of all such peer links as an array.
1066
1067 For more information see https://redis.io/commands/cluster-links
1068 """
1069 return self.execute_command("CLUSTER LINKS", target_nodes=target_node)
1070
1071 def cluster_flushslots(self, target_nodes: "TargetNodesT" | None = None) -> None:
1072 raise NotImplementedError(
1073 "CLUSTER FLUSHSLOTS is intentionally not implemented in the client."
1074 )
1075
1076 def cluster_bumpepoch(self, target_nodes: "TargetNodesT" | None = None) -> None:
1077 raise NotImplementedError(
1078 "CLUSTER BUMPEPOCH is intentionally not implemented in the client."
1079 )
1080
1081 def readonly(self, target_nodes: "TargetNodesT" | None = None) -> ResponseT:
1082 """
1083 Enables read queries.
1084 The command will be sent to the default cluster node if target_nodes is
1085 not specified.
1086
1087 For more information see https://redis.io/commands/readonly
1088 """
1089 if target_nodes == "replicas" or target_nodes == "all":
1090 # read_from_replicas will only be enabled if the READONLY command
1091 # is sent to all replicas
1092 self.read_from_replicas = True
1093 return self.execute_command("READONLY", target_nodes=target_nodes)
1094
1095 def readwrite(self, target_nodes: "TargetNodesT" | None = None) -> ResponseT:
1096 """
1097 Disables read queries.
1098 The command will be sent to the default cluster node if target_nodes is
1099 not specified.
1100
1101 For more information see https://redis.io/commands/readwrite
1102 """
1103 # Reset read from replicas flag
1104 self.read_from_replicas = False
1105 return self.execute_command("READWRITE", target_nodes=target_nodes)
1106
1107 @deprecated_function(
1108 version="7.2.0",
1109 reason="Use client-side caching feature instead.",
1110 )
1111 def client_tracking_on(
1112 self,
1113 clientid: int | None = None,
1114 prefix: Sequence[KeyT] = [],
1115 bcast: bool = False,
1116 optin: bool = False,
1117 optout: bool = False,
1118 noloop: bool = False,
1119 target_nodes: "TargetNodesT" | None = "all",
1120 ) -> ResponseT:
1121 """
1122 Enables the tracking feature of the Redis server, that is used
1123 for server assisted client side caching.
1124
1125 When clientid is provided - in target_nodes only the node that owns the
1126 connection with this id should be provided.
1127 When clientid is not provided - target_nodes can be any node.
1128
1129 For more information see https://redis.io/commands/client-tracking
1130 """
1131 return self.client_tracking(
1132 True,
1133 clientid,
1134 prefix,
1135 bcast,
1136 optin,
1137 optout,
1138 noloop,
1139 target_nodes=target_nodes,
1140 )
1141
1142 @deprecated_function(
1143 version="7.2.0",
1144 reason="Use client-side caching feature instead.",
1145 )
1146 def client_tracking_off(
1147 self,
1148 clientid: int | None = None,
1149 prefix: Sequence[KeyT] = [],
1150 bcast: bool = False,
1151 optin: bool = False,
1152 optout: bool = False,
1153 noloop: bool = False,
1154 target_nodes: "TargetNodesT" | None = "all",
1155 ) -> ResponseT:
1156 """
1157 Disables the tracking feature of the Redis server, that is used
1158 for server assisted client side caching.
1159
1160 When clientid is provided - in target_nodes only the node that owns the
1161 connection with this id should be provided.
1162 When clientid is not provided - target_nodes can be any node.
1163
1164 For more information see https://redis.io/commands/client-tracking
1165 """
1166 return self.client_tracking(
1167 False,
1168 clientid,
1169 prefix,
1170 bcast,
1171 optin,
1172 optout,
1173 noloop,
1174 target_nodes=target_nodes,
1175 )
1176
1177 def hotkeys_start(
1178 self,
1179 metrics: List[HotkeysMetricsTypes],
1180 count: int | None = None,
1181 duration: int | None = None,
1182 sample_ratio: int | None = None,
1183 slots: List[int] | None = None,
1184 **kwargs,
1185 ) -> str | bytes:
1186 """
1187 Cluster client does not support hotkeys command. Please use the non-cluster client.
1188
1189 For more information see https://redis.io/commands/hotkeys-start
1190 """
1191 raise NotImplementedError(
1192 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1193 )
1194
1195 def hotkeys_stop(self, **kwargs) -> str | bytes:
1196 """
1197 Cluster client does not support hotkeys command. Please use the non-cluster client.
1198
1199 For more information see https://redis.io/commands/hotkeys-stop
1200 """
1201 raise NotImplementedError(
1202 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1203 )
1204
1205 def hotkeys_reset(self, **kwargs) -> str | bytes:
1206 """
1207 Cluster client does not support hotkeys command. Please use the non-cluster client.
1208
1209 For more information see https://redis.io/commands/hotkeys-reset
1210 """
1211 raise NotImplementedError(
1212 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1213 )
1214
1215 def hotkeys_get(self, **kwargs) -> list[dict[str | bytes, Any]]:
1216 """
1217 Cluster client does not support hotkeys command. Please use the non-cluster client.
1218
1219 For more information see https://redis.io/commands/hotkeys-get
1220 """
1221 raise NotImplementedError(
1222 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1223 )
1224
1225
1226class AsyncClusterManagementCommands(
1227 ClusterManagementCommands, AsyncManagementCommands
1228):
1229 """
1230 A class for Redis Cluster management commands
1231
1232 The class inherits from Redis's core ManagementCommands class and do the
1233 required adjustments to work with cluster mode
1234 """
1235
1236 async def cluster_delslots(self, *slots: EncodableT) -> List[bool]:
1237 """
1238 Set hash slots as unbound in the cluster.
1239 It determines by it self what node the slot is in and sends it there
1240
1241 Returns a list of the results for each processed slot.
1242
1243 For more information see https://redis.io/commands/cluster-delslots
1244 """
1245 return await asyncio.gather(
1246 *(
1247 asyncio.create_task(self.execute_command("CLUSTER DELSLOTS", slot))
1248 for slot in slots
1249 )
1250 )
1251
1252 @deprecated_function(
1253 version="7.2.0",
1254 reason="Use client-side caching feature instead.",
1255 )
1256 async def client_tracking_on(
1257 self,
1258 clientid: int | None = None,
1259 prefix: Sequence[KeyT] = [],
1260 bcast: bool = False,
1261 optin: bool = False,
1262 optout: bool = False,
1263 noloop: bool = False,
1264 target_nodes: "TargetNodesT" | None = "all",
1265 ) -> ResponseT:
1266 """
1267 Enables the tracking feature of the Redis server, that is used
1268 for server assisted client side caching.
1269
1270 When clientid is provided - in target_nodes only the node that owns the
1271 connection with this id should be provided.
1272 When clientid is not provided - target_nodes can be any node.
1273
1274 For more information see https://redis.io/commands/client-tracking
1275 """
1276 return await self.client_tracking(
1277 True,
1278 clientid,
1279 prefix,
1280 bcast,
1281 optin,
1282 optout,
1283 noloop,
1284 target_nodes=target_nodes,
1285 )
1286
1287 @deprecated_function(
1288 version="7.2.0",
1289 reason="Use client-side caching feature instead.",
1290 )
1291 async def client_tracking_off(
1292 self,
1293 clientid: int | None = None,
1294 prefix: Sequence[KeyT] = [],
1295 bcast: bool = False,
1296 optin: bool = False,
1297 optout: bool = False,
1298 noloop: bool = False,
1299 target_nodes: "TargetNodesT" | None = "all",
1300 ) -> ResponseT:
1301 """
1302 Disables the tracking feature of the Redis server, that is used
1303 for server assisted client side caching.
1304
1305 When clientid is provided - in target_nodes only the node that owns the
1306 connection with this id should be provided.
1307 When clientid is not provided - target_nodes can be any node.
1308
1309 For more information see https://redis.io/commands/client-tracking
1310 """
1311 return await self.client_tracking(
1312 False,
1313 clientid,
1314 prefix,
1315 bcast,
1316 optin,
1317 optout,
1318 noloop,
1319 target_nodes=target_nodes,
1320 )
1321
1322 async def hotkeys_start(
1323 self,
1324 metrics: List[HotkeysMetricsTypes],
1325 count: int | None = None,
1326 duration: int | None = None,
1327 sample_ratio: int | None = None,
1328 slots: List[int] | None = None,
1329 **kwargs,
1330 ) -> Awaitable[str | bytes]:
1331 """
1332 Cluster client does not support hotkeys command. Please use the non-cluster client.
1333
1334 For more information see https://redis.io/commands/hotkeys-start
1335 """
1336 raise NotImplementedError(
1337 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1338 )
1339
1340 async def hotkeys_stop(self, **kwargs) -> Awaitable[str | bytes]:
1341 """
1342 Cluster client does not support hotkeys command. Please use the non-cluster client.
1343
1344 For more information see https://redis.io/commands/hotkeys-stop
1345 """
1346 raise NotImplementedError(
1347 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1348 )
1349
1350 async def hotkeys_reset(self, **kwargs) -> Awaitable[str | bytes]:
1351 """
1352 Cluster client does not support hotkeys command. Please use the non-cluster client.
1353
1354 For more information see https://redis.io/commands/hotkeys-reset
1355 """
1356 raise NotImplementedError(
1357 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1358 )
1359
1360 async def hotkeys_get(self, **kwargs) -> Awaitable[list[dict[str | bytes, Any]]]:
1361 """
1362 Cluster client does not support hotkeys command. Please use the non-cluster client.
1363
1364 For more information see https://redis.io/commands/hotkeys-get
1365 """
1366 raise NotImplementedError(
1367 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1368 )
1369
1370
1371class ClusterDataAccessCommands(DataAccessCommands):
1372 """
1373 A class for Redis Cluster Data Access Commands
1374
1375 The class inherits from Redis's core DataAccessCommand class and do the
1376 required adjustments to work with cluster mode
1377 """
1378
1379 @overload
1380 def stralgo(
1381 self: SyncClientProtocol,
1382 algo: Literal["LCS"],
1383 value1: KeyT,
1384 value2: KeyT,
1385 specific_argument: Literal["strings"] | Literal["keys"] = "strings",
1386 len: bool = False,
1387 idx: bool = False,
1388 minmatchlen: int | None = None,
1389 withmatchlen: bool = False,
1390 **kwargs,
1391 ) -> StralgoResponse: ...
1392
1393 @overload
1394 def stralgo(
1395 self: AsyncClientProtocol,
1396 algo: Literal["LCS"],
1397 value1: KeyT,
1398 value2: KeyT,
1399 specific_argument: Literal["strings"] | Literal["keys"] = "strings",
1400 len: bool = False,
1401 idx: bool = False,
1402 minmatchlen: int | None = None,
1403 withmatchlen: bool = False,
1404 **kwargs,
1405 ) -> Awaitable[StralgoResponse]: ...
1406
1407 def stralgo(
1408 self,
1409 algo: Literal["LCS"],
1410 value1: KeyT,
1411 value2: KeyT,
1412 specific_argument: Literal["strings"] | Literal["keys"] = "strings",
1413 len: bool = False,
1414 idx: bool = False,
1415 minmatchlen: int | None = None,
1416 withmatchlen: bool = False,
1417 **kwargs,
1418 ) -> StralgoResponse | Awaitable[StralgoResponse]:
1419 """
1420 Implements complex algorithms that operate on strings.
1421 Right now the only algorithm implemented is the LCS algorithm
1422 (longest common substring). However new algorithms could be
1423 implemented in the future.
1424
1425 ``algo`` Right now must be LCS
1426 ``value1`` and ``value2`` Can be two strings or two keys
1427 ``specific_argument`` Specifying if the arguments to the algorithm
1428 will be keys or strings. strings is the default.
1429 ``len`` Returns just the len of the match.
1430 ``idx`` Returns the match positions in each string.
1431 ``minmatchlen`` Restrict the list of matches to the ones of a given
1432 minimal length. Can be provided only when ``idx`` set to True.
1433 ``withmatchlen`` Returns the matches with the len of the match.
1434 Can be provided only when ``idx`` set to True.
1435
1436 For more information see https://redis.io/commands/stralgo
1437 """
1438 target_nodes = kwargs.pop("target_nodes", None)
1439 if specific_argument == "strings" and target_nodes is None:
1440 target_nodes = "default-node"
1441 kwargs.update({"target_nodes": target_nodes})
1442 return super().stralgo(
1443 algo,
1444 value1,
1445 value2,
1446 specific_argument,
1447 len,
1448 idx,
1449 minmatchlen,
1450 withmatchlen,
1451 **kwargs,
1452 )
1453
1454 def scan_iter(
1455 self,
1456 match: PatternT | None = None,
1457 count: int | None = None,
1458 _type: str | None = None,
1459 **kwargs,
1460 ) -> Iterator:
1461 # Do the first query with cursor=0 for all nodes
1462 cursors, data = self.scan(match=match, count=count, _type=_type, **kwargs)
1463 yield from data
1464
1465 cursors = {name: cursor for name, cursor in cursors.items() if cursor != 0}
1466 if cursors:
1467 # Get nodes by name
1468 nodes = {name: self.get_node(node_name=name) for name in cursors.keys()}
1469
1470 # Iterate over each node till its cursor is 0
1471 kwargs.pop("target_nodes", None)
1472 while cursors:
1473 for name, cursor in cursors.items():
1474 cur, data = self.scan(
1475 cursor=cursor,
1476 match=match,
1477 count=count,
1478 _type=_type,
1479 target_nodes=nodes[name],
1480 **kwargs,
1481 )
1482 yield from data
1483 cursors[name] = cur[name]
1484
1485 cursors = {
1486 name: cursor for name, cursor in cursors.items() if cursor != 0
1487 }
1488
1489
1490class AsyncClusterDataAccessCommands(
1491 ClusterDataAccessCommands, AsyncDataAccessCommands
1492):
1493 """
1494 A class for Redis Cluster Data Access Commands
1495
1496 The class inherits from Redis's core DataAccessCommand class and do the
1497 required adjustments to work with cluster mode
1498 """
1499
1500 async def scan_iter(
1501 self,
1502 match: PatternT | None = None,
1503 count: int | None = None,
1504 _type: str | None = None,
1505 **kwargs,
1506 ) -> AsyncIterator:
1507 # Do the first query with cursor=0 for all nodes
1508 cursors, data = await self.scan(match=match, count=count, _type=_type, **kwargs)
1509 for value in data:
1510 yield value
1511
1512 cursors = {name: cursor for name, cursor in cursors.items() if cursor != 0}
1513 if cursors:
1514 # Get nodes by name
1515 nodes = {name: self.get_node(node_name=name) for name in cursors.keys()}
1516
1517 # Iterate over each node till its cursor is 0
1518 kwargs.pop("target_nodes", None)
1519 while cursors:
1520 for name, cursor in cursors.items():
1521 cur, data = await self.scan(
1522 cursor=cursor,
1523 match=match,
1524 count=count,
1525 _type=_type,
1526 target_nodes=nodes[name],
1527 **kwargs,
1528 )
1529 for value in data:
1530 yield value
1531 cursors[name] = cur[name]
1532
1533 cursors = {
1534 name: cursor for name, cursor in cursors.items() if cursor != 0
1535 }
1536
1537
1538class ClusterBlessCommands(BlessCommands):
1539 """
1540 A class for Redis Cluster BLESS commands
1541
1542 The class inherits from Redis's core BlessCommands class and does the
1543 required adjustments to work with cluster mode
1544 """
1545
1546 def bless_scan_iter(
1547 self,
1548 flag: BlessFlag,
1549 count: int | None = None,
1550 **kwargs,
1551 ) -> Iterator[bytes | str]:
1552 # Do the first query with cursor=0 for all nodes
1553 cursors, data = self.bless_scan(cursor=0, flag=flag, count=count, **kwargs)
1554 yield from data
1555
1556 cursors = {name: cursor for name, cursor in cursors.items() if cursor != 0}
1557 if cursors:
1558 # Get nodes by name
1559 nodes = {name: self.get_node(node_name=name) for name in cursors.keys()}
1560
1561 # Iterate over each node till its cursor is 0
1562 kwargs.pop("target_nodes", None)
1563 while cursors:
1564 for name, cursor in cursors.items():
1565 cur, data = self.bless_scan(
1566 cursor=cursor,
1567 flag=flag,
1568 count=count,
1569 target_nodes=nodes[name],
1570 **kwargs,
1571 )
1572 yield from data
1573 cursors[name] = cur[name]
1574
1575 cursors = {
1576 name: cursor for name, cursor in cursors.items() if cursor != 0
1577 }
1578
1579
1580class AsyncClusterBlessCommands(ClusterBlessCommands, AsyncBlessCommands):
1581 """
1582 A class for Redis Cluster BLESS commands
1583
1584 The class inherits from Redis's core BlessCommands class and does the
1585 required adjustments to work with cluster mode
1586 """
1587
1588 async def bless_scan_iter(
1589 self,
1590 flag: BlessFlag,
1591 count: int | None = None,
1592 **kwargs,
1593 ) -> AsyncIterator[bytes | str]:
1594 # Do the first query with cursor=0 for all nodes
1595 cursors, data = await self.bless_scan(
1596 cursor=0, flag=flag, count=count, **kwargs
1597 )
1598 for value in data:
1599 yield value
1600
1601 cursors = {name: cursor for name, cursor in cursors.items() if cursor != 0}
1602 if cursors:
1603 # Get nodes by name
1604 nodes = {name: self.get_node(node_name=name) for name in cursors.keys()}
1605
1606 # Iterate over each node till its cursor is 0
1607 kwargs.pop("target_nodes", None)
1608 while cursors:
1609 for name, cursor in cursors.items():
1610 cur, data = await self.bless_scan(
1611 cursor=cursor,
1612 flag=flag,
1613 count=count,
1614 target_nodes=nodes[name],
1615 **kwargs,
1616 )
1617 for value in data:
1618 yield value
1619 cursors[name] = cur[name]
1620
1621 cursors = {
1622 name: cursor for name, cursor in cursors.items() if cursor != 0
1623 }
1624
1625
1626class RedisClusterCommands(
1627 ClusterMultiKeyCommands,
1628 ClusterManagementCommands,
1629 ACLCommands,
1630 PubSubCommands,
1631 ClusterDataAccessCommands,
1632 ScriptCommands,
1633 FunctionCommands,
1634 ClusterBlessCommands,
1635 ModuleCommands,
1636 RedisModuleCommands,
1637):
1638 """
1639 A class for all Redis Cluster commands
1640
1641 For key-based commands, the target node(s) will be internally determined
1642 by the keys' hash slot.
1643 Non-key-based commands can be executed with the 'target_nodes' argument to
1644 target specific nodes. By default, if target_nodes is not specified, the
1645 command will be executed on the default cluster node.
1646
1647 :param target_nodes: type can be one of the following:
1648 - nodes flag: ALL_NODES, PRIMARIES, REPLICAS, RANDOM
1649 - 'ClusterNode'
1650 - 'list(ClusterNodes)'
1651 - 'dict(any:clusterNodes)'
1652
1653 for example:
1654 r.cluster_info(target_nodes=RedisCluster.ALL_NODES)
1655 """
1656
1657
1658class AsyncRedisClusterCommands(
1659 AsyncClusterMultiKeyCommands,
1660 AsyncClusterManagementCommands,
1661 AsyncACLCommands,
1662 PubSubCommands,
1663 AsyncClusterDataAccessCommands,
1664 AsyncScriptCommands,
1665 AsyncFunctionCommands,
1666 AsyncClusterBlessCommands,
1667 AsyncModuleCommands,
1668 AsyncRedisModuleCommands,
1669):
1670 """
1671 A class for all Redis Cluster commands
1672
1673 For key-based commands, the target node(s) will be internally determined
1674 by the keys' hash slot.
1675 Non-key-based commands can be executed with the 'target_nodes' argument to
1676 target specific nodes. By default, if target_nodes is not specified, the
1677 command will be executed on the default cluster node.
1678
1679 :param target_nodes: type can be one of the following:
1680 - nodes flag: ALL_NODES, PRIMARIES, REPLICAS, RANDOM
1681 - 'ClusterNode'
1682 - 'list(ClusterNodes)'
1683 - 'dict(any:clusterNodes)'
1684
1685 for example:
1686 r.cluster_info(target_nodes=RedisCluster.ALL_NODES)
1687 """