From: Arran Cudbard-Bell Date: Thu, 5 Dec 2019 11:57:23 +0000 (+0700) Subject: Add many more tests X-Git-Url: http://git.ipfire.org/cgi-bin/gitweb.cgi?a=commitdiff_plain;h=3bae3bd6b794ea42397a82b60bd468a23a055ce8;p=thirdparty%2Ffreeradius-server.git Add many more tests Add cancel partial state for dealing with partial writes when cancelling Restore in_handler correctly when doing nested handler calls. We *may* need to move this to a stack, but lets not do that unless we have to. When requing requests, free the cancelled ones, don't try and move them to a new connection *sigh* --- diff --git a/src/lib/server/trunk.c b/src/lib/server/trunk.c index 069e6aa50eb..07d2016e153 100644 --- a/src/lib/server/trunk.c +++ b/src/lib/server/trunk.c @@ -58,7 +58,7 @@ static fr_time_t test_time(void) * Allows us to track which */ typedef enum { - FR_TRUNK_REQUEST_UNASSIGNED = 0x0000, //!< Initial state. + FR_TRUNK_REQUEST_UNASSIGNED = 0x0000, //!< Initial state. FR_TRUNK_REQUEST_BACKLOG = 0x0001, //!< In the backlog. FR_TRUNK_REQUEST_PENDING = 0x0002, //!< In the queue of a connection ///< and is pending writing. @@ -70,7 +70,8 @@ typedef enum { FR_TRUNK_REQUEST_CANCEL = 0x0200, //!< A request on a particular socket was cancel. FR_TRUNK_REQUEST_CANCEL_SENT = 0x0400, //!< We've informed the remote server that ///< the request has been cancelled. - FR_TRUNK_REQUEST_CANCEL_COMPLETE= 0x0800, //!< Remote server has acknowledged our cancellation. + FR_TRUNK_REQUEST_CANCEL_PARTIAL = 0x0800, //!< We partially wrote a cancellation request. + FR_TRUNK_REQUEST_CANCEL_COMPLETE= 0x1000, //!< Remote server has acknowledged our cancellation. } fr_trunk_request_state_t; /** All request states @@ -85,6 +86,18 @@ typedef enum { FR_TRUNK_REQUEST_COMPLETE | \ FR_TRUNK_REQUEST_FAILED | \ FR_TRUNK_REQUEST_CANCEL | \ + FR_TRUNK_REQUEST_CANCEL_PARTIAL | \ + FR_TRUNK_REQUEST_CANCEL_SENT | \ + FR_TRUNK_REQUEST_CANCEL_COMPLETE \ +) + +/** All requests in various cancellation states + * + */ +#define FR_TRUNK_REQUEST_CANCEL_ALL \ +(\ + FR_TRUNK_REQUEST_CANCEL | \ + FR_TRUNK_REQUEST_CANCEL_PARTIAL | \ FR_TRUNK_REQUEST_CANCEL_SENT | \ FR_TRUNK_REQUEST_CANCEL_COMPLETE \ ) @@ -173,6 +186,8 @@ struct fr_trunk_connection_s { fr_dlist_head_t cancel; //!< Requests in the cancel state. + fr_trunk_request_t *cancel_partial; //!< Partially written cancellation request. + fr_dlist_head_t cancel_sent; //!< Requests we need to inform a remote server about. bool signalled_full; //!< Connection marked full because of signal. @@ -266,6 +281,7 @@ static fr_table_num_ordered_t const fr_trunk_request_states[] = { { "FAILED", FR_TRUNK_REQUEST_FAILED }, { "CANCEL", FR_TRUNK_REQUEST_CANCEL }, { "CANCEL-SENT", FR_TRUNK_REQUEST_CANCEL_SENT }, + { "CANCEL-PARTIAL", FR_TRUNK_REQUEST_CANCEL_PARTIAL }, { "CANCEL-COMPLETE", FR_TRUNK_REQUEST_CANCEL_COMPLETE} }; static size_t fr_trunk_request_states_len = NUM_ELEMENTS(fr_trunk_request_states); @@ -281,9 +297,9 @@ static fr_table_num_ordered_t const fr_trunk_connection_states[] = { static size_t fr_trunk_connection_states_len = NUM_ELEMENTS(fr_trunk_connection_states); static fr_table_num_ordered_t const fr_trunk_cancellation_reasons[] = { - { "none", FR_TRUNK_CANCEL_REASON_NONE }, - { "signal", FR_TRUNK_CANCEL_REASON_SIGNAL }, - { "move", FR_TRUNK_CANCEL_REASON_MOVE }, + { "FR_TRUNK_CANCEL_REASON_NONE", FR_TRUNK_CANCEL_REASON_NONE }, + { "FR_TRUNK_CANCEL_REASON_SIGNAL", FR_TRUNK_CANCEL_REASON_SIGNAL }, + { "FR_TRUNK_CANCEL_REASON_MOVE", FR_TRUNK_CANCEL_REASON_MOVE }, }; static size_t fr_trunk_cancellation_reasons_len = NUM_ELEMENTS(fr_trunk_cancellation_reasons); @@ -345,10 +361,11 @@ typedef enum { #define DO_REQUEST_CANCEL(_treq, _reason) \ do { \ if ((_treq)->trunk->funcs.request_cancel) { \ + void *prev = (_treq)->trunk->in_handler; \ (_treq)->trunk->in_handler = (void *)(_treq)->trunk->funcs.request_cancel; \ - DEBUG4("Calling funcs.request_cancel(%p, %p, %s, %p)", (_treq)->tconn->conn, (_treq), fr_table_str_by_value(fr_trunk_cancellation_reasons, (_reason), ""), (_treq)->trunk->uctx); \ - (_treq)->trunk->funcs.request_cancel((_treq)->tconn->conn, (_treq), (_reason), (_treq)->trunk->uctx); \ - (_treq)->trunk->in_handler = NULL; \ + DEBUG4("Calling funcs.request_cancel(%p, %p, %p, %s, %p)", (_treq)->tconn->conn, (_treq), (_treq)->preq, fr_table_str_by_value(fr_trunk_cancellation_reasons, (_reason), ""), (_treq)->trunk->uctx); \ + (_treq)->trunk->funcs.request_cancel((_treq)->tconn->conn, (_treq), (_treq)->preq, (_reason), (_treq)->trunk->uctx); \ + (_treq)->trunk->in_handler = prev; \ } \ } while(0) @@ -358,10 +375,11 @@ do { \ #define DO_REQUEST_COMPLETE(_treq) \ do { \ if ((_treq)->trunk->funcs.request_complete) { \ + void *prev = (_treq)->trunk->in_handler; \ DEBUG4("Calling funcs.request_complete(%p, %p, %p, %p)", (_treq)->request, (_treq)->preq, (_treq)->rctx, (_treq)->trunk->uctx); \ (_treq)->trunk->in_handler = (void *)(_treq)->trunk->funcs.request_complete; \ (_treq)->trunk->funcs.request_complete((_treq)->request, (_treq)->preq, (_treq)->rctx, (_treq)->trunk->uctx); \ - (_treq)->trunk->in_handler = NULL; \ + (_treq)->trunk->in_handler = prev; \ } \ } while(0) @@ -371,10 +389,11 @@ do { \ #define DO_REQUEST_FAIL(_treq) \ do { \ if ((_treq)->trunk->funcs.request_fail) { \ + void *prev = (_treq)->trunk->in_handler; \ DEBUG4("Calling funcs.request_fail(%p, %p, %p, %p)", (_treq)->request, (_treq)->preq, (_treq)->rctx, (_treq)->trunk->uctx); \ (_treq)->trunk->in_handler = (void *)(_treq)->trunk->funcs.request_fail; \ (_treq)->trunk->funcs.request_fail((_treq)->request, (_treq)->preq, (_treq)->rctx, (_treq)->trunk->uctx); \ - (_treq)->trunk->in_handler = NULL; \ + (_treq)->trunk->in_handler = prev; \ } \ } while(0) @@ -384,10 +403,11 @@ do { \ #define DO_REQUEST_FREE(_treq) \ do { \ if ((_treq)->trunk->funcs.request_free) { \ + void *prev = (_treq)->trunk->in_handler; \ DEBUG4("Calling funcs.request_free(%p, %p, %p)", (_treq)->request, (_treq)->preq, (_treq)->trunk->uctx); \ (_treq)->trunk->in_handler = (void *)(_treq)->trunk->funcs.request_free; \ (_treq)->trunk->funcs.request_free((_treq)->request, (_treq)->preq, (_treq)->trunk->uctx); \ - (_treq)->trunk->in_handler = NULL; \ + (_treq)->trunk->in_handler = prev; \ } \ } while(0) @@ -396,10 +416,11 @@ do { \ */ #define DO_REQUEST_MUX(_tconn) \ do { \ + void *prev = (_tconn)->trunk->in_handler; \ DEBUG4("Calling funcs.request_mux(%p, %p, %p)", (_tconn), (_tconn)->conn, (_tconn)->trunk->uctx); \ (_tconn)->trunk->in_handler = (void *)(_tconn)->trunk->funcs.request_mux; \ (_tconn)->trunk->funcs.request_mux((_tconn), (_tconn)->conn, (_tconn)->trunk->uctx); \ - (_tconn)->trunk->in_handler = NULL; \ + (_tconn)->trunk->in_handler = prev; \ } while(0) /** Read one or more requests from a connection @@ -407,10 +428,11 @@ do { \ */ #define DO_REQUEST_DEMUX(_tconn) \ do { \ + void *prev = (_tconn)->trunk->in_handler; \ DEBUG4("Calling funcs.request_demux(%p, %p, %p)", (_tconn), (_tconn)->conn, (_tconn)->trunk->uctx); \ (_tconn)->trunk->in_handler = (void *)(_tconn)->trunk->funcs.request_demux; \ (_tconn)->trunk->funcs.request_demux((_tconn), (_tconn)->conn, (_tconn)->trunk->uctx); \ - (_tconn)->trunk->in_handler = NULL; \ + (_tconn)->trunk->in_handler = prev; \ } while(0) /** Write one or more cancellation requests to a connection @@ -419,10 +441,11 @@ do { \ #define DO_REQUEST_CANCEL_MUX(_tconn) \ do { \ if ((_tconn)->trunk->funcs.request_cancel_mux) { \ + void *prev = (_tconn)->trunk->in_handler; \ DEBUG4("Calling funcs.request_cancel_mux(%p, %p, %p)", (_tconn), (_tconn)->conn, (_tconn)->trunk->uctx); \ (_tconn)->trunk->in_handler = (void *)(_tconn)->trunk->funcs.request_cancel_mux; \ (_tconn)->trunk->funcs.request_cancel_mux((_tconn), (_tconn)->conn, (_tconn)->trunk->uctx); \ - (_tconn)->trunk->in_handler = NULL; \ + (_tconn)->trunk->in_handler = prev; \ } \ } while(0) @@ -432,12 +455,13 @@ do { \ #define DO_CONNECTION_NOTIFY(_tconn, _events) \ do { \ if ((_tconn)->trunk->funcs.connection_notify) { \ + void *prev = (_tconn)->trunk->in_handler; \ DEBUG4("Calling funcs.connection_notify(%p, %p, %p, %s, %p)", \ (_tconn), (_tconn)->conn, (_tconn)->trunk->el, \ fr_table_str_by_value(fr_trunk_connection_events, (_events), ""), (_tconn)->trunk->uctx); \ (_tconn)->trunk->in_handler = (void *)(_tconn)->trunk->funcs.connection_notify; \ (_tconn)->trunk->funcs.connection_notify((_tconn), (_tconn)->conn, (_tconn)->trunk->el, (_events), (_tconn)->trunk->uctx); \ - (_tconn)->trunk->in_handler = NULL; \ + (_tconn)->trunk->in_handler = prev; \ } \ } while(0) @@ -446,13 +470,14 @@ do { \ #define IN_REQUEST_DEMUX(_trunk) (((_trunk)->funcs.request_demux) && ((_trunk)->in_handler == (void *)(_trunk)->funcs.request_demux)) #define IN_REQUEST_CANCEL_MUX(_trunk) (((_trunk)->funcs.request_cancel_mux) && ((_trunk)->in_handler == (void *)(_trunk)->funcs.request_cancel_mux)) -/** Remove the current request from the partial slot +/** Remove the current request from the backlog * */ -#define REQUEST_EXTRACT_PARTIAL(_treq) \ +#define REQUEST_EXTRACT_BACKLOG(_treq) \ do { \ - rad_assert((_treq)->tconn->partial == treq); \ - tconn->partial = NULL; \ + int _ret; \ + _ret = fr_heap_extract((_treq)->trunk->backlog, _treq); \ + if (!fr_cond_assert(_ret == 0)) return; \ } while (0) /** Remove the current request from the pending list @@ -465,16 +490,26 @@ do { \ if (!fr_cond_assert(_ret == 0)) return; \ } while (0) -/** Remove the current request from the backlog +/** Remove the current request from the partial slot * */ -#define REQUEST_EXTRACT_BACKLOG(_treq) \ +#define REQUEST_EXTRACT_PARTIAL(_treq) \ do { \ - int _ret; \ - _ret = fr_heap_extract((_treq)->trunk->backlog, _treq); \ - if (!fr_cond_assert(_ret == 0)) return; \ + rad_assert((_treq)->tconn->partial == treq); \ + tconn->partial = NULL; \ +} while (0) + + +/** Remove the current request from the cancel_partial slot + * + */ +#define REQUEST_EXTRACT_CANCEL_PARTIAL(_treq) \ +do { \ + rad_assert((_treq)->tconn->cancel_partial == treq); \ + tconn->cancel_partial = NULL; \ } while (0) + /** Reorder the connections in the active heap * */ @@ -515,9 +550,15 @@ static void trunk_backlog_drain(fr_trunk_t *trunk); /** Remove a request from all connection lists * - * A common function used by init, fail, complete state functions - * to disassociate a request from a connection in preparation for - * freeing or reassignment. + * A common function used by init, fail, complete state functions to disassociate + * a request from a connection in preparation for freeing or reassignment. + * + * Despite its unassuming name, this function is *the* place to put calls to + * functions which need to be called when the number of requests associated with + * a connection changes. + * + * Trunk requests will always be passed to this function before they're removed + * from a connection, even if the requests are being freed. * * @param[in] treq to trigger a state change for. */ @@ -532,10 +573,6 @@ static void trunk_request_remove_from_conn(fr_trunk_request_t *treq) case FR_TRUNK_REQUEST_UNASSIGNED: return; /* Not associated with connection */ - case FR_TRUNK_REQUEST_BACKLOG: - REQUEST_EXTRACT_BACKLOG(treq); - break; - case FR_TRUNK_REQUEST_PENDING: REQUEST_EXTRACT_PENDING(treq); break; @@ -552,6 +589,10 @@ static void trunk_request_remove_from_conn(fr_trunk_request_t *treq) fr_dlist_remove(&tconn->cancel, treq); break; + case FR_TRUNK_REQUEST_CANCEL_PARTIAL: + REQUEST_EXTRACT_CANCEL_PARTIAL(treq); + break; + case FR_TRUNK_REQUEST_CANCEL_SENT: fr_dlist_remove(&tconn->cancel_sent, treq); break; @@ -588,7 +629,7 @@ static void trunk_request_remove_from_conn(fr_trunk_request_t *treq) * * @param[in] treq to trigger a state change for. */ -static void trunk_request_enter_init(fr_trunk_request_t *treq) +static void trunk_request_enter_unassigned(fr_trunk_request_t *treq) { fr_trunk_t *trunk = treq->trunk; @@ -647,13 +688,20 @@ static void trunk_request_enter_backlog(fr_trunk_request_t *treq) * function immediately. * * Likewise, if there's draining connections - * which could be moved back to active. + * which could be moved back to active call + * the trunk manage function. + * + * Remember requests only enter the backlog if + * there's no connections which can service them. */ if ((fr_trunk_connection_count(treq->trunk, FR_TRUNK_CONN_CONNECTING) == 0) || (fr_trunk_connection_count(treq->trunk, FR_TRUNK_CONN_DRAINING) > 0)) trunk_manage(treq->trunk, fr_time()); } /** Transition a request to the pending state, adding it to the backlog of an active connection + * + * All trunk requests being added to a connection get passed to this function. + * All trunk requests being removed from a connection get passed to #trunk_request_remove_from_conn. * * @param[in] treq to trigger a state change for. * @param[in] tconn to enqueue the request on. @@ -693,14 +741,14 @@ static void trunk_request_enter_pending(fr_trunk_request_t *treq, fr_trunk_conne trunk_connection_auto_full(tconn); /* - * Reorder the connection in the heap - * now it has an additional request. + * Reorder the connection in the heap now it has an + * additional request. */ CONN_REORDER(tconn); /* - * We have a new request, see if we need - * to register for I/O events. + * We have a new request, see if we need to register + * for I/O events. */ trunk_connection_event_update(tconn); } @@ -760,6 +808,13 @@ static void trunk_request_enter_sent(fr_trunk_request_t *treq) REQUEST_STATE_TRANSITION(FR_TRUNK_REQUEST_SENT); fr_dlist_insert_tail(&tconn->sent, treq); + + /* + * We just sent a request, we probably need + * to tell the event loop we want to be + * notified if there's data available. + */ + trunk_connection_event_update(tconn); } /** Transition a request to the cancel state, placing it in a connection's cancellation list @@ -817,6 +872,37 @@ static void trunk_request_enter_cancel(fr_trunk_request_t *treq, fr_trunk_cancel if (treq->cancel_reason == FR_TRUNK_CANCEL_REASON_SIGNAL) treq->request = NULL; } +/** Transition a request to the cancel_partial state, placing it in a connection's cancel_partial slot + * + * The request_demux function is then responsible for signalling + * that the cancel request is complete when the remote server + * acknowledges the cancellation request. + * + * @param[in] treq to trigger a state change for. + */ +static void trunk_request_enter_cancel_partial(fr_trunk_request_t *treq) +{ + fr_trunk_connection_t *tconn = treq->tconn; + fr_trunk_t *trunk = treq->trunk; + + if (!fr_cond_assert(!tconn || (tconn->trunk == trunk))) return; + rad_assert(trunk->funcs.request_cancel_mux); + rad_assert(treq->cancel_reason == FR_TRUNK_CANCEL_REASON_SIGNAL); + + switch (treq->state) { + case FR_TRUNK_REQUEST_CANCEL: /* The only valid state cancel_sent can be reached from */ + fr_dlist_remove(&tconn->cancel, treq); + break; + + default: + REQUEST_BAD_STATE_TRANSITION(FR_TRUNK_REQUEST_CANCEL_PARTIAL); + } + + REQUEST_STATE_TRANSITION(FR_TRUNK_REQUEST_CANCEL_PARTIAL); + rad_assert(!tconn->partial); + tconn->cancel_partial = treq; +} + /** Transition a request to the cancel_sent state, placing it in a connection's cancel_sent list * * The request_demux function is then responsible for signalling @@ -827,15 +913,19 @@ static void trunk_request_enter_cancel(fr_trunk_request_t *treq, fr_trunk_cancel */ static void trunk_request_enter_cancel_sent(fr_trunk_request_t *treq) { - fr_trunk_connection_t *tconn = treq->tconn; - fr_trunk_t *trunk = treq->trunk; + fr_trunk_connection_t *tconn = treq->tconn; + fr_trunk_t *trunk = treq->trunk; if (!fr_cond_assert(!tconn || (tconn->trunk == trunk))) return; rad_assert(trunk->funcs.request_cancel_mux); rad_assert(treq->cancel_reason == FR_TRUNK_CANCEL_REASON_SIGNAL); switch (treq->state) { - case FR_TRUNK_REQUEST_CANCEL: /* The only valid state cancel_sent can be reached from */ + case FR_TRUNK_REQUEST_CANCEL_PARTIAL: + REQUEST_EXTRACT_CANCEL_PARTIAL(treq); + break; + + case FR_TRUNK_REQUEST_CANCEL: fr_dlist_remove(&tconn->cancel, treq); break; @@ -863,7 +953,7 @@ static void trunk_request_enter_cancel_complete(fr_trunk_request_t *treq) if (!fr_cond_assert(!treq->request)) return; /* Only a valid state for REQUEST * which have been cancelled */ switch (treq->state) { - case FR_TRUNK_REQUEST_CANCEL_SENT: /* The only valid state cancel_sent can be reached from */ + case FR_TRUNK_REQUEST_CANCEL_SENT: fr_dlist_remove(&tconn->cancel_sent, treq); break; @@ -883,12 +973,19 @@ static void trunk_request_enter_cancel_complete(fr_trunk_request_t *treq) */ static void trunk_request_enter_complete(fr_trunk_request_t *treq) { - fr_trunk_connection_t *tconn = treq->tconn; - fr_trunk_t *trunk = treq->trunk; + fr_trunk_connection_t *tconn = treq->tconn; + fr_trunk_t *trunk = treq->trunk; if (!fr_cond_assert(!tconn || (tconn->trunk == trunk))) return; - trunk_request_remove_from_conn(treq); + switch (treq->state) { + case FR_TRUNK_REQUEST_SENT: + trunk_request_remove_from_conn(treq); + break; + + default: + REQUEST_BAD_STATE_TRANSITION(FR_TRUNK_REQUEST_COMPLETE); + } REQUEST_STATE_TRANSITION(FR_TRUNK_REQUEST_COMPLETE); DO_REQUEST_COMPLETE(treq); @@ -901,12 +998,20 @@ static void trunk_request_enter_complete(fr_trunk_request_t *treq) */ static void trunk_request_enter_failed(fr_trunk_request_t *treq) { - fr_trunk_connection_t *tconn = treq->tconn; - fr_trunk_t *trunk = treq->trunk; + fr_trunk_connection_t *tconn = treq->tconn; + fr_trunk_t *trunk = treq->trunk; if (!fr_cond_assert(!tconn || (tconn->trunk == trunk))) return; - trunk_request_remove_from_conn(treq); + switch (treq->state) { + case FR_TRUNK_REQUEST_BACKLOG: + REQUEST_EXTRACT_BACKLOG(treq); + break; + + default: + trunk_request_remove_from_conn(treq); + break; + } REQUEST_STATE_TRANSITION(FR_TRUNK_REQUEST_FAILED); DO_REQUEST_FAIL(treq); @@ -1024,6 +1129,16 @@ static fr_trunk_enqueue_t trunk_request_enqueue_existing(fr_trunk_request_t *tre /** Shift requests in the specified states onto new connections * + * This function will blindly dequeue any requests in the specified state and get + * them back to the unassigned state, cancelling any sent or partially sent requests. + * + * This function does not check that dequeuing a request in a particular state is a + * sane or sensible thing to do, that's up to the caller! + * + * @param[out] out A list to insert the newly dequeued and unassigned + * requests into. + * @param[in] tconn to dequeue requests from. + * @param[in] states Dequeue request in these states. */ static void trunk_connection_requests_dequeue(fr_dlist_head_t *out, fr_trunk_connection_t *tconn, int states) { @@ -1031,7 +1146,7 @@ static void trunk_connection_requests_dequeue(fr_dlist_head_t *out, fr_trunk_con #define DEQUEUE_ALL(_src_list) \ while ((treq = fr_dlist_head(_src_list))) { \ - trunk_request_enter_init(treq); \ + trunk_request_enter_unassigned(treq); \ fr_dlist_insert_tail(out, treq); \ } @@ -1046,12 +1161,23 @@ static void trunk_connection_requests_dequeue(fr_dlist_head_t *out, fr_trunk_con */ if (states & FR_TRUNK_REQUEST_CANCEL_SENT) DEQUEUE_ALL(&tconn->cancel_sent); + /* + * ....same with cancel partial + */ + if (states & FR_TRUNK_REQUEST_CANCEL_PARTIAL) { + treq = tconn->cancel_partial; + if (treq) { + trunk_request_enter_unassigned(treq); + fr_dlist_insert_tail(out, treq); + } + } + /* * ...and pending. */ if (states & FR_TRUNK_REQUEST_PENDING) { while ((treq = fr_heap_peek(tconn->pending))) { - trunk_request_enter_init(treq); + trunk_request_enter_unassigned(treq); fr_dlist_insert_tail(out, treq); } } @@ -1063,7 +1189,8 @@ static void trunk_connection_requests_dequeue(fr_dlist_head_t *out, fr_trunk_con treq = tconn->partial; if (treq) { trunk_request_enter_cancel(treq, FR_TRUNK_CANCEL_REASON_MOVE); - trunk_request_enter_init(treq); + trunk_request_enter_unassigned(treq); + fr_dlist_insert_tail(out, treq); } } @@ -1073,8 +1200,9 @@ static void trunk_connection_requests_dequeue(fr_dlist_head_t *out, fr_trunk_con if (states & FR_TRUNK_REQUEST_SENT) { while ((treq = fr_dlist_head(&tconn->sent))) { trunk_request_enter_cancel(treq, FR_TRUNK_CANCEL_REASON_MOVE); + trunk_request_enter_unassigned(treq); + fr_dlist_insert_tail(out, treq); } - DEQUEUE_ALL(&tconn->cancel); /* Dequeue the now cancelled requests */ } } @@ -1085,24 +1213,24 @@ static void trunk_connection_requests_dequeue(fr_dlist_head_t *out, fr_trunk_con */ static void trunk_connection_requests_requeue(fr_trunk_connection_t *tconn, int states) { - fr_dlist_head_t to_requeue; - fr_trunk_request_t *treq = NULL; + fr_dlist_head_t to_process; + fr_trunk_request_t *treq = NULL; - fr_dlist_talloc_init(&to_requeue, fr_trunk_request_t, list); + fr_dlist_talloc_init(&to_process, fr_trunk_request_t, list); /* - * Remove requests from the connection + * Remove non-cancelled requests from the connection */ - trunk_connection_requests_dequeue(&to_requeue, tconn, states); + trunk_connection_requests_dequeue(&to_process, tconn, states & ~FR_TRUNK_REQUEST_CANCEL_ALL); /* * Loop over all the requests we gathered * and redistribute them to new connections. */ - while ((treq = fr_dlist_next(&to_requeue, treq))) { + while ((treq = fr_dlist_next(&to_process, treq))) { fr_trunk_request_t *prev; - prev = fr_dlist_remove(&to_requeue, treq); + prev = fr_dlist_remove(&to_process, treq); switch (trunk_request_enqueue_existing(treq)) { case TRUNK_ENQUEUE_OK: break; @@ -1130,6 +1258,24 @@ static void trunk_connection_requests_requeue(fr_trunk_connection_t *tconn, int treq = prev; } + /* + * Deal with the cancelled requests specially we can't + * queue them up again as they were only valid on that + * specific connection. + * + * We just need to run them to completion which, as + * they should already be in the unassigned state, + * just means freeing them. + */ + trunk_connection_requests_dequeue(&to_process, tconn, states & FR_TRUNK_REQUEST_CANCEL_ALL); + while ((treq = fr_dlist_next(&to_process, treq))) { + fr_trunk_request_t *prev; + + prev = fr_dlist_remove(&to_process, treq); + talloc_free(treq); + treq = prev; + } + /* * If the trunk was draining, it wasn't counted * in the requests per connection stats, so @@ -1158,7 +1304,6 @@ static int _trunk_request_free(fr_trunk_request_t *treq) case FR_TRUNK_REQUEST_FAILED: break; - case FR_TRUNK_REQUEST_CANCEL: case FR_TRUNK_REQUEST_CANCEL_SENT: break; @@ -1204,6 +1349,7 @@ static fr_trunk_request_t *trunk_request_alloc(fr_trunk_t *trunk, REQUEST *reque MEM(treq = talloc_zero(request, fr_trunk_request_t)); treq->id = atomic_fetch_add_explicit(&request_counter, 1, memory_order_relaxed); treq->trunk = trunk; + treq->request = request; treq->preq = preq; treq->rctx = rctx; treq->state = FR_TRUNK_REQUEST_UNASSIGNED; @@ -1232,7 +1378,7 @@ void fr_trunk_request_signal_partial(fr_trunk_request_t *treq) "%s can only be called from within request_mux handler", __FUNCTION__)) return; switch (treq->state) { - case FR_TRUNK_REQUEST_SENT: + case FR_TRUNK_REQUEST_PENDING: trunk_request_enter_partial(treq); break; @@ -1251,7 +1397,8 @@ void fr_trunk_request_signal_sent(fr_trunk_request_t *treq) "%s can only be called from within request_mux handler", __FUNCTION__)) return; switch (treq->state) { - case FR_TRUNK_REQUEST_SENT: + case FR_TRUNK_REQUEST_PENDING: + case FR_TRUNK_REQUEST_PARTIAL: trunk_request_enter_sent(treq); break; @@ -1308,13 +1455,68 @@ void fr_trunk_request_signal_cancel(fr_trunk_request_t *treq) "%s cannot be called within a handler", __FUNCTION__)) return; switch (treq->state) { + /* + * We don't call the complete or failed callbacks + * as the request and rctx are no longer viable. + */ case FR_TRUNK_REQUEST_PARTIAL: case FR_TRUNK_REQUEST_SENT: trunk_request_enter_cancel(treq, FR_TRUNK_CANCEL_REASON_SIGNAL); + + switch (treq->state) { + /* + * Need to give this a chance to complete. + */ + case FR_TRUNK_REQUEST_CANCEL_SENT: + break; + + /* + * The cancel callback didn't move + * the request into the FR_TRUNK_REQUEST_CANCEL_SENT + * state so we're done. + */ + case FR_TRUNK_REQUEST_CANCEL: + trunk_request_enter_unassigned(treq); + talloc_free(treq); + break; + + /* + * Shouldn't be in any other state after this + */ + default: + rad_assert(0); + } break; + /* + * For any other state, we just release the request + * from its current connection and free it. + */ default: + trunk_request_enter_unassigned(treq); + talloc_free(treq); + break; + } +} + +/** Signal a partial cancel write + * + * Where there's high load, and the outbound write buffer is full + * + * @param[in] treq to signal state change for. + */ +void fr_trunk_request_signal_cancel_partial(fr_trunk_request_t *treq) +{ + if (!fr_cond_assert_msg(IN_REQUEST_CANCEL_MUX(treq->trunk), + "%s can only be called from within request_cancel_mux handler", __FUNCTION__)) return; + + switch (treq->state) { + case FR_TRUNK_REQUEST_CANCEL: + trunk_request_enter_cancel_partial(treq); break; + + default: + return; } } @@ -1431,6 +1633,7 @@ static inline uint32_t trunk_requests_by_connection_count(fr_trunk_connection_t if (state & FR_TRUNK_REQUEST_PARTIAL) count += tconn->partial ? 1 : 0; if (state & FR_TRUNK_REQUEST_SENT) count += fr_dlist_num_elements(&tconn->sent); if (state & FR_TRUNK_REQUEST_CANCEL) count += fr_dlist_num_elements(&tconn->cancel); + if (state & FR_TRUNK_REQUEST_CANCEL_PARTIAL) count += tconn->cancel_partial ? 1 : 0; if (state & FR_TRUNK_REQUEST_CANCEL_SENT) count += fr_dlist_num_elements(&tconn->cancel_sent); return count; @@ -1680,7 +1883,7 @@ static void trunk_connection_enter_active(fr_trunk_connection_t *tconn) static void _trunk_connection_on_connecting(UNUSED fr_connection_t *conn, UNUSED fr_connection_state_t state, void *uctx) { fr_trunk_connection_t *tconn = talloc_get_type_abort(uctx, fr_trunk_connection_t); - fr_trunk_t *trunk = tconn->trunk; + fr_trunk_t *trunk = tconn->trunk; switch (tconn->state) { case FR_TRUNK_CONN_HALTED: @@ -1718,7 +1921,7 @@ static void _trunk_connection_on_connecting(UNUSED fr_connection_t *conn, UNUSED static void _trunk_connection_on_connected(UNUSED fr_connection_t *conn, UNUSED fr_connection_state_t state, void *uctx) { fr_trunk_connection_t *tconn = talloc_get_type_abort(uctx, fr_trunk_connection_t); - fr_trunk_t *trunk = tconn->trunk; + fr_trunk_t *trunk = tconn->trunk; /* * If a connection was just connected, @@ -1750,9 +1953,9 @@ static void _trunk_connection_on_connected(UNUSED fr_connection_t *conn, UNUSED static void _trunk_connection_on_closed(UNUSED fr_connection_t *conn, UNUSED fr_connection_state_t state, void *uctx) { fr_trunk_connection_t *tconn = talloc_get_type_abort(uctx, fr_trunk_connection_t); - fr_trunk_t *trunk = tconn->trunk; - int ret; - bool need_requeue = false; + fr_trunk_t *trunk = tconn->trunk; + int ret; + bool need_requeue = false; switch (tconn->state) { case FR_TRUNK_CONN_HALTED: /* Failed during handle initialisation */ @@ -1974,6 +2177,7 @@ static int trunk_connection_spawn(fr_trunk_t *trunk, fr_time_t now) MEM(tconn->pending = fr_heap_talloc_create(tconn, trunk->funcs.request_prioritise, fr_trunk_request_t, heap_id)); fr_dlist_talloc_init(&tconn->sent, fr_trunk_request_t, list); + fr_dlist_talloc_init(&tconn->cancel, fr_trunk_request_t, list); fr_dlist_talloc_init(&tconn->cancel_sent, fr_trunk_request_t, list); /* @@ -2026,16 +2230,23 @@ static int trunk_connection_spawn(fr_trunk_t *trunk, fr_time_t now) * The request was cancelled and we don't need to wait, clean it * up immediately. * - * @param[out] treq_out The dequeued cancellation request. - * @param[in] tconn Connection to drain cancellation request from. + * @param[out] preq associated with the trunk request. + * @param[in] tconn Connection to drain cancellation request from. */ -fr_trunk_request_t *fr_trunk_connection_pop_cancellation(fr_trunk_connection_t *tconn) +fr_trunk_request_t *fr_trunk_connection_pop_cancellation(void **preq, fr_trunk_connection_t *tconn) { + fr_trunk_request_t *treq; + if (!fr_cond_assert_msg(IN_REQUEST_CANCEL_MUX(tconn->trunk), "%s can only be called from within request_cancel_mux handler", __FUNCTION__)) return NULL; - return fr_dlist_head(&tconn->cancel); + treq = tconn->cancel_partial ? tconn->cancel_partial : fr_dlist_head(&tconn->cancel); + if (!treq) return NULL; + + if (preq) *preq = treq->preq; + + return treq; } /** Pop a request off a connection's pending queue @@ -2066,16 +2277,25 @@ fr_trunk_request_t *fr_trunk_connection_pop_cancellation(fr_trunk_connection_t * * - #fr_trunk_request_signal_sent * Successfully sent a request. * + * @param[out] preq associated with the trunk request. + * @param[out] rctx associated with the trunk request. * @param[in] tconn to pop a request from. */ -fr_trunk_request_t *fr_trunk_connection_pop_request(fr_trunk_connection_t *tconn) +fr_trunk_request_t *fr_trunk_connection_pop_request(void **preq, void **rctx, fr_trunk_connection_t *tconn) { + fr_trunk_request_t *treq; + if (!fr_cond_assert_msg(IN_REQUEST_MUX(tconn->trunk), "%s can only be called from within request_mux handler", __FUNCTION__)) return NULL; - if (tconn->partial) return tconn->partial; - return fr_heap_peek(tconn->pending); + treq = tconn->partial ? tconn->partial : fr_heap_peek(tconn->pending); + if (!treq) return NULL; + + if (preq) *preq = treq->preq; + if (rctx) *rctx = treq->rctx; + + return treq; } /** Signal that a trunk connection is writable @@ -2086,9 +2306,13 @@ fr_trunk_request_t *fr_trunk_connection_pop_request(fr_trunk_connection_t *tconn */ void fr_trunk_connection_signal_writable(fr_trunk_connection_t *tconn) { + fr_trunk_t *trunk = tconn->trunk; + if (!fr_cond_assert_msg(!IN_HANDLER(tconn->trunk), "%s cannot be called within a handler", __FUNCTION__)) return; + DEBUG4("[%" PRIu64 "] Signalled writable", fr_connection_get_id(tconn->conn)); + trunk_connection_writable(tconn); } @@ -2100,9 +2324,13 @@ void fr_trunk_connection_signal_writable(fr_trunk_connection_t *tconn) */ void fr_trunk_connection_signal_readable(fr_trunk_connection_t *tconn) { + fr_trunk_t *trunk = tconn->trunk; + if (!fr_cond_assert_msg(!IN_HANDLER(tconn->trunk), "%s cannot be called within a handler", __FUNCTION__)) return; + DEBUG4("[%" PRIu64 "] Signalled readable", fr_connection_get_id(tconn->conn)); + trunk_connection_readable(tconn); } @@ -2667,7 +2895,8 @@ static int8_t _trunk_connection_order_by_shortest_queue(void const *one, void co */ static int _trunk_free(fr_trunk_t *trunk) { - fr_connection_t *tconn; + fr_connection_t *tconn; + fr_trunk_request_t *treq; DEBUG4("Trunk free %p", trunk); @@ -2694,6 +2923,11 @@ static int _trunk_free(fr_trunk_t *trunk) while ((tconn = fr_dlist_head(&trunk->failed))) talloc_free(tconn); while ((tconn = fr_dlist_head(&trunk->draining))) talloc_free(tconn); + /* + * Free any requests left in the backlog + */ + while ((treq = fr_heap_peek(trunk->backlog))) trunk_request_enter_failed(treq); + return 0; } @@ -2794,6 +3028,16 @@ static void *dummy_uctx = NULL; #include #include +typedef struct { + fr_trunk_request_t *treq; //!< Trunk request. + bool cancelled; //!< Seen by the cancelled callback. + bool completed; //!< Seen by the complete callback. + bool failed; //!< Seen by the failed callback. + bool freed; //!< Seen by the free callback. + bool signal_partial; //!< Muxer should signal that this request is partially written. + bool signal_cancel_partial; //!< Muxer should signal that this request is partially cancelled. +} test_proto_request_t; + #define DEBUG_LVL_SET if (test_verbose_level__ >= 3) fr_debug_lvl = L_DBG_LVL_4 + 1 static void _conn_close(void *h, UNUSED void *uctx) @@ -2820,15 +3064,10 @@ static void _conn_io_loopback(fr_event_list_t *el, int fd, int flags, void *uctx rad_assert(fd == our_h[1]); slen = read(fd, buff, sizeof(buff)); + printf("Received %zu bytes of data, sending it back\n", slen); write(our_h[1], buff, (size_t)slen); } -static void _conn_io_demux(fr_event_list_t *el, int fd, int flags, void *uctx) -{ - fr_trunk_connection_t *tconn = talloc_get_type_abort(uctx, fr_trunk_connection_t); - fr_trunk_connection_signal_readable(tconn); -} - /** Insert I/O handlers that loop any data back round * */ @@ -2836,8 +3075,10 @@ static fr_connection_state_t _conn_open(fr_event_list_t *el, void *h, void *uctx { int *our_h = talloc_get_type_abort(h, int); + /* + * This always needs to be inserted + */ fr_event_fd_insert(our_h, el, our_h[1], _conn_io_loopback, NULL, NULL, our_h); - fr_event_fd_insert(our_h, el, our_h[0], _conn_io_demux, NULL, NULL, uctx); return FR_CONNECTION_STATE_CONNECTED; } @@ -2902,7 +3143,7 @@ static void test_socket_pair_alloc_then_free(void) TEST_CHECK(events == 0); /* I/O events should have been cleared */ talloc_free(trunk); - talloc_free(el); + talloc_free(ctx); } static void test_socket_pair_alloc_then_reconnect_then_free(void) @@ -2948,7 +3189,7 @@ static void test_socket_pair_alloc_then_reconnect_then_free(void) TEST_CHECK(events == 0); /* I/O events should have been cleared */ talloc_free(trunk); - talloc_free(el); + talloc_free(ctx); } static fr_connection_state_t _conn_init_no_signal(void **h_out, fr_connection_t *conn, void *uctx) @@ -3027,7 +3268,7 @@ static void test_socket_pair_alloc_then_connect_timeout(void) TEST_CHECK(events == 0); /* I/O events should have been cleared */ talloc_free(trunk); - talloc_free(el); + talloc_free(ctx); } static fr_connection_t *test_setup_socket_pair_1s_reconnection_delay_alloc(fr_trunk_connection_t *tconn, @@ -3101,25 +3342,184 @@ static void test_socket_pair_alloc_then_reconnect_check_delay(void) TEST_CHECK(events == 2); /* Should have a pending I/O event and a timer */ talloc_free(trunk); - talloc_free(el); + talloc_free(ctx); } static void test_mux(fr_trunk_connection_t *tconn, fr_connection_t *conn, void *uctx) { fr_trunk_request_t *treq; + void *preq; size_t count = 0; + int fd = *(talloc_get_type_abort(fr_connection_get_handle(conn), int)); - while ((treq = fr_trunk_connection_pop_request(tconn))) { + while ((treq = fr_trunk_connection_pop_request(&preq, NULL, tconn))) { + test_proto_request_t *our_preq = preq; count++; + + /* + * Simulate a partial write + */ + if (our_preq->signal_partial) { + fr_trunk_request_signal_partial(treq); + our_preq->signal_partial = false; + break; + } + + write(fd, &preq, talloc_array_length(preq)); + fr_trunk_request_signal_sent(treq); + } + TEST_CHECK(count > 0); +} + +static void test_cancel_mux(fr_trunk_connection_t *tconn, fr_connection_t *conn, void *uctx) +{ + fr_trunk_request_t *treq; + void *preq; + size_t count = 0; + int fd = *(talloc_get_type_abort(fr_connection_get_handle(conn), int)); + + /* + * For cancellation we just do + */ + while ((treq = fr_trunk_connection_pop_cancellation(&preq, tconn))) { + test_proto_request_t *our_preq = preq; + count++; + + /* + * Simulate a partial cancel write + */ + if (our_preq->signal_cancel_partial) { + fr_trunk_request_signal_cancel_partial(treq); + our_preq->signal_cancel_partial = false; + break; + } + + write(fd, &preq, talloc_array_length(preq)); + fr_trunk_request_signal_cancel_sent(treq); } TEST_CHECK(count > 0); } static void test_demux(fr_trunk_connection_t *tconn, fr_connection_t *conn, void *uctx) { + int fd = *(talloc_get_type_abort(fr_connection_get_handle(conn), int)); + test_proto_request_t *preq; + + TEST_CHECK(read(fd, &preq, sizeof(preq)) == sizeof(preq)); + + talloc_get_type_abort(preq, test_proto_request_t); + + /* + * Demuxer can handle both normal requests and cancelled ones + */ + switch (preq->treq->state) { + case FR_TRUNK_REQUEST_CANCEL_SENT: + fr_trunk_request_signal_cancel_complete(preq->treq); + break; + + case FR_TRUNK_REQUEST_SENT: + fr_trunk_request_signal_complete(preq->treq); + break; + + default: + rad_assert(0); + break; + } +} + +static void _conn_io_error(UNUSED fr_event_list_t *el, UNUSED int fd, UNUSED int flags, + UNUSED int fd_errno, void *uctx) +{ + fr_trunk_connection_t *tconn = talloc_get_type_abort(uctx, fr_trunk_connection_t); + fr_trunk_connection_signal_reconnect(tconn); +} +static void _conn_io_read(UNUSED fr_event_list_t *el, UNUSED int fd, UNUSED int flags, void *uctx) +{ + fr_trunk_connection_t *tconn = talloc_get_type_abort(uctx, fr_trunk_connection_t); + fr_trunk_connection_signal_readable(tconn); } +static void _conn_io_write(UNUSED fr_event_list_t *el, UNUSED int fd, UNUSED int flags, void *uctx) +{ + fr_trunk_connection_t *tconn = talloc_get_type_abort(uctx, fr_trunk_connection_t); + fr_trunk_connection_signal_writable(tconn); +} + +static void _conn_notify(fr_trunk_connection_t *tconn, fr_connection_t *conn, + fr_event_list_t *el, + fr_trunk_connection_event_t notify_on, void *uctx) +{ + int fd = *(talloc_get_type_abort(fr_connection_get_handle(conn), int)); + + switch (notify_on) { + case FR_TRUNK_CONN_EVENT_NONE: + fr_event_fd_delete(el, fd, FR_EVENT_FILTER_IO); + break; + + case FR_TRUNK_CONN_EVENT_READ: + fr_event_fd_insert(conn, el, fd, _conn_io_read, NULL, _conn_io_error, tconn); + break; + + case FR_TRUNK_CONN_EVENT_WRITE: + fr_event_fd_insert(conn, el, fd, NULL, _conn_io_write, _conn_io_error, tconn); + break; + + case FR_TRUNK_CONN_EVENT_BOTH: + fr_event_fd_insert(conn, el, fd, _conn_io_read, _conn_io_write, _conn_io_error, tconn); + break; + + default: + rad_assert(0); + } +} + +static void test_request_cancel(UNUSED fr_connection_t *conn, UNUSED fr_trunk_request_t *treq, void *preq, + UNUSED fr_trunk_cancel_reason_t reason, UNUSED void *uctx) +{ + test_proto_request_t *our_preq = talloc_get_type_abort(preq, test_proto_request_t); + + our_preq->cancelled = true; +} + +static void test_request_complete(REQUEST *request, void *preq, void *rctx, void *uctx) +{ + test_proto_request_t *our_preq = talloc_get_type_abort(preq, test_proto_request_t); + + our_preq->completed = true; +} + +static void test_request_fail(REQUEST *request, void *preq, void *rctx, void *uctx) +{ + test_proto_request_t *our_preq = talloc_get_type_abort(preq, test_proto_request_t); + + our_preq->failed = true; +} + +static void test_request_free(REQUEST *request, void *preq, void *uctx) +{ + test_proto_request_t *our_preq = talloc_get_type_abort(preq, test_proto_request_t); + + our_preq->freed = true; +} + +/** Write a failure result to the rctx so that the API client is aware that the request failed + * + * This function should free any memory not bound to the lifetime of the rctx + * or request, or that was allocated explicitly to prepare for the REQUEST * + * being used by a trunk, this may include library request handles and + * (partially-)encoded packets. + * + * The rctx should be modified in such a way that indicates to the API client + * that the request could not be sent using the trunk. + * + * After this function returns the REQUEST * will be marked as runnable. + * + * @note If a cancel function is provided, this function should be used to remove + * active requests from any request/response matching, not the fail function. + */ +typedef void (*fr_trunk_request_fail_t)(REQUEST *request, void *preq, void *rctx, void *uctx); + static void test_enqueue_basic(void) { TALLOC_CTX *ctx = talloc_init("test"); @@ -3132,13 +3532,17 @@ static void test_enqueue_basic(void) }; fr_trunk_io_funcs_t io_funcs = { .connection_alloc = test_setup_socket_pair_connection_alloc, + .connection_notify = _conn_notify, .request_prioritise = fr_pointer_cmp, .request_mux = test_mux, .request_demux = test_demux, - + .request_cancel = test_request_cancel, + .request_cancel_mux = test_cancel_mux, + .request_complete = test_request_complete, + .request_fail = test_request_fail, + .request_free = test_request_free }; - char *test_data = "Testing123"; - char *preq; + test_proto_request_t *preq; fr_trunk_request_t *treq; fr_trunk_enqueue_t rcode; REQUEST *request; @@ -3150,7 +3554,12 @@ static void test_enqueue_basic(void) trunk = fr_trunk_alloc(ctx, el, "test_socket_pair", false, &conf, &io_funcs, &dummy_uctx); - preq = talloc_strdup(NULL, test_data); + /* + * Our preq is a pointer to the trunk + * request so we don't have to manage + * a tree of requests and responses. + */ + preq = talloc_zero(NULL, test_proto_request_t); /* * The trunk is active, but there's no @@ -3161,7 +3570,8 @@ static void test_enqueue_basic(void) * backlog. */ request = request_alloc(ctx); - rcode = fr_trunk_request_enqueue(&treq, trunk, request, preq, ctx); + rcode = fr_trunk_request_enqueue(&treq, trunk, request, preq, NULL); + preq->treq = treq; TEST_CHECK(rcode == TRUNK_ENQUEUE_IN_BACKLOG); TEST_CHECK(fr_trunk_requests_by_state_count(trunk, FR_TRUNK_REQUEST_BACKLOG) == 1); @@ -3186,8 +3596,222 @@ static void test_enqueue_basic(void) TEST_CHECK(fr_trunk_requests_by_state_count(trunk, FR_TRUNK_REQUEST_PENDING) == 0); TEST_CHECK(fr_trunk_connection_count(trunk, FR_TRUNK_CONN_ALL) == 1); + /* + * Allow the connection to establish (again) + */ + events = fr_event_corral(el, test_time_base, false); + fr_event_service(el); + + TEST_CHECK(fr_trunk_requests_by_state_count(trunk, FR_TRUNK_REQUEST_BACKLOG) == 0); + TEST_CHECK(fr_trunk_requests_by_state_count(trunk, FR_TRUNK_REQUEST_PENDING) == 1); + + /* + * Should now be active and have a write event + * inserted into the event loop. + */ + TEST_CHECK(fr_trunk_connection_count(trunk, FR_TRUNK_CONN_ACTIVE) == 1); + + /* + * Trunk should be signalled the connection is + * writable. + * + * We should then: + * - Pop a request from the pending queue. + * - Write the request to the socket pair + */ + events = fr_event_corral(el, test_time_base, false); + fr_event_service(el); + + TEST_CHECK(fr_trunk_requests_by_state_count(trunk, FR_TRUNK_REQUEST_SENT) == 1); + + /* + * Gives the loopback function a chance + * to read the data, and write it back. + */ + events = fr_event_corral(el, test_time_base, false); + fr_event_service(el); + + /* + * Trunk should be signalled the connection is + * readable. + * + * We should then: + * - Read the (looped back) response. + * - Signal the trunk that the connection is readable. + */ + events = fr_event_corral(el, test_time_base, false); + fr_event_service(el); + + TEST_CHECK(preq->completed == true); + TEST_CHECK(preq->failed == false); + TEST_CHECK(preq->cancelled == false); + TEST_CHECK(preq->freed == true); + talloc_free(preq); + + talloc_free(trunk); + talloc_free(ctx); +} + +static void test_enqueue_cancellation_points(void) +{ + TALLOC_CTX *ctx = talloc_init("test"); + fr_trunk_t *trunk; + fr_event_list_t *el; + int events; + fr_trunk_conf_t conf = { + .min_connections = 1, + .manage_interval = NSEC * 0.5 + }; + fr_trunk_io_funcs_t io_funcs = { + .connection_alloc = test_setup_socket_pair_connection_alloc, + .connection_notify = _conn_notify, + .request_prioritise = fr_pointer_cmp, + .request_mux = test_mux, + .request_demux = test_demux, + .request_cancel = test_request_cancel, + .request_cancel_mux = test_cancel_mux, + .request_complete = test_request_complete, + .request_fail = test_request_fail, + .request_free = test_request_free + }; + test_proto_request_t *preq; + fr_trunk_request_t *treq; + REQUEST *request; + + DEBUG_LVL_SET; + + el = fr_event_list_alloc(ctx, NULL, NULL); + fr_event_list_set_time_func(el, test_time); + + request = request_alloc(ctx); + + trunk = fr_trunk_alloc(ctx, el, "test_socket_pair", false, &conf, &io_funcs, &dummy_uctx); + preq = talloc_zero(NULL, test_proto_request_t); + fr_trunk_request_enqueue(&treq, trunk, request, preq, NULL); + + /* + * Cancellation in backlog via trunk free + */ + talloc_free(trunk); + TEST_CHECK(preq->completed == false); + TEST_CHECK(preq->failed == true); + TEST_CHECK(preq->cancelled == false); + TEST_CHECK(preq->freed == true); + talloc_free(preq); + + /* + * Cancellation in backlog via signal + */ + trunk = fr_trunk_alloc(ctx, el, "test_socket_pair", false, &conf, &io_funcs, &dummy_uctx); + preq = talloc_zero(NULL, test_proto_request_t); + fr_trunk_request_enqueue(&treq, trunk, request, preq, NULL); + fr_trunk_request_signal_cancel(treq); + TEST_CHECK(fr_trunk_requests_by_state_count(trunk, FR_TRUNK_REQUEST_ALL) == 0); + + TEST_CHECK(preq->completed == false); + TEST_CHECK(preq->failed == false); /* Request/rctx not guaranteed after signal, so can't call fail */ + TEST_CHECK(preq->cancelled == false); + TEST_CHECK(preq->freed == true); + talloc_free(preq); + talloc_free(trunk); + + + /* + * Cancellation in partial via trunk free + */ + trunk = fr_trunk_alloc(ctx, el, "test_socket_pair", false, &conf, &io_funcs, &dummy_uctx); + preq = talloc_zero(NULL, test_proto_request_t); + preq->signal_partial = true; + fr_trunk_request_enqueue(&treq, trunk, request, preq, NULL); + + events = fr_event_corral(el, test_time_base, false); /* Connect the connection */ + fr_event_service(el); + + events = fr_event_corral(el, test_time_base, false); /* Send the request */ + fr_event_service(el); + + TEST_CHECK(fr_trunk_requests_by_state_count(trunk, FR_TRUNK_REQUEST_PARTIAL)); + talloc_free(trunk); - talloc_free(el); + + TEST_CHECK(preq->completed == false); + TEST_CHECK(preq->failed == true); + TEST_CHECK(preq->cancelled == true); + TEST_CHECK(preq->freed == true); + talloc_free(preq); + + /* + * Cancellation in partial via signal + */ + trunk = fr_trunk_alloc(ctx, el, "test_socket_pair", false, &conf, &io_funcs, &dummy_uctx); + preq = talloc_zero(NULL, test_proto_request_t); + preq->signal_partial = true; + fr_trunk_request_enqueue(&treq, trunk, request, preq, NULL); + + events = fr_event_corral(el, test_time_base, false); /* Connect the connection */ + fr_event_service(el); + + events = fr_event_corral(el, test_time_base, false); /* Send the request */ + fr_event_service(el); + + TEST_CHECK(fr_trunk_requests_by_state_count(trunk, FR_TRUNK_REQUEST_PARTIAL) == 1); + fr_trunk_request_signal_cancel(treq); + TEST_CHECK(fr_trunk_requests_by_state_count(trunk, FR_TRUNK_REQUEST_ALL) == 0); + + TEST_CHECK(preq->completed == false); + TEST_CHECK(preq->failed == false); /* Request/rctx not guaranteed after signal, so can't call fail */ + TEST_CHECK(preq->cancelled == true); + TEST_CHECK(preq->freed == true); + talloc_free(preq); + talloc_free(trunk); + + /* + * Cancellation in sent via trunk free + */ + trunk = fr_trunk_alloc(ctx, el, "test_socket_pair", false, &conf, &io_funcs, &dummy_uctx); + preq = talloc_zero(NULL, test_proto_request_t); + fr_trunk_request_enqueue(&treq, trunk, request, preq, NULL); + + events = fr_event_corral(el, test_time_base, false); /* Connect the connection */ + fr_event_service(el); + + events = fr_event_corral(el, test_time_base, false); /* Send the request */ + fr_event_service(el); + + TEST_CHECK(fr_trunk_requests_by_state_count(trunk, FR_TRUNK_REQUEST_SENT) == 1); + talloc_free(trunk); + + TEST_CHECK(preq->completed == false); + TEST_CHECK(preq->failed == true); + TEST_CHECK(preq->cancelled == true); + TEST_CHECK(preq->freed == true); + talloc_free(preq); + + /* + * Cancellation in sent via signal + */ + trunk = fr_trunk_alloc(ctx, el, "test_socket_pair", false, &conf, &io_funcs, &dummy_uctx); + preq = talloc_zero(NULL, test_proto_request_t); + fr_trunk_request_enqueue(&treq, trunk, request, preq, NULL); + + events = fr_event_corral(el, test_time_base, false); /* Connect the connection */ + fr_event_service(el); + + events = fr_event_corral(el, test_time_base, false); /* Send the request */ + fr_event_service(el); + + TEST_CHECK(fr_trunk_requests_by_state_count(trunk, FR_TRUNK_REQUEST_SENT) == 1); + fr_trunk_request_signal_cancel(treq); + TEST_CHECK(fr_trunk_requests_by_state_count(trunk, FR_TRUNK_REQUEST_ALL) == 0); + + TEST_CHECK(preq->completed == false); + TEST_CHECK(preq->failed == false); /* Request/rctx not guaranteed after signal, so can't call fail */ + TEST_CHECK(preq->cancelled == true); + TEST_CHECK(preq->freed == true); + talloc_free(preq); + talloc_free(trunk); + + talloc_free(ctx); } @@ -3208,7 +3832,7 @@ TEST_LIST = { * Basic enqueue/dequeue */ { "Enqueue - Basic", test_enqueue_basic }, - + { "Enqueue - Cancellation points", test_enqueue_cancellation_points }, { NULL } }; #endif diff --git a/src/lib/server/trunk.h b/src/lib/server/trunk.h index f37531c3617..ae3af8015ae 100644 --- a/src/lib/server/trunk.h +++ b/src/lib/server/trunk.h @@ -278,11 +278,12 @@ typedef void (*fr_trunk_request_cancel_mux_t)(fr_trunk_connection_t *tconn, fr_c * writable. * * @param[in] conn to remove request from. - * @param[in] treq Trunk request to remove. + * @param[in] treq Trunk request to cancel. + * @param[in] preq Preq to cancel. * @param[in] reason Why the request was cancelled. * @param[in] uctx User context data passed to #fr_trunk_alloc. */ -typedef void (*fr_trunk_request_cancel_t)(fr_connection_t *conn, fr_trunk_request_t *treq, +typedef void (*fr_trunk_request_cancel_t)(fr_connection_t *conn, fr_trunk_request_t *treq, void *preq, fr_trunk_cancel_reason_t reason, void *uctx); /** Write a successful result to the rctx so that the API client is aware of the result @@ -384,6 +385,8 @@ void fr_trunk_request_signal_fail(fr_trunk_request_t *treq); void fr_trunk_request_signal_cancel(fr_trunk_request_t *treq); +void fr_trunk_request_signal_cancel_partial(fr_trunk_request_t *treq); + void fr_trunk_request_signal_cancel_sent(fr_trunk_request_t *treq); void fr_trunk_request_signal_cancel_complete(fr_trunk_request_t *treq); @@ -400,9 +403,9 @@ int fr_trunk_request_enqueue(fr_trunk_request_t **treq, fr_trunk_t *trunk, REQU /** @name Dequeue protocol requests and cancellations * @{ */ -fr_trunk_request_t *fr_trunk_connection_pop_cancellation(fr_trunk_connection_t *tconn); +fr_trunk_request_t *fr_trunk_connection_pop_cancellation(void **preq, fr_trunk_connection_t *tconn); -fr_trunk_request_t *fr_trunk_connection_pop_request(fr_trunk_connection_t *tconn); +fr_trunk_request_t *fr_trunk_connection_pop_request(void **preq, void **rctx, fr_trunk_connection_t *tconn); /** @} */ /** @name Connection state signalling