* 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.
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
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 \
)
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.
{ "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);
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);
#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), "<INVALID>"), (_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), "<INVALID>"), (_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)
#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)
#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)
#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)
*/
#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
*/
#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
#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)
#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), "<INVALID>"), (_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)
#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
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
*
*/
/** 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.
*/
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;
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;
*
* @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;
* 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.
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);
}
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
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
*/
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;
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;
*/
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);
*/
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);
/** 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)
{
#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); \
}
*/
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);
}
}
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);
}
}
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 */
}
}
*/
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;
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
case FR_TRUNK_REQUEST_FAILED:
break;
- case FR_TRUNK_REQUEST_CANCEL:
case FR_TRUNK_REQUEST_CANCEL_SENT:
break;
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;
"%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;
"%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;
"%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;
}
}
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;
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:
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,
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 */
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);
/*
* 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
* - #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
*/
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);
}
*/
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);
}
*/
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);
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;
}
#include <sys/types.h>
#include <sys/socket.h>
+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)
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
*
*/
{
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;
}
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)
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)
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,
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");
};
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;
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
* 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);
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);
}
* Basic enqueue/dequeue
*/
{ "Enqueue - Basic", test_enqueue_basic },
-
+ { "Enqueue - Cancellation points", test_enqueue_cancellation_points },
{ NULL }
};
#endif