diff options
Diffstat (limited to 'src/libskabus')
| -rw-r--r-- | src/libskabus/deps-lib/skabus | 4 | ||||
| -rw-r--r-- | src/libskabus/skabus_rpc_get.c | 10 | ||||
| -rw-r--r-- | src/libskabus/skabus_rpc_qlist.c | 12 | ||||
| -rw-r--r-- | src/libskabus/skabus_rpc_qlist_ack.c | 15 | ||||
| -rw-r--r-- | src/libskabus/skabus_rpc_r_notimpl.c | 2 | ||||
| -rw-r--r-- | src/libskabus/skabus_rpc_release.c | 6 | ||||
| -rw-r--r-- | src/libskabus/skabus_rpc_reply.c | 13 | ||||
| -rw-r--r-- | src/libskabus/skabus_rpc_reply_async.c | 17 | ||||
| -rw-r--r-- | src/libskabus/skabus_rpc_replyv.c | 15 | ||||
| -rw-r--r-- | src/libskabus/skabus_rpc_replyv_async.c | 19 | ||||
| -rw-r--r-- | src/libskabus/skabus_rpc_update.c | 1 |
11 files changed, 87 insertions, 27 deletions
diff --git a/src/libskabus/deps-lib/skabus b/src/libskabus/deps-lib/skabus index 423b672..21582c7 100644 --- a/src/libskabus/deps-lib/skabus +++ b/src/libskabus/deps-lib/skabus @@ -11,11 +11,15 @@ skabus_rpc_interface_unregister.o skabus_rpc_interface_unregister_async.o skabus_rpc_interface_zero.o skabus_rpc_qinfo_zero.o +skabus_rpc_qlist.o +skabus_rpc_qlist_ack.o skabus_rpc_r_notimpl.o skabus_rpc_rcancel_ignore.o skabus_rpc_release.o skabus_rpc_reply.o skabus_rpc_replyv.o +skabus_rpc_reply_async.o +skabus_rpc_replyv_async.o skabus_rpc_rinfo_pack.o skabus_rpc_rinfo_unpack.o skabus_rpc_rinfo_zero.o diff --git a/src/libskabus/skabus_rpc_get.c b/src/libskabus/skabus_rpc_get.c index d780211..132597a 100644 --- a/src/libskabus/skabus_rpc_get.c +++ b/src/libskabus/skabus_rpc_get.c @@ -2,6 +2,7 @@ #include <stdint.h> #include <errno.h> +#include <skalibs/error.h> #include <skalibs/gensetdyn.h> #include <skalibs/avltree.h> #include <skabus/rpc.h> @@ -10,9 +11,14 @@ int skabus_rpc_get (skabus_rpc_t *a, uint64_t serial, int *result, unixmessage_t { uint32_t id ; skabus_rpc_qinfo_t *p ; - if (!avltree_search(&a->qmap, &serial, &id)) return 0 ; + if (!avltree_search(&a->qmap, &serial, &id)) + { + if (errno == ESRCH) errno = EINVAL ; + return -1 ; + } p = GENSETDYN_P(skabus_rpc_qinfo_t, &a->q, id) ; - if (p->status) return (errno = p->status, 0) ; + if (p->status) + return error_isagain(p->status) ? 0 : (errno = p->status, -1) ; *result = (int)(unsigned char)p->result ; *m = p->message ; return 1 ; diff --git a/src/libskabus/skabus_rpc_qlist.c b/src/libskabus/skabus_rpc_qlist.c new file mode 100644 index 0000000..3f11a53 --- /dev/null +++ b/src/libskabus/skabus_rpc_qlist.c @@ -0,0 +1,12 @@ +/* ISC license. */ + +#include <skalibs/uint64.h> +#include <skalibs/genalloc.h> +#include <skabus/rpc.h> + +size_t skabus_rpc_qlist (skabus_rpc_t *a, uint64_t **list) +{ + uint64_t n = genalloc_len(uint64_t, &a->qlist) ; + if (n) *list = genalloc_s(uint64_t, &a->qlist) ; + return n ; +} diff --git a/src/libskabus/skabus_rpc_qlist_ack.c b/src/libskabus/skabus_rpc_qlist_ack.c new file mode 100644 index 0000000..894f530 --- /dev/null +++ b/src/libskabus/skabus_rpc_qlist_ack.c @@ -0,0 +1,15 @@ +/* ISC license. */ + +#include <string.h> +#include <skalibs/uint64.h> +#include <skalibs/genalloc.h> +#include <skabus/rpc.h> + +void skabus_rpc_qlist_ack (skabus_rpc_t *a, size_t n) +{ + uint64_t len = genalloc_len(uint64_t, &a->qlist) ; + uint64_t *p = genalloc_s(uint64_t, &a->qlist) ; + if (n > len) n = len ; + memmove(p, p + n * sizeof(uint64_t), (len - n) * sizeof(uint64_t)) ; + genalloc_setlen(uint64_t, &a->qlist, len - n) ; +} diff --git a/src/libskabus/skabus_rpc_r_notimpl.c b/src/libskabus/skabus_rpc_r_notimpl.c index 8ec3922..e36450f 100644 --- a/src/libskabus/skabus_rpc_r_notimpl.c +++ b/src/libskabus/skabus_rpc_r_notimpl.c @@ -7,5 +7,5 @@ int skabus_rpc_r_notimpl (skabus_rpc_t *a, skabus_rpc_rinfo_t const *info, unixm { (void)m ; (void)data ; - return skabus_rpc_reply(a, info->serial, ENOSYS, "", 0) ; + return skabus_rpc_reply(a, info->serial, ENOSYS, "", 0, 0, 0) ; } diff --git a/src/libskabus/skabus_rpc_release.c b/src/libskabus/skabus_rpc_release.c index d0a6875..80ad13e 100644 --- a/src/libskabus/skabus_rpc_release.c +++ b/src/libskabus/skabus_rpc_release.c @@ -12,7 +12,11 @@ int skabus_rpc_release (skabus_rpc_t *a, uint64_t serial) { uint32_t id ; skabus_rpc_qinfo_t *p ; - if (!avltree_search(&a->qmap, &serial, &id)) return 0 ; + if (!avltree_search(&a->qmap, &serial, &id)) + { + if (errno == ESRCH) errno = EINVAL ; + return 0 ; + } p = GENSETDYN_P(skabus_rpc_qinfo_t, &a->q, id) ; if (p->status) return (errno = p->status, 0) ; alloc_free(p->message.s) ; diff --git a/src/libskabus/skabus_rpc_reply.c b/src/libskabus/skabus_rpc_reply.c index 491772b..c1b72b2 100644 --- a/src/libskabus/skabus_rpc_reply.c +++ b/src/libskabus/skabus_rpc_reply.c @@ -1,17 +1,10 @@ /* ISC license. */ -#include <sys/uio.h> -#include <skalibs/uint64.h> -#include <skalibs/unixmessage.h> #include <skalibs/skaclient.h> #include <skabus/rpc.h> -int skabus_rpc_reply_withfds (skabus_rpc_t *a, uint64_t serial, char result, char const *s, size_t len, int const *fds, unsigned int nfds, unsigned char const *bits) +int skabus_rpc_reply_withfds (skabus_rpc_t *a, uint64_t serial, char result, char const *s, size_t len, int const *fds, unsigned int nfds, unsigned char const *bits, tain_t const *deadline, tain_t *stamp) { - char pack[10] = "R" ; - struct iovec v[2] = { { .iov_base = pack, .iov_len = 10 }, { .iov_base = (char *)s, .iov_len = len } } ; - unixmessage_v_t m = { .v = v, .vlen = 2, .fds = (int *)fds, .nfds = nfds } ; - uint64_pack_big(pack+1, serial) ; - pack[9] = result ; - return skaclient_aputv_and_close(&a->connection, &m, bits) ; + return skabus_rpc_reply_withfds_async(a, serial, result, s, len, fds, nfds, bits) + && skaclient_timed_aflush(&a->connection, deadline, stamp) ; } diff --git a/src/libskabus/skabus_rpc_reply_async.c b/src/libskabus/skabus_rpc_reply_async.c new file mode 100644 index 0000000..a6b106c --- /dev/null +++ b/src/libskabus/skabus_rpc_reply_async.c @@ -0,0 +1,17 @@ +/* ISC license. */ + +#include <sys/uio.h> +#include <skalibs/uint64.h> +#include <skalibs/unixmessage.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> + +int skabus_rpc_reply_withfds_async (skabus_rpc_t *a, uint64_t serial, char result, char const *s, size_t len, int const *fds, unsigned int nfds, unsigned char const *bits) +{ + char pack[10] = "R" ; + struct iovec v[2] = { { .iov_base = pack, .iov_len = 10 }, { .iov_base = (char *)s, .iov_len = len } } ; + unixmessage_v_t m = { .v = v, .vlen = 2, .fds = (int *)fds, .nfds = nfds } ; + uint64_pack_big(pack+1, serial) ; + pack[9] = result ; + return skaclient_aputv_and_close(&a->connection, &m, bits) ; +} diff --git a/src/libskabus/skabus_rpc_replyv.c b/src/libskabus/skabus_rpc_replyv.c index ef727e9..c972ce4 100644 --- a/src/libskabus/skabus_rpc_replyv.c +++ b/src/libskabus/skabus_rpc_replyv.c @@ -1,19 +1,10 @@ /* ISC license. */ -#include <sys/uio.h> -#include <skalibs/uint64.h> -#include <skalibs/unixmessage.h> #include <skalibs/skaclient.h> #include <skabus/rpc.h> -int skabus_rpc_replyv_withfds (skabus_rpc_t *a, uint64_t serial, char result, struct iovec const *v, unsigned int vlen, int const *fds, unsigned int nfds, unsigned char const *bits) +int skabus_rpc_replyv_withfds (skabus_rpc_t *a, uint64_t serial, char result, struct iovec const *v, unsigned int vlen, int const *fds, unsigned int nfds, unsigned char const *bits, tain_t const *deadline, tain_t *stamp) { - char pack[10] = "R" ; - struct iovec vv[vlen + 1] ; - unixmessage_v_t m = { .v = vv, .vlen = vlen+1, .fds = (int *)fds, .nfds = nfds } ; - vv[0].iov_base = pack ; vv[0].iov_len = 10 ; - for (unsigned int i = 0 ; i < vlen ; i++) vv[1+i] = v[i] ; - uint64_pack_big(pack+1, serial) ; - pack[9] = result ; - return skaclient_aputv_and_close(&a->connection, &m, bits) ; + return skabus_rpc_replyv_withfds_async(a, serial, result, v, vlen, fds, nfds, bits) + && skaclient_timed_aflush(&a->connection, deadline, stamp) ; } diff --git a/src/libskabus/skabus_rpc_replyv_async.c b/src/libskabus/skabus_rpc_replyv_async.c new file mode 100644 index 0000000..7bd4397 --- /dev/null +++ b/src/libskabus/skabus_rpc_replyv_async.c @@ -0,0 +1,19 @@ +/* ISC license. */ + +#include <sys/uio.h> +#include <skalibs/uint64.h> +#include <skalibs/unixmessage.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> + +int skabus_rpc_replyv_withfds_async (skabus_rpc_t *a, uint64_t serial, char result, struct iovec const *v, unsigned int vlen, int const *fds, unsigned int nfds, unsigned char const *bits) +{ + char pack[10] = "R" ; + struct iovec vv[vlen + 1] ; + unixmessage_v_t m = { .v = vv, .vlen = vlen+1, .fds = (int *)fds, .nfds = nfds } ; + vv[0].iov_base = pack ; vv[0].iov_len = 10 ; + for (unsigned int i = 0 ; i < vlen ; i++) vv[1+i] = v[i] ; + uint64_pack_big(pack+1, serial) ; + pack[9] = result ; + return skaclient_aputv_and_close(&a->connection, &m, bits) ; +} diff --git a/src/libskabus/skabus_rpc_update.c b/src/libskabus/skabus_rpc_update.c index d727de1..8a24ec3 100644 --- a/src/libskabus/skabus_rpc_update.c +++ b/src/libskabus/skabus_rpc_update.c @@ -115,6 +115,5 @@ static int handler (unixmessage_t const *m, void *x) int skabus_rpc_update (skabus_rpc_t *a) { - genalloc_setlen(uint64_t, &a->qlist, 0) ; return skaclient_update(&a->connection, &handler, a) ; } |
