static int debug_lvl = 0;
static int kq_master, kq_worker;
+static fr_atomic_queue_t *aq_master, *aq_worker;
static int max_messages = 10;
static int max_outstanding = 1;
static bool touch_memory = false;
kq_worker = kqueue();
rad_assert(kq_worker >= 0);
- channel = fr_channel_create(autofree, kq_master, kq_worker);
+ aq_master = fr_atomic_queue_create(autofree, 128);
+ rad_assert(aq_master != NULL);
+
+ aq_worker = fr_atomic_queue_create(autofree, 128);
+ rad_assert(aq_worker != NULL);
+
+ channel = fr_channel_create(autofree,
+ kq_master, aq_master,
+ kq_worker, aq_worker);
if (!channel) {
fprintf(stderr, "channel_test: Failed to create channel\n");
exit(1);
typedef struct fr_channel_end_t {
int kq; //!< the kqueue associated with the channel
+ fr_atomic_queue_t *aq_control; //!< the control channel queue
+
int num_outstanding; //!< number of outstanding requests with no reply
size_t num_signals;
* - NULL on error
* - channel on success
*/
-fr_channel_t *fr_channel_create(TALLOC_CTX *ctx, int kq_master, int kq_worker)
+fr_channel_t *fr_channel_create(TALLOC_CTX *ctx, int kq_master, fr_atomic_queue_t *aq_master,
+ int kq_worker, fr_atomic_queue_t *aq_worker)
{
fr_time_t when;
fr_channel_t *ch;
}
ch->end[TO_WORKER].kq = kq_worker;
+ ch->end[TO_WORKER].aq_control = aq_worker;
ch->end[FROM_WORKER].kq = kq_master;
+ ch->end[FROM_WORKER].aq_control = aq_master;
/*
* Initialize all of the timers to now.
FR_CHANNEL_CLOSE,
} fr_channel_event_t;
-fr_channel_t *fr_channel_create(TALLOC_CTX *ctx, int kq_master, int kq_worker);
+fr_channel_t *fr_channel_create(TALLOC_CTX *ctx, int kq_master, fr_atomic_queue_t *aq_master,
+ int kq_worker, fr_atomic_queue_t *aq_worker);
int fr_channel_send_request(fr_channel_t *ch, fr_channel_data_t *cm, fr_channel_data_t **p_reply) CC_HINT(nonnull);
fr_channel_data_t *fr_channel_recv_request(fr_channel_t *ch) CC_HINT(nonnull);