Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/redis/commands/cluster.py: 60%

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

337 statements  

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 AsyncDataAccessCommands, 

44 AsyncFunctionCommands, 

45 AsyncManagementCommands, 

46 AsyncModuleCommands, 

47 AsyncScriptCommands, 

48 DataAccessCommands, 

49 FunctionCommands, 

50 HotkeysMetricsTypes, 

51 ManagementCommands, 

52 ModuleCommands, 

53 PubSubCommands, 

54 ScriptCommands, 

55) 

56from .helpers import list_or_args 

57from .redismodules import AsyncRedisModuleCommands, RedisModuleCommands 

58 

59if TYPE_CHECKING: 

60 from redis.asyncio.cluster import TargetNodesT 

61 from redis.cluster import LoadBalancingStrategy 

62 

63# DEPRECATED - no longer consulted by the default metadata routing, and it will be removed in a future release. 

64# 

65# Replica safety is now decided from command metadata by the metadata resolver, 

66# so this set no longer describes what commands are safe to execute on replicas. 

67# It is kept as a public attribute only so an external caller reading it keeps working; 

68# editing it changes nothing. 

69# 

70# To change replica safety, edit ``redis.commands.metadata._STATIC_COMMAND_METADATA`` or pass 

71# a ``metadata_resolver`` to the client. Nothing here. 

72READ_COMMANDS = frozenset( 

73 [ 

74 # Bit Operations 

75 "BITCOUNT", 

76 "BITFIELD_RO", 

77 "BITPOS", 

78 # Scripting 

79 "EVAL_RO", 

80 "EVALSHA_RO", 

81 "FCALL_RO", 

82 # Key Operations 

83 "DBSIZE", 

84 "DIGEST", 

85 "DUMP", 

86 "EXISTS", 

87 "EXPIRETIME", 

88 "PEXPIRETIME", 

89 "KEYS", 

90 "SCAN", 

91 "PTTL", 

92 "RANDOMKEY", 

93 "TTL", 

94 "TYPE", 

95 # String Operations 

96 "GET", 

97 "GETBIT", 

98 "GETRANGE", 

99 "MGET", 

100 "STRLEN", 

101 "LCS", 

102 # Geo Operations 

103 "GEODIST", 

104 "GEOHASH", 

105 "GEOPOS", 

106 "GEOSEARCH", 

107 # Hash Operations 

108 "HEXISTS", 

109 "HGET", 

110 "HGETALL", 

111 "HKEYS", 

112 "HLEN", 

113 "HMGET", 

114 "HSTRLEN", 

115 "HVALS", 

116 "HRANDFIELD", 

117 "HEXPIRETIME", 

118 "HPEXPIRETIME", 

119 "HTTL", 

120 "HPTTL", 

121 "HSCAN", 

122 # List Operations 

123 "LINDEX", 

124 "LPOS", 

125 "LLEN", 

126 "LRANGE", 

127 # Set Operations 

128 "SCARD", 

129 "SDIFF", 

130 "SDIFFCARD", 

131 "SINTER", 

132 "SINTERCARD", 

133 "SISMEMBER", 

134 "SMISMEMBER", 

135 "SMEMBERS", 

136 "SRANDMEMBER", 

137 "SUNION", 

138 "SUNIONCARD", 

139 "SSCAN", 

140 # Sorted Set Operations 

141 "ZCARD", 

142 "ZCOUNT", 

143 "ZDIFF", 

144 "ZINTER", 

145 "ZINTERCARD", 

146 "ZLEXCOUNT", 

147 "ZMSCORE", 

148 "ZRANDMEMBER", 

149 "ZRANGE", 

150 "ZRANGEBYLEX", 

151 "ZRANGEBYSCORE", 

152 "ZRANK", 

153 "ZREVRANGE", 

154 "ZREVRANGEBYLEX", 

155 "ZREVRANGEBYSCORE", 

156 "ZREVRANK", 

157 "ZSCAN", 

158 "ZSCORE", 

159 "ZUNION", 

160 # Stream Operations 

161 "XLEN", 

162 "XPENDING", 

163 "XRANGE", 

164 "XREAD", 

165 "XREVRANGE", 

166 # JSON Module 

167 "JSON.ARRINDEX", 

168 "JSON.ARRLEN", 

169 "JSON.GET", 

170 "JSON.MGET", 

171 "JSON.OBJKEYS", 

172 "JSON.OBJLEN", 

173 "JSON.RESP", 

174 "JSON.STRLEN", 

175 "JSON.TYPE", 

176 # RediSearch Module 

177 "FT.ALIASLIST", 

178 "FT.EXPLAIN", 

179 "FT.INFO", 

180 "FT.PROFILE", 

181 "FT.SEARCH", 

182 ] 

183) 

184 

185 

186class ClusterMultiKeyCommands(ClusterCommandsProtocol): 

187 """ 

188 A class containing commands that handle more than one key 

189 """ 

190 

191 # Read-routing configuration, which the cluster clients all set on the instance. The 

192 # defaults are here for a host class that mixes these commands in without one - see 

193 # ``_is_replica_safe`` below, which is the same fallback for the same audience - so 

194 # that ``_execute_pipeline_by_slot`` reads a value rather than raising 

195 # ``AttributeError``. Typed under TYPE_CHECKING because ``redis.cluster`` imports this 

196 # module. 

197 read_from_replicas: bool = False 

198 load_balancing_strategy: LoadBalancingStrategy | None = None 

199 

200 def _partition_keys_by_slot(self, keys: Iterable[KeyT]) -> Dict[int, List[KeyT]]: 

201 """Split keys into a dictionary that maps a slot to a list of keys.""" 

202 

203 slots_to_keys = {} 

204 for key in keys: 

205 slot = key_slot(self.encoder.encode(key)) 

206 slots_to_keys.setdefault(slot, []).append(key) 

207 

208 return slots_to_keys 

209 

210 def _partition_pairs_by_slot( 

211 self, mapping: Mapping[AnyKeyT, EncodableT] 

212 ) -> Dict[int, List[EncodableT]]: 

213 """Split pairs into a dictionary that maps a slot to a list of pairs.""" 

214 

215 slots_to_pairs = {} 

216 for pair in mapping.items(): 

217 slot = key_slot(self.encoder.encode(pair[0])) 

218 slots_to_pairs.setdefault(slot, []).extend(pair) 

219 

220 return slots_to_pairs 

221 

222 def _is_replica_safe(self, command_name: str) -> bool: 

223 warnings.warn( 

224 "Using READ_COMMANDS for replica-safe checks is deprecated and will be removed in a future release. " 

225 "Use a MetadataResolver instead.", 

226 DeprecationWarning, 

227 stacklevel=2, 

228 ) 

229 return isinstance(command_name, str) and command_name.upper() in READ_COMMANDS 

230 

231 def _execute_pipeline_by_slot( 

232 self, command: str, slots_to_args: Mapping[int, Iterable[EncodableT]] 

233 ) -> List[Any]: 

234 replica_safe = ( 

235 self.read_from_replicas or self.load_balancing_strategy is not None 

236 ) and self._is_replica_safe(command) 

237 pipe = self.pipeline() 

238 [ 

239 pipe.execute_command( 

240 command, 

241 *slot_args, 

242 target_nodes=[ 

243 self.nodes_manager.get_node_from_slot( 

244 slot, 

245 replica_safe, 

246 self.load_balancing_strategy if replica_safe else None, 

247 ) 

248 ], 

249 ) 

250 for slot, slot_args in slots_to_args.items() 

251 ] 

252 return pipe.execute() 

253 

254 def _reorder_keys_by_command( 

255 self, 

256 keys: Iterable[KeyT], 

257 slots_to_args: Mapping[int, Iterable[EncodableT]], 

258 responses: Iterable[Any], 

259 ) -> List[Any]: 

260 results = { 

261 k: v 

262 for slot_values, response in zip(slots_to_args.values(), responses) 

263 for k, v in zip(slot_values, response) 

264 } 

265 return [results[key] for key in keys] 

266 

267 def mget_nonatomic(self, keys: KeysT, *args: KeyT) -> List[Any | None]: 

268 """ 

269 Splits the keys into different slots and then calls MGET 

270 for the keys of every slot. This operation will not be atomic 

271 if keys belong to more than one slot. 

272 

273 Returns a list of values ordered identically to ``keys`` 

274 

275 For more information see https://redis.io/commands/mget 

276 """ 

277 

278 # Concatenate all keys into a list 

279 keys = list_or_args(keys, args) 

280 

281 # Split keys into slots 

282 slots_to_keys = self._partition_keys_by_slot(keys) 

283 

284 # Execute commands using a pipeline 

285 res = self._execute_pipeline_by_slot("MGET", slots_to_keys) 

286 

287 # Reorder keys in the order the user provided & return 

288 return self._reorder_keys_by_command(keys, slots_to_keys, res) 

289 

290 def mset_nonatomic(self, mapping: Mapping[AnyKeyT, EncodableT]) -> List[bool]: 

291 """ 

292 Sets key/values based on a mapping. Mapping is a dictionary of 

293 key/value pairs. Both keys and values should be strings or types that 

294 can be cast to a string via str(). 

295 

296 Splits the keys into different slots and then calls MSET 

297 for the keys of every slot. This operation will not be atomic 

298 if keys belong to more than one slot. 

299 

300 For more information see https://redis.io/commands/mset 

301 """ 

302 

303 # Partition the keys by slot 

304 slots_to_pairs = self._partition_pairs_by_slot(mapping) 

305 

306 # Execute commands using a pipeline & return list of replies 

307 return self._execute_pipeline_by_slot("MSET", slots_to_pairs) 

308 

309 def _split_command_across_slots(self, command: str, *keys: KeyT) -> int: 

310 """ 

311 Runs the given command once for the keys 

312 of each slot. Returns the sum of the return values. 

313 """ 

314 

315 # Partition the keys by slot 

316 slots_to_keys = self._partition_keys_by_slot(keys) 

317 

318 # Sum up the reply from each command 

319 return sum(self._execute_pipeline_by_slot(command, slots_to_keys)) 

320 

321 @overload 

322 def exists(self: SyncClientProtocol, *keys: KeyT) -> int: ... 

323 

324 @overload 

325 def exists(self: AsyncClientProtocol, *keys: KeyT) -> Awaitable[int]: ... 

326 

327 def exists(self, *keys: KeyT) -> int | Awaitable[int]: 

328 """ 

329 Returns the number of ``names`` that exist in the 

330 whole cluster. The keys are first split up into slots 

331 and then an EXISTS command is sent for every slot 

332 

333 For more information see https://redis.io/commands/exists 

334 """ 

335 return self._split_command_across_slots("EXISTS", *keys) 

336 

337 @overload 

338 def delete(self: SyncClientProtocol, *keys: KeyT) -> int: ... 

339 

340 @overload 

341 def delete(self: AsyncClientProtocol, *keys: KeyT) -> Awaitable[int]: ... 

342 

343 def delete(self, *keys: KeyT) -> int | Awaitable[int]: 

344 """ 

345 Deletes the given keys in the cluster. 

346 The keys are first split up into slots 

347 and then an DEL command is sent for every slot 

348 

349 Non-existent keys are ignored. 

350 Returns the number of keys that were deleted. 

351 

352 For more information see https://redis.io/commands/del 

353 """ 

354 return self._split_command_across_slots("DEL", *keys) 

355 

356 @overload 

357 def touch(self: SyncClientProtocol, *keys: KeyT) -> int: ... 

358 

359 @overload 

360 def touch(self: AsyncClientProtocol, *keys: KeyT) -> Awaitable[int]: ... 

361 

362 def touch(self, *keys: KeyT) -> int | Awaitable[int]: 

363 """ 

364 Updates the last access time of given keys across the 

365 cluster. 

366 

367 The keys are first split up into slots 

368 and then an TOUCH command is sent for every slot 

369 

370 Non-existent keys are ignored. 

371 Returns the number of keys that were touched. 

372 

373 For more information see https://redis.io/commands/touch 

374 """ 

375 return self._split_command_across_slots("TOUCH", *keys) 

376 

377 @overload 

378 def unlink(self: SyncClientProtocol, *keys: KeyT) -> int: ... 

379 

380 @overload 

381 def unlink(self: AsyncClientProtocol, *keys: KeyT) -> Awaitable[int]: ... 

382 

383 def unlink(self, *keys: KeyT) -> int | Awaitable[int]: 

384 """ 

385 Remove the specified keys in a different thread. 

386 

387 The keys are first split up into slots 

388 and then an TOUCH command is sent for every slot 

389 

390 Non-existent keys are ignored. 

391 Returns the number of keys that were unlinked. 

392 

393 For more information see https://redis.io/commands/unlink 

394 """ 

395 return self._split_command_across_slots("UNLINK", *keys) 

396 

397 

398class AsyncClusterMultiKeyCommands(ClusterMultiKeyCommands): 

399 """ 

400 A class containing commands that handle more than one key 

401 """ 

402 

403 async def mget_nonatomic(self, keys: KeysT, *args: KeyT) -> List[Any | None]: 

404 """ 

405 Splits the keys into different slots and then calls MGET 

406 for the keys of every slot. This operation will not be atomic 

407 if keys belong to more than one slot. 

408 

409 Returns a list of values ordered identically to ``keys`` 

410 

411 For more information see https://redis.io/commands/mget 

412 """ 

413 

414 # Concatenate all keys into a list 

415 keys = list_or_args(keys, args) 

416 

417 # Split keys into slots 

418 slots_to_keys = self._partition_keys_by_slot(keys) 

419 

420 # Execute commands using a pipeline 

421 res = await self._execute_pipeline_by_slot("MGET", slots_to_keys) 

422 

423 # Reorder keys in the order the user provided & return 

424 return self._reorder_keys_by_command(keys, slots_to_keys, res) 

425 

426 async def mset_nonatomic(self, mapping: Mapping[AnyKeyT, EncodableT]) -> List[bool]: 

427 """ 

428 Sets key/values based on a mapping. Mapping is a dictionary of 

429 key/value pairs. Both keys and values should be strings or types that 

430 can be cast to a string via str(). 

431 

432 Splits the keys into different slots and then calls MSET 

433 for the keys of every slot. This operation will not be atomic 

434 if keys belong to more than one slot. 

435 

436 For more information see https://redis.io/commands/mset 

437 """ 

438 

439 # Partition the keys by slot 

440 slots_to_pairs = self._partition_pairs_by_slot(mapping) 

441 

442 # Execute commands using a pipeline & return list of replies 

443 return await self._execute_pipeline_by_slot("MSET", slots_to_pairs) 

444 

445 async def _split_command_across_slots(self, command: str, *keys: KeyT) -> int: 

446 """ 

447 Runs the given command once for the keys 

448 of each slot. Returns the sum of the return values. 

449 """ 

450 

451 # Partition the keys by slot 

452 slots_to_keys = self._partition_keys_by_slot(keys) 

453 

454 # Sum up the reply from each command 

455 return sum(await self._execute_pipeline_by_slot(command, slots_to_keys)) 

456 

457 async def _is_replica_safe(self, command_name: str) -> bool: 

458 warnings.warn( 

459 "Using READ_COMMANDS for replica-safe checks is deprecated and will be removed in a future release. " 

460 "Use an AsyncMetadataResolver instead.", 

461 DeprecationWarning, 

462 stacklevel=2, 

463 ) 

464 return isinstance(command_name, str) and command_name.upper() in READ_COMMANDS 

465 

466 async def _execute_pipeline_by_slot( 

467 self, command: str, slots_to_args: Mapping[int, Iterable[EncodableT]] 

468 ) -> List[Any]: 

469 if self._initialize: 

470 await self.initialize() 

471 replica_safe = ( 

472 self.read_from_replicas or self.load_balancing_strategy is not None 

473 ) and await self._is_replica_safe(command) 

474 pipe = self.pipeline() 

475 [ 

476 pipe.execute_command( 

477 command, 

478 *slot_args, 

479 target_nodes=[ 

480 self.nodes_manager.get_node_from_slot( 

481 slot, 

482 replica_safe, 

483 self.load_balancing_strategy if replica_safe else None, 

484 ) 

485 ], 

486 ) 

487 for slot, slot_args in slots_to_args.items() 

488 ] 

489 return await pipe.execute() 

490 

491 

492class ClusterManagementCommands(ManagementCommands): 

493 """ 

494 A class for Redis Cluster management commands 

495 

496 The class inherits from Redis's core ManagementCommands class and do the 

497 required adjustments to work with cluster mode 

498 """ 

499 

500 def slaveof(self, *args, **kwargs) -> NoReturn: 

501 """ 

502 Make the server a replica of another instance, or promote it as master. 

503 

504 For more information see https://redis.io/commands/slaveof 

505 """ 

506 raise RedisClusterException("SLAVEOF is not supported in cluster mode") 

507 

508 def replicaof(self, *args, **kwargs) -> NoReturn: 

509 """ 

510 Make the server a replica of another instance, or promote it as master. 

511 

512 For more information see https://redis.io/commands/replicaof 

513 """ 

514 raise RedisClusterException("REPLICAOF is not supported in cluster mode") 

515 

516 def swapdb(self, *args, **kwargs) -> NoReturn: 

517 """ 

518 Swaps two Redis databases. 

519 

520 For more information see https://redis.io/commands/swapdb 

521 """ 

522 raise RedisClusterException("SWAPDB is not supported in cluster mode") 

523 

524 @overload 

525 def cluster_myid( 

526 self: SyncClientProtocol, target_node: "TargetNodesT" 

527 ) -> bytes | str: ... 

528 

529 @overload 

530 def cluster_myid( 

531 self: AsyncClientProtocol, target_node: "TargetNodesT" 

532 ) -> Awaitable[bytes | str]: ... 

533 

534 def cluster_myid(self, target_node: "TargetNodesT") -> (bytes | str) | Awaitable[ 

535 bytes | str 

536 ]: 

537 """ 

538 Returns the node's id. 

539 

540 :target_node: 'ClusterNode' 

541 The node to execute the command on 

542 

543 For more information check https://redis.io/commands/cluster-myid/ 

544 """ 

545 return self.execute_command("CLUSTER MYID", target_nodes=target_node) 

546 

547 @overload 

548 def cluster_addslots( 

549 self: SyncClientProtocol, target_node: "TargetNodesT", *slots: EncodableT 

550 ) -> bool: ... 

551 

552 @overload 

553 def cluster_addslots( 

554 self: AsyncClientProtocol, target_node: "TargetNodesT", *slots: EncodableT 

555 ) -> Awaitable[bool]: ... 

556 

557 def cluster_addslots( 

558 self, target_node: "TargetNodesT", *slots: EncodableT 

559 ) -> bool | Awaitable[bool]: 

560 """ 

561 Assign new hash slots to receiving node. Sends to specified node. 

562 

563 :target_node: 'ClusterNode' 

564 The node to execute the command on 

565 

566 For more information see https://redis.io/commands/cluster-addslots 

567 """ 

568 return self.execute_command( 

569 "CLUSTER ADDSLOTS", *slots, target_nodes=target_node 

570 ) 

571 

572 @overload 

573 def cluster_addslotsrange( 

574 self: SyncClientProtocol, target_node: "TargetNodesT", *slots: EncodableT 

575 ) -> bool: ... 

576 

577 @overload 

578 def cluster_addslotsrange( 

579 self: AsyncClientProtocol, target_node: "TargetNodesT", *slots: EncodableT 

580 ) -> Awaitable[bool]: ... 

581 

582 def cluster_addslotsrange( 

583 self, target_node: "TargetNodesT", *slots: EncodableT 

584 ) -> bool | Awaitable[bool]: 

585 """ 

586 Similar to the CLUSTER ADDSLOTS command. 

587 The difference between the two commands is that ADDSLOTS takes a list of slots 

588 to assign to the node, while ADDSLOTSRANGE takes a list of slot ranges 

589 (specified by start and end slots) to assign to the node. 

590 

591 :target_node: 'ClusterNode' 

592 The node to execute the command on 

593 

594 For more information see https://redis.io/commands/cluster-addslotsrange 

595 """ 

596 return self.execute_command( 

597 "CLUSTER ADDSLOTSRANGE", *slots, target_nodes=target_node 

598 ) 

599 

600 @overload 

601 def cluster_countkeysinslot(self: SyncClientProtocol, slot_id: int) -> int: ... 

602 

603 @overload 

604 def cluster_countkeysinslot( 

605 self: AsyncClientProtocol, slot_id: int 

606 ) -> Awaitable[int]: ... 

607 

608 def cluster_countkeysinslot(self, slot_id: int) -> int | Awaitable[int]: 

609 """ 

610 Return the number of local keys in the specified hash slot 

611 Send to node based on specified slot_id 

612 

613 For more information see https://redis.io/commands/cluster-countkeysinslot 

614 """ 

615 return self.execute_command("CLUSTER COUNTKEYSINSLOT", slot_id) 

616 

617 @overload 

618 def cluster_count_failure_report(self: SyncClientProtocol, node_id: str) -> int: ... 

619 

620 @overload 

621 def cluster_count_failure_report( 

622 self: AsyncClientProtocol, node_id: str 

623 ) -> Awaitable[int]: ... 

624 

625 def cluster_count_failure_report(self, node_id: str) -> int | Awaitable[int]: 

626 """ 

627 Return the number of failure reports active for a given node 

628 Sends to a random node 

629 

630 For more information see https://redis.io/commands/cluster-count-failure-reports 

631 """ 

632 return self.execute_command("CLUSTER COUNT-FAILURE-REPORTS", node_id) 

633 

634 def cluster_delslots(self, *slots: EncodableT) -> List[bool]: 

635 """ 

636 Set hash slots as unbound in the cluster. 

637 It determines by it self what node the slot is in and sends it there 

638 

639 Returns a list of the results for each processed slot. 

640 

641 For more information see https://redis.io/commands/cluster-delslots 

642 """ 

643 return [self.execute_command("CLUSTER DELSLOTS", slot) for slot in slots] 

644 

645 @overload 

646 def cluster_delslotsrange(self: SyncClientProtocol, *slots: EncodableT) -> bool: ... 

647 

648 @overload 

649 def cluster_delslotsrange( 

650 self: AsyncClientProtocol, *slots: EncodableT 

651 ) -> Awaitable[bool]: ... 

652 

653 def cluster_delslotsrange(self, *slots: EncodableT) -> bool | Awaitable[bool]: 

654 """ 

655 Similar to the CLUSTER DELSLOTS command. 

656 The difference is that CLUSTER DELSLOTS takes a list of hash slots to remove 

657 from the node, while CLUSTER DELSLOTSRANGE takes a list of slot ranges to remove 

658 from the node. 

659 

660 For more information see https://redis.io/commands/cluster-delslotsrange 

661 """ 

662 return self.execute_command("CLUSTER DELSLOTSRANGE", *slots) 

663 

664 @overload 

665 def cluster_failover( 

666 self: SyncClientProtocol, 

667 target_node: "TargetNodesT", 

668 option: str | None = None, 

669 ) -> bool: ... 

670 

671 @overload 

672 def cluster_failover( 

673 self: AsyncClientProtocol, 

674 target_node: "TargetNodesT", 

675 option: str | None = None, 

676 ) -> Awaitable[bool]: ... 

677 

678 def cluster_failover( 

679 self, target_node: "TargetNodesT", option: str | None = None 

680 ) -> bool | Awaitable[bool]: 

681 """ 

682 Forces a slave to perform a manual failover of its master 

683 Sends to specified node 

684 

685 :target_node: 'ClusterNode' 

686 The node to execute the command on 

687 

688 For more information see https://redis.io/commands/cluster-failover 

689 """ 

690 if option: 

691 if option.upper() not in ["FORCE", "TAKEOVER"]: 

692 raise RedisError( 

693 f"Invalid option for CLUSTER FAILOVER command: {option}" 

694 ) 

695 else: 

696 return self.execute_command( 

697 "CLUSTER FAILOVER", option, target_nodes=target_node 

698 ) 

699 else: 

700 return self.execute_command("CLUSTER FAILOVER", target_nodes=target_node) 

701 

702 @overload 

703 def cluster_info( 

704 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None 

705 ) -> dict[str, str]: ... 

706 

707 @overload 

708 def cluster_info( 

709 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None 

710 ) -> Awaitable[dict[str, str]]: ... 

711 

712 def cluster_info( 

713 self, target_nodes: "TargetNodesT" | None = None 

714 ) -> dict[str, str] | Awaitable[dict[str, str]]: 

715 """ 

716 Provides info about Redis Cluster node state. 

717 The command will be sent to a random node in the cluster if no target 

718 node is specified. 

719 

720 For more information see https://redis.io/commands/cluster-info 

721 """ 

722 return self.execute_command("CLUSTER INFO", target_nodes=target_nodes) 

723 

724 @overload 

725 def cluster_keyslot(self: SyncClientProtocol, key: str) -> int: ... 

726 

727 @overload 

728 def cluster_keyslot(self: AsyncClientProtocol, key: str) -> Awaitable[int]: ... 

729 

730 def cluster_keyslot(self, key: str) -> int | Awaitable[int]: 

731 """ 

732 Returns the hash slot of the specified key 

733 Sends to random node in the cluster 

734 

735 For more information see https://redis.io/commands/cluster-keyslot 

736 """ 

737 return self.execute_command("CLUSTER KEYSLOT", key) 

738 

739 @overload 

740 def cluster_meet( 

741 self: SyncClientProtocol, 

742 host: str, 

743 port: int, 

744 target_nodes: "TargetNodesT" | None = None, 

745 ) -> bool: ... 

746 

747 @overload 

748 def cluster_meet( 

749 self: AsyncClientProtocol, 

750 host: str, 

751 port: int, 

752 target_nodes: "TargetNodesT" | None = None, 

753 ) -> Awaitable[bool]: ... 

754 

755 def cluster_meet( 

756 self, host: str, port: int, target_nodes: "TargetNodesT" | None = None 

757 ) -> bool | Awaitable[bool]: 

758 """ 

759 Force a node cluster to handshake with another node. 

760 Sends to specified node. 

761 

762 For more information see https://redis.io/commands/cluster-meet 

763 """ 

764 return self.execute_command( 

765 "CLUSTER MEET", host, port, target_nodes=target_nodes 

766 ) 

767 

768 @overload 

769 def cluster_nodes(self: SyncClientProtocol) -> dict[str, ClusterNodeDetail]: ... 

770 

771 @overload 

772 def cluster_nodes( 

773 self: AsyncClientProtocol, 

774 ) -> Awaitable[dict[str, ClusterNodeDetail]]: ... 

775 

776 def cluster_nodes( 

777 self, 

778 ) -> dict[str, ClusterNodeDetail] | Awaitable[dict[str, ClusterNodeDetail]]: 

779 """ 

780 Get Cluster config for the node. 

781 Sends to random node in the cluster 

782 

783 For more information see https://redis.io/commands/cluster-nodes 

784 """ 

785 return self.execute_command("CLUSTER NODES") 

786 

787 @overload 

788 def cluster_replicate( 

789 self: SyncClientProtocol, target_nodes: "TargetNodesT", node_id: str 

790 ) -> bool: ... 

791 

792 @overload 

793 def cluster_replicate( 

794 self: AsyncClientProtocol, target_nodes: "TargetNodesT", node_id: str 

795 ) -> Awaitable[bool]: ... 

796 

797 def cluster_replicate( 

798 self, target_nodes: "TargetNodesT", node_id: str 

799 ) -> bool | Awaitable[bool]: 

800 """ 

801 Reconfigure a node as a slave of the specified master node 

802 

803 For more information see https://redis.io/commands/cluster-replicate 

804 """ 

805 return self.execute_command( 

806 "CLUSTER REPLICATE", node_id, target_nodes=target_nodes 

807 ) 

808 

809 @overload 

810 def cluster_reset( 

811 self: SyncClientProtocol, 

812 soft: bool = True, 

813 target_nodes: "TargetNodesT" | None = None, 

814 ) -> bool: ... 

815 

816 @overload 

817 def cluster_reset( 

818 self: AsyncClientProtocol, 

819 soft: bool = True, 

820 target_nodes: "TargetNodesT" | None = None, 

821 ) -> Awaitable[bool]: ... 

822 

823 def cluster_reset( 

824 self, soft: bool = True, target_nodes: "TargetNodesT" | None = None 

825 ) -> bool | Awaitable[bool]: 

826 """ 

827 Reset a Redis Cluster node 

828 

829 If 'soft' is True then it will send 'SOFT' argument 

830 If 'soft' is False then it will send 'HARD' argument 

831 

832 For more information see https://redis.io/commands/cluster-reset 

833 """ 

834 return self.execute_command( 

835 "CLUSTER RESET", b"SOFT" if soft else b"HARD", target_nodes=target_nodes 

836 ) 

837 

838 @overload 

839 def cluster_save_config( 

840 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None 

841 ) -> bool: ... 

842 

843 @overload 

844 def cluster_save_config( 

845 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None 

846 ) -> Awaitable[bool]: ... 

847 

848 def cluster_save_config( 

849 self, target_nodes: "TargetNodesT" | None = None 

850 ) -> bool | Awaitable[bool]: 

851 """ 

852 Forces the node to save cluster state on disk 

853 

854 For more information see https://redis.io/commands/cluster-saveconfig 

855 """ 

856 return self.execute_command("CLUSTER SAVECONFIG", target_nodes=target_nodes) 

857 

858 @overload 

859 def cluster_get_keys_in_slot( 

860 self: SyncClientProtocol, slot: int, num_keys: int 

861 ) -> list[bytes | str]: ... 

862 

863 @overload 

864 def cluster_get_keys_in_slot( 

865 self: AsyncClientProtocol, slot: int, num_keys: int 

866 ) -> Awaitable[list[bytes | str]]: ... 

867 

868 def cluster_get_keys_in_slot( 

869 self, slot: int, num_keys: int 

870 ) -> list[bytes | str] | Awaitable[list[bytes | str]]: 

871 """ 

872 Returns the number of keys in the specified cluster slot 

873 

874 For more information see https://redis.io/commands/cluster-getkeysinslot 

875 """ 

876 return self.execute_command("CLUSTER GETKEYSINSLOT", slot, num_keys) 

877 

878 @overload 

879 def cluster_set_config_epoch( 

880 self: SyncClientProtocol, epoch: int, target_nodes: "TargetNodesT" | None = None 

881 ) -> bool: ... 

882 

883 @overload 

884 def cluster_set_config_epoch( 

885 self: AsyncClientProtocol, 

886 epoch: int, 

887 target_nodes: "TargetNodesT" | None = None, 

888 ) -> Awaitable[bool]: ... 

889 

890 def cluster_set_config_epoch( 

891 self, epoch: int, target_nodes: "TargetNodesT" | None = None 

892 ) -> bool | Awaitable[bool]: 

893 """ 

894 Set the configuration epoch in a new node 

895 

896 For more information see https://redis.io/commands/cluster-set-config-epoch 

897 """ 

898 return self.execute_command( 

899 "CLUSTER SET-CONFIG-EPOCH", epoch, target_nodes=target_nodes 

900 ) 

901 

902 @overload 

903 def cluster_setslot( 

904 self: SyncClientProtocol, 

905 target_node: "TargetNodesT", 

906 node_id: str, 

907 slot_id: int, 

908 state: str, 

909 ) -> bool: ... 

910 

911 @overload 

912 def cluster_setslot( 

913 self: AsyncClientProtocol, 

914 target_node: "TargetNodesT", 

915 node_id: str, 

916 slot_id: int, 

917 state: str, 

918 ) -> Awaitable[bool]: ... 

919 

920 def cluster_setslot( 

921 self, target_node: "TargetNodesT", node_id: str, slot_id: int, state: str 

922 ) -> bool | Awaitable[bool]: 

923 """ 

924 Bind an hash slot to a specific node 

925 

926 :target_node: 'ClusterNode' 

927 The node to execute the command on 

928 

929 For more information see https://redis.io/commands/cluster-setslot 

930 """ 

931 if state.upper() in ("IMPORTING", "NODE", "MIGRATING"): 

932 return self.execute_command( 

933 "CLUSTER SETSLOT", slot_id, state, node_id, target_nodes=target_node 

934 ) 

935 elif state.upper() == "STABLE": 

936 raise RedisError('For "stable" state please use cluster_setslot_stable') 

937 else: 

938 raise RedisError(f"Invalid slot state: {state}") 

939 

940 @overload 

941 def cluster_setslot_stable(self: SyncClientProtocol, slot_id: int) -> bool: ... 

942 

943 @overload 

944 def cluster_setslot_stable( 

945 self: AsyncClientProtocol, slot_id: int 

946 ) -> Awaitable[bool]: ... 

947 

948 def cluster_setslot_stable(self, slot_id: int) -> bool | Awaitable[bool]: 

949 """ 

950 Clears migrating / importing state from the slot. 

951 It determines by it self what node the slot is in and sends it there. 

952 

953 For more information see https://redis.io/commands/cluster-setslot 

954 """ 

955 return self.execute_command("CLUSTER SETSLOT", slot_id, "STABLE") 

956 

957 @overload 

958 def cluster_replicas( 

959 self: SyncClientProtocol, 

960 node_id: str, 

961 target_nodes: "TargetNodesT" | None = None, 

962 ) -> dict[str, ClusterNodeDetail]: ... 

963 

964 @overload 

965 def cluster_replicas( 

966 self: AsyncClientProtocol, 

967 node_id: str, 

968 target_nodes: "TargetNodesT" | None = None, 

969 ) -> Awaitable[dict[str, ClusterNodeDetail]]: ... 

970 

971 def cluster_replicas( 

972 self, node_id: str, target_nodes: "TargetNodesT" | None = None 

973 ) -> dict[str, ClusterNodeDetail] | Awaitable[dict[str, ClusterNodeDetail]]: 

974 """ 

975 Provides a list of replica nodes replicating from the specified primary 

976 target node. 

977 

978 For more information see https://redis.io/commands/cluster-replicas 

979 """ 

980 return self.execute_command( 

981 "CLUSTER REPLICAS", node_id, target_nodes=target_nodes 

982 ) 

983 

984 @overload 

985 def cluster_slots( 

986 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None 

987 ) -> list[Any]: ... 

988 

989 @overload 

990 def cluster_slots( 

991 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None 

992 ) -> Awaitable[list[Any]]: ... 

993 

994 def cluster_slots( 

995 self, target_nodes: "TargetNodesT" | None = None 

996 ) -> list[Any] | Awaitable[list[Any]]: 

997 """ 

998 Get array of Cluster slot to node mappings 

999 

1000 For more information see https://redis.io/commands/cluster-slots 

1001 """ 

1002 return self.execute_command("CLUSTER SLOTS", target_nodes=target_nodes) 

1003 

1004 @overload 

1005 def cluster_shards( 

1006 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None 

1007 ) -> ClusterShardsResponse: ... 

1008 

1009 @overload 

1010 def cluster_shards( 

1011 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None 

1012 ) -> Awaitable[ClusterShardsResponse]: ... 

1013 

1014 def cluster_shards( 

1015 self, target_nodes: "TargetNodesT" | None = None 

1016 ) -> ClusterShardsResponse | Awaitable[ClusterShardsResponse]: 

1017 """ 

1018 Returns details about the shards of the cluster. 

1019 

1020 For more information see https://redis.io/commands/cluster-shards 

1021 """ 

1022 return self.execute_command("CLUSTER SHARDS", target_nodes=target_nodes) 

1023 

1024 @overload 

1025 def cluster_myshardid( 

1026 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None 

1027 ) -> bytes | str: ... 

1028 

1029 @overload 

1030 def cluster_myshardid( 

1031 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None 

1032 ) -> Awaitable[bytes | str]: ... 

1033 

1034 def cluster_myshardid(self, target_nodes: "TargetNodesT" | None = None) -> ( 

1035 bytes | str 

1036 ) | Awaitable[bytes | str]: 

1037 """ 

1038 Returns the shard ID of the node. 

1039 

1040 For more information see https://redis.io/commands/cluster-myshardid/ 

1041 """ 

1042 return self.execute_command("CLUSTER MYSHARDID", target_nodes=target_nodes) 

1043 

1044 @overload 

1045 def cluster_links( 

1046 self: SyncClientProtocol, target_node: "TargetNodesT" 

1047 ) -> ClusterLinksResponse: ... 

1048 

1049 @overload 

1050 def cluster_links( 

1051 self: AsyncClientProtocol, target_node: "TargetNodesT" 

1052 ) -> Awaitable[ClusterLinksResponse]: ... 

1053 

1054 def cluster_links( 

1055 self, target_node: "TargetNodesT" 

1056 ) -> ClusterLinksResponse | Awaitable[ClusterLinksResponse]: 

1057 """ 

1058 Each node in a Redis Cluster maintains a pair of long-lived TCP link with each 

1059 peer in the cluster: One for sending outbound messages towards the peer and one 

1060 for receiving inbound messages from the peer. 

1061 

1062 This command outputs information of all such peer links as an array. 

1063 

1064 For more information see https://redis.io/commands/cluster-links 

1065 """ 

1066 return self.execute_command("CLUSTER LINKS", target_nodes=target_node) 

1067 

1068 def cluster_flushslots(self, target_nodes: "TargetNodesT" | None = None) -> None: 

1069 raise NotImplementedError( 

1070 "CLUSTER FLUSHSLOTS is intentionally not implemented in the client." 

1071 ) 

1072 

1073 def cluster_bumpepoch(self, target_nodes: "TargetNodesT" | None = None) -> None: 

1074 raise NotImplementedError( 

1075 "CLUSTER BUMPEPOCH is intentionally not implemented in the client." 

1076 ) 

1077 

1078 def readonly(self, target_nodes: "TargetNodesT" | None = None) -> ResponseT: 

1079 """ 

1080 Enables read queries. 

1081 The command will be sent to the default cluster node if target_nodes is 

1082 not specified. 

1083 

1084 For more information see https://redis.io/commands/readonly 

1085 """ 

1086 if target_nodes == "replicas" or target_nodes == "all": 

1087 # read_from_replicas will only be enabled if the READONLY command 

1088 # is sent to all replicas 

1089 self.read_from_replicas = True 

1090 return self.execute_command("READONLY", target_nodes=target_nodes) 

1091 

1092 def readwrite(self, target_nodes: "TargetNodesT" | None = None) -> ResponseT: 

1093 """ 

1094 Disables read queries. 

1095 The command will be sent to the default cluster node if target_nodes is 

1096 not specified. 

1097 

1098 For more information see https://redis.io/commands/readwrite 

1099 """ 

1100 # Reset read from replicas flag 

1101 self.read_from_replicas = False 

1102 return self.execute_command("READWRITE", target_nodes=target_nodes) 

1103 

1104 @deprecated_function( 

1105 version="7.2.0", 

1106 reason="Use client-side caching feature instead.", 

1107 ) 

1108 def client_tracking_on( 

1109 self, 

1110 clientid: int | None = None, 

1111 prefix: Sequence[KeyT] = [], 

1112 bcast: bool = False, 

1113 optin: bool = False, 

1114 optout: bool = False, 

1115 noloop: bool = False, 

1116 target_nodes: "TargetNodesT" | None = "all", 

1117 ) -> ResponseT: 

1118 """ 

1119 Enables the tracking feature of the Redis server, that is used 

1120 for server assisted client side caching. 

1121 

1122 When clientid is provided - in target_nodes only the node that owns the 

1123 connection with this id should be provided. 

1124 When clientid is not provided - target_nodes can be any node. 

1125 

1126 For more information see https://redis.io/commands/client-tracking 

1127 """ 

1128 return self.client_tracking( 

1129 True, 

1130 clientid, 

1131 prefix, 

1132 bcast, 

1133 optin, 

1134 optout, 

1135 noloop, 

1136 target_nodes=target_nodes, 

1137 ) 

1138 

1139 @deprecated_function( 

1140 version="7.2.0", 

1141 reason="Use client-side caching feature instead.", 

1142 ) 

1143 def client_tracking_off( 

1144 self, 

1145 clientid: int | None = None, 

1146 prefix: Sequence[KeyT] = [], 

1147 bcast: bool = False, 

1148 optin: bool = False, 

1149 optout: bool = False, 

1150 noloop: bool = False, 

1151 target_nodes: "TargetNodesT" | None = "all", 

1152 ) -> ResponseT: 

1153 """ 

1154 Disables the tracking feature of the Redis server, that is used 

1155 for server assisted client side caching. 

1156 

1157 When clientid is provided - in target_nodes only the node that owns the 

1158 connection with this id should be provided. 

1159 When clientid is not provided - target_nodes can be any node. 

1160 

1161 For more information see https://redis.io/commands/client-tracking 

1162 """ 

1163 return self.client_tracking( 

1164 False, 

1165 clientid, 

1166 prefix, 

1167 bcast, 

1168 optin, 

1169 optout, 

1170 noloop, 

1171 target_nodes=target_nodes, 

1172 ) 

1173 

1174 def hotkeys_start( 

1175 self, 

1176 metrics: List[HotkeysMetricsTypes], 

1177 count: int | None = None, 

1178 duration: int | None = None, 

1179 sample_ratio: int | None = None, 

1180 slots: List[int] | None = None, 

1181 **kwargs, 

1182 ) -> str | bytes: 

1183 """ 

1184 Cluster client does not support hotkeys command. Please use the non-cluster client. 

1185 

1186 For more information see https://redis.io/commands/hotkeys-start 

1187 """ 

1188 raise NotImplementedError( 

1189 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client." 

1190 ) 

1191 

1192 def hotkeys_stop(self, **kwargs) -> str | bytes: 

1193 """ 

1194 Cluster client does not support hotkeys command. Please use the non-cluster client. 

1195 

1196 For more information see https://redis.io/commands/hotkeys-stop 

1197 """ 

1198 raise NotImplementedError( 

1199 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client." 

1200 ) 

1201 

1202 def hotkeys_reset(self, **kwargs) -> str | bytes: 

1203 """ 

1204 Cluster client does not support hotkeys command. Please use the non-cluster client. 

1205 

1206 For more information see https://redis.io/commands/hotkeys-reset 

1207 """ 

1208 raise NotImplementedError( 

1209 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client." 

1210 ) 

1211 

1212 def hotkeys_get(self, **kwargs) -> list[dict[str | bytes, Any]]: 

1213 """ 

1214 Cluster client does not support hotkeys command. Please use the non-cluster client. 

1215 

1216 For more information see https://redis.io/commands/hotkeys-get 

1217 """ 

1218 raise NotImplementedError( 

1219 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client." 

1220 ) 

1221 

1222 

1223class AsyncClusterManagementCommands( 

1224 ClusterManagementCommands, AsyncManagementCommands 

1225): 

1226 """ 

1227 A class for Redis Cluster management commands 

1228 

1229 The class inherits from Redis's core ManagementCommands class and do the 

1230 required adjustments to work with cluster mode 

1231 """ 

1232 

1233 async def cluster_delslots(self, *slots: EncodableT) -> List[bool]: 

1234 """ 

1235 Set hash slots as unbound in the cluster. 

1236 It determines by it self what node the slot is in and sends it there 

1237 

1238 Returns a list of the results for each processed slot. 

1239 

1240 For more information see https://redis.io/commands/cluster-delslots 

1241 """ 

1242 return await asyncio.gather( 

1243 *( 

1244 asyncio.create_task(self.execute_command("CLUSTER DELSLOTS", slot)) 

1245 for slot in slots 

1246 ) 

1247 ) 

1248 

1249 @deprecated_function( 

1250 version="7.2.0", 

1251 reason="Use client-side caching feature instead.", 

1252 ) 

1253 async def client_tracking_on( 

1254 self, 

1255 clientid: int | None = None, 

1256 prefix: Sequence[KeyT] = [], 

1257 bcast: bool = False, 

1258 optin: bool = False, 

1259 optout: bool = False, 

1260 noloop: bool = False, 

1261 target_nodes: "TargetNodesT" | None = "all", 

1262 ) -> ResponseT: 

1263 """ 

1264 Enables the tracking feature of the Redis server, that is used 

1265 for server assisted client side caching. 

1266 

1267 When clientid is provided - in target_nodes only the node that owns the 

1268 connection with this id should be provided. 

1269 When clientid is not provided - target_nodes can be any node. 

1270 

1271 For more information see https://redis.io/commands/client-tracking 

1272 """ 

1273 return await self.client_tracking( 

1274 True, 

1275 clientid, 

1276 prefix, 

1277 bcast, 

1278 optin, 

1279 optout, 

1280 noloop, 

1281 target_nodes=target_nodes, 

1282 ) 

1283 

1284 @deprecated_function( 

1285 version="7.2.0", 

1286 reason="Use client-side caching feature instead.", 

1287 ) 

1288 async def client_tracking_off( 

1289 self, 

1290 clientid: int | None = None, 

1291 prefix: Sequence[KeyT] = [], 

1292 bcast: bool = False, 

1293 optin: bool = False, 

1294 optout: bool = False, 

1295 noloop: bool = False, 

1296 target_nodes: "TargetNodesT" | None = "all", 

1297 ) -> ResponseT: 

1298 """ 

1299 Disables the tracking feature of the Redis server, that is used 

1300 for server assisted client side caching. 

1301 

1302 When clientid is provided - in target_nodes only the node that owns the 

1303 connection with this id should be provided. 

1304 When clientid is not provided - target_nodes can be any node. 

1305 

1306 For more information see https://redis.io/commands/client-tracking 

1307 """ 

1308 return await self.client_tracking( 

1309 False, 

1310 clientid, 

1311 prefix, 

1312 bcast, 

1313 optin, 

1314 optout, 

1315 noloop, 

1316 target_nodes=target_nodes, 

1317 ) 

1318 

1319 async def hotkeys_start( 

1320 self, 

1321 metrics: List[HotkeysMetricsTypes], 

1322 count: int | None = None, 

1323 duration: int | None = None, 

1324 sample_ratio: int | None = None, 

1325 slots: List[int] | None = None, 

1326 **kwargs, 

1327 ) -> Awaitable[str | bytes]: 

1328 """ 

1329 Cluster client does not support hotkeys command. Please use the non-cluster client. 

1330 

1331 For more information see https://redis.io/commands/hotkeys-start 

1332 """ 

1333 raise NotImplementedError( 

1334 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client." 

1335 ) 

1336 

1337 async def hotkeys_stop(self, **kwargs) -> Awaitable[str | bytes]: 

1338 """ 

1339 Cluster client does not support hotkeys command. Please use the non-cluster client. 

1340 

1341 For more information see https://redis.io/commands/hotkeys-stop 

1342 """ 

1343 raise NotImplementedError( 

1344 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client." 

1345 ) 

1346 

1347 async def hotkeys_reset(self, **kwargs) -> Awaitable[str | bytes]: 

1348 """ 

1349 Cluster client does not support hotkeys command. Please use the non-cluster client. 

1350 

1351 For more information see https://redis.io/commands/hotkeys-reset 

1352 """ 

1353 raise NotImplementedError( 

1354 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client." 

1355 ) 

1356 

1357 async def hotkeys_get(self, **kwargs) -> Awaitable[list[dict[str | bytes, Any]]]: 

1358 """ 

1359 Cluster client does not support hotkeys command. Please use the non-cluster client. 

1360 

1361 For more information see https://redis.io/commands/hotkeys-get 

1362 """ 

1363 raise NotImplementedError( 

1364 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client." 

1365 ) 

1366 

1367 

1368class ClusterDataAccessCommands(DataAccessCommands): 

1369 """ 

1370 A class for Redis Cluster Data Access Commands 

1371 

1372 The class inherits from Redis's core DataAccessCommand class and do the 

1373 required adjustments to work with cluster mode 

1374 """ 

1375 

1376 @overload 

1377 def stralgo( 

1378 self: SyncClientProtocol, 

1379 algo: Literal["LCS"], 

1380 value1: KeyT, 

1381 value2: KeyT, 

1382 specific_argument: Literal["strings"] | Literal["keys"] = "strings", 

1383 len: bool = False, 

1384 idx: bool = False, 

1385 minmatchlen: int | None = None, 

1386 withmatchlen: bool = False, 

1387 **kwargs, 

1388 ) -> StralgoResponse: ... 

1389 

1390 @overload 

1391 def stralgo( 

1392 self: AsyncClientProtocol, 

1393 algo: Literal["LCS"], 

1394 value1: KeyT, 

1395 value2: KeyT, 

1396 specific_argument: Literal["strings"] | Literal["keys"] = "strings", 

1397 len: bool = False, 

1398 idx: bool = False, 

1399 minmatchlen: int | None = None, 

1400 withmatchlen: bool = False, 

1401 **kwargs, 

1402 ) -> Awaitable[StralgoResponse]: ... 

1403 

1404 def stralgo( 

1405 self, 

1406 algo: Literal["LCS"], 

1407 value1: KeyT, 

1408 value2: KeyT, 

1409 specific_argument: Literal["strings"] | Literal["keys"] = "strings", 

1410 len: bool = False, 

1411 idx: bool = False, 

1412 minmatchlen: int | None = None, 

1413 withmatchlen: bool = False, 

1414 **kwargs, 

1415 ) -> StralgoResponse | Awaitable[StralgoResponse]: 

1416 """ 

1417 Implements complex algorithms that operate on strings. 

1418 Right now the only algorithm implemented is the LCS algorithm 

1419 (longest common substring). However new algorithms could be 

1420 implemented in the future. 

1421 

1422 ``algo`` Right now must be LCS 

1423 ``value1`` and ``value2`` Can be two strings or two keys 

1424 ``specific_argument`` Specifying if the arguments to the algorithm 

1425 will be keys or strings. strings is the default. 

1426 ``len`` Returns just the len of the match. 

1427 ``idx`` Returns the match positions in each string. 

1428 ``minmatchlen`` Restrict the list of matches to the ones of a given 

1429 minimal length. Can be provided only when ``idx`` set to True. 

1430 ``withmatchlen`` Returns the matches with the len of the match. 

1431 Can be provided only when ``idx`` set to True. 

1432 

1433 For more information see https://redis.io/commands/stralgo 

1434 """ 

1435 target_nodes = kwargs.pop("target_nodes", None) 

1436 if specific_argument == "strings" and target_nodes is None: 

1437 target_nodes = "default-node" 

1438 kwargs.update({"target_nodes": target_nodes}) 

1439 return super().stralgo( 

1440 algo, 

1441 value1, 

1442 value2, 

1443 specific_argument, 

1444 len, 

1445 idx, 

1446 minmatchlen, 

1447 withmatchlen, 

1448 **kwargs, 

1449 ) 

1450 

1451 def scan_iter( 

1452 self, 

1453 match: PatternT | None = None, 

1454 count: int | None = None, 

1455 _type: str | None = None, 

1456 **kwargs, 

1457 ) -> Iterator: 

1458 # Do the first query with cursor=0 for all nodes 

1459 cursors, data = self.scan(match=match, count=count, _type=_type, **kwargs) 

1460 yield from data 

1461 

1462 cursors = {name: cursor for name, cursor in cursors.items() if cursor != 0} 

1463 if cursors: 

1464 # Get nodes by name 

1465 nodes = {name: self.get_node(node_name=name) for name in cursors.keys()} 

1466 

1467 # Iterate over each node till its cursor is 0 

1468 kwargs.pop("target_nodes", None) 

1469 while cursors: 

1470 for name, cursor in cursors.items(): 

1471 cur, data = self.scan( 

1472 cursor=cursor, 

1473 match=match, 

1474 count=count, 

1475 _type=_type, 

1476 target_nodes=nodes[name], 

1477 **kwargs, 

1478 ) 

1479 yield from data 

1480 cursors[name] = cur[name] 

1481 

1482 cursors = { 

1483 name: cursor for name, cursor in cursors.items() if cursor != 0 

1484 } 

1485 

1486 

1487class AsyncClusterDataAccessCommands( 

1488 ClusterDataAccessCommands, AsyncDataAccessCommands 

1489): 

1490 """ 

1491 A class for Redis Cluster Data Access Commands 

1492 

1493 The class inherits from Redis's core DataAccessCommand class and do the 

1494 required adjustments to work with cluster mode 

1495 """ 

1496 

1497 async def scan_iter( 

1498 self, 

1499 match: PatternT | None = None, 

1500 count: int | None = None, 

1501 _type: str | None = None, 

1502 **kwargs, 

1503 ) -> AsyncIterator: 

1504 # Do the first query with cursor=0 for all nodes 

1505 cursors, data = await self.scan(match=match, count=count, _type=_type, **kwargs) 

1506 for value in data: 

1507 yield value 

1508 

1509 cursors = {name: cursor for name, cursor in cursors.items() if cursor != 0} 

1510 if cursors: 

1511 # Get nodes by name 

1512 nodes = {name: self.get_node(node_name=name) for name in cursors.keys()} 

1513 

1514 # Iterate over each node till its cursor is 0 

1515 kwargs.pop("target_nodes", None) 

1516 while cursors: 

1517 for name, cursor in cursors.items(): 

1518 cur, data = await self.scan( 

1519 cursor=cursor, 

1520 match=match, 

1521 count=count, 

1522 _type=_type, 

1523 target_nodes=nodes[name], 

1524 **kwargs, 

1525 ) 

1526 for value in data: 

1527 yield value 

1528 cursors[name] = cur[name] 

1529 

1530 cursors = { 

1531 name: cursor for name, cursor in cursors.items() if cursor != 0 

1532 } 

1533 

1534 

1535class RedisClusterCommands( 

1536 ClusterMultiKeyCommands, 

1537 ClusterManagementCommands, 

1538 ACLCommands, 

1539 PubSubCommands, 

1540 ClusterDataAccessCommands, 

1541 ScriptCommands, 

1542 FunctionCommands, 

1543 ModuleCommands, 

1544 RedisModuleCommands, 

1545): 

1546 """ 

1547 A class for all Redis Cluster commands 

1548 

1549 For key-based commands, the target node(s) will be internally determined 

1550 by the keys' hash slot. 

1551 Non-key-based commands can be executed with the 'target_nodes' argument to 

1552 target specific nodes. By default, if target_nodes is not specified, the 

1553 command will be executed on the default cluster node. 

1554 

1555 :param target_nodes: type can be one of the following: 

1556 - nodes flag: ALL_NODES, PRIMARIES, REPLICAS, RANDOM 

1557 - 'ClusterNode' 

1558 - 'list(ClusterNodes)' 

1559 - 'dict(any:clusterNodes)' 

1560 

1561 for example: 

1562 r.cluster_info(target_nodes=RedisCluster.ALL_NODES) 

1563 """ 

1564 

1565 

1566class AsyncRedisClusterCommands( 

1567 AsyncClusterMultiKeyCommands, 

1568 AsyncClusterManagementCommands, 

1569 AsyncACLCommands, 

1570 PubSubCommands, 

1571 AsyncClusterDataAccessCommands, 

1572 AsyncScriptCommands, 

1573 AsyncFunctionCommands, 

1574 AsyncModuleCommands, 

1575 AsyncRedisModuleCommands, 

1576): 

1577 """ 

1578 A class for all Redis Cluster commands 

1579 

1580 For key-based commands, the target node(s) will be internally determined 

1581 by the keys' hash slot. 

1582 Non-key-based commands can be executed with the 'target_nodes' argument to 

1583 target specific nodes. By default, if target_nodes is not specified, the 

1584 command will be executed on the default cluster node. 

1585 

1586 :param target_nodes: type can be one of the following: 

1587 - nodes flag: ALL_NODES, PRIMARIES, REPLICAS, RANDOM 

1588 - 'ClusterNode' 

1589 - 'list(ClusterNodes)' 

1590 - 'dict(any:clusterNodes)' 

1591 

1592 for example: 

1593 r.cluster_info(target_nodes=RedisCluster.ALL_NODES) 

1594 """