fr_heap_t *workers; //!< heap of workers
fr_heap_t *done_workers; //!< heap of done workers
+
+ uint32_t num_transports; //!< how many transport layers we have
+ fr_transport_t **transports; //!< array of active transports.
};
return NULL;
}
- sw->worker = fr_worker_create(ctx);
+ sw->worker = fr_worker_create(ctx, sc->num_transports, sc->transports);
if (!sw->worker) {
talloc_free(ctx);
goto fail;
* - fr_schedule_t new scheduler
*/
fr_schedule_t *fr_schedule_create(TALLOC_CTX *ctx, int max_inputs, int max_workers,
+ uint32_t num_transports, fr_transport_t **transports,
fr_schedule_thread_instantiate_t worker_thread_instantiate,
void *worker_thread_ctx)
{
sc->worker_instantiate_ctx = worker_thread_ctx;
sc->running = true;
+ sc->num_transports = num_transports;
+ sc->transports = transports;
/*
* No inputs or workers, we're single threaded mode.
typedef int (*fr_schedule_thread_instantiate_t)(void *ctx);
fr_schedule_t *fr_schedule_create(TALLOC_CTX *ctx, int max_inputs, int max_workers,
+ uint32_t num_transports, fr_transport_t **transports,
fr_schedule_thread_instantiate_t worker_thread_instantiate,
void *worker_thread_ctx);
/* schedulers are async, so there's no fr_schedule_run() */
fr_time_tracking_t tracking; //!< how much time the worker has spent doing things.
+ uint32_t num_transports; //!< how many transport layers we have
fr_transport_t **transports; //!< array of active transports.
fr_channel_t *channel[1]; //!< list of channels
* Receive a message to the worker queue, and decode it
* to a to a request.
*/
+ rad_assert(cd->transport <= worker->num_transports);
rad_assert(worker->transports[cd->transport] != NULL);
request = worker->transports[cd->transport]->recv_request(worker->transports[cd->transport], cd->ctx, ctx, cd->m.data, cd->m.data_size);
* - NULL on error
* - fr_worker_t on success
*/
-fr_worker_t *fr_worker_create(TALLOC_CTX *ctx)
+fr_worker_t *fr_worker_create(TALLOC_CTX *ctx, uint32_t num_transports, fr_transport_t **transports)
{
fr_worker_t *worker;
+ if (!num_transports || !transports) return NULL;
+
worker = talloc_zero(ctx, fr_worker_t);
worker->el = fr_event_list_create(worker, fr_worker_idle, worker);
return NULL;
}
+ worker->num_transports = num_transports;
+ worker->transports = transports;
+
return worker;
}
*/
typedef struct fr_worker_t fr_worker_t;
-fr_worker_t *fr_worker_create(TALLOC_CTX *ctx);
+fr_worker_t *fr_worker_create(TALLOC_CTX *ctx, uint32_t num_transports, fr_transport_t **transports);
void fr_worker_destroy(fr_worker_t *worker) CC_HINT(nonnull);
int fr_worker_kq(fr_worker_t *worker) CC_HINT(nonnull);
void fr_worker(fr_worker_t *worker) CC_HINT(nonnull);