diff options
Diffstat (limited to 'src/protocol/survey0')
| -rw-r--r-- | src/protocol/survey0/respond.c | 18 | ||||
| -rw-r--r-- | src/protocol/survey0/survey.c | 19 | ||||
| -rw-r--r-- | src/protocol/survey0/xrespond.c | 26 | ||||
| -rw-r--r-- | src/protocol/survey0/xsurvey.c | 25 |
4 files changed, 66 insertions, 22 deletions
diff --git a/src/protocol/survey0/respond.c b/src/protocol/survey0/respond.c index e553f6ce..7738a8b7 100644 --- a/src/protocol/survey0/respond.c +++ b/src/protocol/survey0/respond.c @@ -290,6 +290,15 @@ resp0_sock_close(void *arg) } static void +resp0_pipe_stop(void *arg) +{ + resp0_pipe *p = arg; + + nni_aio_stop(p->aio_send); + nni_aio_stop(p->aio_recv); +} + +static void resp0_pipe_fini(void *arg) { resp0_pipe *p = arg; @@ -344,12 +353,15 @@ resp0_pipe_start(void *arg) } static void -resp0_pipe_stop(void *arg) +resp0_pipe_close(void *arg) { resp0_pipe *p = arg; resp0_sock *s = p->psock; resp0_ctx * ctx; + nni_aio_close(p->aio_send); + nni_aio_close(p->aio_recv); + nni_mtx_lock(&s->mtx); while ((ctx = nni_list_first(&p->sendq)) != NULL) { nni_aio *aio; @@ -369,9 +381,6 @@ resp0_pipe_stop(void *arg) } nni_idhash_remove(s->pipes, p->id); nni_mtx_unlock(&s->mtx); - - nni_aio_stop(p->aio_send); - nni_aio_stop(p->aio_recv); } static void @@ -626,6 +635,7 @@ static nni_proto_pipe_ops resp0_pipe_ops = { .pipe_init = resp0_pipe_init, .pipe_fini = resp0_pipe_fini, .pipe_start = resp0_pipe_start, + .pipe_close = resp0_pipe_close, .pipe_stop = resp0_pipe_stop, }; diff --git a/src/protocol/survey0/survey.c b/src/protocol/survey0/survey.c index e725d2b3..51bce0c8 100644 --- a/src/protocol/survey0/survey.c +++ b/src/protocol/survey0/survey.c @@ -284,6 +284,16 @@ surv0_sock_close(void *arg) } static void +surv0_pipe_stop(void *arg) +{ + surv0_pipe *p = arg; + + nni_aio_stop(p->aio_getq); + nni_aio_stop(p->aio_send); + nni_aio_stop(p->aio_recv); +} + +static void surv0_pipe_fini(void *arg) { surv0_pipe *p = arg; @@ -338,14 +348,14 @@ surv0_pipe_start(void *arg) } static void -surv0_pipe_stop(void *arg) +surv0_pipe_close(void *arg) { surv0_pipe *p = arg; surv0_sock *s = p->sock; - 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); @@ -532,6 +542,7 @@ static nni_proto_pipe_ops surv0_pipe_ops = { .pipe_init = surv0_pipe_init, .pipe_fini = surv0_pipe_fini, .pipe_start = surv0_pipe_start, + .pipe_close = surv0_pipe_close, .pipe_stop = surv0_pipe_stop, }; diff --git a/src/protocol/survey0/xrespond.c b/src/protocol/survey0/xrespond.c index 7aaed6da..bcbbcbc7 100644 --- a/src/protocol/survey0/xrespond.c +++ b/src/protocol/survey0/xrespond.c @@ -63,7 +63,6 @@ xresp0_sock_fini(void *arg) { xresp0_sock *s = arg; - nni_aio_stop(s->aio_getq); nni_aio_fini(s->aio_getq); nni_idhash_fini(s->pipes); nni_mtx_fini(&s->mtx); @@ -107,7 +106,18 @@ xresp0_sock_close(void *arg) { xresp0_sock *s = arg; - nni_aio_abort(s->aio_getq, NNG_ECLOSED); + nni_aio_close(s->aio_getq); +} + +static void +xresp0_pipe_stop(void *arg) +{ + xresp0_pipe *p = arg; + + nni_aio_stop(p->aio_putq); + nni_aio_stop(p->aio_getq); + nni_aio_stop(p->aio_send); + nni_aio_stop(p->aio_recv); } static void @@ -170,16 +180,17 @@ xresp0_pipe_start(void *arg) } static void -xresp0_pipe_stop(void *arg) +xresp0_pipe_close(void *arg) { xresp0_pipe *p = arg; xresp0_sock *s = p->psock; + nni_aio_close(p->aio_putq); + nni_aio_close(p->aio_getq); + nni_aio_close(p->aio_send); + nni_aio_close(p->aio_recv); + nni_msgq_close(p->sendq); - nni_aio_stop(p->aio_putq); - nni_aio_stop(p->aio_getq); - nni_aio_stop(p->aio_send); - nni_aio_stop(p->aio_recv); nni_mtx_lock(&s->mtx); nni_idhash_remove(s->pipes, p->id); @@ -366,6 +377,7 @@ static nni_proto_pipe_ops xresp0_pipe_ops = { .pipe_init = xresp0_pipe_init, .pipe_fini = xresp0_pipe_fini, .pipe_start = xresp0_pipe_start, + .pipe_close = xresp0_pipe_close, .pipe_stop = xresp0_pipe_stop, }; diff --git a/src/protocol/survey0/xsurvey.c b/src/protocol/survey0/xsurvey.c index cf311b15..47ebef3c 100644 --- a/src/protocol/survey0/xsurvey.c +++ b/src/protocol/survey0/xsurvey.c @@ -61,7 +61,6 @@ xsurv0_sock_fini(void *arg) { xsurv0_sock *s = arg; - nni_aio_stop(s->aio_getq); nni_aio_fini(s->aio_getq); nni_mtx_fini(&s->mtx); NNI_FREE_STRUCT(s); @@ -104,7 +103,18 @@ xsurv0_sock_close(void *arg) { xsurv0_sock *s = arg; - nni_aio_abort(s->aio_getq, NNG_ECLOSED); + nni_aio_close(s->aio_getq); +} + +static void +xsurv0_pipe_stop(void *arg) +{ + xsurv0_pipe *p = arg; + + nni_aio_stop(p->aio_getq); + nni_aio_stop(p->aio_send); + nni_aio_stop(p->aio_recv); + nni_aio_stop(p->aio_putq); } static void @@ -166,15 +176,15 @@ xsurv0_pipe_start(void *arg) } static void -xsurv0_pipe_stop(void *arg) +xsurv0_pipe_close(void *arg) { xsurv0_pipe *p = arg; xsurv0_sock *s = p->psock; - nni_aio_stop(p->aio_getq); - nni_aio_stop(p->aio_send); - nni_aio_stop(p->aio_recv); - nni_aio_stop(p->aio_putq); + nni_aio_close(p->aio_getq); + nni_aio_close(p->aio_send); + nni_aio_close(p->aio_recv); + nni_aio_close(p->aio_putq); nni_msgq_close(p->sendq); @@ -338,6 +348,7 @@ static nni_proto_pipe_ops xsurv0_pipe_ops = { .pipe_init = xsurv0_pipe_init, .pipe_fini = xsurv0_pipe_fini, .pipe_start = xsurv0_pipe_start, + .pipe_close = xsurv0_pipe_close, .pipe_stop = xsurv0_pipe_stop, }; |
