diff options
| author | Garrett D'Amore <garrett@damore.org> | 2018-08-08 00:16:48 +0300 |
|---|---|---|
| committer | Garrett D'Amore <garrett@damore.org> | 2018-08-08 08:24:32 +0300 |
| commit | 1a752ba26a7c65c386ae5431b2615ace25db9521 (patch) | |
| tree | fbc1a586cff0f408626a04bddee591cb727a57a4 /src | |
| parent | 53a6e9683589988971751256c4cd96cba79d07ab (diff) | |
| download | nng-1a752ba26a7c65c386ae5431b2615ace25db9521.tar.gz nng-1a752ba26a7c65c386ae5431b2615ace25db9521.tar.bz2 nng-1a752ba26a7c65c386ae5431b2615ace25db9521.zip | |
fixes #630 IPC start elimination
fixes #615 IPC close on Windows leaves handle open
This reintroduces the code to eradicate the separate transport
start function for IPC, as a incremental step towards the full
fix for 599 and 208. It also addresses 615, by including revised
logic for the handling of close.
Diffstat (limited to 'src')
| -rw-r--r-- | src/transport/ipc/ipc.c | 765 |
1 files changed, 330 insertions, 435 deletions
diff --git a/src/transport/ipc/ipc.c b/src/transport/ipc/ipc.c index ee72d6d8..c40880c0 100644 --- a/src/transport/ipc/ipc.c +++ b/src/transport/ipc/ipc.c @@ -20,17 +20,21 @@ // Windows named pipes. Other platforms could use other mechanisms, // but all implementations on the platform must use the same mechanism. -typedef struct ipctran_pipe ipctran_pipe; -typedef struct ipctran_dialer ipctran_dialer; -typedef struct ipctran_listener ipctran_listener; +typedef struct ipctran_pipe ipctran_pipe; +typedef struct ipctran_ep ipctran_ep; // ipc_pipe is one end of an IPC connection. struct ipctran_pipe { - nni_ipc_conn *conn; - uint16_t peer; - uint16_t proto; - size_t rcvmax; - nni_sockaddr sa; + nni_ipc_conn * conn; + uint16_t peer; + uint16_t proto; + size_t rcvmax; + bool closed; + nni_sockaddr sa; + ipctran_ep * ep; + nni_list_node node; + nni_atomic_flag reaped; + nni_reap_item reap; uint8_t txhead[1 + sizeof(uint64_t)]; uint8_t rxhead[1 + sizeof(uint64_t)]; @@ -41,32 +45,25 @@ struct ipctran_pipe { nni_list recvq; nni_list sendq; - nni_aio *user_negaio; + nni_aio *useraio; nni_aio *txaio; nni_aio *rxaio; - nni_aio *negaio; + nni_aio *negoaio; + nni_aio *connaio; nni_msg *rxmsg; nni_mtx mtx; }; -struct ipctran_dialer { - nni_sockaddr sa; - nni_ipc_dialer *dialer; - uint16_t proto; - size_t rcvmax; - nni_aio * aio; - nni_aio * user_aio; - nni_mtx mtx; -}; - -struct ipctran_listener { +struct ipctran_ep { + nni_mtx mtx; nni_sockaddr sa; - nni_ipc_listener *listener; - uint16_t proto; size_t rcvmax; - nni_aio * aio; - nni_aio * user_aio; - nni_mtx mtx; + uint16_t proto; + nni_list pipes; + bool fini; + nni_ipc_dialer * dialer; + nni_ipc_listener *listener; + nni_reap_item reap; }; static void ipctran_pipe_send_start(ipctran_pipe *); @@ -74,8 +71,8 @@ static void ipctran_pipe_recv_start(ipctran_pipe *); static void ipctran_pipe_send_cb(void *); static void ipctran_pipe_recv_cb(void *); static void ipctran_pipe_nego_cb(void *); -static void ipctran_dialer_cb(void *); -static void ipctran_listener_cb(void *); +static void ipctran_pipe_conn_cb(void *); +static void ipctran_ep_fini(void *); static int ipctran_init(void) @@ -93,9 +90,14 @@ ipctran_pipe_close(void *arg) { ipctran_pipe *p = arg; + nni_mtx_lock(&p->mtx); + p->closed = true; + nni_mtx_unlock(&p->mtx); + nni_aio_close(p->rxaio); nni_aio_close(p->txaio); - nni_aio_close(p->negaio); + nni_aio_close(p->negoaio); + nni_aio_close(p->connaio); nni_ipc_conn_close(p->conn); } @@ -107,17 +109,32 @@ ipctran_pipe_stop(void *arg) nni_aio_stop(p->rxaio); nni_aio_stop(p->txaio); - nni_aio_stop(p->negaio); + nni_aio_stop(p->negoaio); + nni_aio_stop(p->connaio); } static void ipctran_pipe_fini(void *arg) { ipctran_pipe *p = arg; + ipctran_ep * ep; + if (p == NULL) { + return; + } + ipctran_pipe_stop(p); + if ((ep = p->ep) != NULL) { + nni_mtx_lock(&ep->mtx); + nni_list_remove(&ep->pipes, p); + if (ep->fini && nni_list_empty(&ep->pipes)) { + nni_reap(&ep->reap, ipctran_ep_fini, ep); + } + nni_mtx_unlock(&ep->mtx); + } nni_aio_fini(p->rxaio); nni_aio_fini(p->txaio); - nni_aio_fini(p->negaio); + nni_aio_fini(p->negoaio); + nni_aio_fini(p->connaio); if (p->conn != NULL) { nni_ipc_conn_fini(p->conn); } @@ -128,8 +145,16 @@ ipctran_pipe_fini(void *arg) NNI_FREE_STRUCT(p); } +static void +ipctran_pipe_reap(ipctran_pipe *p) +{ + if (!nni_atomic_flag_test_and_set(&p->reaped)) { + nni_reap(&p->reap, ipctran_pipe_fini, p); + } +} + static int -ipctran_pipe_init(ipctran_pipe **pipep, void *conn) +ipctran_pipe_init(ipctran_pipe **pipep, ipctran_ep *ep) { ipctran_pipe *p; int rv; @@ -140,51 +165,78 @@ ipctran_pipe_init(ipctran_pipe **pipep, void *conn) nni_mtx_init(&p->mtx); if (((rv = nni_aio_init(&p->txaio, ipctran_pipe_send_cb, p)) != 0) || ((rv = nni_aio_init(&p->rxaio, ipctran_pipe_recv_cb, p)) != 0) || - ((rv = nni_aio_init(&p->negaio, ipctran_pipe_nego_cb, p)) != 0)) { + ((rv = nni_aio_init(&p->connaio, ipctran_pipe_conn_cb, p)) != 0) || + ((rv = nni_aio_init(&p->negoaio, ipctran_pipe_nego_cb, p)) != 0)) { ipctran_pipe_fini(p); return (rv); } nni_aio_list_init(&p->sendq); nni_aio_list_init(&p->recvq); + nni_atomic_flag_reset(&p->reaped); + nni_list_append(&ep->pipes, p); + + p->proto = ep->proto; + p->rcvmax = ep->rcvmax; + p->sa = ep->sa; + p->ep = ep; - p->conn = conn; -#if 0 - p->proto = ep->proto; - p->rcvmax = ep->rcvmax; - p->sa.s_ipc.sa_family = NNG_AF_IPC; - p->sa = ep->sa; -#endif *pipep = p; return (0); } static void -ipctran_pipe_nego_cancel(nni_aio *aio, int rv) +ipctran_pipe_conn_cb(void *arg) { - ipctran_pipe *p = nni_aio_get_prov_data(aio); + ipctran_pipe *p = arg; + nni_aio * aio = p->connaio; + nni_iov iov; + int rv; - nni_mtx_lock(&p->mtx); - if (p->user_negaio != aio) { - nni_mtx_unlock(&p->mtx); + nni_mtx_lock(&p->ep->mtx); + if ((rv = nni_aio_result(aio)) != 0) { + nni_aio *uaio = p->useraio; + p->useraio = NULL; + nni_mtx_unlock(&p->ep->mtx); + if (uaio != NULL) { + nni_aio_finish_error(uaio, rv); + } return; } - p->user_negaio = NULL; - nni_mtx_unlock(&p->mtx); + p->conn = nni_aio_get_output(aio, 0); + p->txhead[0] = 0; + p->txhead[1] = 'S'; + p->txhead[2] = 'P'; + p->txhead[3] = 0; + NNI_PUT16(&p->txhead[4], p->proto); + NNI_PUT16(&p->txhead[6], 0); - nni_aio_abort(p->negaio, rv); - nni_aio_finish_error(aio, rv); + p->gotrxhead = 0; + p->gottxhead = 0; + p->wantrxhead = 8; + p->wanttxhead = 8; + iov.iov_len = 8; + iov.iov_buf = &p->txhead[0]; + nni_aio_set_iov(p->negoaio, 1, &iov); + nni_ipc_conn_send(p->conn, p->negoaio); + nni_mtx_unlock(&p->ep->mtx); } static void ipctran_pipe_nego_cb(void *arg) { ipctran_pipe *p = arg; - nni_aio * aio = p->negaio; + nni_aio * aio = p->negoaio; + nni_aio * uaio; int rv; - nni_mtx_lock(&p->mtx); + nni_mtx_lock(&p->ep->mtx); + if ((uaio = p->useraio) == NULL) { + nni_mtx_unlock(&p->ep->mtx); + ipctran_pipe_reap(p); + return; + } if ((rv = nni_aio_result(aio)) != 0) { - goto done; + goto error; } // We start transmitting before we receive. @@ -193,7 +245,6 @@ ipctran_pipe_nego_cb(void *arg) } else if (p->gotrxhead < p->wantrxhead) { p->gotrxhead += nni_aio_count(aio); } - if (p->gottxhead < p->wanttxhead) { nni_iov iov; iov.iov_len = p->wanttxhead - p->gottxhead; @@ -201,7 +252,7 @@ ipctran_pipe_nego_cb(void *arg) nni_aio_set_iov(aio, 1, &iov); // send it down... nni_ipc_conn_send(p->conn, aio); - nni_mtx_unlock(&p->mtx); + nni_mtx_unlock(&p->ep->mtx); return; } if (p->gotrxhead < p->wantrxhead) { @@ -210,7 +261,7 @@ ipctran_pipe_nego_cb(void *arg) iov.iov_buf = &p->rxhead[p->gotrxhead]; nni_aio_set_iov(aio, 1, &iov); nni_ipc_conn_recv(p->conn, aio); - nni_mtx_unlock(&p->mtx); + nni_mtx_unlock(&p->ep->mtx); return; } // We have both sent and received the headers. Lets check the @@ -219,17 +270,21 @@ ipctran_pipe_nego_cb(void *arg) (p->rxhead[2] != 'P') || (p->rxhead[3] != 0) || (p->rxhead[6] != 0) || (p->rxhead[7] != 0)) { rv = NNG_EPROTO; - goto done; + goto error; } NNI_GET16(&p->rxhead[4], p->peer); + p->useraio = NULL; + nni_mtx_unlock(&p->ep->mtx); + nni_aio_set_output(uaio, 0, p); + nni_aio_finish(uaio, 0, 0); + return; -done: - if ((aio = p->user_negaio) != NULL) { - p->user_negaio = NULL; - nni_aio_finish(aio, rv, 0); - } - nni_mtx_unlock(&p->mtx); +error: + p->useraio = NULL; + nni_mtx_unlock(&p->ep->mtx); + nni_aio_finish_error(uaio, rv); + ipctran_pipe_reap(p); } static void @@ -243,17 +298,18 @@ ipctran_pipe_send_cb(void *arg) nni_aio * txaio = p->txaio; nni_mtx_lock(&p->mtx); - aio = nni_list_first(&p->sendq); - if ((rv = nni_aio_result(txaio)) != 0) { // Intentionally we do not queue up another transfer. // There's an excellent chance that the pipe is no longer // usable, with a partial transfer. // The protocol should see this error, and close the // pipe itself, we hope. - nni_aio_list_remove(aio); + + while ((aio = nni_list_first(&p->sendq)) != NULL) { + nni_aio_list_remove(aio); + nni_aio_finish_error(aio, rv); + } nni_mtx_unlock(&p->mtx); - nni_aio_finish_error(aio, rv); return; } @@ -265,6 +321,7 @@ ipctran_pipe_send_cb(void *arg) return; } + aio = nni_list_first(&p->sendq); nni_aio_list_remove(aio); ipctran_pipe_send_start(p); @@ -294,7 +351,7 @@ ipctran_pipe_recv_cb(void *arg) // Error on receive. This has to cause an error back // to the user. Also, if we had allocated an rxmsg, lets // toss it. - goto recv_error; + goto error; } n = nni_aio_count(rxaio); @@ -315,7 +372,7 @@ ipctran_pipe_recv_cb(void *arg) // Check to make sure we got msg type 1. if (p->rxhead[0] != 1) { rv = NNG_EPROTO; - goto recv_error; + goto error; } // We should have gotten a message header. @@ -325,7 +382,7 @@ ipctran_pipe_recv_cb(void *arg) // the caller will shut down the pipe. if ((len > p->rcvmax) && (p->rcvmax > 0)) { rv = NNG_EMSGSIZE; - goto recv_error; + goto error; } // Note that all IO on this pipe is blocked behind this @@ -334,7 +391,7 @@ ipctran_pipe_recv_cb(void *arg) // transmits to proceed normally. In practice this is // unlikely to be much of an issue though. if ((rv = nni_msg_alloc(&p->rxmsg, (size_t) len)) != 0) { - goto recv_error; + goto error; } if (len != 0) { @@ -354,20 +411,22 @@ ipctran_pipe_recv_cb(void *arg) // Otherwise we got a message read completely. Let the user know the // good news. + aio = nni_list_first(&p->recvq); nni_aio_list_remove(aio); msg = p->rxmsg; p->rxmsg = NULL; - if (!nni_list_empty(&p->recvq)) { - ipctran_pipe_recv_start(p); - } + ipctran_pipe_recv_start(p); nni_mtx_unlock(&p->mtx); nni_aio_set_msg(aio, msg); nni_aio_finish_synch(aio, 0, nni_msg_len(msg)); return; -recv_error: - nni_aio_list_remove(aio); +error: + while ((aio = nni_list_first(&p->recvq)) != NULL) { + nni_aio_list_remove(aio); + nni_aio_finish_error(aio, rv); + } msg = p->rxmsg; p->rxmsg = NULL; // Intentionally, we do not queue up another receive. @@ -375,7 +434,6 @@ recv_error: nni_mtx_unlock(&p->mtx); nni_msg_free(msg); - nni_aio_finish_error(aio, rv); } static void @@ -412,6 +470,13 @@ ipctran_pipe_send_start(ipctran_pipe *p) nni_iov iov[3]; uint64_t len; + if (p->closed) { + while ((aio = nni_list_first(&p->sendq)) != NULL) { + nni_list_remove(&p->sendq, aio); + nni_aio_finish_error(aio, NNG_ECLOSED); + } + return; + } if ((aio = nni_list_first(&p->sendq)) == NULL) { return; } @@ -494,6 +559,18 @@ ipctran_pipe_recv_start(ipctran_pipe *p) nni_iov iov; NNI_ASSERT(p->rxmsg == NULL); + if (p->closed) { + nni_aio *aio; + while ((aio = nni_list_first(&p->recvq)) != NULL) { + nni_list_remove(&p->recvq, aio); + nni_aio_finish_error(aio, NNG_ECLOSED); + } + return; + } + if (nni_list_empty(&p->recvq)) { + return; + } + // Schedule a read of the IPC header. rxaio = p->rxaio; iov.iov_buf = p->rxhead; @@ -513,7 +590,11 @@ ipctran_pipe_recv(void *arg, nni_aio *aio) return; } nni_mtx_lock(&p->mtx); - + if (p->closed) { + nni_mtx_unlock(&p->mtx); + nni_aio_finish_error(aio, NNG_ECLOSED); + return; + } if ((rv = nni_aio_schedule(aio, ipctran_pipe_recv_cancel, p)) != 0) { nni_mtx_unlock(&p->mtx); nni_aio_finish_error(aio, rv); @@ -527,43 +608,6 @@ ipctran_pipe_recv(void *arg, nni_aio *aio) nni_mtx_unlock(&p->mtx); } -static void -ipctran_pipe_start(void *arg, nni_aio *aio) -{ - ipctran_pipe *p = arg; - nni_aio * negaio; - nni_iov iov; - int rv; - - if (nni_aio_begin(aio) != 0) { - return; - } - nni_mtx_lock(&p->mtx); - if ((rv = nni_aio_schedule(aio, ipctran_pipe_nego_cancel, p)) != 0) { - nni_mtx_unlock(&p->mtx); - nni_aio_finish_error(aio, rv); - return; - } - p->txhead[0] = 0; - p->txhead[1] = 'S'; - p->txhead[2] = 'P'; - p->txhead[3] = 0; - NNI_PUT16(&p->txhead[4], p->proto); - NNI_PUT16(&p->txhead[6], 0); - - p->user_negaio = aio; - p->gotrxhead = 0; - p->gottxhead = 0; - p->wantrxhead = 8; - p->wanttxhead = 8; - negaio = p->negaio; - iov.iov_len = 8; - iov.iov_buf = &p->txhead[0]; - nni_aio_set_iov(negaio, 1, &iov); - nni_ipc_conn_send(p->conn, negaio); - nni_mtx_unlock(&p->mtx); -} - static uint16_t ipctran_pipe_peer(void *arg) { @@ -628,368 +672,217 @@ ipctran_pipe_get_peer_zoneid(void *arg, void *buf, size_t *szp, nni_opt_type t) } static void -ipctran_dialer_fini(void *arg) +ipctran_pipe_conn_cancel(nni_aio *aio, int rv) { - ipctran_dialer *d = arg; + ipctran_pipe *p = nni_aio_get_prov_data(aio); - nni_aio_stop(d->aio); - if (d->dialer != NULL) { - nni_ipc_dialer_fini(d->dialer); + nni_mtx_lock(&p->ep->mtx); + if (aio == p->useraio) { + nni_aio_close(p->negoaio); + nni_aio_close(p->connaio); + p->useraio = NULL; + nni_aio_finish_error(aio, rv); } - nni_aio_fini(d->aio); - nni_mtx_fini(&d->mtx); - NNI_FREE_STRUCT(d); + nni_mtx_unlock(&p->ep->mtx); } static void -ipctran_dialer_close(void *arg) +ipctran_ep_fini(void *arg) { - ipctran_dialer *d = arg; + ipctran_ep *ep = arg; - nni_aio_close(d->aio); - nni_ipc_dialer_close(d->dialer); -} - -static int -ipctran_dialer_init(void **dp, nni_url *url, nni_sock *sock) -{ - ipctran_dialer *d; - int rv; - size_t sz; - - if ((d = NNI_ALLOC_STRUCT(d)) == NULL) { - return (NNG_ENOMEM); - } - nni_mtx_init(&d->mtx); - - sz = sizeof(d->sa.s_ipc.sa_path); - d->sa.s_ipc.sa_family = NNG_AF_IPC; - - if (nni_strlcpy(d->sa.s_ipc.sa_path, url->u_path, sz) >= sz) { - ipctran_dialer_fini(d); - return (NNG_EADDRINVAL); - } - - if (((rv = nni_ipc_dialer_init(&d->dialer)) != 0) || - ((rv = nni_aio_init(&d->aio, ipctran_dialer_cb, d)) != 0)) { - ipctran_dialer_fini(d); - return (rv); - } - - d->proto = nni_sock_proto_id(sock); - - *dp = d; - return (0); -} - -static void -ipctran_dialer_cb(void *arg) -{ - ipctran_dialer *d = arg; - ipctran_pipe * p; - nni_ipc_conn * conn; - nni_aio * aio; - int rv; - - nni_mtx_lock(&d->mtx); - aio = d->user_aio; - rv = nni_aio_result(d->aio); - - if (aio == NULL) { - nni_mtx_unlock(&d->mtx); - if (rv == 0) { - conn = nni_aio_get_output(d->aio, 0); - nni_ipc_conn_fini(conn); - } + nni_mtx_lock(&ep->mtx); + ep->fini = true; + if (!nni_list_empty(&ep->pipes)) { + nni_mtx_unlock(&ep->mtx); return; } - - if (rv != 0) { - d->user_aio = NULL; - nni_mtx_unlock(&d->mtx); - nni_aio_finish_error(aio, rv); - return; + if (ep->dialer != NULL) { + nni_ipc_dialer_fini(ep->dialer); } - - d->user_aio = NULL; - conn = nni_aio_get_output(d->aio, 0); - NNI_ASSERT(conn != NULL); - if ((rv = ipctran_pipe_init(&p, conn)) != 0) { - nni_mtx_unlock(&d->mtx); - nni_ipc_conn_fini(conn); - nni_aio_finish_error(aio, rv); - return; + if (ep->listener != NULL) { + nni_ipc_listener_fini(ep->listener); } - - p->proto = d->proto; - p->rcvmax = d->rcvmax; - p->sa = d->sa; - nni_mtx_unlock(&d->mtx); - - nni_aio_set_output(aio, 0, p); - nni_aio_finish(aio, 0, 0); + nni_mtx_fini(&ep->mtx); + NNI_FREE_STRUCT(ep); } static void -ipctran_dialer_cancel(nni_aio *aio, int rv) +ipctran_ep_close(void *arg) { - ipctran_dialer *d = nni_aio_get_prov_data(aio); + ipctran_ep * ep = arg; + ipctran_pipe *p; - nni_mtx_lock(&d->mtx); - if (d->user_aio != aio) { - nni_mtx_unlock(&d->mtx); - return; + nni_mtx_lock(&ep->mtx); + NNI_LIST_FOREACH (&ep->pipes, p) { + nni_aio_close(p->negoaio); + nni_aio_close(p->connaio); + nni_aio_close(p->txaio); + nni_aio_close(p->rxaio); } - d->user_aio = NULL; - nni_mtx_unlock(&d->mtx); - - nni_aio_abort(d->aio, rv); - nni_aio_finish_error(aio, rv); -} - -static void -ipctran_dialer_connect(void *arg, nni_aio *aio) -{ - ipctran_dialer *d = arg; - int rv; - - if (nni_aio_begin(aio) != 0) { - return; + if (ep->dialer != NULL) { + nni_ipc_dialer_close(ep->dialer); } - nni_mtx_lock(&d->mtx); - NNI_ASSERT(d->user_aio == NULL); - - if ((rv = nni_aio_schedule(aio, ipctran_dialer_cancel, d)) != 0) { - nni_mtx_unlock(&d->mtx); - nni_aio_finish_error(aio, rv); - return; + if (ep->listener != NULL) { + nni_ipc_listener_close(ep->listener); } - d->user_aio = aio; - - nni_ipc_dialer_dial(d->dialer, &d->sa, d->aio); - nni_mtx_unlock(&d->mtx); + nni_mtx_unlock(&ep->mtx); } static int -ipctran_dialer_get_recvmaxsz(void *arg, void *v, size_t *szp, nni_opt_type t) +ipctran_ep_init_dialer(void **dp, nni_url *url, nni_sock *sock) { - ipctran_dialer *d = arg; - int rv; - nni_mtx_lock(&d->mtx); - rv = nni_copyout_size(d->rcvmax, v, szp, t); - nni_mtx_unlock(&d->mtx); - return (rv); -} + ipctran_ep *ep; + int rv; + size_t sz; -static int -ipctran_dialer_set_recvmaxsz( - void *arg, const void *v, size_t sz, nni_opt_type t) -{ - ipctran_dialer *d = arg; - size_t val; - int rv; - if ((rv = nni_copyin_size(&val, v, sz, 0, NNI_MAXSZ, t)) == 0) { - nni_mtx_lock(&d->mtx); - d->rcvmax = val; - nni_mtx_unlock(&d->mtx); + if ((ep = NNI_ALLOC_STRUCT(ep)) == NULL) { + return (NNG_ENOMEM); } - return (rv); -} + nni_mtx_init(&ep->mtx); + NNI_LIST_INIT(&ep->pipes, ipctran_pipe, node); -static void -ipctran_listener_fini(void *arg) -{ - ipctran_listener *l = arg; + sz = sizeof(ep->sa.s_ipc.sa_path); + ep->sa.s_ipc.sa_family = NNG_AF_IPC; + ep->proto = nni_sock_proto_id(sock); - nni_aio_stop(l->aio); - if (l->listener != NULL) { - nni_ipc_listener_fini(l->listener); + if (nni_strlcpy(ep->sa.s_ipc.sa_path, url->u_path, sz) >= sz) { + ipctran_ep_fini(ep); + return (NNG_EADDRINVAL); } - nni_aio_fini(l->aio); - nni_mtx_fini(&l->mtx); - NNI_FREE_STRUCT(l); + + if ((rv = nni_ipc_dialer_init(&ep->dialer)) != 0) { + ipctran_ep_fini(ep); + return (rv); + } + + *dp = ep; + return (0); } static int -ipctran_listener_init(void **lp, nni_url *url, nni_sock *sock) +ipctran_ep_init_listener(void **dp, nni_url *url, nni_sock *sock) { - ipctran_listener *l; - int rv; - size_t sz; + ipctran_ep *ep; + int rv; + size_t sz; - if ((l = NNI_ALLOC_STRUCT(l)) == NULL) { + if ((ep = NNI_ALLOC_STRUCT(ep)) == NULL) { return (NNG_ENOMEM); } - nni_mtx_init(&l->mtx); + nni_mtx_init(&ep->mtx); + NNI_LIST_INIT(&ep->pipes, ipctran_pipe, node); - sz = sizeof(l->sa.s_ipc.sa_path); - l->sa.s_ipc.sa_family = NNG_AF_IPC; + sz = sizeof(ep->sa.s_ipc.sa_path); + ep->sa.s_ipc.sa_family = NNG_AF_IPC; + ep->proto = nni_sock_proto_id(sock); - if (nni_strlcpy(l->sa.s_ipc.sa_path, url->u_path, sz) >= sz) { - ipctran_listener_fini(l); + if (nni_strlcpy(ep->sa.s_ipc.sa_path, url->u_path, sz) >= sz) { + ipctran_ep_fini(ep); return (NNG_EADDRINVAL); } - if (((rv = nni_ipc_listener_init(&l->listener)) != 0) || - ((rv = nni_aio_init(&l->aio, ipctran_listener_cb, l)) != 0)) { - ipctran_listener_fini(l); + if ((rv = nni_ipc_listener_init(&ep->listener)) != 0) { + ipctran_ep_fini(ep); return (rv); } - l->proto = nni_sock_proto_id(sock); - - *lp = l; + *dp = ep; return (0); } static void -ipctran_listener_close(void *arg) +ipctran_ep_connect(void *arg, nni_aio *aio) { - ipctran_listener *l = arg; + ipctran_ep * ep = arg; + ipctran_pipe *p = NULL; + int rv; - nni_aio_close(l->aio); - nni_ipc_listener_close(l->listener); + if (nni_aio_begin(aio) != 0) { + return; + } + nni_mtx_lock(&ep->mtx); + if (((rv = ipctran_pipe_init(&p, ep)) != 0) || + ((rv = nni_aio_schedule(aio, ipctran_pipe_conn_cancel, p)) != 0)) { + nni_mtx_unlock(&ep->mtx); + nni_aio_finish_error(aio, rv); + ipctran_pipe_reap(p); + return; + } + p->useraio = aio; + nni_ipc_dialer_dial(ep->dialer, &p->sa, p->connaio); + nni_mtx_unlock(&ep->mtx); } static int -ipctran_listener_bind(void *arg) +ipctran_ep_get_recvmaxsz(void *arg, void *v, size_t *szp, nni_opt_type t) { - ipctran_listener *l = arg; - int rv; - - nni_mtx_lock(&l->mtx); - rv = nni_ipc_listener_listen(l->listener, &l->sa); - nni_mtx_unlock(&l->mtx); + ipctran_ep *ep = arg; + int rv; + nni_mtx_lock(&ep->mtx); + rv = nni_copyout_size(ep->rcvmax, v, szp, t); + nni_mtx_unlock(&ep->mtx); return (rv); } -static void -ipctran_listener_cb(void *arg) +static int +ipctran_ep_set_recvmaxsz(void *arg, const void *v, size_t sz, nni_opt_type t) { - ipctran_listener *l = arg; - nni_aio * aio; - int rv; - ipctran_pipe * p = NULL; - nni_ipc_conn * conn; - - nni_mtx_lock(&l->mtx); - rv = nni_aio_result(l->aio); - aio = l->user_aio; - l->user_aio = NULL; - - if (aio == NULL) { - nni_mtx_unlock(&l->mtx); - if (rv == 0) { - conn = nni_aio_get_output(l->aio, 0); - nni_ipc_conn_fini(conn); - } - return; - } - - if (rv != 0) { - nni_mtx_unlock(&l->mtx); - nni_aio_finish_error(aio, rv); - return; - } - - conn = nni_aio_get_output(l->aio, 0); - NNI_ASSERT(conn != NULL); - - // Attempt to allocate the parent pipe. If this fails we'll - // drop the connection (ENOMEM probably). - if ((rv = ipctran_pipe_init(&p, conn)) != 0) { - nni_mtx_unlock(&l->mtx); - nni_ipc_conn_fini(conn); - nni_aio_finish_error(aio, rv); - return; + ipctran_ep *ep = arg; + size_t val; + int rv; + if ((rv = nni_copyin_size(&val, v, sz, 0, NNI_MAXSZ, t)) == 0) { + nni_mtx_lock(&ep->mtx); + ep->rcvmax = val; + nni_mtx_unlock(&ep->mtx); } - p->proto = l->proto; - p->rcvmax = l->rcvmax; - p->sa = l->sa; - nni_mtx_unlock(&l->mtx); - - nni_aio_set_output(aio, 0, p); - nni_aio_finish(aio, 0, 0); + return (rv); } -static void -ipctran_listener_cancel(nni_aio *aio, int rv) +static int +ipctran_ep_bind(void *arg) { - ipctran_listener *l = nni_aio_get_prov_data(aio); - - NNI_ASSERT(rv != 0); - nni_mtx_lock(&l->mtx); - if (l->user_aio != aio) { - nni_mtx_unlock(&l->mtx); - return; - } - l->user_aio = NULL; - nni_mtx_unlock(&l->mtx); + ipctran_ep *ep = arg; + int rv; - nni_aio_abort(l->aio, rv); - nni_aio_finish_error(aio, rv); + nni_mtx_lock(&ep->mtx); + rv = nni_ipc_listener_listen(ep->listener, &ep->sa); + nni_mtx_unlock(&ep->mtx); + return (rv); } static void -ipctran_listener_accept(void *arg, nni_aio *aio) +ipctran_ep_accept(void *arg, nni_aio *aio) { - ipctran_listener *l = arg; - int rv; + ipctran_ep * ep = arg; + ipctran_pipe *p = NULL; + int rv; if (nni_aio_begin(aio) != 0) { return; } - nni_mtx_lock(&l->mtx); - NNI_ASSERT(l->user_aio == NULL); - - if ((rv = nni_aio_schedule(aio, ipctran_listener_cancel, l)) != 0) { - nni_mtx_unlock(&l->mtx); + nni_mtx_lock(&ep->mtx); + if (((rv = ipctran_pipe_init(&p, ep)) != 0) || + ((rv = nni_aio_schedule(aio, ipctran_pipe_conn_cancel, p)) != 0)) { + nni_mtx_unlock(&ep->mtx); nni_aio_finish_error(aio, rv); + ipctran_pipe_reap(p); return; } - l->user_aio = aio; - - nni_ipc_listener_accept(l->listener, l->aio); - nni_mtx_unlock(&l->mtx); + p->useraio = aio; + nni_ipc_listener_accept(ep->listener, p->connaio); + nni_mtx_unlock(&ep->mtx); } static int -ipctran_listener_set_recvmaxsz( - void *arg, const void *data, size_t sz, nni_opt_type t) +ipctran_ep_get_locaddr(void *arg, void *buf, size_t *szp, nni_opt_type t) { - ipctran_listener *l = arg; - size_t val; - int rv; - - if ((rv = nni_copyin_size(&val, data, sz, 0, NNI_MAXSZ, t)) == 0) { - nni_mtx_lock(&l->mtx); - l->rcvmax = val; - nni_mtx_unlock(&l->mtx); - } - return (rv); -} + ipctran_ep *ep = arg; + int rv; -static int -ipctran_listener_get_recvmaxsz( - void *arg, void *data, size_t *szp, nni_opt_type t) -{ - ipctran_listener *l = arg; - int rv; - nni_mtx_lock(&l->mtx); - rv = nni_copyout_size(l->rcvmax, data, szp, t); - nni_mtx_unlock(&l->mtx); - return (rv); -} - -static int -ipctran_listener_get_locaddr(void *arg, void *buf, size_t *szp, nni_opt_type t) -{ - ipctran_listener *l = arg; - int rv; - - nni_mtx_lock(&l->mtx); - rv = nni_copyout_sockaddr(&l->sa, buf, szp, t); - nni_mtx_unlock(&l->mtx); + nni_mtx_lock(&ep->mtx); + rv = nni_copyout_sockaddr(&ep->sa, buf, szp, t); + nni_mtx_unlock(&ep->mtx); return (rv); } @@ -1000,19 +893,18 @@ ipctran_check_recvmaxsz(const void *data, size_t sz, nni_opt_type t) } static int -ipctran_listener_set_perms( - void *arg, const void *data, size_t sz, nni_opt_type t) +ipctran_ep_set_perms(void *arg, const void *data, size_t sz, nni_opt_type t) { - ipctran_listener *l = arg; - int val; - int rv; + ipctran_ep *ep = arg; + int val; + int rv; // Probably we could further limit this -- most systems don't have // meaningful chmod beyond the lower 9 bits. if ((rv = nni_copyin_int(&val, data, sz, 0, 0x7FFFFFFF, t)) == 0) { - nni_mtx_lock(&l->mtx); - rv = nni_ipc_listener_set_permissions(l->listener, val); - nni_mtx_unlock(&l->mtx); + nni_mtx_lock(&ep->mtx); + rv = nni_ipc_listener_set_permissions(ep->listener, val); + nni_mtx_unlock(&ep->mtx); } return (rv); } @@ -1024,18 +916,17 @@ ipctran_check_perms(const void *data, size_t sz, nni_opt_type t) } static int -ipctran_listener_set_sec_desc( - void *arg, const void *data, size_t sz, nni_opt_type t) +ipctran_ep_set_sec_desc(void *arg, const void *data, size_t sz, nni_opt_type t) { - ipctran_listener *l = arg; - void * ptr; - int rv; + ipctran_ep *ep = arg; + void * ptr; + int rv; if ((rv = nni_copyin_ptr(&ptr, data, sz, t)) == 0) { - nni_mtx_lock(&l->mtx); - rv = - nni_ipc_listener_set_security_descriptor(l->listener, ptr); - nni_mtx_unlock(&l->mtx); + nni_mtx_lock(&ep->mtx); + rv = nni_ipc_listener_set_security_descriptor( + ep->listener, ptr); + nni_mtx_unlock(&ep->mtx); } return (rv); } @@ -1085,7 +976,6 @@ static nni_tran_option ipctran_pipe_options[] = { static nni_tran_pipe_ops ipctran_pipe_ops = { .p_fini = ipctran_pipe_fini, - .p_start = ipctran_pipe_start, .p_stop = ipctran_pipe_stop, .p_send = ipctran_pipe_send, .p_recv = ipctran_pipe_recv, @@ -1094,45 +984,50 @@ static nni_tran_pipe_ops ipctran_pipe_ops = { .p_options = ipctran_pipe_options, }; -static nni_tran_option ipctran_dialer_options[] = { +static nni_tran_option ipctran_ep_dialer_options[] = { { .o_name = NNG_OPT_RECVMAXSZ, .o_type = NNI_TYPE_SIZE, - .o_get = ipctran_dialer_get_recvmaxsz, - .o_set = ipctran_dialer_set_recvmaxsz, + .o_get = ipctran_ep_get_recvmaxsz, + .o_set = ipctran_ep_set_recvmaxsz, .o_chk = ipctran_check_recvmaxsz, }, + { + .o_name = NNG_OPT_LOCADDR, + .o_type = NNI_TYPE_SOCKADDR, + .o_get = ipctran_ep_get_locaddr, + }, // terminate list { .o_name = NULL, }, }; -static nni_tran_option ipctran_listener_options[] = { +static nni_tran_option ipctran_ep_listener_options[] = { { .o_name = NNG_OPT_RECVMAXSZ, .o_type = NNI_TYPE_SIZE, - .o_get = ipctran_listener_get_recvmaxsz, - .o_set = ipctran_listener_set_recvmaxsz, + .o_get = ipctran_ep_get_recvmaxsz, + .o_set = ipctran_ep_set_recvmaxsz, .o_chk = ipctran_check_recvmaxsz, }, { .o_name = NNG_OPT_LOCADDR, .o_type = NNI_TYPE_SOCKADDR, - .o_get = ipctran_listener_get_locaddr, + .o_get = ipctran_ep_get_locaddr, }, { .o_name = NNG_OPT_IPC_SECURITY_DESCRIPTOR, .o_type = NNI_TYPE_POINTER, .o_get = NULL, - .o_set = ipctran_listener_set_sec_desc, + .o_set = ipctran_ep_set_sec_desc, .o_chk = ipctran_check_sec_desc, }, { .o_name = NNG_OPT_IPC_PERMISSIONS, .o_type = NNI_TYPE_INT32, .o_get = NULL, - .o_set = ipctran_listener_set_perms, + .o_set = ipctran_ep_set_perms, .o_chk = ipctran_check_perms, }, // terminate list @@ -1142,20 +1037,20 @@ static nni_tran_option ipctran_listener_options[] = { }; static nni_tran_dialer_ops ipctran_dialer_ops = { - .d_init = ipctran_dialer_init, - .d_fini = ipctran_dialer_fini, - .d_connect = ipctran_dialer_connect, - .d_close = ipctran_dialer_close, - .d_options = ipctran_dialer_options, + .d_init = ipctran_ep_init_dialer, + .d_fini = ipctran_ep_fini, + .d_connect = ipctran_ep_connect, + .d_close = ipctran_ep_close, + .d_options = ipctran_ep_dialer_options, }; static nni_tran_listener_ops ipctran_listener_ops = { - .l_init = ipctran_listener_init, - .l_fini = ipctran_listener_fini, - .l_bind = ipctran_listener_bind, - .l_accept = ipctran_listener_accept, - .l_close = ipctran_listener_close, - .l_options = ipctran_listener_options, + .l_init = ipctran_ep_init_listener, + .l_fini = ipctran_ep_fini, + .l_bind = ipctran_ep_bind, + .l_accept = ipctran_ep_accept, + .l_close = ipctran_ep_close, + .l_options = ipctran_ep_listener_options, }; static nni_tran ipc_tran = { |
