The FreeRADIUS server $Id: f3670dba8951ca10eb4948feb3dc3db9423a334f $
Loading...
Searching...
No Matches
pipeline.c
Go to the documentation of this file.
1/*
2 * This program 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: e0178ad7c432a6aeaa19c12bae65a4edc9952514 $
19 * @file lib/redis/pipeline.c
20 * @brief Functions for pipelining commands.
21 *
22 * @copyright 2019 The FreeRADIUS server project
23 * @copyright 2019 Network RADIUS SAS (legal@networkradius.com)
24 *
25 * @author Arran Cudbard-Bell (a.cudbardb@freeradius.org)
26 */
27
28#include <freeradius-devel/server/connection.h>
29#include <freeradius-devel/server/trunk.h>
30
31#include "pipeline.h"
32#include "cluster_async.h"
33#include "io.h"
34
35
36/** The thread local free list
37 *
38 * Any entries remaining in the list will be freed when the thread is joined
39 */
41
42typedef enum {
43 FR_REDIS_COMMAND_NORMAL = 0, //!< A normal, non-transactional command.
44 FR_REDIS_COMMAND_TRANSACTION_START, //!< Start of a transaction block. Either WATCH or MULTI.
45 ///< if a transaction is started with WATCH, then multi
46 ///< is not marked up as a transaction start.
47 FR_REDIS_COMMAND_TRANSACTION_END //!< End of a transaction block. Either EXEC or DISCARD.
48 ///< If this command fails with
49 ///< MOVED or ASK, all commands back to the previous
50 ///< MULTI command must be requeued.
52
53typedef enum {
54 FR_REDIS_COMMAND_FMT_EXPANDED = 0, //!< A command as a single string
55 FR_REDIS_COMMAND_FMT_ARGV, //!< A command as an argv array
56 FR_REDIS_COMMAND_FMT_PREFORMATTED //!< A command preformatted with redisCommandFormat
58
59/** Represents a single command
60 *
61 */
63 fr_redis_command_set_t *cmds; //!< Command set this entry belongs to.
64 fr_dlist_t entry; //!< Entry in the command buffer.
65
66 fr_redis_command_type_t type; //!< Redis command type.
67 fr_redis_command_fmt_t fmt; //!< Redis command format.
68
69 union {
70 struct{
71 char const *str; //!< The command string.
72 size_t str_len; //!< Length of the command string.
73 };
74 struct {
75 size_t argc; //!< Number of argv arguments.
76 char const **argv; //!< Arguments for the redis command.
77 size_t *argv_len; //!< Lengths of the arguments.
78 };
79 };
80
81 uint64_t sqn; //!< The sequence number of the command. This is only
82 ///< valid for a specific handle, and is unique within
83 ///< the handle.
84
85 fr_redis_command_complete_t complete; //!< Callback to process result from this command.
86
87 void *rctx; //!< To be passed to the callback.
88};
89
90/** Represents a collection of pipelined commands
91 *
92 * Commands MUST map to the same cluster node if using clustering.
93 */
96
97 fr_redis_async_rcode_t rcode; //!< Code from last error returned.
98 bool autofree; //!< Should the command set be freed when it is complete
99
100 char *next_node_ip; //!< IP address of node from MOVED / ASK reply
101 uint16_t next_node_port; //!< Port of node from MOVED / ASK reply
102
103 /** @name Command state lists
104 * @{
105 */
106 fr_dlist_head_t pending; //!< Commands yet to be sent.
107 fr_dlist_head_t sent; //!< Commands sent.
108 fr_dlist_head_t completed; //!< Commands complete with replies.
109 /** @} */
110
111 uint8_t redirected; //!< How many times this command set was redirected.
112 uint32_t max_redirects; //!< How many times this command set may be redirected.
113 ///< Copied from the cluster configuration when the
114 ///< command set is enqueued.
115
116 /** @name Request state
117 *
118 * treq and request are duplicated here with the trunk code.
119 * The reason for this, is because a fr_command_set_t, may need to be transferred
120 * between trunks when redirects are being followed, and so we need this information
121 * encapsulated within the command set, not just within the trunk.
122 * @{
123 */
124 trunk_request_t *treq; //!< Trunk request this command set is associated with.
125 request_t *request; //!< Request this commands set is associated with (if any).
126 void *rctx; //!< Resume context to write results to.
127 /** @} */
128
129 /** @name Callback functions
130 * @{
131 */
132 fr_redis_command_set_complete_t complete; //!< Notify the creator of the command set
133 ///< that the command set has executed to
134 ///< to completion. We have results for
135 ///< all commands.
136
137 fr_redis_command_set_fail_t fail; //!< Notify the creator of the command set
138 ///< that the command set failed to execute
139 ///< to completion. Partial results will
140 ///< be available.
141 /** @} */
142
143 /** @name Command set transaction stats
144 *
145 * We do these checks as REDIS commands from a great number of requests may pipeline
146 * requests on the same connection and leaving a transaction open would be fairly
147 * catastrophic, potentially causing errors across all future command sets set to
148 * the connection.
149 * @{
150 */
151 bool txn_watch; //!< Transaction was started with a watch statement.
152 uint16_t txn_start; //!< Number of times a transaction block was started
153 ///< in this command set.
154 uint16_t txn_end; //!< The number of times a transaction block ended
155 ///< in this command set.
156
157 /** @} */
158
159 bool blocking; //!< This command set contains one or more commands
160 ///< which block the client (e.g. WAIT)
161};
162
164 fr_redis_io_conf_t const *io_conf; //!< Redis I/O configuration. Specifies how to connect
165 ///< to the host this trunk is used to communicate with.
166 trunk_t *trunk; //!< Trunk containing all the connections to a specific
167 ///< host.
168 fr_redis_ct_t *rtcluster; //!< Cluster this trunk belongs to.
169
170 fr_redis_trunk_active_t active; //!< Callback to run when the trunk becomes active.
171 void *active_uctx; //!< Uctx to pass to active callback.
172};
173
174/** Free any free requests when the thread is joined
175 *
176 */
178{
179 fr_dlist_head_t *list = talloc_get_type_abort(arg, fr_dlist_head_t);
181
182 /*
183 * See the destructor for why this works
184 */
185 while ((cmds = fr_dlist_head(list))) if (talloc_free(cmds) < 0) return -1;
186 return talloc_free(list);
187}
188
189/** Free a command set
190 *
191 */
193{
195 (likely(!fr_dlist_entry_in_list(&cmds->entry)))) return 0; /* Keep a buffer of 1024 */
196
197 /*
198 * Freed from the free list....
199 */
201 fr_dlist_entry_unlink(&cmds->entry); /* Don't trust the list head to be available */
202 return 0;
203 }
204
205 /*
206 * It is possible for a command set to be freed while its trunk request
207 * is still inflight.
208 * This is an edge case such as shutting down the server when scripts
209 * are still being loaded, since the script loading done on redis trunk
210 * startup are not run through requests, so there isn't a cancellation
211 * path.
212 */
213 if (cmds->treq) {
214 switch (cmds->treq->state) {
219 break;
220
221 default:
222 break;
223
224 }
225 }
226
227 talloc_free_children(cmds);
228 memset(cmds, 0, sizeof(*cmds));
229
231
232 return -1; /* Prevent the free */
233}
234
235/** Allocate a new command set
236 *
237 * This is a set of commands that the calling module wants to execute
238 * on the redis server in sequence.
239 *
240 * Control will be returned to the caller via the registered complete
241 * and fail functions.
242 *
243 * @param[in] ctx to bind the command set's lifetime to.
244 * @param[in] request to pass to places that need it.
245 * @param[in] complete Function to call when all commands have been processed.
246 * @param[in] fail Function to call if the command set was not executed
247 * or was partially executed.
248 * @param[in] rctx Resume context to pass to complete and fail functions.
249 * @param[in] autofree Should the command set be freed when completed.
250 * @return A new or refurbished command set.
251 */
253 request_t *request,
256 void *rctx, bool autofree)
257
258{
260 fr_dlist_head_t *free_list;
261
262#define COMMAND_PRE_ALLOC_COUNT 8 //!< How much room we pre-allocate for commands.
263#define COMMAND_PRE_ALLOC_LEN 64 //!< How much we allocate for each command string.
264
265 /*
266 * Initialise the free list
267 */
269 MEM(free_list = talloc(NULL, fr_dlist_head_t));
270 fr_dlist_init(free_list, fr_redis_command_set_t, entry);
272 } else {
273 free_list = command_set_free_list;
274 }
275
276 /*
277 * Pull an element out of the free list
278 * or allocate a new one.
279 */
280 cmds = fr_dlist_pop_head(free_list);
281 if (!cmds) {
286 talloc_set_destructor(cmds, _redis_command_set_free);
288 }
289
293 cmds->request = request;
294 cmds->complete = complete;
295 cmds->fail = fail;
296 cmds->rctx = rctx;
297 cmds->autofree = autofree;
298
299 if (ctx) talloc_link_ctx(ctx, cmds);
300
301 return cmds;
302}
303
305 fr_redis_command_set_t *cmds, char const *cmd)
306{
307 /*
308 * Transaction sanity checks.
309 *
310 * Because commands from many different requests share the same connection
311 * we need to ensure that transaction blocks aren't left dangling and
312 * that the commands are all in the right order.
313 *
314 * We try very hard to do this without incurring a performance penalty
315 * for non-transactional commands.
316 */
317 switch (tolower(cmd[0])) {
318 case 'm':
319 if (tolower(cmd[1]) != 'u') break;
320 if (strncasecmp(cmd, "multi", sizeof("multi") - 1) != 0) break;
321 /*
322 * There should only ever be a difference of
323 * 1 between txn starts and txn ends.
324 */
325 if ((cmds->txn_end < cmds->txn_start) && ((cmds->txn_start - cmds->txn_end) > 1)) {
326 ROPTIONAL(REDEBUG, ERROR, "Too many consecutive \"MULTI\" commands");
328 }
329 /*
330 * If we have a watch before the MULTI,
331 * that's marked as the start of the transaction
332 * block.
333 */
335 cmds->txn_start++; /* Yes MULTI increments start, not WATCH */
336 break;
337
338 case 'e':
339 if (tolower(cmd[1]) != 'x') break;
340 if (strncasecmp(cmd, "exec", sizeof("exec") - 1) != 0) break;
341 goto txn_end;
342
343 /*
344 * It's useful to allow discard as it allows command syntax checks
345 * to be performed against the REDIS server without actually
346 * executing the commands.
347 */
348 case 'd':
349 if (tolower(cmd[1]) != 'i') break;
350 if (strncasecmp(cmd, "discard", sizeof("discard") - 1) != 0) break;
351 txn_end:
352 if (cmds->txn_start <= cmds->txn_end) {
353 ROPTIONAL(REDEBUG, ERROR, "Transaction not started, missing \"MULTI\" command");
355 }
357 cmds->txn_end++;
358 break;
359
360 case 'w':
361 if (tolower(cmd[1]) != 'a') break;
362
363 if (strncasecmp(cmd, "wait", sizeof("wait") - 1) == 0) {
364 cmds->blocking = true;
365 break;
366 }
367
368 if (strncasecmp(cmd, "watch", sizeof("watch") - 1) != 0) break;
369 if (cmds->txn_watch) {
370 ROPTIONAL(REDEBUG, ERROR, "Too many consecutive \"WATCH\" commands");
372 }
373 if (cmds->txn_start > cmds->txn_end) {
374 ROPTIONAL(REDEBUG, ERROR, "\"WATCH\" can only be used before \"MULTI\"");
376 }
378
379 default:
380 break;
381 }
382
384}
385
386/** Add a literal command to the command set
387 *
388 * The command must either be entirely static, or parented by the command set.
389 *
390 * @note Caller should disallow "SUBSCRIBE" et al, if they're not appropriate.
391 * As subscribing to a stream where we're not expecting it would break
392 * things, badly.
393 *
394 * @param[in] cmds Command set to add command to.
395 * @param[in] cmd_str A fully expanded/formatted command to send to redis.
396 * Must be static, or have the same lifetime as the
397 * command set (allocated with the command set as the parent).
398 * @param[in] complete Callback to run when this command completes
399 * @param[in] rctx to pass to `complete`
400 * @return
401 * - FR_REDIS_PIPELINE_BAD_CMDS if a bad command sequence is enqueued.
402 * - FR_REDIS_PIPELINE_OK if command was enqueued successfully.
403 */
405 fr_redis_command_complete_t complete, void *rctx)
406{
407 request_t *request = cmds->request;
410
412
413 MEM(cmd = talloc_zero(cmds, fr_redis_command_t));
414 cmd->cmds = cmds;
415 cmd->type = type;
416 cmd->str = cmd_str;
417 cmd->complete = complete;
418 cmd->rctx = rctx;
420 fr_dlist_insert_tail(&cmds->pending, cmd);
421
423}
424
425/** Add a command with arguments to the command set
426 *
427 * The command and arguments must either be entirely static, or parented by the command set.
428 *
429 * @param[in] cmds Command set to add command to.
430 * @param[in] argc Number of arguments.
431 * @param[in] argv Redis command arguments.
432 * @param[in] argv_len Length of the command arguments.
433 * @param[in] complete Callback to run when this command completes
434 * @param[in] rctx to pass to `complete`
435 * @return
436 * - FR_REDIS_PIPELINE_BAD_CMDS if a bad command sequence is enqueued.
437 * - FR_REDIS_PIPELINE_OK if command was enqueued successfully.
438 */
440 char const **argv, size_t *argv_len,
441 fr_redis_command_complete_t complete, void *rctx)
442{
443 request_t *request = cmds->request;
446
448
449 MEM(cmd = talloc_zero(cmds, fr_redis_command_t));
450 cmd->cmds = cmds;
451 cmd->type = type;
452 cmd->argc = argc;
453 cmd->argv = argv;
454 cmd->argv_len = argv_len;
455 cmd->complete = complete;
456 cmd->rctx = rctx;
458 fr_dlist_insert_tail(&cmds->pending, cmd);
459
461}
462
463/** Add an preformatted command to the command set as formatted by redisCommandFormat or it's variants
464 *
465 * The command must either be entirely static, or parented by the command set.
466 *
467 * @note Caller should disallow "SUBSCRIBE" et al, if they're not appropriate.
468 * As subscribing to a stream where we're not expecting it would break
469 * things, badly.
470 *
471 * @param[in] cmds Command set to add command to.
472 * @param[in] cmd_str A fully formatted command to send to redis.
473 * Must be static, or have the same lifetime as the
474 * command set (allocated with the command set as the parent).
475 * @param[in] cmd_len The length of cmd_str (as returned by redisCommandForamt)
476 * @param[in] complete Callback to run when this command completes
477 * @param[in] rctx to pass to `complete`
478 * @return
479 * - FR_REDIS_PIPELINE_BAD_CMDS if a bad command sequence is enqueued.
480 * - FR_REDIS_PIPELINE_OK if command was enqueued successfully.
481 */
483 size_t cmd_len,
484 fr_redis_command_complete_t complete, void *rctx)
485{
486 request_t *request = cmds->request;
489 char const *p = cmd_str, *end;
490
491 /*
492 * Preformatted Redis commands start *<n>\r\n$<n>\r\n<cmd>. Verify that is what we have.
493 */
494 end = p + cmd_len;
495 if (*p++ != '*') {
496 error:
497 ERROR("Incorrect Redis command format");
499 }
500 while (isdigit(*p) && (p < end)) p++;
501 if (*p++ != '\r') goto error;
502 if (*p++ != '\n') goto error;
503 if (*p++ != '$') goto error;
504 while (isdigit(*p) && (p < end)) p++;
505 if (*p++ != '\r') goto error;
506 if (*p++ != '\n') goto error;
507
509
510 MEM(cmd = talloc_zero(cmds, fr_redis_command_t));
511 cmd->cmds = cmds;
512 cmd->type = type;
513 cmd->str = cmd_str;
514 cmd->str_len = cmd_len;
515 cmd->complete = complete;
516 cmd->rctx = rctx;
518 fr_dlist_insert_tail(&cmds->pending, cmd);
519
521}
522
523/** Enqueue a command set on a specific trunk
524 *
525 * The command set may be passed around several trunks before it is complete.
526 * This is to allow it to follow MOVED and ASK responses.
527 *
528 * @param[in] rtrunk to enqueue command set on.
529 * @param[in] cmds Command set to enqueue.
530 * @return
531 * - FR_REDIS_PIPELINE_OK if commands were immediately enqueued or placed in the backlog.
532 * - FR_REDIS_PIPELINE_DST_UNAVAILABLE if the REDIS host is unreachable.
533 * - FR_REDIS_PIPELINE_FAIL any other general error.
534 */
536{
537 if (cmds->txn_start != cmds->txn_end) {
538 ERROR("Refusing to enqueue - Unbalanced transaction start/stop commands");
540 }
541
542 /*
543 * Record the limit, rather than a pointer to the cluster, so that the command set does not
544 * depend on the cluster still being around when a redirect is received.
545 */
547
548 switch (trunk_request_enqueue(&cmds->treq, rtrunk->trunk, cmds->request, cmds, cmds->rctx)) {
549 case TRUNK_ENQUEUE_OK:
551 if (cmds->blocking) trunk_request_mark_blocking(cmds->treq);
553
556
557 default:
559 }
560}
561
562/** Cancel a command set
563 *
564 * @param[in] cmds Command set to cancel.
565 */
567{
568 if (cmds->treq) {
570 cmds->treq = NULL;
571 }
573}
574
575/** Convert a MOVED / ASK reply into an address and port
576 *
577 */
578static int redis_addr_from_redirect(TALLOC_CTX *ctx, char **addr, uint16_t *port, redisReply *redirect)
579{
580 unsigned long key;
581 fr_sbuff_t sbuff;
582 fr_ipaddr_t ipaddr;
584
585 if (!redirect || (redirect->type != REDIS_REPLY_ERROR)) return -1;
586
587 fr_sbuff_init_in(&sbuff, redirect->str, redirect->len);
588
591 fr_strerror_const("No '-MOVED' or '-ASK' log_prefix");
592 return -1;
593 }
594
595 if (fr_sbuff_out(&key, &sbuff) < 0) {
596 fr_strerror_const("Failed to parse key slot from MOVED / ASK reply");
597 return -1;
598 };
599 if (key >= KEY_SLOTS) {
600 fr_strerror_printf("Key %lu outside of redis slot range", key);
601 return -1;
602 }
603
604 if (!fr_sbuff_next_if_char(&sbuff, ' ')) {
605 fr_strerror_const("Missing key/host separator");
606 return -1;
607 }
608
609 if (fr_inet_pton_port(&ipaddr, port, fr_sbuff_current(&sbuff), fr_sbuff_remaining(&sbuff),
610 AF_UNSPEC, true, true) < 0) {
611 return -1;
612 }
613 fr_assert(ipaddr.af);
614
615 *addr = talloc_strdup(ctx, fr_inet_ntop(buff, sizeof(buff), &ipaddr));
616
617 return 0;
618}
619
620/** Callback for for receiving Redis replies
621 *
622 * This is called by hiredis for each response is receives. privData is set to the
623 * fr_command_set
624 *
625 * @note Called only from hiredis, not the trunk itself.
626 *
627 * @param[in] ac The async context the command was enqueued on.
628 * @param[in] vreply redisReply containing the result of the command.
629 * @param[in] privdata fr_redis_command_t that was sent to the Redis server.
630 * The fr_redis_command_t contains a pointer to the
631 * fr_redis_command_set_t which holds the treq which
632 * we use to signal that we have responses for all
633 * commands.
634 */
635static void _redis_pipeline_demux(struct redisAsyncContext *ac, void *vreply, void *privdata)
636{
639 connection_t *conn = talloc_get_type_abort(ac->ev.data, connection_t);
640 fr_redis_handle_t *h = talloc_get_type_abort(conn->h, fr_redis_handle_t);
641 redisReply *reply = vreply;
642
643 /*
644 * If we're already disconnecting, then ignore the response.
645 * Testing has shown this callback can be called by hiredis after the connection
646 * has errored. At that point privdata is no longer valid so there's nothing
647 * that can be done.
648 */
649 if (h->freeing) return;
650
651 /*
652 * First check if we should ignore the response
653 */
655 DEBUG4("Ignoring response with SQN %"PRIu64, (h->rsp_sqn - 1)); /* Already incremented */
656 return;
657 }
658
659 cmd = talloc_get_type_abort(privdata, fr_redis_command_t);
660 cmds = cmd->cmds;
661
662 /*
663 * The trunk request has already failed, nothing more to do.
664 */
665 if (cmds->rcode == REDIS_ASYNC_RCODE_FAIL) return;
666
667 fr_dlist_remove(&cmds->sent, cmd);
668 fr_dlist_insert_tail(&cmds->completed, cmd);
669
670 if (!reply) {
672 error:
673 /*
674 * Mark remaining sent commands to be ignored and fail the treq
675 */
676 fr_dlist_foreach(&cmds->sent, fr_redis_command_t, sent_cmd) {
677 fr_redis_connection_ignore_response(h, sent_cmd->sqn);
678 }
679
680 /*
681 * No reply, move unfinished "sent" commands to the "completed" list. This allows the
682 * "fail" callback to see the full list.
683 *
684 * @todo - do we want to have a "failed" list?
685 */
686 fr_dlist_move(&cmds->completed, &cmds->sent);
687
688 /*
689 * Only REDIS_ASYNC_RCODE_ERROR is really a failure.
690 */
691 if (cmds->rcode == REDIS_ASYNC_RCODE_ERROR) {
693 } else {
695 }
696 cmds->treq = NULL;
697 return;
698 }
699
700 /*
701 * If the reply was an error, look for known types.
702 */
703 if (reply->type == REDIS_REPLY_ERROR) {
704 request_t *request = cmds->request;
705
706 fr_assert_msg(reply->str, "Error response contained no error string");
707
708 if (strncmp(REDIS_ERROR_MOVED_STR, reply->str, sizeof(REDIS_ERROR_MOVED_STR) - 1) == 0) {
709 ROPTIONAL(RWARN, WARN, "Server returned %s", reply->str);
711 goto redirect;
712 } else if (strncmp(REDIS_ERROR_ASK_STR, reply->str, sizeof(REDIS_ERROR_ASK_STR) - 1) == 0) {
713 ROPTIONAL(RWARN, WARN, "Server returned %s", reply->str);
715 redirect:
716 cmds->redirected++;
717 if (cmds->redirected >= cmds->max_redirects) {
718 ROPTIONAL(REDEBUG, ERROR, "Redirected too many times (%u > max_redirects %u)",
719 cmds->redirected, cmds->max_redirects);
721 goto error;
722 }
723
724 if (redis_addr_from_redirect(cmds, &cmds->next_node_ip, &cmds->next_node_port, reply) < 0) {
726 }
727 } else if (strncmp(REDIS_ERROR_TRY_AGAIN_STR, reply->str, sizeof(REDIS_ERROR_TRY_AGAIN_STR) - 1) == 0) {
728 ROPTIONAL(RWARN, WARN, "Server returned %s", reply->str);
730 } else if (strncmp(REDIS_ERROR_NO_SCRIPT_STR, reply->str, sizeof(REDIS_ERROR_NO_SCRIPT_STR) - 1) == 0) {
731 ROPTIONAL(RWARN, WARN, "Server returned %s", reply->str);
733 } else {
734 fr_strerror_printf("Server error: %s", reply->str);
736 }
737 goto error;
738 }
739
740 if (cmd->complete) cmd->complete(cmds->request, cmd, reply, cmd->rctx);
742
743 /*
744 * Check is the command set is complete,
745 * and if it is, tell the trunk the treq
746 * is complete.
747 */
748 if ((fr_dlist_num_elements(&cmds->pending) == 0) &&
749 (fr_dlist_num_elements(&cmds->sent) == 0)) {
751 cmds->treq = NULL;
752 }
753}
754
755CC_NO_UBSAN(function) /* UBSAN: false positive - public vs private connection_t trips --fsanitize=function */
757 connection_conf_t const *conf,
758 char const *log_prefix, void *uctx)
759{
760 fr_redis_trunk_t *rtrunk = talloc_get_type_abort(uctx, fr_redis_trunk_t);
761
762 return fr_redis_connection_alloc(tconn, el, conf, rtrunk->io_conf,
763#ifdef HAVE_REDIS_SSL
764 fr_redis_ct_ssl_ctx(rtrunk->rtcluster),
765#endif
766 log_prefix);
767}
768
769/** Enqueue one or more command sets onto a redis handle
770 *
771 * Because the trunk is in always writable mode, _redis_pipeline_mux
772 * will be called any time trunk_request_enqueue is called, so there'll only
773 * ever be one command to dequeue.
774 *
775 * @param[in] el Event list for trunk events. Unused.
776 * @param[in] tconn Trunk connection holding the commands to enqueue.
777 * @param[in] conn Connection handle containing the fr_redis_handle_t.
778 * @param[in] uctx fr_redis_cluster_t. Unused.
779 */
780CC_NO_UBSAN(function) /* UBSAN: false positive - public vs private connection_t trips --fsanitize=function */
782 connection_t *conn, UNUSED void *uctx)
783{
784 trunk_request_t *treq;
787 fr_redis_handle_t *h = talloc_get_type_abort(conn->h, fr_redis_handle_t);
788 request_t *request;
789 int ret;
790
791 while (trunk_connection_pop_request(&treq, tconn) == 0) {
792 cmds = talloc_get_type_abort(treq->preq, fr_redis_command_set_t);
793 request = treq->request;
794 while ((cmd = fr_dlist_head(&cmds->pending))) {
795 /*
796 * If this fails it probably means the connection
797 * is disconnecting, but if that's happening then
798 * we shouldn't be enqueueing new requests?
799 */
800 switch (cmd->fmt) {
802 if (DEBUG_ENABLED3) {
803 size_t i;
804 ROPTIONAL(RDEBUG3, DEBUG3, "Sending Redis argv command");
805 for (i = 0; i < cmd->argc; i++) {
806 ROPTIONAL(RDEBUG3, DEBUG3, " %pV",
807 fr_box_strvalue_len(cmd->argv[i], cmd->argv_len[i]));
808 }
809
810 }
811 ret = redisAsyncCommandArgv(h->ac, _redis_pipeline_demux, cmd, cmd->argc,
812 cmd->argv, cmd->argv_len);
813 break;
814
816 ROPTIONAL(RDEBUG3, DEBUG3, "Sending Redis command %s", cmd->str);
817 ret = redisAsyncCommand(h->ac, _redis_pipeline_demux, cmd, cmd->str);
818 break;
819
821 ROPTIONAL(RDEBUG3, DEBUG3, "Sending Redis formatted command %s", cmd->str);
822 ret = redisAsyncFormattedCommand(h->ac, _redis_pipeline_demux, cmd,
823 cmd->str, cmd->str_len);
824 break;
825 }
826
827 if (unlikely(ret != REDIS_OK)) {
828 ROPTIONAL(REDEBUG, ERROR, "Unexpected error queueing REDIS command");
829
830 while ((cmd = fr_dlist_head(&cmds->sent))) {
832 fr_dlist_remove(&cmds->sent, cmd);
833 fr_dlist_insert_tail(&cmds->pending, cmd);
834 }
836 cmds->treq = NULL;
837 return;
838 }
840 fr_dlist_remove(&cmds->pending, cmd);
841 fr_dlist_insert_tail(&cmds->sent, cmd);
842 }
844 }
845}
846
847/** Deal with cancellation of sent requests
848 *
849 * We can't actually signal redis to not process the request, so depending
850 * on why the commands were cancelled, we either tell the handle to ignore
851 * them, or move them back into the pending list.
852 */
854 trunk_cancel_reason_t reason, UNUSED void *uctx)
855{
856 fr_redis_command_set_t *cmds = talloc_get_type_abort(preq, fr_redis_command_set_t);
857 fr_redis_handle_t *h = conn->h;
858
859 /*
860 * How we cancel is very different depending
861 * on _WHY_ we're cancelling.
862 */
863 switch (reason) {
864 /*
865 * Cancel is only called for requests that
866 * have been sent, and only when the connection
867 * is about to be closed for some reason.
868 *
869 * We don't need to tell the handle to ignore
870 * the responses, we just need to get the
871 * command set back into the correct state for
872 * execution by another handle.
873 */
876 fr_dlist_move(&cmds->pending, &cmds->sent);
877 return;
878
879 /*
880 * If the request was cancelled due to a signal
881 * we'll have a response coming back for a
882 * request, pctx and rctx that no longer exist.
883 * Tell the handle to signal that the response
884 * should be ignored when it's received.
885 *
886 * Free will take care of cleaning up the
887 * pending commands.
888 */
890 {
892
893 /*
894 * Only connected connections will get replies that
895 * need to be ignored.
896 */
897 if (conn->state != CONNECTION_STATE_CONNECTED) return;
898
899 for (cmd = fr_dlist_head(&cmds->sent);
900 cmd;
901 cmd = fr_dlist_next(&cmds->sent, cmd)) {
903 }
904 }
905 return;
906
908 fr_assert(0);
909 return;
910 }
911}
912
913/** Signal the API client that we got a complete set of responses to a command set
914 *
915 */
917 UNUSED void *rctx, UNUSED void *uctx)
918{
919 fr_redis_command_set_t *cmds = talloc_get_type_abort(preq, fr_redis_command_set_t);
920
921 if (cmds->complete) cmds->complete(cmds->request, &cmds->completed, cmds->rctx);
923}
924
925/** Signal the API client that we failed enqueuing the commands
926 *
927 */
928static void _redis_pipeline_command_set_fail(UNUSED request_t *request, void *preq, UNUSED void *rctx,
929 UNUSED trunk_request_state_t state, UNUSED void *uctx)
930{
931 fr_redis_command_set_t *cmds = talloc_get_type_abort(preq, fr_redis_command_set_t);
932
934 if (cmds->fail) cmds->fail(cmds->request, &cmds->completed, cmds->rctx);
936}
937
938/** Free the command set
939 *
940 */
941static void _redis_pipeline_command_set_free(UNUSED request_t *request, void *preq,
942 UNUSED void *uctx)
943{
944 fr_redis_command_set_t *cmds = talloc_get_type_abort(preq, fr_redis_command_set_t);
945
946 if (cmds->autofree) {
947 talloc_free(cmds);
948 return;
949 }
950
951 /*
952 * Once this callback completes, the trunk either frees the treq or
953 * returns the treq to the trunk's `free_requests` list. The
954 * caller-allocated command set outlives the treq. We clear `cmds->treq`
955 * because after the callback returns, `cmds->treq` is left dangling.
956 * If we do NOT clear `cmds->treq`, the command set destructor and
957 * `fr_redis_command_set_cancel()` would access the freed treq through
958 * `cmds->treq`.
959 *
960 * We do NOT clear any other field. The caller needs to read `rcode`
961 * after the request resumes to pick the result path. A MOVED / ASK
962 * reply makes the caller resend the command set to a new node.
963 * `fr_redis_command_set_reset()` needs `completed` to move the
964 * completed commands back to `pending` for the resend.
965 * `fr_redis_command_set_next_node()` needs `next_node_ip` and
966 * `next_node_port` to name the new node.
967 */
968 cmds->treq = NULL;
969}
970
971CC_NO_UBSAN(function) /* UBSAN: false positive - public vs private trunk_t trips --fsanitize=function */
972static void _redis_trunk_active(UNUSED trunk_t *trunk, UNUSED trunk_state_t prev, UNUSED trunk_state_t state, void *uctx)
973{
974 fr_redis_trunk_t *rtcluster = talloc_get_type_abort(uctx, fr_redis_trunk_t);
975
976 rtcluster->active(rtcluster, rtcluster->active_uctx);
977}
978
979/** Allocate a new trunk
980 *
981 * @param[in] rtcluster to allocate the trunk for.
982 * @param[in] io_conf Describing the connection to a single REDIS host.
983 * @param[in] trigger_args Pairs to pass to trigger requests, if triggers are enabled.
984 * @param[in] active Callback to run when the trunk becomes active.
985 * @param[in] active_uctx Uctx to pass to active callback.
986 * @param[in] active_oneshot Should the call back be run just once.
987 * @return
988 * - On success, a new fr_redis_trunk_t which can be used for pipelining commands.
989 * - NULL on failure.
990 */
992 fr_pair_list_t *trigger_args, fr_redis_trunk_active_t active,
993 void *active_uctx, bool active_oneshot)
994{
995 fr_redis_trunk_t *rtrunk;
998 .request_mux = _redis_pipeline_mux,
999 /* demux called directly by hiredis */
1000 .request_cancel = _redis_pipeline_command_set_cancel,
1001 .request_complete = _redis_pipeline_command_set_complete,
1002 .request_fail = _redis_pipeline_command_set_fail,
1003 .request_free = _redis_pipeline_command_set_free
1004 };
1005
1006 MEM(rtrunk = talloc(rtcluster, fr_redis_trunk_t));
1007 *rtrunk = (fr_redis_trunk_t) {
1008 .io_conf = io_conf,
1009 .rtcluster = rtcluster,
1010 .active = active,
1011 .active_uctx = active_uctx,
1012 };
1013 rtrunk->trunk = trunk_alloc(rtrunk, fr_redis_ct_el(rtcluster), &io_funcs, fr_redis_ct_trunk_conf(rtcluster),
1014 io_conf->log_prefix, rtrunk, false, trigger_args);
1015 if (!rtrunk->trunk) {
1016 talloc_free(rtrunk);
1017 return NULL;
1018 }
1019
1020 if (active) trunk_add_watch(rtrunk->trunk, TRUNK_STATE_ACTIVE, _redis_trunk_active, active_oneshot, rtrunk);
1021
1022 return rtrunk;
1023}
1024
1026{
1027 switch(cmd->type) {
1030 return cmd->str;
1031
1033 return cmd->argv[0];
1034 }
1035 return NULL;
1036}
1037
1038/** Extract the rcode from a command set
1039 */
1044
1045/** Set the rcode for a command set
1046 */
1051
1052/** Extract the next node address and port from a command set
1053 */
1055{
1056 ioconf->hostname = cmds->next_node_ip;
1057 ioconf->port = cmds->next_node_port;
1058}
1059
1060/** Reset a command set to it's state before enqueuing
1061 *
1062 * For use when handling MOVED / ASK where the command set needs to be sent
1063 * to another node.
1064 */
1066{
1067 fr_redis_command_t *cmd;
1068
1069 /*
1070 * Move sent and completed commands back to the pending list
1071 * Popping from the tail of sent, then completed and inserting
1072 * into the head of pending ensures pending is back in the
1073 * original sequence.
1074 */
1075 while ((cmd = fr_dlist_pop_tail(&cmds->sent))) {
1076 fr_dlist_insert_head(&cmds->pending, cmd);
1077 }
1078 while ((cmd = fr_dlist_pop_tail(&cmds->completed))) {
1079 fr_dlist_insert_head(&cmds->pending, cmd);
1080 }
1081
1082 TALLOC_FREE(cmds->next_node_ip);
1083 cmds->next_node_port = 0;
1084 cmds->treq = NULL;
1085
1086 return 0;
1087}
1088
1089/** Reinitialise a command set so that it can be used again
1090 *
1091 * Frees the completed commands, and returns the command set to the state it was in when it was created via
1092 * fr_redis_command_set_alloc(). The command set keeps the request, the resume ctx, the callbacks, and the
1093 * autofree setting it was allocated with, so that the caller can add a new set of commands and enqueue the
1094 * command set again.
1095 *
1096 * @note The caller must have finished reading the results before calling this function. The completed
1097 * commands are freed here, as is the address of the node named by a MOVED / ASK reply.
1098 * fr_redis_command_set_rcode() reports success after this function returns, whatever the previous
1099 * execution did.
1100 *
1101 * @param[in] cmds Command set to reinitialise.
1102 * @return
1103 * - 0 on success.
1104 * - -1 if the command set is still executing, in which case nothing is changed.
1105 */
1107{
1108 /*
1109 * @todo - the run-time checks should likely be asserts, and this function should probably return
1110 * "void".
1111 */
1112 if (fr_dlist_num_elements(&cmds->pending) > 0) return -1;
1113 if (fr_dlist_num_elements(&cmds->sent) > 0) return -1;
1114
1115 /*
1116 * Free the commands, rather than just emptying the list. The commands are allocated from the
1117 * command set's pool, so leaving the commands allocated means the command set grows every time
1118 * the command set is reused.
1119 *
1120 * The reply is passed to the per command callback as the reply arrives, and is not stored on the
1121 * command, so freeing a command here cannot free a reply which the caller is still reading.
1122 */
1124
1125 /*
1126 * Clear the execution state.
1127 */
1128 TALLOC_FREE(cmds->next_node_ip);
1129 cmds->next_node_port = 0;
1130 cmds->treq = NULL;
1132 cmds->redirected = 0;
1133
1134 /*
1135 * Clear the command state. The transaction counts are checked when the command set is enqueued,
1136 * so stale counts from the previous execution would be applied to the new set of commands.
1137 */
1138 cmds->txn_watch = false;
1139 cmds->txn_start = 0;
1140 cmds->txn_end = 0;
1141 cmds->blocking = false;
1142
1143 return 0;
1144}
#define _Thread_local
Definition atexit.h:213
#define fr_atexit_thread_local(_name, _free, _uctx)
Definition atexit.h:224
#define FALL_THROUGH
clang 10 doesn't recognised the FALL-THROUGH comment anymore
Definition build.h:391
#define CC_NO_UBSAN(_sanitize)
Definition build.h:503
#define unlikely(_x)
Definition build.h:455
#define UNUSED
Definition build.h:384
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.
trunk_conf_t const * fr_redis_ct_trunk_conf(fr_redis_ct_t *rtcluster)
fr_event_list_t * fr_redis_ct_el(fr_redis_ct_t *rtcluster)
fr_redis_trunk_active_t active
Callback to run when the trunk becomes active.
Thread local state for a cluster.
Redis asynchronous cluster management.
#define KEY_SLOTS
Maximum number of keyslots (should not change).
TALLOC_CTX * autofree
Definition common.c:29
#define fr_assert_msg(_x, _msg,...)
Calls panic_action ifndef NDEBUG, else logs error and causes the server to exit immediately with code...
Definition debug.h:248
#define MEM(x)
Definition debug.h:38
#define ERROR(fmt,...)
Definition dhcpclient.c:40
#define fr_dlist_init(_head, _type, _field)
Initialise the head structure of a doubly linked list.
Definition dlist.h:242
static void * fr_dlist_head(fr_dlist_head_t const *list_head)
Return the HEAD item of a list or NULL if the list is empty.
Definition dlist.h:468
#define fr_dlist_foreach(_list_head, _type, _iter)
Iterate over the contents of a list.
Definition dlist.h:98
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_talloc_free(fr_dlist_head_t *head)
Free all items in a doubly linked list (with talloc)
Definition dlist.h:892
static void fr_dlist_entry_unlink(fr_dlist_t *entry)
Remove an item from the dlist when we don't have access to the head.
Definition dlist.h:128
static unsigned int fr_dlist_num_elements(fr_dlist_head_t const *head)
Return the number of elements in the dlist.
Definition dlist.h:921
static void * fr_dlist_pop_tail(fr_dlist_head_t *list_head)
Remove the tail item in a list.
Definition dlist.h:670
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
static int fr_dlist_move(fr_dlist_head_t *list_dst, fr_dlist_head_t *list_src)
Merge two lists, inserting the source at the tail of the destination.
Definition dlist.h:745
#define fr_dlist_talloc_init(_head, _type, _field)
Initialise the head structure of a doubly linked list.
Definition dlist.h:257
static int fr_dlist_insert_head(fr_dlist_head_t *list_head, void *ptr)
Insert an item into the head of a list.
Definition dlist.h:320
static void fr_dlist_entry_init(fr_dlist_t *entry)
Initialise a linked list without metadata.
Definition dlist.h:120
static void * fr_dlist_next(fr_dlist_head_t const *list_head, void const *ptr)
Get the next item in a list.
Definition dlist.h:537
Head of a doubly linked list.
Definition dlist.h:51
Entry in a doubly linked list.
Definition dlist.h:41
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
int af
Address family.
Definition inet.h:64
IPv4/6 prefix.
void unlang_interpret_mark_runnable(request_t *request)
Mark a request as resumable.
Definition interpret.c:2008
connection_t * fr_redis_connection_alloc(TALLOC_CTX *ctx, fr_event_list_t *el, connection_conf_t const *conn_conf, fr_redis_io_conf_t const *io_conf, char const *log_prefix)
Allocate an async redis I/O connection.
Definition io.c:610
fr_redis_sqn_t rsp_sqn
Current redis response number.
Definition io.h:95
redisAsyncContext * ac
Async handle for hiredis.
Definition io.h:84
static void fr_redis_connection_ignore_response(fr_redis_handle_t *h, fr_redis_sqn_t sqn)
Ignore a response with a specific sequence number.
Definition io.h:116
char const * hostname
Definition io.h:51
char const * log_prefix
Definition io.h:57
static bool fr_redis_connection_process_response(fr_redis_handle_t *h)
Update the response sequence number and check if we should ignore the response.
Definition io.h:147
uint16_t port
Definition io.h:52
bool freeing
Ensure that redisAsyncFree doesn't cause a callback loop.
Definition io.h:78
static fr_redis_sqn_t fr_redis_connection_sent_request(fr_redis_handle_t *h)
Tell the handle we sent a command, and get the SQN that command was assigned.
Definition io.h:106
Store I/O state.
Definition io.h:75
#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 RWARN(fmt,...)
Definition log.h:314
#define DEBUG4(_fmt,...)
Definition log.h:272
#define DEBUG_ENABLED3
True if global debug level 1-3 messages are enabled.
Definition log.h:264
Stores all information relating to an event list.
Definition event.c:377
unsigned short uint16_t
unsigned int uint32_t
unsigned char uint8_t
int strncasecmp(char *s1, char *s2, int n)
Definition missing.c:35
static const trunk_io_funcs_t io_funcs
Definition bio.c:2719
Function prototypes and datatypes for the REST (HTTP) transport.
fr_redis_async_rcode_t rcode
Code from last error returned.
Definition pipeline.c:97
char const * fr_redis_command_get_cmd(fr_redis_command_t *cmd)
Definition pipeline.c:1025
request_t * request
Request this commands set is associated with (if any).
Definition pipeline.c:125
static int _command_set_free_list_free_on_exit(void *arg)
Free any free requests when the thread is joined.
Definition pipeline.c:177
fr_redis_command_complete_t complete
Callback to process result from this command.
Definition pipeline.c:85
fr_dlist_head_t sent
Commands sent.
Definition pipeline.c:107
fr_redis_async_rcode_t fr_redis_command_set_rcode(fr_redis_command_set_t *cmds)
Extract the rcode from a command set.
Definition pipeline.c:1040
uint16_t txn_start
Number of times a transaction block was started in this command set.
Definition pipeline.c:152
fr_redis_pipeline_status_t fr_redis_command_preformatted_add(fr_redis_command_set_t *cmds, char const *cmd_str, size_t cmd_len, fr_redis_command_complete_t complete, void *rctx)
Add an preformatted command to the command set as formatted by redisCommandFormat or it's variants.
Definition pipeline.c:482
fr_redis_io_conf_t const * io_conf
Redis I/O configuration.
Definition pipeline.c:164
static void _redis_pipeline_command_set_cancel(connection_t *conn, void *preq, trunk_cancel_reason_t reason, UNUSED void *uctx)
Deal with cancellation of sent requests.
Definition pipeline.c:853
trunk_t * trunk
Trunk containing all the connections to a specific host.
Definition pipeline.c:166
fr_dlist_t entry
Entry in the command buffer.
Definition pipeline.c:64
static connection_t * _redis_pipeline_connection_alloc(trunk_connection_t *tconn, fr_event_list_t *el, connection_conf_t const *conf, char const *log_prefix, void *uctx)
Definition pipeline.c:756
void * active_uctx
Uctx to pass to active callback.
Definition pipeline.c:171
char * next_node_ip
IP address of node from MOVED / ASK reply.
Definition pipeline.c:100
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
#define COMMAND_PRE_ALLOC_COUNT
void * rctx
Resume context to write results to.
Definition pipeline.c:126
static void _redis_pipeline_command_set_fail(UNUSED request_t *request, void *preq, UNUSED void *rctx, UNUSED trunk_request_state_t state, UNUSED void *uctx)
Signal the API client that we failed enqueuing the commands.
Definition pipeline.c:928
uint64_t sqn
The sequence number of the command.
Definition pipeline.c:81
fr_redis_command_set_complete_t complete
Notify the creator of the command set that the command set has executed to to completion.
Definition pipeline.c:132
fr_redis_command_fmt_t fmt
Redis command format.
Definition pipeline.c:67
fr_redis_command_set_fail_t fail
Notify the creator of the command set that the command set failed to execute to completion.
Definition pipeline.c:137
fr_redis_trunk_active_t active
Callback to run when the trunk becomes active.
Definition pipeline.c:170
uint16_t txn_end
The number of times a transaction block ended in this command set.
Definition pipeline.c:154
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
fr_redis_ct_t * rtcluster
Cluster this trunk belongs to.
Definition pipeline.c:168
static void _redis_trunk_active(UNUSED trunk_t *trunk, UNUSED trunk_state_t prev, UNUSED trunk_state_t state, void *uctx)
Definition pipeline.c:972
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
static void _redis_pipeline_demux(struct redisAsyncContext *ac, void *vreply, void *privdata)
Callback for for receiving Redis replies.
Definition pipeline.c:635
bool autofree
Should the command set be freed when it is complete.
Definition pipeline.c:98
#define COMMAND_PRE_ALLOC_LEN
fr_redis_command_set_t * fr_redis_command_set_alloc(TALLOC_CTX *ctx, request_t *request, fr_redis_command_set_complete_t complete, fr_redis_command_set_fail_t fail, void *rctx, bool autofree)
Allocate a new command set.
Definition pipeline.c:252
static int redis_addr_from_redirect(TALLOC_CTX *ctx, char **addr, uint16_t *port, redisReply *redirect)
Convert a MOVED / ASK reply into an address and port.
Definition pipeline.c:578
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
static fr_redis_pipeline_status_t redis_command_transaction_check(request_t *request, fr_redis_command_type_t *type, fr_redis_command_set_t *cmds, char const *cmd)
Definition pipeline.c:304
static void _redis_pipeline_command_set_free(UNUSED request_t *request, void *preq, UNUSED void *uctx)
Free the command set.
Definition pipeline.c:941
static _Thread_local fr_dlist_head_t * command_set_free_list
The thread local free list.
Definition pipeline.c:40
uint16_t next_node_port
Port of node from MOVED / ASK reply.
Definition pipeline.c:101
void fr_redis_command_set_cancel(fr_redis_command_set_t *cmds)
Cancel a command set.
Definition pipeline.c:566
fr_redis_pipeline_status_t fr_redis_command_argv_add(fr_redis_command_set_t *cmds, size_t argc, char const **argv, size_t *argv_len, fr_redis_command_complete_t complete, void *rctx)
Add a command with arguments to the command set.
Definition pipeline.c:439
void * rctx
To be passed to the callback.
Definition pipeline.c:87
fr_redis_command_type_t type
Redis command type.
Definition pipeline.c:66
static int _redis_command_set_free(fr_redis_command_set_t *cmds)
Free a command set.
Definition pipeline.c:192
static void _redis_pipeline_command_set_complete(UNUSED request_t *request, void *preq, UNUSED void *rctx, UNUSED void *uctx)
Signal the API client that we got a complete set of responses to a command set.
Definition pipeline.c:916
bool txn_watch
Transaction was started with a watch statement.
Definition pipeline.c:151
fr_dlist_head_t completed
Commands complete with replies.
Definition pipeline.c:108
fr_redis_command_type_t
Definition pipeline.c:42
@ FR_REDIS_COMMAND_TRANSACTION_START
Start of a transaction block.
Definition pipeline.c:44
@ FR_REDIS_COMMAND_NORMAL
A normal, non-transactional command.
Definition pipeline.c:43
@ FR_REDIS_COMMAND_TRANSACTION_END
End of a transaction block.
Definition pipeline.c:47
uint8_t redirected
How many times this command set was redirected.
Definition pipeline.c:111
fr_redis_pipeline_status_t fr_redis_command_literal_add(fr_redis_command_set_t *cmds, char const *cmd_str, fr_redis_command_complete_t complete, void *rctx)
Add a literal command to the command set.
Definition pipeline.c:404
bool blocking
This command set contains one or more commands which block the client (e.g.
Definition pipeline.c:159
fr_redis_trunk_t * fr_redis_trunk_alloc(fr_redis_ct_t *rtcluster, fr_redis_io_conf_t const *io_conf, fr_pair_list_t *trigger_args, fr_redis_trunk_active_t active, void *active_uctx, bool active_oneshot)
Allocate a new trunk.
Definition pipeline.c:991
trunk_request_t * treq
Trunk request this command set is associated with.
Definition pipeline.c:124
int fr_redis_command_set_clear(fr_redis_command_set_t *cmds)
Reinitialise a command set so that it can be used again.
Definition pipeline.c:1106
uint32_t max_redirects
How many times this command set may be redirected.
Definition pipeline.c:112
static void _redis_pipeline_mux(UNUSED fr_event_list_t *el, trunk_connection_t *tconn, connection_t *conn, UNUSED void *uctx)
Enqueue one or more command sets onto a redis handle.
Definition pipeline.c:781
fr_dlist_head_t pending
Commands yet to be sent.
Definition pipeline.c:106
fr_redis_command_set_t * cmds
Command set this entry belongs to.
Definition pipeline.c:63
fr_redis_command_fmt_t
Definition pipeline.c:53
@ FR_REDIS_COMMAND_FMT_ARGV
A command as an argv array.
Definition pipeline.c:55
@ FR_REDIS_COMMAND_FMT_PREFORMATTED
A command preformatted with redisCommandFormat.
Definition pipeline.c:56
@ FR_REDIS_COMMAND_FMT_EXPANDED
A command as a single string.
Definition pipeline.c:54
Represents a single command.
Definition pipeline.c:62
Represents a collection of pipelined commands.
Definition pipeline.c:94
Redis asynchronous command pipelining.
void(* fr_redis_command_set_complete_t)(request_t *request, fr_dlist_head_t *completed, void *rctx)
Do something meaningful with the replies to the commands previously issued.
Definition pipeline.h:65
void(* fr_redis_command_complete_t)(request_t *request, fr_redis_command_t *cmd, redisReply *reply, void *rctx)
Process the reply from a single command.
Definition pipeline.h:60
void(* fr_redis_command_set_fail_t)(request_t *request, fr_dlist_head_t *completed, void *rctx)
Write a failure result to the rctx so that the module is aware that the request failed.
Definition pipeline.h:70
struct fr_redis_trunk_s fr_redis_trunk_t
Definition pipeline.h:53
fr_redis_pipeline_status_t
Definition pipeline.h:43
@ FR_REDIS_PIPELINE_OK
No failure.
Definition pipeline.h:44
@ FR_REDIS_PIPELINE_BAD_CMDS
Malformed command set.
Definition pipeline.h:45
@ FR_REDIS_PIPELINE_DST_UNAVAILABLE
Cluster or host is down.
Definition pipeline.h:46
@ FR_REDIS_PIPELINE_FAIL
Generic failure.
Definition pipeline.h:48
void(* fr_redis_trunk_active_t)(fr_redis_trunk_t *rtrunk, void *uctx)
Definition pipeline.h:55
#define fr_assert(_expr)
Definition rad_assert.h:37
#define REDEBUG(fmt,...)
#define WARN(fmt,...)
static rs_t * conf
Definition radsniff.c:52
#define REDIS_ERROR_TRY_AGAIN_STR
Definition base.h:49
fr_redis_async_rcode_t
Definition base.h:80
@ REDIS_ASYNC_RCODE_MOVE
Attempt operation on an alternative node with remap.
Definition base.h:88
@ REDIS_ASYNC_RCODE_ERROR
Unrecoverable error.
Definition base.h:82
@ REDIS_ASYNC_RCODE_ASK
Attempt operation on an alternative node.
Definition base.h:87
@ 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_NO_SCRIPT
Script doesn't exist.
Definition base.h:89
@ REDIS_ASYNC_RCODE_SUCCESS
Operation was successful.
Definition base.h:81
#define REDIS_ERROR_MOVED_STR
Definition base.h:47
#define REDIS_ERROR_ASK_STR
Definition base.h:48
#define REDIS_ERROR_NO_SCRIPT_STR
Definition base.h:50
bool fr_sbuff_next_if_char(fr_sbuff_t *sbuff, char c)
Return true if the current char matches, and if it does, advance.
Definition sbuff.c:2247
#define fr_sbuff_adv_past_str_literal(_sbuff, _needle)
#define fr_sbuff_current(_sbuff_or_marker)
#define fr_sbuff_out(_out, _in)
#define fr_sbuff_init_in(_out, _start, _len_or_end)
#define fr_sbuff_remaining(_sbuff_or_marker)
@ CONNECTION_STATE_CONNECTED
File descriptor is open (ready for writing).
Definition connection.h:54
static char buff[sizeof("18446744073709551615")+3]
Definition size_tests.c:37
fr_aka_sim_id_type_t type
int talloc_link_ctx(TALLOC_CTX *parent, TALLOC_CTX *child)
Link two different parent and child contexts, so the child is freed before the parent.
Definition talloc.c:168
#define talloc_zero_pooled_object(_ctx, _type, _num_subobjects, _total_subobjects_size)
Definition talloc.h:208
#define talloc_strdup(_ctx, _str)
Definition talloc.h:149
void trunk_request_signal_fail(trunk_request_t *treq)
Signal that a trunk request failed.
Definition trunk.c:2196
trunk_watch_entry_t * trunk_add_watch(trunk_t *trunk, trunk_state_t state, trunk_watch_t watch, bool oneshot, void const *uctx)
Add a watch entry to the trunk state list.
Definition trunk.c:915
trunk_enqueue_t trunk_request_enqueue(trunk_request_t **treq_out, trunk_t *trunk, request_t *request, void *preq, void *rctx)
Enqueue a request that needs data written to the trunk.
Definition trunk.c:2657
int trunk_connection_pop_request(trunk_request_t **treq_out, trunk_connection_t *tconn)
Pop a request off a connection's pending queue.
Definition trunk.c:3979
void trunk_request_signal_cancel(trunk_request_t *treq)
Cancel a trunk request.
Definition trunk.c:2216
trunk_t * trunk_alloc(TALLOC_CTX *ctx, fr_event_list_t *el, trunk_io_funcs_t const *funcs, trunk_conf_t const *conf, char const *log_prefix, void const *uctx, bool delay_start, fr_pair_list_t *trigger_args)
Allocate a new collection of connections.
Definition trunk.c:5124
void trunk_request_mark_blocking(trunk_request_t *treq)
Mark a trunk request as one which will block the connection until it is completed.
Definition trunk.c:2859
void trunk_request_signal_sent(trunk_request_t *treq)
Signal that the request was written to a connection successfully.
Definition trunk.c:2114
void trunk_request_signal_complete(trunk_request_t *treq)
Signal that a trunk request is complete.
Definition trunk.c:2158
Associates request queues with a connection.
Definition trunk.c:137
Wraps a normal request.
Definition trunk.c:99
Main trunk management handle.
Definition trunk.c:219
trunk_connection_alloc_t connection_alloc
Allocate a new connection_t.
Definition trunk.h:747
trunk_cancel_reason_t
Reasons for a request being cancelled.
Definition trunk.h:55
@ TRUNK_CANCEL_REASON_NONE
Request has not been cancelled.
Definition trunk.h:56
@ TRUNK_CANCEL_REASON_SIGNAL
Request cancelled due to a signal.
Definition trunk.h:57
@ TRUNK_CANCEL_REASON_REQUEUE
A previously sent request is being requeued.
Definition trunk.h:59
@ TRUNK_CANCEL_REASON_MOVE
Request cancelled because it's being moved.
Definition trunk.h:58
trunk_state_t
Definition trunk.h:62
@ TRUNK_STATE_ACTIVE
Trunk has at least one active connection which can service requests.
Definition trunk.h:64
@ TRUNK_ENQUEUE_DST_UNAVAILABLE
Destination is down.
Definition trunk.h:163
@ TRUNK_ENQUEUE_OK
Operation was successful.
Definition trunk.h:160
@ TRUNK_ENQUEUE_IN_BACKLOG
Request should be enqueued in backlog.
Definition trunk.h:159
trunk_request_state_t
Used for sanity checks and to simplify freeing.
Definition trunk.h:171
@ TRUNK_REQUEST_STATE_BACKLOG
In the backlog.
Definition trunk.h:177
@ TRUNK_REQUEST_STATE_PENDING
In the queue of a connection and is pending writing.
Definition trunk.h:178
@ TRUNK_REQUEST_STATE_SENT
Was written to a socket. Waiting for a response.
Definition trunk.h:182
I/O functions to pass to trunk_alloc.
Definition trunk.h:746
static fr_event_list_t * el
#define fr_strerror_printf(_fmt,...)
Log to thread local error buffer.
Definition strerror.h:64
#define fr_strerror_const(_msg)
Definition strerror.h:223
#define fr_box_strvalue_len(_val, _len)
Definition value.h:334