* saying "I took saved the data, but the socket wasn't ready, so you
* need to call me again at a later point".
*
- * @param[in] sockfd the file descriptor to use
* @param[in] io_ctx the context for this function
* @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
* - <0 on error
* - >=0 length of the data read or written.
*/
-typedef ssize_t (*fr_io_data_t)(int sockfd, void *io_ctx, uint8_t *buffer, size_t buffer_len);
+typedef ssize_t (*fr_io_data_t)(void *io_ctx, uint8_t *buffer, size_t buffer_len);
/** Handle a close or error on the socket.
*
* before "close". On normal finish, the "close" function will be
* called.
*
- * @param[in] sockfd the file descriptor to use
* @param[in] io_ctx the context for this function
* @return
* - 0 on success
* - <0 on error
*/
-typedef int (*fr_io_signal_t)(int sockfd, void *io_ctx);
+typedef int (*fr_io_signal_t)(void *io_ctx);
/** Process a request through the transport async state machine.
*
typedef struct fr_network_socket_t {
fr_dlist_t entry;
- int fd; //!< the file descriptor
void *ctx; //!< transport context
fr_io_op_t *transport; //!< the transport
ssize_t data_size;
fr_channel_data_t *cd;
- rad_assert(s->fd == sockfd);
+ rad_assert(s->transport->fd(s->ctx) == sockfd);
fr_log(nr->log, L_DBG, "network read");
* network side knows that it needs to close the
* connection.
*/
- data_size = s->transport->read(sockfd, s->ctx, cd->m.data, cd->m.rb_size);
+ data_size = s->transport->read(s->ctx, cd->m.data, cd->m.rb_size);
if (data_size == 0) {
fr_log(nr->log, L_DBG_ERR, "got no data from transport read");
* fr_io_op_t.
*/
if (data_size < 0) {
- fr_log(nr->log, L_DBG_ERR, "error from transport read on socket %d", s->fd);
+ fr_log(nr->log, L_DBG_ERR, "error from transport read on socket %d", sockfd);
- (void) fr_event_fd_delete(nr->el, s->fd);
+ (void) fr_event_fd_delete(nr->el, sockfd);
fr_dlist_remove(&s->entry);
talloc_free(s);
return;
* Initialize the rest of the fields of the channel data.
*/
cd->m.when = fr_time();
- cd->io->fd = sockfd;
cd->priority = 0;
cd->io->ctx = s->ctx;
cd->io->op = s->transport;
* @param sockfd the socket which is ready to write
* @param ctx the network socket context.
*/
-static void fr_network_write(UNUSED fr_event_list_t *el, int sockfd, void *ctx)
+static void fr_network_write(UNUSED fr_event_list_t *el, UNUSED int sockfd, void *ctx)
{
fr_network_socket_t *s = ctx;
- if (s->transport->flush(sockfd, s->ctx) < 0) {
- s->transport->error(sockfd, s->ctx);
+ if (s->transport->flush(s->ctx) < 0) {
+ s->transport->error(s->ctx);
talloc_free(s);
}
}
* @param sockfd the socket which has a fatal error.
* @param ctx the network socket context.
*/
-static void fr_network_error(UNUSED fr_event_list_t *el, int sockfd, void *ctx)
+static void fr_network_error(UNUSED fr_event_list_t *el, UNUSED int sockfd, void *ctx)
{
fr_network_socket_t *s = ctx;
- s->transport->error(sockfd, s->ctx);
+ s->transport->error(s->ctx);
talloc_free(s);
}
{
fr_network_t *nr = talloc_parent(s);
- fr_event_fd_delete(nr->el, s->fd);
+ fr_event_fd_delete(nr->el, s->transport->fd(s->ctx));
fr_dlist_remove(&s->entry);
if (s->transport->close) {
- s->transport->close(s->fd, s->ctx);
+ s->transport->close(s->ctx);
} else {
- close(s->fd);
+ close(s->transport->fd(s->ctx));
}
return 0;
*/
static void fr_network_socket_callback(void *ctx, void const *data, size_t data_size, UNUSED fr_time_t now)
{
+ int fd;
fr_network_t *nr = ctx;
fr_network_socket_t *s;
fr_event_fd_handler_t write_fn, error_fn;
if (s->transport->error) error_fn = fr_network_error;
- if (fr_event_fd_insert(nr->el, s->fd, fr_network_read, write_fn, error_fn, s) < 0) {
+ fd = s->transport->fd(s->ctx);
+
+ if (fr_event_fd_insert(nr->el, fd, fr_network_read, write_fn, error_fn, s) < 0) {
fr_log(nr->log, L_ERR, "Failed adding new socket to event loop: %s", fr_strerror());
talloc_free(s);
return;
fr_dlist_insert_head(&nr->sockets, &s->entry);
- fr_log(nr->log, L_DBG, "Using new socket with FD %d", s->fd);
+ fr_log(nr->log, L_DBG, "Using new socket with FD %d", fd);
}
* the reply is a NAK, don't write it to the
* network.
*/
- rcode = io->op->write(io->fd, io->ctx, cd->m.data, cd->m.data_size);
+ rcode = io->op->write(io->ctx, cd->m.data, cd->m.data_size);
if (rcode < 0) {
fr_dlist_t *entry;
* Don't call close, as that will be done
* in the destructor.
*/
- if (io->op->error) io->op->error(io->fd, io->ctx);
+ if (io->op->error) io->op->error(io->ctx);
/*
- * Find the socket which matches this
- * file descriptor.
+ * Find the fr_network_socket_t which
+ * matches this message.
*
- * @todo - put them into a binary tree
- * based on FD. That way we can handle
- * tens of thousands without walking a
- * linked list.
+ * @todo - put them into a binary tree.
+ * That way we can handle tens of
+ * thousands without walking a linked
+ * list.
*/
for (entry = FR_DLIST_FIRST(nr->sockets);
entry != NULL;
fr_network_socket_t *s;
s = fr_ptr_to_type(fr_network_socket_t, entry, entry);
- if (s->fd == io->fd) {
+ if (s->ctx == io->ctx) {
talloc_free(s);
break;
}
/** Add a socket to a network
*
* @param nr the network
- * @param fd the file descriptor for the socket
* @param ctx the context for the transport
* @param transport the transport
*/
-int fr_network_socket_add(fr_network_t *nr, int fd, void *ctx, fr_io_op_t *transport)
+int fr_network_socket_add(fr_network_t *nr, void *ctx, fr_io_op_t *transport)
{
fr_network_socket_t m;
memset(&m, 0, sizeof(m));
- m.fd = fd;
m.ctx = ctx;
m.transport = transport;
int fr_network_destroy(fr_network_t *nr) CC_HINT(nonnull);
void fr_network(fr_network_t *nr) CC_HINT(nonnull);
-int fr_network_socket_add(fr_network_t *nr, int fd, void *ctx, fr_io_op_t *transport) CC_HINT(nonnull);
+int fr_network_socket_add(fr_network_t *nr, void *ctx, fr_io_op_t *transport) CC_HINT(nonnull);
int fr_network_worker_add(fr_network_t *nr, fr_worker_t *worker) CC_HINT(nonnull);
#ifdef __cplusplus
/** Add a socket to a scheduler.
*
* @param sc the scheduler
- * @param fd the file descriptor for the socket
* @param ctx the context for the transport
* @param transport the transport
* @return
* - NULL on error
* - the fr_network_t that the socket was added to.
*/
-fr_network_t *fr_schedule_socket_add(fr_schedule_t *sc, int fd, void *ctx, fr_io_op_t *transport)
+fr_network_t *fr_schedule_socket_add(fr_schedule_t *sc, void *ctx, fr_io_op_t *transport)
{
- if (fr_network_socket_add(sc->sn->rc, fd, ctx, transport) < 0) {
+ if (fr_network_socket_add(sc->sn->rc, ctx, transport) < 0) {
return NULL;
}
int fr_schedule_destroy(fr_schedule_t *sc);
int fr_schedule_get_worker_kq(fr_schedule_t *sc);
-fr_network_t *fr_schedule_socket_add(fr_schedule_t *sc, int fd, void *ctx, fr_io_op_t *transport) CC_HINT(nonnull);
+fr_network_t *fr_schedule_socket_add(fr_schedule_t *sc, void *ctx, fr_io_op_t *transport) CC_HINT(nonnull);
#ifdef __cplusplus
}
return -1;
}
+ /*
+ * Add port_name
+ */
+
app_io = (fr_app_io_t const *) module->common;
if (app_io->instantiate(io_cs, io_ctx) < 0) {
cf_log_err_cs(cs, "Failed instantiating 'transport = %s'", value);
CONF_PARSER_TERMINATOR
};
-static ssize_t mod_read(int sockfd, void *ctx, uint8_t *buffer, size_t buffer_len)
+static ssize_t mod_read(void *ctx, uint8_t *buffer, size_t buffer_len)
{
ssize_t data_size;
size_t packet_len;
pc->salen = sizeof(pc->src);
- data_size = recvfrom(sockfd, buffer, buffer_len, 0, (struct sockaddr *) &pc->src, &pc->salen);
+ data_size = recvfrom(pc->sockfd, buffer, buffer_len, 0, (struct sockaddr *) &pc->src, &pc->salen);
if (data_size <= 0) return data_size;
packet_len = data_size;
}
-static ssize_t mod_write(int sockfd, void *ctx, uint8_t *buffer, size_t buffer_len)
+static ssize_t mod_write(void *ctx, uint8_t *buffer, size_t buffer_len)
{
ssize_t data_size;
fr_packet_ctx_t *pc = ctx;
/*
* @todo - do more stuff
*/
- data_size = sendto(sockfd, buffer, buffer_len, 0, (struct sockaddr *) &pc->src, pc->salen);
+ data_size = sendto(pc->sockfd, buffer, buffer_len, 0, (struct sockaddr *) &pc->src, pc->salen);
if (data_size <= 0) return data_size;
/*
return 10;
}
-static ssize_t test_read(int sockfd, void *ctx, uint8_t *buffer, size_t buffer_len)
+static ssize_t test_read(void *ctx, uint8_t *buffer, size_t buffer_len)
{
ssize_t data_size;
fr_packet_ctx_t *pc = ctx;
pc->salen = sizeof(pc->src);
- data_size = recvfrom(sockfd, buffer, buffer_len, 0, (struct sockaddr *) &pc->src, &pc->salen);
+ data_size = recvfrom(pc->sockfd, buffer, buffer_len, 0, (struct sockaddr *) &pc->src, &pc->salen);
if (data_size <= 0) return data_size;
/*
}
-static ssize_t test_write(int sockfd, void *ctx, uint8_t *buffer, size_t buffer_len)
+static ssize_t test_write(void *ctx, uint8_t *buffer, size_t buffer_len)
{
ssize_t data_size;
fr_packet_ctx_t *pc = ctx;
pc->salen = sizeof(pc->src);
- data_size = sendto(sockfd, buffer, buffer_len, 0, (struct sockaddr *) &pc->src, pc->salen);
+ data_size = sendto(pc->sockfd, buffer, buffer_len, 0, (struct sockaddr *) &pc->src, pc->salen);
if (data_size <= 0) return data_size;
/*
return data_size;
}
+static int test_fd(void *ctx)
+{
+ fr_packet_ctx_t *pc = ctx;
+
+ return pc->sockfd;
+}
+
static fr_io_op_t transport = {
.name = "schedule-test",
.decode = test_decode,
.encode = test_encode,
.nak = test_nak,
+ .fd = test_fd,
};
static void NEVER_RETURNS usage(void)
packet_ctx.sockfd = sockfd;
- (void) fr_schedule_socket_add(sched, sockfd, &packet_ctx, &transport);
+ (void) fr_schedule_socket_add(sched, &packet_ctx, &transport);
sleep(10);