summaryrefslogtreecommitdiff
path: root/src/protocol/pubsub0
diff options
context:
space:
mode:
authorGarrett D'Amore <garrett@damore.org>2018-05-15 01:47:12 -0700
committerGitHub <noreply@github.com>2018-05-15 01:47:12 -0700
commit1d033484ee1a2ec26d3eead073e7bc0f889ffdf4 (patch)
tree15d3897d405cb0beb1ada6270ecf70241451ca70 /src/protocol/pubsub0
parent16b4c4019c7b7904de171c588ed8c72ca732d2cf (diff)
downloadnng-1d033484ee1a2ec26d3eead073e7bc0f889ffdf4.tar.gz
nng-1d033484ee1a2ec26d3eead073e7bc0f889ffdf4.tar.bz2
nng-1d033484ee1a2ec26d3eead073e7bc0f889ffdf4.zip
fixes #419 want to nni_aio_stop without blocking (#428)
* fixes #419 want to nni_aio_stop without blocking This actually introduces an nni_aio_close() API that causes nni_aio_begin to return NNG_ECLOSED, while scheduling a callback on the AIO to do an NNG_ECLOSED as well. This should be called in non-blocking close() contexts instead of nni_aio_stop(), and the cases where we call nni_aio_fini() multiple times are updated updated to add nni_aio_stop() calls on all "interlinked" aios before finalizing them. Furthermore, we call nni_aio_close() as soon as practical in the close path. This closes an annoying race condition where the callback from a lower subsystem could wind up rescheduling an operation that we wanted to abort.
Diffstat (limited to 'src/protocol/pubsub0')
-rw-r--r--src/protocol/pubsub0/pub.c23
-rw-r--r--src/protocol/pubsub0/sub.c16
2 files changed, 30 insertions, 9 deletions
diff --git a/src/protocol/pubsub0/pub.c b/src/protocol/pubsub0/pub.c
index 45f4b7d9..4db48754 100644
--- a/src/protocol/pubsub0/pub.c
+++ b/src/protocol/pubsub0/pub.c
@@ -61,7 +61,6 @@ pub0_sock_fini(void *arg)
{
pub0_sock *s = arg;
- nni_aio_stop(s->aio_getq);
nni_aio_fini(s->aio_getq);
nni_mtx_fini(&s->mtx);
NNI_FREE_STRUCT(s);
@@ -103,13 +102,24 @@ pub0_sock_close(void *arg)
{
pub0_sock *s = arg;
- nni_aio_abort(s->aio_getq, NNG_ECLOSED);
+ nni_aio_close(s->aio_getq);
+}
+
+static void
+pub0_pipe_stop(void *arg)
+{
+ pub0_pipe *p = arg;
+
+ nni_aio_stop(p->aio_getq);
+ nni_aio_stop(p->aio_send);
+ nni_aio_stop(p->aio_recv);
}
static void
pub0_pipe_fini(void *arg)
{
pub0_pipe *p = arg;
+
nni_aio_fini(p->aio_getq);
nni_aio_fini(p->aio_send);
nni_aio_fini(p->aio_recv);
@@ -164,14 +174,14 @@ pub0_pipe_start(void *arg)
}
static void
-pub0_pipe_stop(void *arg)
+pub0_pipe_close(void *arg)
{
pub0_pipe *p = arg;
pub0_sock *s = p->pub;
- nni_aio_stop(p->aio_getq);
- nni_aio_stop(p->aio_send);
- nni_aio_stop(p->aio_recv);
+ nni_aio_close(p->aio_getq);
+ nni_aio_close(p->aio_send);
+ nni_aio_close(p->aio_recv);
nni_msgq_close(p->sendq);
@@ -290,6 +300,7 @@ static nni_proto_pipe_ops pub0_pipe_ops = {
.pipe_init = pub0_pipe_init,
.pipe_fini = pub0_pipe_fini,
.pipe_start = pub0_pipe_start,
+ .pipe_close = pub0_pipe_close,
.pipe_stop = pub0_pipe_stop,
};
diff --git a/src/protocol/pubsub0/sub.c b/src/protocol/pubsub0/sub.c
index b41b33ea..c244b0ad 100644
--- a/src/protocol/pubsub0/sub.c
+++ b/src/protocol/pubsub0/sub.c
@@ -99,6 +99,15 @@ sub0_sock_close(void *arg)
}
static void
+sub0_pipe_stop(void *arg)
+{
+ sub0_pipe *p = arg;
+
+ nni_aio_stop(p->aio_putq);
+ nni_aio_stop(p->aio_recv);
+}
+
+static void
sub0_pipe_fini(void *arg)
{
sub0_pipe *p = arg;
@@ -139,12 +148,12 @@ sub0_pipe_start(void *arg)
}
static void
-sub0_pipe_stop(void *arg)
+sub0_pipe_close(void *arg)
{
sub0_pipe *p = arg;
- nni_aio_stop(p->aio_putq);
- nni_aio_stop(p->aio_recv);
+ nni_aio_close(p->aio_putq);
+ nni_aio_close(p->aio_recv);
}
static void
@@ -338,6 +347,7 @@ static nni_proto_pipe_ops sub0_pipe_ops = {
.pipe_init = sub0_pipe_init,
.pipe_fini = sub0_pipe_fini,
.pipe_start = sub0_pipe_start,
+ .pipe_close = sub0_pipe_close,
.pipe_stop = sub0_pipe_stop,
};