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

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

367 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 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 """