union {
struct {
- fr_time_t *start_time; //!< time original request started (network -> worker)
+ fr_time_t *recv_time; //!< time original request was received (network -> worker)
fr_dlist_t list; //!< list of unprocessed packets for the worker
} request;
*
* @param[in] instance the context for this function
* @param[out] packet_ctx Where to write a newly allocated packet_ctx struct containing request specific data.
+ * @param[in,out] recv_time A pointer to a time when the packet was received
* @param[in,out] buffer the buffer where the raw packet will be written to (or read from)
* @param[in] buffer_len the length of the buffer
* @return
* - <0 on error
* - >=0 length of the data read or written.
*/
-typedef ssize_t (*fr_io_data_read_t)(void const *instance, void **packet_ctx, uint8_t *buffer, size_t buffer_len);
+typedef ssize_t (*fr_io_data_read_t)(void const *instance, void **packet_ctx, fr_time_t **recv_time, uint8_t *buffer, size_t buffer_len);
/** Write a socket.
*
}
-static fr_time_t start_time = 0;
-
-
/** Read a packet from the network.
*
* @param[in] el the event list.
fr_network_t *nr = talloc_parent(s);
ssize_t data_size;
fr_channel_data_t *cd;
+ fr_time_t *recv_time;
rad_assert(s->listen->app_io->fd(s->listen->app_io_instance) == sockfd);
* network side knows that it needs to close the
* connection.
*/
- data_size = s->listen->app_io->read(s->listen->app_io_instance, &cd->packet_ctx, cd->m.data, cd->m.rb_size);
+ data_size = s->listen->app_io->read(s->listen->app_io_instance, &cd->packet_ctx, &recv_time, cd->m.data, cd->m.rb_size);
if (data_size == 0) {
fr_log(nr->log, L_DBG_ERR, "got no data from transport read");
/*
* Initialize the rest of the fields of the channel data.
*/
- cd->m.when = fr_time();
+ if (recv_time) {
+ cd->m.when = *recv_time;
+ } else {
+ cd->m.when = fr_time();
+ }
cd->priority = 0;
cd->listen = s->listen;
- cd->request.start_time = &start_time; /* @todo - set by transport */
-
- start_time = cd->m.when;
+ cd->request.recv_time = recv_time;
(void) fr_message_alloc(s->ms, &cd->m, data_size);
* This message has asynchronously aged out while it was
* in the queue. Delete it, and go get another one.
*/
- if (cd->request.start_time && (cd->m.when != *cd->request.start_time)) {
- fr_log(worker->log, L_DBG, "\t%sIGNORING old message", worker->name);
+ if (cd->request.recv_time && (cd->m.when != *cd->request.recv_time)) {
+ fr_log(worker->log, L_DBG, "\t%sIGNORING old message: was %zd now %zd", worker->name,
+ *cd->request.recv_time, cd->m.when);
fr_worker_nak(worker, cd, fr_time());
cd = NULL;
}
* processing this message.
*/
request->async->channel = cd->channel.ch;
- request->async->original_recv_time = cd->request.start_time;
+ request->async->original_recv_time = cd->request.recv_time;
request->async->recv_time = cd->m.when;
request->async->el = worker->el;
request->number = worker->number++;
/*
* Hoist run-time checks here.
*/
- if (!cd->request.start_time) request->async->original_recv_time = &request->async->recv_time;
+ if (!cd->request.recv_time) request->async->original_recv_time = &request->async->recv_time;
/*
* We're done with this message.
return address->client;
}
-static ssize_t mod_read(void const *instance, void **packet_ctx, uint8_t *buffer, size_t buffer_len)
+static ssize_t mod_read(void const *instance, void **packet_ctx, fr_time_t **recv_time, uint8_t *buffer, size_t buffer_len)
{
proto_radius_udp_t const *inst = talloc_get_type_abort(instance, proto_radius_udp_t);
}
*packet_ctx = track;
+ *recv_time = &track->timestamp;
return packet_len;
}
* The original packet has changed. Suppress the write,
* as the client will never accept the response.
*/
- if (track->timestamp > request_time) return buffer_len;
+ if (track->timestamp != request_time) return buffer_len;
/*
* Figure out when we've sent the reply.
return 0;
}
-static ssize_t test_read(void const *ctx, UNUSED void **packet_ctx, uint8_t *buffer, size_t buffer_len)
+static fr_time_t start_time;
+
+static ssize_t test_read(void const *ctx, UNUSED void **packet_ctx, fr_time_t **recv_time, uint8_t *buffer, size_t buffer_len)
{
ssize_t data_size;
fr_listen_test_t *io_ctx = talloc_get_type_abort(ctx, fr_listen_test_t);
tpc.id = buffer[1];
memcpy(tpc.vector, buffer + 4, sizeof(tpc.vector));
+ start_time = fr_time();
+ *recv_time = &start_time;
+
return data_size;
}