The FreeRADIUS server $Id: f3670dba8951ca10eb4948feb3dc3db9423a334f $
Loading...
Searching...
No Matches
cluster_async.c
Go to the documentation of this file.
1/*
2 * This program is is free software; you can redistribute it and/or modify
3 * it under the terms of the GNU General Public License as published by
4 * the Free Software Foundation; either version 2 of the License, or (at
5 * your option) any later version.
6 *
7 * This program is distributed in the hope that it will be useful,
8 * but WITHOUT ANY WARRANTY; without even the implied warranty of
9 * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
10 * GNU General Public License for more details.
11 *
12 * You should have received a copy of the GNU General Public License
13 * along with this program; if not, write to the Free Software
14 * Foundation, Inc., 51 Franklin St, Fifth Floor, Boston, MA 02110-1301, USA
15 */
16
17/**
18 * $Id: 9f28cd72cbd682d534c702097e81a7bb384d1f90 $
19 * @file cluster_async.c
20 * @brief conf functions for interacting asynchronously with Redis cluster via Hiredis.
21 *
22 * @author Arran Cudbard-Bell (a.cudbardb@freeradius.org)
23 *
24 * @copyright 2026 Network RADIUS (legal@networkradius.com)
25 *
26 * Overview
27 * ========
28 *
29 * Read and understand this http://redis.io/topics/cluster-spec first, else the text below
30 * will not be useful.
31 *
32 * Using the cluster's public API
33 * ------------------------------
34 *
35 * The cluster requires use of a coordinator to fetch the cluster map.
36 *
37 * Any module using a Redis cluster should register a coordinator which uses
38 * the `redis` virtual server.
39 *
40 * In `mod_coord_attach`, #fr_redis_ct_map_bootstrap can be used to initiate
41 * the loading of the cluster map. Typically this should not be called if
42 * the pool start is set to zero.
43 * In that case, the first attempt to enqueue a command set will indicate
44 * that the cluster map need to be bootstrapped.
45 *
46 * At runtime the function #fr_redis_async_cmd_start is used to enqueue a set of
47 * commands, and the statis it returns should be checked with the macro
48 * REDIS_ASYNC_START_RCODE_PROCESS to initiate the cluster bootstrap or get
49 * an updated map if needed.
50 *
51 * With calling Redis using it's async API, the majority of results processing has to
52 * be done in a callback called by hiredis - the `redisReply` structure is freed
53 * after the callback is called.
54 *
55 * This callback is associated with the individual commands in a Redis command set
56 * as they are added to the command set with the fr_redis_command_*_add functions.
57 *
58 * Structures
59 * ----------
60 *
61 * This code maintains a series structures for efficient lookup and lockless operations.
62 *
63 * The important ones are:
64 * - An array of #fr_redis_cluster_node_t. These are pre-allocated on startup and are
65 * never added to, or removed from.
66 * - An #fr_fifo_t. This contains the queue of nodes that may be re-used.
67 * - An #fr_rb_tree_t. This contains a tree of nodes which are active. The tree is built on IP
68 * address and port.
69 *
70 * Each #fr_redis_cluster_node_t contains a master ID, and an array of slave IDs. The IDs are array
71 * indexes in the fr_redis_cluster_t.node array. We use 8bit unsigned integers instead of
72 * pointers to save space. Using pointers, the node[] array would need 784K, using IDs
73 * it uses 112K. Still not light on memory, but a bit more acceptable.
74 * Currently the key_slot array is shadowed by key_slot_pending, used to stage new key_slot
75 * mappings. This doubles the memory used. We may want to consider allocating key_slot_pending
76 * only during re-mappings and freeing it after.
77 *
78 * Mapping/Remapping the cluster
79 * -----------------------------
80 *
81 * On startup, and during cluster operation, a remap may be performed. A remap involves
82 * the following steps:
83 *
84 * 1. Request the cluster map from the coordinator.
85 * 2. The coordinator:
86 * a. Checks to see when it last fetched the cluster map. If it was less than 1 second ago,
87 * replies with the most recently fetched data.
88 * b. Executes the Redis 'cluster info' on all known nodes to establish which nodes believe
89 * they can see a working cluster, from the `cluster_state` response and which nodes have
90 * the most up to date representation of the cluster, from the `cluster_current_epoch`
91 * response.
92 * c. Nodes reporting the cluster is OK are issued the Redis 'cluster slots' or 'cluster shards'
93 * command depending on the Redis server version.
94 * d. Validating the result of this command. We need to do extensive validation to
95 * avoid SEGV on invalid data, due to the way libhiredis presents the result.
96 * e. Return the cluster map to the workers.
97 * 4. Determining the intersection between nodes described in the result, and those already
98 * in our #fr_rb_tree_t.
99 * 5. Creating trunk connections to nodes that were in the result, but not in the tree.
100 * 6. Mapping keyslot ranges to nodes in the key_slot_pending array.
101 * 7. Verifying there are no holes in the ranges (if there are, we roll back and error out).
102 * 8. Applying the new keyslot ranges.
103 * 9. Removing nodes no longer used by the key slots, and adding them back to the free
104 * nodes queue.
105 *
106 * #fr_redis_ct_map_get is used to request an updated map from the coordinator and
107 * #fr_redis_ct_map_update is used to process the message received from the coordinator to update
108 * the thread local copy of the cluster map.
109 *
110 * The cluster client can continue to operate, albeit inefficiently, with a stale cluster map
111 * by following '-ASK' and '-MOVE' redirects.
112 *
113 * Remaps are limited to one per second. If any operation sets the remap_needed flag, or
114 * attempts a remap directly, the remap may be skipped if one occurred recently.
115 *
116 *
117 * Processing '-ASK' and '-MOVE' redirects
118 * ---------------------------------------
119 *
120 * Resume functions which are run after an async Redis command set has completed should
121 * fetch the rcode with #fr_redis_command_set_rcode.
122 * If the rcode indicates MOVE or ASK, then #fr_redis_async_cmd_redirect should be used
123 * to re-enqueue the command set on the indicated node. In addition, if the response
124 * was MOVE, then #fr_redis_ct_map_get should be used to initiate a refresh of the cluster map.
125 *
126 * The data from '-MOVE' responses, is not used to alter the cluster map. That is only done
127 * on successful remap.
128 *
129 *
130 * Processing '-TRYAGAIN'
131 * ----------------------
132 *
133 * If the cluster is in a state of flux, a node may return '-TRYAGAIN' to indicated that we
134 * should attempt the operation again. #fr_redis_async_cmd_resend can be used to re-enqueue
135 * the command set.
136 *
137 */
138
139#include <freeradius-devel/util/debug.h>
140
141#include "config.h"
142#include "attrs.h"
143#include "cluster_async.h"
144#include "crc16.h"
145
146#ifndef WITH_TLS
147# undef HAVE_REDIS_SSL
148#endif
149
150#ifdef HAVE_REDIS_SSL
151#include <freeradius-devel/tls/strerror.h>
152#include <hiredis/hiredis_ssl.h>
153#endif
154
155#define MAX_REPLICAS 5 //!< Maximum number of replicas associated
156 //!< with a keyslot.
158 uint8_t replica[MAX_REPLICAS]; //!< Array of ids of replica nodes
159 uint8_t num_replicas; //!< Number of replica nodes
160 uint8_t master; //!< id of the master node.
161};
162
163typedef enum {
164 CLUSTER_INIT = 0, //!< Cluster has been initialised.
165 CLUSTER_MAP_FETCHING, //!< The cluster map is currently being fetched.
166 CLUSTER_READY, //!< The cluster is available to handle requests.
167 CLUSTER_FAIL, //!< The coordinator reported a failed cluster map update.
169
170/** Thread local state for a cluster
171 *
172 */
174 uint16_t cluster_id; //!< Number assigned to the cluster by coordinator.
176 trunk_conf_t const *tconf; //!< Configuration for all trunks in the cluster.
177 bool delay_start; //!< Prevent connections from spawning immediately.
178 fr_redis_conf_t const *conf; //!< Redis configuration for the cluster.
179 CONF_SECTION const *tls_cs; //!< TLS CONF_SECTION
180
181 fr_redis_trunk_active_t active; //!< Callback to run when the trunk becomes active.
182 void *active_uctx; //!< Uctx to pass to active callback.
183 bool active_oneshot; //!< Should the callback only be called once.
184 uint8_t max_redirects; //!< maximum number of redirects.
185
186#ifdef HAVE_REDIS_SSL
187 SSL_CTX *ssl_ctx; //!< SSL context.
188#endif
189
190 fr_redis_ct_node_t *node; //!< Array of nodes in this cluster.
191 fr_fifo_t *free_nodes; //!< Nodes not currently active.
192 fr_rb_tree_t *used_nodes; //!< Active nodes.
193
195
196 fr_redis_ct_state_t state; //!< State of the cluster.
197 fr_dlist_head_t pend_cmds; //!< Commands awaiting cluster map.
198 fr_dlist_head_t pend_reqs; //!< Requests awaiting cluster map.
199 fr_time_t map_updated; //!< Time the cluster last updated.
200};
201
203 fr_rb_node_t rbnode; //!< Entry in the tree of used nodes
204 char name[INET6_ADDRSTRLEN];
205 uint8_t id; //!< Array offset in the array of available nodes.
206
207 bool is_active; //!< Is this node currently active.
208 bool is_master; //!< Is this node currently a master.
209
210 fr_redis_ct_t *rtcluster; //!< Cluster this node belongs to
211 fr_redis_io_conf_t ioconf; //!< Connection config for this node.
212 fr_redis_trunk_t *trunk; //!< Trunk connection to this node.
213 fr_pair_list_t trigger_args; //!< Pairs to pass to trigger functions.
214};
215
216/** Structure for holding the state of an async redis command set.
217 *
218 */
220 request_t *request; //!< Request this command set relates to.
221 fr_redis_ct_t *rtcluster; //!< Cluster this command set is running on.
222 fr_redis_trunk_t *rtrunk; //!< Trunk the command set is currently running on.
223 fr_redis_command_set_t *cmds; //!< Command set to run.
224 uint8_t const *key; //!< Key used to identify key slot.
225 size_t key_len; //!< Length of key.
226 fr_redis_ct_key_slot_t const *key_slot; //!< Key slot identified from the command key.
227 bool read_only; //!< Should this command be run read only.
228 uint8_t replica_no; //!< Current replica number being used.
229 fr_dlist_t entry; //!< Entry in the list of commands waiting for a cluster remap.
230 fr_redis_ct_node_t *node; //!< Specific node to run command set on.
231};
232
233/** Structure to record that a request is waiting for the cluster map.
234 *
235 */
236typedef struct {
237 request_t *request; //!< The request waiting for the map.
238 fr_dlist_t entry; //!< Entry in the list of pending requests.
239 fr_redis_ct_t *rtcluster; //!< Cluster the request is waiting for.
241
242#define CONFIGURE_NODE(_node, _addr, _port) \
243do { \
244 _node->ioconf = (fr_redis_io_conf_t) { \
245 .port = _port, \
246 .database = rtcluster->conf->database, \
247 .username = rtcluster->conf->username, \
248 .password = rtcluster->conf->password, \
249 .use_tls = rtcluster->conf->use_tls, \
250 }; \
251 _node->ioconf.hostname = talloc_strdup(rtcluster, _addr); \
252 _node->ioconf.log_prefix = talloc_asprintf(rtcluster, "%s %s:%d", rtcluster->conf->log_prefix, \
253 _addr, _node->ioconf.port); \
254 if (rtcluster->conf->trunk_conf.conn_triggers) { \
255 module_trigger_args_build(rtcluster, &_node->trigger_args, NULL, \
256 &(module_trigger_args_t) { \
257 .module = rtcluster->conf->module_name, \
258 .name = rtcluster->conf->inst_name, \
259 .server = _addr, \
260 .port = _node->ioconf.port \
261 }); \
262 } \
263 _node->trunk = fr_redis_trunk_alloc(rtcluster, &_node->ioconf, &_node->trigger_args, rtcluster->active, \
264 rtcluster->active_uctx, rtcluster->active_oneshot); \
265 if (!_node->trunk) goto error; \
266} while (0)
267
268/** Resolve key to key slot index
269 *
270 * Identical to the example implementation, except it uses memchr which will
271 * be faster, and isn't so needlessly complex.
272 *
273 * @param[in] key to resolve.
274 * @param[in] key_len length of key.
275 * @return key slot index for the key.
276 */
277static uint16_t cluster_key_hash(uint8_t const *key, size_t key_len)
278{
279 uint8_t *p, *q;
280
281 p = memchr(key, '{', key_len);
282 if (!p) {
283 all:
284 return fr_crc16_xmodem(key, key_len) & (KEY_SLOTS - 1);
285 }
286
287 q = memchr(p, '}', key_len - (p - key)); /* look for } after { */
288 if (!q || (q == p + 1)) goto all; /* no } or {}, hash everything */
289
290 p++; /* skip '{' */
291
292 return fr_crc16_xmodem(p, q - p) & (KEY_SLOTS - 1); /* hash stuff between { and } */
293}
294
295/** Resolve key to key slot
296 *
297 * @param[in] rtcluster to resolve the key slot in.
298 * @param[in] request Current request (for debugging).
299 * @param[in] key to resolve.
300 * @param[in] key_len length of key.
301 * @return key slot for the key.
302 */
304 uint8_t const *key, size_t key_len)
305{
306 fr_redis_ct_key_slot_t *key_slot;
307
308 if (!key || (key_len == 0)) {
309 key_slot = &rtcluster->key_slot[(uint16_t)(fr_rand() & (KEY_SLOTS - 1))];
310 ROPTIONAL(RDEBUG2, DEBUG2, "Key rand() -> slot %zu", key_slot - rtcluster->key_slot);
311
312 return key_slot;
313 }
314
315 /*
316 * Avoid CRC16 if we're operating with one cluster node or
317 * without clustering.
318 */
319 if (fr_rb_num_elements(rtcluster->used_nodes) > 1) {
320 key_slot = &rtcluster->key_slot[cluster_key_hash(key, key_len)];
321 ROPTIONAL(RDEBUG2, DEBUG2, "Key \"%pV\" -> slot %zu",
322 fr_box_strvalue_len((char const *)key, key_len), key_slot - rtcluster->key_slot);
323
324 return key_slot;
325 }
326 ROPTIONAL(RDEBUG3, DEBUG3, "Single node available, skipping key selection");
327
328 return &rtcluster->key_slot[0];
329}
330
331/** Return the master node that would be used for a particular key slot
332 *
333 * @param[in] rtcluster To resolve key slot in.
334 * @param[in] key_slot to resolve to node.
335 * @return
336 * - The current master node.
337 * - NULL if no master node is currently assigned to a particular key slot.
338 */
340 fr_redis_ct_key_slot_t const *key_slot)
341{
342 return &rtcluster->node[key_slot->master];
343}
344
345/** Return the replica node that would be used for a particular key slot
346 *
347 * @param[in] rtcluster To resolve key slot in.
348 * @param[in] key_slot To resolve to node.
349 * @param[in] replica_num 0..n.
350 * @return
351 * - A replica node.
352 * - NULL if no replica node is assigned, or is at the specific key slot.
353 *
354 */
356 fr_redis_ct_key_slot_t const *key_slot, uint8_t replica_num)
357{
358 if (replica_num >= key_slot->num_replicas) return NULL; /* No replica available */
359
360 return &rtcluster->node[key_slot->replica[replica_num]];
361}
362
363/** Return the ipaddr of a particular node
364 *
365 * @param[in] node to get ip address from.
366 * @return
367 * - IP address of node
368 * - NULL on failure (node is NULL).
369 */
370char const * fr_redis_ct_ipaddr(fr_redis_ct_node_t const *node)
371{
372 if (!node) return NULL;
373
374 return node->ioconf.hostname;
375}
376
377/** Return the port of a particular node
378 *
379 * @param[out] out Port of the node.
380 * @param[in] node to get ip address from.
381 * @return
382 * - 0 on success.
383 * - -1 on failure (node is NULL).
384 */
386{
387 if (!node) return -1;
388
389 *out = node->ioconf.port;
390
391 return 0;
392}
393
394/** Enqueue a command set on a node identified by the key.
395 *
396 */
398{
399 fr_redis_ct_t *rtcluster = cmd->rtcluster;
401 bool dst_unavail = false;
402
403 if (likely(!cmd->node)) cmd->key_slot = fr_redis_ct_slot_by_key(rtcluster, cmd->request, cmd->key, cmd->key_len);
404
405 if (unlikely(cmd->node != NULL)) {
406 cmd->rtrunk = cmd->node->trunk;
407 }
408 /*
409 * Read only commands start on the first replica, if there are any.
410 */
411 else if (cmd->read_only && cmd->key_slot->num_replicas) {
412 cmd->rtrunk = rtcluster->node[cmd->key_slot->replica[0]].trunk;
413 } else {
414 cmd->rtrunk = rtcluster->node[cmd->key_slot->master].trunk;
415 }
416
417 if (unlikely(!cmd->rtrunk)) return REDIS_ASYNC_RCODE_ERROR;
418
419again:
420 ret = redis_command_set_enqueue(cmd->rtrunk, cmd->cmds);
421
422 switch (ret) {
424 /*
425 * If one or more nodes reported failed to enqueue with
426 * destination unavailable, tell the caller that the cluster
427 * map should be updated.
428 */
430
432 if (cmd->node) return REDIS_ASYNC_RCODE_ERROR;
433 dst_unavail = true;
434 if (cmd->replica_no < cmd->key_slot->num_replicas) {
435 cmd->rtrunk = rtcluster->node[cmd->key_slot->replica[cmd->replica_no]].trunk;
436 cmd->replica_no++;
437 goto again;
438 }
439 /*
440 * Read only commands can also try the master node.
441 * Non-read only first tried the master.
442 */
443 if (cmd->read_only && (cmd->rtrunk != rtcluster->node[cmd->key_slot->master].trunk)) {
444 cmd->rtrunk = rtcluster->node[cmd->key_slot->master].trunk;
445 goto again;
446 }
448
449 default:
451 }
452
453}
454
456{
457 if (!fr_dlist_entry_in_list(&cmd->entry)) return 0;
459 return 0;
460}
461
462/** Start running a command set on an async redis cluster
463 *
464 * @param ctx to allocate tracking structure.
465 * @param request current request.
466 * @param rcode Where to write the result code.
467 * @param rtcluster to start the command set on
468 * @param key to identify the cluster slot.
469 * @param key_len Length of key.
470 * @param cmds Command set to run.
471 * @param read_only Should the command set be run on read only nodes.
472 * @param node Specific node to run the command set on.
473 * @return The async redis command
474 */
476 fr_redis_ct_t *rtcluster, uint8_t const *key, size_t key_len,
477 fr_redis_command_set_t *cmds, bool read_only, fr_redis_ct_node_t *node)
478{
480
481 MEM(cmd = talloc(ctx, fr_redis_async_cmd_t));
482
483 *cmd = (fr_redis_async_cmd_t) {
484 .request = request,
485 .rtcluster = rtcluster,
486 .cmds = cmds,
487 .read_only = read_only,
488 .key = key,
489 .key_len = key_len,
490 .node = node,
491 };
492
493 switch (rtcluster->state) {
494 case CLUSTER_INIT:
495 /*
496 * If the cluster has not bootstrapped, that must be done first.
497 */
498 fr_dlist_insert_tail(&rtcluster->pend_cmds, cmd);
499 talloc_set_destructor(cmd, _fr_redis_async_cmd_free);
501 break;
502
504 fr_dlist_insert_tail(&rtcluster->pend_cmds, cmd);
505 talloc_set_destructor(cmd, _fr_redis_async_cmd_free);
507 break;
508
509 case CLUSTER_FAIL:
510 /*
511 * The coordinator reported a failed cluster.
512 */
513 *rcode = REDIS_ASYNC_RCODE_FAIL;
514 fr_strerror_const("Cluster failed");
515 talloc_free(cmd);
516 return NULL;
517
518 default:
519 *rcode = fr_redis_async_cmd_enqueue(cmd);
520 break;
521 }
522
523 return cmd;
524}
525
526/** Cancel a Redis async command.
527 *
528 */
533
534/** Fetch the redis trunk a command is associated with.
535 *
536 */
541
542/** Fetch the cluster node a command was last sent to
543 */
545 if (!cmd->key_slot) return NULL;
546 if (cmd->replica_no == 0) return &cmd->rtcluster->node[cmd->key_slot->master];
547 return &cmd->rtcluster->node[cmd->key_slot->replica[cmd->replica_no - 1]];
548}
549
550/** Re-submit a redis async command set on a different node
551 *
552 * Using the node returned by a MOVED / ASK response.
553 * @param cmd Async command set to redirect
554 * @return fr_redis_async_rcode_t
555 */
557{
558 fr_redis_ct_node_t find, *cluster_node;
559
561
562 fr_rb_find((void **)&cluster_node, cmd->rtcluster->used_nodes, &find);
563 if (!cluster_node) {
564 ERROR("Asked to redirect to a node not in the current cluster map");
566 }
567
569 cmd->node = cluster_node;
570 return fr_redis_async_cmd_enqueue(cmd);
571}
572
573/** Re-submit a redis async command set
574 *
575 * To be used following TRYAGAIN responses
576 * @param cmd Async command set to redirect
577 * @return fr_redis_async_rcode_t
578 */
584
585/** Compare two redis nodes to check equality
586 *
587 * @param[in] one first node.
588 * @param[in] two second node.
589 * @return CMP(one, two)
590 */
591static fr_cmp_ret_t _cluster_thread_node_cmp(void const *one, void const *two)
592{
593 fr_redis_ct_node_t const *a = one;
594 fr_redis_ct_node_t const *b = two;
595 int ret;
596
597 ret = strcmp(a->ioconf.hostname, b->ioconf.hostname);
598 if (ret != 0) return CMP(ret, 0);
599
600 return CMP(a->ioconf.port, b->ioconf.port);
601}
602
603#ifdef HAVE_REDIS_SSL
604static int _redis_cluster_thread_free(fr_redis_ct_t *rtcluster)
605{
606 if (rtcluster->ssl_ctx) SSL_CTX_free(rtcluster->ssl_ctx);
607 return 0;
608}
609#endif
610
611/** Allocate per-thread, per-cluster instance
612 *
613 * This structure represents all the connections for a given thread for a given cluster.
614 * The structures holds the trunk connections to talk to each cluster member.
615 *
616 */
618 fr_redis_trunk_active_t active, void *active_uctx, bool active_oneshot)
619{
620 fr_redis_ct_t *rtcluster;
621 trunk_conf_t *our_tconf;
622 uint8_t i;
623 uint32_t s, num_nodes;
624
625 MEM(rtcluster = talloc_zero(ctx, fr_redis_ct_t));
626 *rtcluster = (fr_redis_ct_t) {
627 .el = el,
628 .conf = conf,
629 .tls_cs = tls_cs,
630 .active = active,
631 .active_uctx = active_uctx,
632 .active_oneshot = active_oneshot
633 };
634 MEM(our_tconf = talloc_memdup(rtcluster, &conf->trunk_conf, sizeof(conf->trunk_conf)));
635 our_tconf->always_writable = true;
636
637 rtcluster->tconf = our_tconf;
640
641 if (conf->max_nodes == UINT8_MAX) {
642 ERROR("%s - Maximum number of connected nodes allowed is %i", conf->log_prefix, UINT8_MAX - 1);
643 error:
644 talloc_free(rtcluster);
645 return NULL;
646 }
647
648 if (conf->max_nodes == 0) {
649 ERROR("%s - Minimum number of nodes allowed is 1", conf->log_prefix);
650 goto error;
651 }
652
653 MEM(rtcluster->node = talloc_zero_array(rtcluster, fr_redis_ct_node_t, conf->max_nodes + 1));
654 MEM(rtcluster->used_nodes = fr_rb_inline_alloc(rtcluster, fr_redis_ct_node_t, rbnode, _cluster_thread_node_cmp, NULL));
655 MEM(rtcluster->free_nodes = fr_fifo_create(rtcluster, conf->max_nodes, NULL));
656
657 /*
658 * Node id 0 is reserved, so we can detect misconfigured
659 * clusters.
660 */
661 for (i = 1; i <= conf->max_nodes; i++) {
662 rtcluster->node[i].id = i;
663 rtcluster->node[i].rtcluster = rtcluster;
664 fr_pair_list_init(&rtcluster->node[i].trigger_args);
665
666 /* Push them all into the queue */
667 fr_fifo_push(rtcluster->free_nodes, &rtcluster->node[i]);
668 }
669
670 if (conf->use_tls) {
671#ifdef HAVE_REDIS_SSL
672 fr_tls_conf_t *tls_conf;
673 if (!tls_cs) {
674 ERROR("%s - Missing TLS configuration", conf->log_prefix);
675 goto error;
676 }
677
678 tls_conf = fr_tls_conf_parse_client(tls_cs);
679 if (!tls_conf) {
680 ERROR("%s - Failed to parse TLS configuration", conf->log_prefix);
681 goto error;
682 }
683
684 rtcluster->ssl_ctx = fr_tls_ctx_alloc(tls_conf, true);
685 if (!rtcluster->ssl_ctx) {
686 ERROR("%s - Failed to allocate SSL context", conf->log_prefix);
687 goto error;
688 }
689 talloc_set_destructor(rtcluster, _redis_cluster_thread_free);
690#else
691 WARN("%s - No redis SSL support, ignoring \"use_tls = yes\"", conf->log_prefix);
692#endif
693 }
694
695 /*
696 * cmds->redirected is uint8_t, so we limit the number of redirects here.
697 */
698 if (rtcluster->conf->max_redirects > UINT8_MAX) {
699 rtcluster->max_redirects = UINT8_MAX;
700 } else {
701 rtcluster->max_redirects = rtcluster->conf->max_redirects;
702 }
703
704 if (conf->use_cluster_map) return rtcluster;
705
706 /*
707 * If we are not using a cluster map, just configure nodes from
708 * the bootstrap list and distribute them through the key slots.
709 */
710 for (s = 0; s < talloc_array_length(conf->hostname); s++) {
711 fr_redis_ct_node_t *cluster_node;
712 fr_ipaddr_t addr;
713 uint16_t port;
715
716 cluster_node = fr_fifo_pop(rtcluster->free_nodes);
717 if (!cluster_node) {
718 ERROR("Reached maximum connected nodes");
719 goto error;
720 }
721 if (fr_inet_pton_port(&addr, &port,
722 conf->hostname[s], talloc_strlen(conf->hostname[s]), AF_UNSPEC, true, true) < 0) {
723 PERROR("Failed parsing %s", conf->hostname[s]);
724 goto error;
725 }
726 if (port == 0) port = conf->port;
727 fr_inet_ntop(buff, sizeof(buff), &addr);
728 CONFIGURE_NODE(cluster_node, buff, port);
729 fr_rb_insert(rtcluster->used_nodes, cluster_node);
730 cluster_node->is_active = true;
731 cluster_node->is_master = true;
732 }
733
734 num_nodes = fr_rb_num_elements(rtcluster->used_nodes);
735 if (!num_nodes) {
736 ERROR("%s - No bootstrap servers configured", conf->log_prefix);
737 goto error;
738 }
739
740 for (s = 0; s < KEY_SLOTS; s++) rtcluster->key_slot[s].master = (s % (uint16_t) num_nodes) + 1;
741
742 rtcluster->state = CLUSTER_READY;
743
744 return rtcluster;
745}
746
748{
749 return rtcluster->el;
750}
751
753{
754 return rtcluster->tconf;
755}
756
757/** How many times a command set running on this cluster may be redirected
758 *
759 * @param[in] rtcluster to return the configured limit for.
760 * @return the value of the max_redirects configuration item.
761 */
763{
764 return rtcluster->max_redirects;
765}
766
767#ifdef HAVE_REDIS_SSL
768SSL_CTX *fr_redis_ct_ssl_ctx(fr_redis_ct_t *rtcluster)
769{
770 return rtcluster->ssl_ctx;
771}
772#endif
773
774/** Update a Redis cluster map from a pair list returned from a coordinator
775 *
776 * @param rtcluster Cluster to update
777 * @param list pairs sent by a coordinator
778 * @return
779 * - 0 om success
780 * - -1 on error
781 */
783{
784 fr_pair_t *vp, *shard = NULL, *slot, *start, *end, *node, *role, *node_ip, *node_port;
785 uint16_t i;
786 uint8_t r = 0;
787 uint8_t rollback[UINT8_MAX]; // Set of nodes to re-add to the queue on failure.
788 bool active[UINT8_MAX]; // Set of nodes active in the new cluster map.
789 bool master[UINT8_MAX]; // Master nodes.
790
791 fr_redis_ct_node_t find, *cluster_node;
792 fr_redis_ct_key_slot_t tmp_slot;
793 fr_redis_ct_key_slot_t key_slot_pending[KEY_SLOTS];
795 fr_redis_ct_pend_req_t *pend_req;
796
797#define SET_INACTIVE(_node) \
798do { \
799 (_node)->is_active = false; \
800 (_node)->is_master = false; \
801 talloc_const_free((_node)->ioconf.log_prefix); \
802 (_node)->ioconf.log_prefix = NULL; \
803 TALLOC_FREE((_node)->trunk); \
804 fr_pair_list_free(&(_node)->trigger_args); \
805 fr_rb_delete(rtcluster->used_nodes, _node); \
806 fr_fifo_push(rtcluster->free_nodes, _node); \
807} while (0)
808
809#define SET_ACTIVE(_node) \
810do { \
811 fr_rb_insert(rtcluster->used_nodes, _node); \
812 fr_fifo_pop(rtcluster->free_nodes); \
813 (_node)->is_active = true; \
814 active[(_node)->id] = true; \
815 rollback[r++] = (_node)->id; \
816} while (0)
817
819 if (unlikely(!vp)) {
820 ERROR("Missing cluster ID");
821 return -1;
822 }
823 if (rtcluster->cluster_id == 0) rtcluster->cluster_id = vp->vp_uint16;
824
825 if (rtcluster->cluster_id != vp->vp_uint16) {
826 ERROR("Got map for cluster ID %d, expected ID %d", vp->vp_uint16, rtcluster->cluster_id);
827 return -1;
828 }
829
830 DEBUG3("Updating cluster %d", rtcluster->cluster_id);
831
832 memset(&key_slot_pending, 0, sizeof(key_slot_pending));
833 memset(active, 0, sizeof(active));
834 memset(master, 0, sizeof(master));
835
836 while ((shard = fr_pair_find_by_da(list, shard, attr_redis_shard))) {
837 cluster_node = NULL;
838 memset(&tmp_slot, 0, sizeof(fr_redis_ct_key_slot_t));
839 node = NULL;
840 while ((node = fr_pair_find_by_da(&shard->vp_group, node, attr_redis_node))) {
841 role = fr_pair_find_by_da(&node->vp_group, NULL, attr_redis_node_role);
842 if (unlikely(!role)) continue;
843 if (role->vp_uint8 == 1) {
844 DEBUG3("Master node %pP", node);
845
846 node_ip = fr_pair_find_by_da(&node->vp_group, NULL, attr_redis_node_endpoint);
847 if (unlikely(!node_ip)) continue;
848 find.ioconf.hostname = node_ip->vp_strvalue;
849 node_port = fr_pair_find_by_da(&node->vp_group, NULL, attr_redis_node_port);
850 if (unlikely(!node_port)) continue;
851 find.ioconf.port = node_port->vp_uint16;
852
853 fr_rb_find((void **)&cluster_node, rtcluster->used_nodes, &find);
854 break;
855 }
856 }
857
858 if (!node) {
859 ERROR("Missing master node");
860 error:
861 for (i = 0; i < r; i++) SET_INACTIVE(&rtcluster->node[rollback[i]]);
862 return -1;
863 }
864
865 if (!cluster_node) {
866 cluster_node = fr_fifo_peek(rtcluster->free_nodes);
867 if (!cluster_node) {
868 out_of_nodes:
869 fr_strerror_const("Reached maximum connected nodes");
870 goto error;
871 }
872 CONFIGURE_NODE(cluster_node, find.ioconf.hostname, find.ioconf.port);
873 SET_ACTIVE(cluster_node);
874 } else {
875 active[cluster_node->id] = true;
876 }
877 master[cluster_node->id] = true;
878 tmp_slot.master = cluster_node->id;
879
880 node = NULL;
881 while ((node = fr_pair_find_by_da(&shard->vp_group, node, attr_redis_node))) {
882 role = fr_pair_find_by_da(&node->vp_group, NULL, attr_redis_node_role);
883 if (unlikely(!role)) continue;
884 if (tmp_slot.num_replicas >= MAX_REPLICAS) break;
885 if (role->vp_uint8 != 2) continue;
886
887 DEBUG3("Replica node %pP", node);
888 node_ip = fr_pair_find_by_da(&node->vp_group, NULL, attr_redis_node_endpoint);
889 if (unlikely(!node_ip)) continue;
890 find.ioconf.hostname = node_ip->vp_strvalue;
891 node_port = fr_pair_find_by_da(&node->vp_group, NULL, attr_redis_node_port);
892 if (unlikely(!node_port)) continue;
893 find.ioconf.port = node_port->vp_uint16;
894
895 fr_rb_find((void **)&cluster_node, rtcluster->used_nodes, &find);
896
897 if (cluster_node) {
898 tmp_slot.replica[tmp_slot.num_replicas++] = cluster_node->id;
899 active[cluster_node->id] = true;
900 continue;
901 }
902
903 cluster_node = fr_fifo_peek(rtcluster->free_nodes);
904 if (!cluster_node) goto out_of_nodes;
905
906 CONFIGURE_NODE(cluster_node, find.ioconf.hostname, find.ioconf.port);
907 tmp_slot.replica[tmp_slot.num_replicas++] = cluster_node->id;
908 SET_ACTIVE(cluster_node);
909 }
910
911 slot = NULL;
912 while ((slot = fr_pair_find_by_da(&shard->vp_group, slot, attr_redis_slot))) {
913 start = fr_pair_find_by_da(&slot->vp_group, NULL, attr_redis_slot_start);
914 if (unlikely(!start)) {
915 ERROR("Missing slot start");
916 goto error;
917 }
918 if (unlikely(start->vp_uint16 >= KEY_SLOTS)) {
919 ERROR("Value of %d for slot start greater than expected maximum %d",
920 start->vp_uint16, KEY_SLOTS);
921 goto error;
922 }
923 end = fr_pair_find_by_da(&slot->vp_group, NULL, attr_redis_slot_end);
924 if (unlikely(!end)) {
925 ERROR("Missing slot end");
926 goto error;
927 }
928 if (unlikely(end->vp_uint16 >= KEY_SLOTS)) {
929 ERROR("Value of %d for slot end greater than expected maximum %d",
930 end->vp_uint16, KEY_SLOTS);
931 goto error;
932 }
933 if (unlikely(end->vp_uint16 < start->vp_uint16)) {
934 ERROR("Value of %d for slot end less than value of %d for slot start",
935 end->vp_uint16, start->vp_uint16);
936 goto error;
937 }
938 DEBUG4("Setting nodes for slots %d to %d", start->vp_uint16, end->vp_uint16);
939 for (i = start->vp_uint16; i <= end->vp_uint16; i++) {
940 memcpy(&key_slot_pending[i], &tmp_slot, sizeof(*key_slot_pending));
941 }
942 }
943 }
944
945 memcpy(&rtcluster->key_slot, &key_slot_pending, sizeof(rtcluster->key_slot));
946
947 /*
948 * Anything not in the active set of nodes gets
949 * added back into the queue, to be re-used.
950 *
951 * We start at 1, as node 0 is reserved.
952 */
953 for (i = 1; i <= rtcluster->conf->max_nodes; i++) {
954#ifndef NDEBUG
955 fr_redis_ct_node_t *found;
956
957 if (rtcluster->node[i].is_active) {
958 /* Sanity check for duplicates that are active */
959 fr_rb_find((void **)&found, rtcluster->used_nodes, &rtcluster->node[i]);
960 fr_assert(found);
961 fr_assert(found->is_active);
962 fr_assert(found->id == i);
963 }
964#endif
965
966 if (!active[i] && rtcluster->node[i].is_active) {
967 SET_INACTIVE(&rtcluster->node[i]);
968
969 /*
970 * Only change the masters once we've successfully
971 * remapped the cluster.
972 */
973 } else if (master[i]) {
974 rtcluster->node[i].is_master = true;
975 } else {
976 rtcluster->node[i].is_master = false;
977 }
978 }
979
980 rtcluster->state = CLUSTER_READY;
981 rtcluster->map_updated = fr_time();
982
983 /*
984 * Enqueue any commands which were waiting for the cluster remap.
985 */
986 while ((cmd = fr_dlist_pop_head(&rtcluster->pend_cmds))) {
988 }
989
990 /*
991 * Resume any requests which were waiting for the cluster remap.
992 */
993 while ((pend_req = fr_dlist_pop_head(&rtcluster->pend_reqs))) {
995 talloc_free(pend_req);
996 }
997
998 return 0;
999}
1000
1001/** Process a cluster map fail message from the coordinator.
1002 *
1003 * @param rtcluster Cluster to update
1004 * @param list pairs sent by a coordinator
1005 */
1007{
1009 fr_redis_ct_pend_req_t *pend_req;
1010
1011 DEBUG3("Cluster %d failed", rtcluster->cluster_id);
1012 rtcluster->state = CLUSTER_FAIL;
1013
1014 /*
1015 * Inform any pending requests that the cluster has failed.
1016 */
1017 while ((cmd = fr_dlist_pop_head(&rtcluster->pend_cmds))) {
1020 }
1021
1022 /*
1023 * Resume any requests which were waiting for the cluster remap.
1024 */
1025 while ((pend_req = fr_dlist_pop_head(&rtcluster->pend_reqs))) {
1027 talloc_free(pend_req);
1028 }
1029
1030 return 0;
1031}
1032
1033/** Initiate bootstrapping of the cluster map
1034 *
1035 * To be used when a module first wants to fetch a cluster map
1036 *
1037 * @param rtcluster Cluster to fetch map for
1038 * @param cw Coord worker to launch request
1039 * @param coord_pair_reg Coord pair registration
1040 * @return
1041 * - 0 on success.
1042 * - -1 on failure.
1043 */
1045{
1046 fr_redis_conf_t const *conf = rtcluster->conf;
1047 fr_pair_list_t list;
1048 fr_pair_t *vp;
1049 TALLOC_CTX *local = talloc_new(NULL);
1050 int ret;
1051 size_t i;
1052
1053 fr_pair_list_init(&list);
1055 if (!vp) {
1056 error:
1057 talloc_free(local);
1058 return -1;
1059 }
1060
1061 if (fr_pair_append_by_da(local, &vp, &list, attr_redis_log_prefix) < 0) goto error;
1062 if (fr_value_box_strdup(vp, &vp->data, NULL, conf->log_prefix, false) < 0) goto error;
1063
1064 fr_pair_list_append_by_da(local, vp, &list, attr_redis_max_nodes, conf->max_nodes, false);
1065 if (!vp) goto error;
1066
1067 for (i = 0; i < talloc_array_length(conf->hostname); i++) {
1068 if (fr_pair_append_by_da(local, &vp, &list, attr_redis_bootstrap_node) < 0) goto error;
1069 if (fr_value_box_strdup(vp, &vp->data, NULL, conf->hostname[i], false) < 0) goto error;
1070 }
1071
1072 fr_pair_list_append_by_da(local, vp, &list, attr_redis_bootstrap_port, conf->port, false);
1073 if (!vp) goto error;
1074
1075 if (conf->password) {
1076 if (fr_pair_append_by_da(local, &vp, &list, attr_redis_password) < 0) goto error;
1077 if (fr_value_box_strdup(vp, &vp->data, NULL, conf->password, false) < 0) goto error;
1078 if (conf->username) {
1079 if (fr_pair_append_by_da(local, &vp, &list, attr_redis_username) < 0) goto error;
1080 if (fr_value_box_strdup(vp, &vp->data, NULL, conf->username, false) < 0) goto error;
1081 }
1082 }
1083
1084 if (conf->use_tls) {
1085 uintptr_t tls_conf = (uintptr_t)rtcluster->tls_cs;
1086 fr_pair_list_append_by_da(local, vp, &list, attr_redis_use_tls, false, false);
1087 if (!vp) goto error;
1088 fr_pair_list_append_by_da(local, vp, &list, attr_redis_tls_conf, (uint64_t)tls_conf, false);
1089 if (!vp) goto error;
1090 }
1091
1092 ret = fr_worker_to_coord_pair_send(cw, coord_pair_reg, &list);
1093 talloc_free(local);
1094
1095 if (ret < 0) return -1;
1096 rtcluster->state = CLUSTER_MAP_FETCHING;
1097
1098 return 0;
1099}
1100
1101/** Initiate updating of the cluster map
1102 *
1103 * To be used when a command returns MOVED
1104 */
1106 fr_coord_pair_reg_t *coord_pair_reg, bool force)
1107{
1108 fr_pair_list_t list;
1109 fr_pair_t *vp;
1110 TALLOC_CTX *local;
1111 int ret;
1112
1113 if (rtcluster->cluster_id == 0) return REDIS_ASYNC_RCODE_BOOTSTRAP;
1114
1115 /*
1116 * The update request has already been sent.
1117 */
1118 if ((rtcluster->state == CLUSTER_MAP_FETCHING) && !force) return REDIS_ASYNC_RCODE_SUCCESS;
1119
1120 /*
1121 * If the cluster was updated less than 1 sec ago, don't ask.
1122 */
1123 if (((fr_time_to_sec(fr_time()) == fr_time_to_sec(rtcluster->map_updated))) && !force) return REDIS_ASYNC_RCODE_SUCCESS;
1124
1125 DEBUG3("Requesting updated map for cluster %d", rtcluster->cluster_id);
1126
1127 local = talloc_new(NULL);
1128 fr_pair_list_init(&list);
1130 if (!vp) {
1131 error:
1132 talloc_free(local);
1134 }
1135
1136 fr_pair_list_append_by_da(local, vp, &list, attr_redis_cluster_id, rtcluster->cluster_id, false);
1137 if (!vp) goto error;
1138
1139 if (force) {
1140 fr_pair_list_append_by_da(local, vp, &list, attr_redis_force_update, true, false);
1141 }
1142
1143 ret = fr_worker_to_coord_pair_send(cw, coord_pair_reg, &list);
1144 talloc_free(local);
1145
1146 if (ret < 0) return REDIS_ASYNC_RCODE_ERROR;
1147 rtcluster->state = CLUSTER_MAP_FETCHING;
1148
1150}
1151
1153{
1154 fr_redis_ct_node_t find, *found;
1155
1156 find.ioconf.hostname = ioconf->hostname;
1157 find.ioconf.port = ioconf->port;
1158 fr_rb_find((void **)&found, rtcluster->used_nodes, &find);
1159 return found;
1160}
1161
1163 fr_redis_ct_t *rtcluster, bool is_master, bool is_replica)
1164{
1165 uint64_t in_use = fr_rb_num_elements(rtcluster->used_nodes);
1167 fr_redis_ct_node_t *node;
1168 uint8_t count = 0;
1169 fr_redis_io_conf_t *found;
1170
1171 switch (rtcluster->state) {
1172 case CLUSTER_INIT:
1174
1177
1178 default:
1179 break;
1180 }
1181
1182 if (in_use == 0) {
1183 *out = NULL;
1184 *count_out = 0;
1186 }
1187
1188 found = talloc_zero_array(ctx, fr_redis_io_conf_t, in_use);
1189 if (!found) {
1190 fr_strerror_const("Out of memory");
1192 }
1193
1194 for (node = fr_rb_iter_init_inorder(rtcluster->used_nodes, &iter);
1195 node;
1196 node = fr_rb_iter_next_inorder(rtcluster->used_nodes, &iter)) {
1197 if ((is_master && node->is_master) || (is_replica && !node->is_master)) found[count++] = node->ioconf;
1198 }
1199
1200 if (count == 0) {
1201 *out = NULL;
1202 talloc_free(found);
1203 } else {
1204 *out = found;
1205 }
1206 *count_out = count;
1208}
1209
1210/** Ensure pending request is removed from the list on freeing.
1211 */
1213{
1214 if (!fr_dlist_entry_in_list(&pend_req->entry)) return 0;
1215 fr_dlist_remove(&pend_req->rtcluster->pend_reqs, pend_req);
1216 return 0;
1217}
1218
1219/** Add a request to the list of those waiting for the cluster map
1220 *
1221 */
1222void fr_redis_ct_request_yield(TALLOC_CTX *ctx, fr_redis_ct_t *rtcluster, request_t *request)
1223{
1224 fr_redis_ct_pend_req_t *pend_req;
1225
1226 MEM(pend_req = talloc_zero(ctx, fr_redis_ct_pend_req_t));
1227 pend_req->request = request;
1228 pend_req->rtcluster = rtcluster;
1229 fr_dlist_insert_tail(&rtcluster->pend_reqs, pend_req);
1230 talloc_set_destructor(pend_req, _fr_redis_ct_pend_req_free);
1231}
#define CMP(_a, _b)
Same as CMP_PREFER_SMALLER use when you don't really care about ordering, you just want an ordering.
Definition build.h:113
#define unlikely(_x)
Definition build.h:455
#define UNUSED
Definition build.h:384
A section grouping multiple CONF_PAIR.
Definition cf_priv.h:106
uint8_t fr_redis_ct_max_redirects(fr_redis_ct_t *rtcluster)
How many times a command set running on this cluster may be redirected.
fr_redis_conf_t const * conf
Redis configuration for the cluster.
static int _fr_redis_async_cmd_free(fr_redis_async_cmd_t *cmd)
fr_redis_ct_t * rtcluster
Cluster the request is waiting for.
fr_redis_ct_state_t state
State of the cluster.
#define SET_ACTIVE(_node)
bool read_only
Should this command be run read only.
void fr_redis_ct_request_yield(TALLOC_CTX *ctx, fr_redis_ct_t *rtcluster, request_t *request)
Add a request to the list of those waiting for the cluster map.
int fr_redis_ct_map_bootstrap(fr_redis_ct_t *rtcluster, fr_coord_worker_t *cw, fr_coord_pair_reg_t *coord_pair_reg)
Initiate bootstrapping of the cluster map.
request_t * request
The request waiting for the map.
fr_rb_tree_t * used_nodes
Active nodes.
request_t * request
Request this command set relates to.
bool active_oneshot
Should the callback only be called once.
size_t key_len
Length of key.
fr_redis_async_cmd_t * fr_redis_async_cmd_start(TALLOC_CTX *ctx, request_t *request, fr_redis_async_rcode_t *rcode, fr_redis_ct_t *rtcluster, uint8_t const *key, size_t key_len, fr_redis_command_set_t *cmds, bool read_only, fr_redis_ct_node_t *node)
Start running a command set on an async redis cluster.
uint16_t cluster_id
Number assigned to the cluster by coordinator.
uint8_t max_redirects
maximum number of redirects.
static uint16_t cluster_key_hash(uint8_t const *key, size_t key_len)
Resolve key to key slot index.
trunk_conf_t const * tconf
Configuration for all trunks in the cluster.
fr_redis_async_rcode_t fr_redis_ct_node_addr_by_role(TALLOC_CTX *ctx, fr_redis_io_conf_t *out[], uint8_t *count_out, fr_redis_ct_t *rtcluster, bool is_master, bool is_replica)
fr_dlist_head_t pend_reqs
Requests awaiting cluster map.
trunk_conf_t const * fr_redis_ct_trunk_conf(fr_redis_ct_t *rtcluster)
fr_redis_trunk_t * rtrunk
Trunk the command set is currently running on.
fr_redis_ct_node_t * node
Array of nodes in this cluster.
fr_redis_async_rcode_t fr_redis_ct_map_get(fr_redis_ct_t *rtcluster, fr_coord_worker_t *cw, fr_coord_pair_reg_t *coord_pair_reg, bool force)
Initiate updating of the cluster map.
fr_dlist_t entry
Entry in the list of pending requests.
fr_redis_ct_node_t * node
Specific node to run command set on.
void fr_redis_async_cmd_cancel(fr_redis_async_cmd_t *cmd)
Cancel a Redis async command.
uint8_t num_replicas
Number of replica nodes.
fr_redis_ct_key_slot_t key_slot[KEY_SLOTS]
fr_redis_ct_node_t * fr_redis_ct_node_by_addr(fr_redis_ct_t *rtcluster, fr_redis_io_conf_t *ioconf)
fr_redis_async_rcode_t fr_redis_async_cmd_resend(fr_redis_async_cmd_t *cmd)
Re-submit a redis async command set.
fr_redis_trunk_t * trunk
Trunk connection to this node.
#define MAX_REPLICAS
Maximum number of replicas associated with a keyslot.
fr_redis_ct_node_t const * fr_redis_ct_replica(fr_redis_ct_t *rtcluster, fr_redis_ct_key_slot_t const *key_slot, uint8_t replica_num)
Return the replica node that would be used for a particular key slot.
fr_redis_ct_node_t * fr_redis_async_cmd_node(fr_redis_async_cmd_t *cmd)
Fetch the cluster node a command was last sent to.
fr_redis_trunk_t * fr_redis_async_cmd_trunk(fr_redis_async_cmd_t *cmd)
Fetch the redis trunk a command is associated with.
#define SET_INACTIVE(_node)
fr_redis_ct_key_slot_t const * fr_redis_ct_slot_by_key(fr_redis_ct_t *rtcluster, request_t *request, uint8_t const *key, size_t key_len)
Resolve key to key slot.
CONF_SECTION const * tls_cs
TLS CONF_SECTION.
fr_rb_node_t rbnode
Entry in the tree of used nodes.
#define CONFIGURE_NODE(_node, _addr, _port)
fr_pair_list_t trigger_args
Pairs to pass to trigger functions.
fr_redis_ct_key_slot_t const * key_slot
Key slot identified from the command key.
static int _fr_redis_ct_pend_req_free(fr_redis_ct_pend_req_t *pend_req)
Ensure pending request is removed from the list on freeing.
fr_redis_async_rcode_t fr_redis_async_cmd_redirect(fr_redis_async_cmd_t *cmd)
Re-submit a redis async command set on a different node.
uint8_t master
id of the master node.
uint8_t replica_no
Current replica number being used.
fr_redis_ct_t * rtcluster
Cluster this command set is running on.
fr_time_t map_updated
Time the cluster last updated.
fr_redis_ct_t * fr_redis_ct_alloc(TALLOC_CTX *ctx, CONF_SECTION *tls_cs, fr_event_list_t *el, fr_redis_conf_t *conf, fr_redis_trunk_active_t active, void *active_uctx, bool active_oneshot)
Allocate per-thread, per-cluster instance.
fr_redis_ct_state_t
@ CLUSTER_FAIL
The coordinator reported a failed cluster map update.
@ CLUSTER_MAP_FETCHING
The cluster map is currently being fetched.
@ CLUSTER_READY
The cluster is available to handle requests.
@ CLUSTER_INIT
Cluster has been initialised.
static fr_cmp_ret_t _cluster_thread_node_cmp(void const *one, void const *two)
Compare two redis nodes to check equality.
uint8_t replica[MAX_REPLICAS]
Array of ids of replica nodes.
uint8_t const * key
Key used to identify key slot.
char name[INET6_ADDRSTRLEN]
fr_redis_ct_t * rtcluster
Cluster this node belongs to.
fr_fifo_t * free_nodes
Nodes not currently active.
int fr_redis_ct_port(uint16_t *out, fr_redis_ct_node_t const *node)
Return the port of a particular node.
bool is_master
Is this node currently a master.
fr_redis_ct_node_t const * fr_redis_ct_master(fr_redis_ct_t *rtcluster, fr_redis_ct_key_slot_t const *key_slot)
Return the master node that would be used for a particular key slot.
fr_redis_command_set_t * cmds
Command set to run.
bool delay_start
Prevent connections from spawning immediately.
fr_event_list_t * fr_redis_ct_el(fr_redis_ct_t *rtcluster)
fr_event_list_t * el
bool is_active
Is this node currently active.
fr_dlist_t entry
Entry in the list of commands waiting for a cluster remap.
fr_redis_trunk_active_t active
Callback to run when the trunk becomes active.
fr_dlist_head_t pend_cmds
Commands awaiting cluster map.
int fr_redis_ct_map_update(fr_redis_ct_t *rtcluster, fr_pair_list_t const *list)
Update a Redis cluster map from a pair list returned from a coordinator.
void * active_uctx
Uctx to pass to active callback.
int fr_redis_ct_map_fail(fr_redis_ct_t *rtcluster, UNUSED fr_pair_list_t const *list)
Process a cluster map fail message from the coordinator.
char const * fr_redis_ct_ipaddr(fr_redis_ct_node_t const *node)
Return the ipaddr of a particular node.
uint8_t id
Array offset in the array of available nodes.
fr_redis_io_conf_t ioconf
Connection config for this node.
static fr_redis_async_rcode_t fr_redis_async_cmd_enqueue(fr_redis_async_cmd_t *cmd)
Enqueue a command set on a node identified by the key.
Structure for holding the state of an async redis command set.
Structure to record that a request is waiting for the cluster map.
Thread local state for a cluster.
Redis asynchronous cluster management.
#define KEY_SLOTS
Maximum number of keyslots (should not change).
struct fr_redis_async_cmd_s fr_redis_async_cmd_t
The worker end of worker <-> coordinator communication.
Definition coord.c:73
int fr_worker_to_coord_pair_send(fr_coord_worker_t *cw, fr_coord_pair_reg_t *coord_pair_reg, fr_pair_list_t *list)
Send a pair list from a worker to a coordinator.
Definition coord_pair.c:824
struct fr_coord_pair_reg_s fr_coord_pair_reg_t
Definition coord_pair.h:32
uint16_t fr_crc16_xmodem(uint8_t const *in, size_t in_len)
CRC16 implementation according to CCITT standards.
Definition crc16.c:91
#define MEM(x)
Definition debug.h:38
#define ERROR(fmt,...)
Definition dhcpclient.c:40
static void * fr_dlist_remove(fr_dlist_head_t *list_head, void *ptr)
Remove an item from the list.
Definition dlist.h:620
static bool fr_dlist_entry_in_list(fr_dlist_t const *entry)
Check if a list entry is part of a list.
Definition dlist.h:145
static void * fr_dlist_pop_head(fr_dlist_head_t *list_head)
Remove the head item in a list.
Definition dlist.h:654
static int fr_dlist_insert_tail(fr_dlist_head_t *list_head, void *ptr)
Insert an item into the tail of a list.
Definition dlist.h:360
#define fr_dlist_talloc_init(_head, _type, _field)
Initialise the head structure of a doubly linked list.
Definition dlist.h:257
Head of a doubly linked list.
Definition dlist.h:51
Entry in a doubly linked list.
Definition dlist.h:41
void * fr_fifo_peek(fr_fifo_t *fi)
Examine the next element that would be popped.
Definition fifo.c:158
int fr_fifo_push(fr_fifo_t *fi, void *data)
Push data onto the fifo.
Definition fifo.c:111
void * fr_fifo_pop(fr_fifo_t *fi)
Pop data off of the fifo.
Definition fifo.c:135
#define fr_fifo_create(_ctx, _max_entries, _node_free)
Creates a fifo.
Definition fifo.h:66
talloc_free(hp)
int fr_inet_pton_port(fr_ipaddr_t *out, uint16_t *port_out, char const *value, ssize_t inlen, int af, bool resolve, bool mask)
Parses IPv4/6 address + port, to fr_ipaddr_t and integer (port)
Definition inet.c:944
char * fr_inet_ntop(char out[static FR_IPADDR_STRLEN], size_t outlen, fr_ipaddr_t const *addr)
Print the address portion of a fr_ipaddr_t.
Definition inet.c:1025
#define FR_IPADDR_STRLEN
Like INET6_ADDRSTRLEN but includes space for the textual Zone ID.
Definition inet.h:89
IPv4/6 prefix.
void unlang_interpret_mark_runnable(request_t *request)
Mark a request as resumable.
Definition interpret.c:2008
fr_dict_attr_t const * attr_redis_node_role
Definition redis.c:74
fr_dict_attr_t const * attr_redis_bootstrap_node
Definition redis.c:62
fr_dict_attr_t const * attr_redis_use_tls
Definition redis.c:76
fr_dict_attr_t const * attr_redis_slot_end
Definition redis.c:70
fr_dict_attr_t const * attr_redis_force_update
Definition redis.c:75
fr_dict_attr_t const * attr_redis_packet_type
Definition redis.c:59
fr_dict_attr_t const * attr_redis_log_prefix
Definition redis.c:60
fr_dict_attr_t const * attr_redis_node
Definition redis.c:71
fr_dict_attr_t const * attr_redis_node_port
Definition redis.c:73
fr_dict_attr_t const * attr_redis_slot
Definition redis.c:68
fr_dict_attr_t const * attr_redis_slot_start
Definition redis.c:69
fr_dict_attr_t const * attr_redis_max_nodes
Definition redis.c:61
fr_dict_attr_t const * attr_redis_password
Definition redis.c:65
fr_dict_attr_t const * attr_redis_bootstrap_port
Definition redis.c:63
fr_dict_attr_t const * attr_redis_cluster_id
Definition redis.c:66
fr_dict_attr_t const * attr_redis_tls_conf
Definition redis.c:77
fr_dict_attr_t const * attr_redis_username
Definition redis.c:64
fr_dict_attr_t const * attr_redis_node_endpoint
Definition redis.c:72
fr_dict_attr_t const * attr_redis_shard
Definition redis.c:67
char const * hostname
Definition io.h:51
uint16_t port
Definition io.h:52
#define PERROR(_fmt,...)
Definition log.h:233
#define DEBUG3(_fmt,...)
Definition log.h:271
#define ROPTIONAL(_l_request, _l_global, _fmt,...)
Use different logging functions depending on whether request is NULL or not.
Definition log.h:545
#define RDEBUG3(fmt,...)
Definition log.h:360
#define DEBUG4(_fmt,...)
Definition log.h:272
#define fr_time()
Definition event.c:60
Stores all information relating to an event list.
Definition event.c:377
unsigned short uint16_t
unsigned int uint32_t
unsigned char uint8_t
#define UINT8_MAX
fr_cmp_ret_t
Result of an ordering comparison.
Definition misc.h:50
int fr_pair_append_by_da(TALLOC_CTX *ctx, fr_pair_t **out, fr_pair_list_t *list, fr_dict_attr_t const *da)
Alloc a new fr_pair_t (and append)
Definition pair.c:1417
fr_pair_t * fr_pair_find_by_da(fr_pair_list_t const *list, fr_pair_t const *prev, fr_dict_attr_t const *da)
Find the first pair with a matching da.
Definition pair.c:708
void fr_pair_list_init(fr_pair_list_t *list)
Initialise a pair list header.
Definition pair.c:47
void fr_redis_command_set_rcode_set(fr_redis_command_set_t *cmds, fr_redis_async_rcode_t rcode)
Set the rcode for a command set.
Definition pipeline.c:1047
fr_redis_pipeline_status_t redis_command_set_enqueue(fr_redis_trunk_t *rtrunk, fr_redis_command_set_t *cmds)
Enqueue a command set on a specific trunk.
Definition pipeline.c:535
int fr_redis_command_set_reset(fr_redis_command_set_t *cmds)
Reset a command set to it's state before enqueuing.
Definition pipeline.c:1065
void fr_redis_command_set_next_node(fr_redis_command_set_t *cmds, fr_redis_io_conf_t *ioconf)
Extract the next node address and port from a command set.
Definition pipeline.c:1054
void fr_redis_command_set_cancel(fr_redis_command_set_t *cmds)
Cancel a command set.
Definition pipeline.c:566
Represents a collection of pipelined commands.
Definition pipeline.c:94
fr_redis_pipeline_status_t
Definition pipeline.h:43
@ FR_REDIS_PIPELINE_OK
No failure.
Definition pipeline.h:44
@ FR_REDIS_PIPELINE_DST_UNAVAILABLE
Cluster or host is down.
Definition pipeline.h:46
void(* fr_redis_trunk_active_t)(fr_redis_trunk_t *rtrunk, void *uctx)
Definition pipeline.h:55
VQP attributes.
#define fr_assert(_expr)
Definition rad_assert.h:37
#define RDEBUG2(fmt,...)
#define DEBUG2(fmt,...)
#define WARN(fmt,...)
static rs_t * conf
Definition radsniff.c:52
uint32_t fr_rand(void)
Return a 32-bit random number.
Definition rand.c:104
uint32_t fr_rb_num_elements(fr_rb_tree_t *tree)
Return how many nodes there are in a tree.
Definition rb.c:807
int fr_rb_find(void **found, fr_rb_tree_t const *tree, void const *data)
Find an element in the tree, returning the data, not the node.
Definition rb.c:586
void * fr_rb_iter_init_inorder(fr_rb_tree_t *tree, fr_rb_iter_inorder_t *iter)
Initialise an in-order iterator.
Definition rb.c:850
int fr_rb_insert(fr_rb_tree_t *tree, void const *data)
Insert data into a tree.
Definition rb.c:637
void * fr_rb_iter_next_inorder(UNUSED fr_rb_tree_t *tree, fr_rb_iter_inorder_t *iter)
Return the next node.
Definition rb.c:876
#define fr_rb_inline_alloc(_ctx, _type, _field, _data_cmp, _data_free)
Allocs a red black tree.
Definition rb.h:269
Iterator structure for in-order traversal of an rbtree.
Definition rb.h:319
The main red black tree structure.
Definition rb.h:71
uint8_t max_nodes
Maximum number of cluster nodes to connect to.
Definition base.h:124
fr_redis_async_rcode_t
Definition base.h:80
@ REDIS_ASYNC_RCODE_BOOTSTRAP
The caller should issue a request to bootstrap the cluster map.
Definition base.h:83
@ REDIS_ASYNC_RCODE_ERROR
Unrecoverable error.
Definition base.h:82
@ REDIS_ASYNC_RCODE_GETMAP
The caller should issue a request to update the cluster map.
Definition base.h:84
@ REDIS_ASYNC_RCODE_FAIL
The command set trunk request has been failed.
Definition base.h:90
@ REDIS_ASYNC_RCODE_TRY_AGAIN
Try the operation again.
Definition base.h:86
@ REDIS_ASYNC_RCODE_SUCCESS
Operation was successful.
Definition base.h:81
uint32_t max_redirects
Maximum number of times we can be redirected.
Definition base.h:125
@ FR_REDIS_CLUSTER_MAP_BOOTSTRAP
Definition base.h:102
@ FR_REDIS_CLUSTER_MAP_GET
Definition base.h:103
struct fr_redis_ct_s fr_redis_ct_t
Definition base.h:56
Configuration parameters for a redis connection.
Definition base.h:114
static char buff[sizeof("18446744073709551615")+3]
Definition size_tests.c:37
return count
Definition module.c:155
fr_pair_t * vp
Stores an attribute, a value and various bits of other data.
Definition pair.h:68
static size_t talloc_strlen(char const *s)
Returns the length of a talloc array containing a string.
Definition talloc.h:143
static int64_t fr_time_to_sec(fr_time_t when)
Convert an fr_time_t (internal time) to number of sec since the unix epoch (wallclock time)
Definition time.h:731
"server local" time.
Definition time.h:69
bool always_writable
Set to true if our ability to write requests to a connection handle is not dependent on the state of ...
Definition trunk.h:281
Common configuration parameters for a trunk.
Definition trunk.h:234
static fr_event_list_t * el
#define fr_pair_list_append_by_da(_ctx, _vp, _list, _attr, _val, _tainted)
Append a pair to a list, assigning its value.
Definition pair.h:306
#define fr_strerror_const(_msg)
Definition strerror.h:223
int fr_value_box_strdup(TALLOC_CTX *ctx, fr_value_box_t *dst, fr_dict_attr_t const *enumv, char const *src, bool tainted)
Copy a nul terminated string to a fr_value_box_t.
Definition value.c:4648
#define fr_box_strvalue_len(_val, _len)
Definition value.h:334
static fr_sbuff_err_t char ** out
Definition value.h:1062