aboutsummaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
authorGarrett D'Amore <garrett@damore.org>2018-08-08 00:16:48 +0300
committerGarrett D'Amore <garrett@damore.org>2018-08-08 08:24:32 +0300
commit1a752ba26a7c65c386ae5431b2615ace25db9521 (patch)
treefbc1a586cff0f408626a04bddee591cb727a57a4 /src
parent53a6e9683589988971751256c4cd96cba79d07ab (diff)
downloadnng-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.c765
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 = {