diff options
Diffstat (limited to 'src')
| -rw-r--r-- | src/include/skabus/rpc.h | 31 | ||||
| -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 | ||||
| -rw-r--r-- | src/rpc/skabus-rpcd.c | 5 |
13 files changed, 112 insertions, 38 deletions
diff --git a/src/include/skabus/rpc.h b/src/include/skabus/rpc.h index 733a1bf..a96586a 100644 --- a/src/include/skabus/rpc.h +++ b/src/include/skabus/rpc.h @@ -69,14 +69,20 @@ extern skabus_rpc_interface_t const skabus_rpc_interface_zero ; extern skabus_rpc_r_func_t skabus_rpc_r_notimpl ; extern skabus_rpc_rcancel_func_t skabus_rpc_rcancel_ignore ; -extern int skabus_rpc_reply_withfds (skabus_rpc_t *, uint64_t, char, char const *, size_t, int const *, unsigned int, unsigned char const *) ; -#define skabus_rpc_reply(a, serial, result, s, len) skabus_rpc_reply_withfds(a, serial, result, s, len, 0, 0, unixmessage_bits_closenone) -extern int skabus_rpc_replyv_withfds (skabus_rpc_t *, uint64_t, char, struct iovec const *, unsigned int, int const *, unsigned int, unsigned char const *) ; -#define skabus_rpc_replyv(a, serial, result, v, vlen) skabus_rpc_replyv_withfds(a, serial, result, v, vlen, 0, 0, unixmessage_bits_closenone) +extern int skabus_rpc_reply_withfds_async (skabus_rpc_t *, uint64_t, char, char const *, size_t, int const *, unsigned int, unsigned char const *) ; +#define skabus_rpc_reply_async(a, serial, result, s, len) skabus_rpc_reply_withfds(a, serial, result, s, len, 0, 0, unixmessage_bits_closenone) +extern int skabus_rpc_replyv_withfds_async (skabus_rpc_t *, uint64_t, char, struct iovec const *, unsigned int, int const *, unsigned int, unsigned char const *) ; +#define skabus_rpc_replyv_async(a, serial, result, v, vlen) skabus_rpc_replyv_withfds(a, serial, result, v, vlen, 0, 0, unixmessage_bits_closenone) -#define skabus_rpc_rfd(a) skaclient_afd(&(a)->connection) -#define skabus_rpc_riswritable skaclient_aiswritable(&(a)->connection) -#define skabus_rpc_rflush(a) skaclient_aflush(&(a)->connection) +extern int skabus_rpc_reply_withfds (skabus_rpc_t *, uint64_t, char, char const *, size_t, int const *, unsigned int, unsigned char const *, tain_t const *, tain_t *) ; +#define skabus_rpc_reply(a, serial, result, s, len, deadline, stamp) skabus_rpc_reply_withfds(a, serial, result, s, len, 0, 0, unixmessage_bits_closenone, deadline, stamp) +extern int skabus_rpc_replyv_withfds (skabus_rpc_t *, uint64_t, char, struct iovec const *, unsigned int, int const *, unsigned int, unsigned char const *, tain_t const *, tain_t *) ; +#define skabus_rpc_replyv(a, serial, result, v, vlen, deadline, stamp) skabus_rpc_replyv_withfds(a, serial, result, v, vlen, 0, 0, unixmessage_bits_closenone, deadline, stamp) + +#define skabus_rpc_reply_withfds_g(a, serial, result, s, len, fds, nfds, bits, deadline) skabus_rpc_reply_withfds(a, serial, result, s, len, fds, nfds, bits, (deadline), &STAMP) +#define skabus_rpc_reply_g(a, serial, result, s, len, deadline) skabus_rpc_reply(a, serial, result, s, len, (deadline), &STAMP) +#define skabus_rpc_replyv_withfds_g(a, serial, result, v, vlen, fds, nfds, bits, deadline) skabus_rpc_replyv_withfds(a, serial, result, v, vlen, fds, nfds, bits, (deadline), &STAMP) +#define skabus_rpc_replyv_g(a, serial, result, v, vlen, deadline) skabus_rpc_replyv(a, serial, result, v, vlen, (deadline), &STAMP) /* Internal client interface storage */ @@ -153,9 +159,12 @@ extern void skabus_rpc_end (skabus_rpc_t *) ; /* Getting results */ +#define skabus_rpc_fd(a) skaclient_fd(&(a)->connection) extern int skabus_rpc_update (skabus_rpc_t *) ; +extern size_t skabus_rpc_qlist (skabus_rpc_t *, uint64_t **) ; extern int skabus_rpc_get (skabus_rpc_t *, uint64_t, int *, unixmessage_t *) ; extern int skabus_rpc_release (skabus_rpc_t *, uint64_t) ; +extern void skabus_rpc_qlist_ack(skabus_rpc_t *, size_t) ; /* Registering an interface */ @@ -191,7 +200,7 @@ extern int skabus_rpc_send_withfds_async (skabus_rpc_t *, char const *, char con #define skabus_rpc_send_async(a, ifname, s, len, limit, r) skabus_rpc_send_withfds_async(a, ifname, s, len, 0, 0, unixmessage_bits_closenone, limit, r) extern uint64_t skabus_rpc_send_withfds (skabus_rpc_t *, char const *, char const *, size_t, int const *, unsigned int, unsigned char const *, tain_t const *, tain_t const *, tain_t *) ; -#define skabus_rpc_send_withfds_g(a, ifname, s, len, fds, nfds, bits, limit, deadline) skabus_rpc_send_withfds(a, ifname, s, len, fds, nfds, limit, (deadline), &STAMP) +#define skabus_rpc_send_withfds_g(a, ifname, s, len, fds, nfds, bits, limit, deadline) skabus_rpc_send_withfds(a, ifname, s, len, fds, nfds, bits, limit, (deadline), &STAMP) #define skabus_rpc_send(a, ifname, s, len, limit, deadline, stamp) skabus_rpc_send_withfds(a, ifname, s, len, 0, 0, unixmessage_bits_closenone, limit, deadline, stamp) #define skabus_rpc_send_g(a, ifname, s, len, limit, deadline) skabus_rpc_send(a, ifname, s, len, limit, (deadline), &STAMP) @@ -199,7 +208,7 @@ extern int skabus_rpc_sendv_withfds_async (skabus_rpc_t *, char const *, struct #define skabus_rpc_sendv_async(a, ifname, v, vlen, limit, r) skabus_rpc_sendv_withfds_async(a, ifname, v, vlen, 0, 0, unixmessage_bits_closenone, limit, r) extern uint64_t skabus_rpc_sendv_withfds (skabus_rpc_t *, char const *, struct iovec const *, unsigned int, int const *, unsigned int, unsigned char const *, tain_t const *, tain_t const *, tain_t *) ; -#define skabus_rpc_sendv_withfds_g(a, ifname, v, vlen, fds, nfds, bits, limit, deadline) skabus_rpc_sendv_withfds(a, ifname, v, vlen, fds, nfds, limit, (deadline), &STAMP) +#define skabus_rpc_sendv_withfds_g(a, ifname, v, vlen, fds, nfds, bits, limit, deadline) skabus_rpc_sendv_withfds(a, ifname, v, vlen, fds, nfds, bits, limit, (deadline), &STAMP) #define skabus_rpc_sendv(a, ifname, v, vlen, limit, deadline, stamp) skabus_rpc_sendv_withfds(a, ifname, v, vlen, 0, 0, unixmessage_bits_closenone, limit, deadline, stamp) #define skabus_rpc_sendv_g(a, ifname, v, vlen, limit, deadline) skabus_rpc_sendv(a, ifname, v, vlen, limit, (deadline), &STAMP) @@ -207,7 +216,7 @@ extern int skabus_rpc_sendpm_withfds_async (skabus_rpc_t *, char const *, char c #define skabus_rpc_sendpm_async(a, cname, s, len, limit, r) skabus_rpc_sendpm_withfds_async(a, cname, s, len, 0, 0, unixmessage_bits_closenone, limit, r) extern uint64_t skabus_rpc_sendpm_withfds (skabus_rpc_t *, char const *, char const *, size_t, int const *, unsigned int, unsigned char const *, tain_t const *, tain_t const *, tain_t *) ; -#define skabus_rpc_sendpm_withfds_g(a, cname, s, len, fds, nfds, bits, limit, deadline) skabus_rpc_sendpm_withfds(a, cname, s, len, fds, nfds, limit, (deadline), &STAMP) +#define skabus_rpc_sendpm_withfds_g(a, cname, s, len, fds, nfds, bits, limit, deadline) skabus_rpc_sendpm_withfds(a, cname, s, len, fds, nfds, bits, limit, (deadline), &STAMP) #define skabus_rpc_sendpm(a, cname, s, len, limit, deadline, stamp) skabus_rpc_sendpm_withfds(a, cname, s, len, 0, 0, unixmessage_bits_closenone, limit, deadline, stamp) #define skabus_rpc_sendpm_g(a, cname, s, len, limit, deadline) skabus_rpc_sendpm(a, cname, s, len, limit, (deadline), &STAMP) @@ -215,7 +224,7 @@ extern int skabus_rpc_sendvpm_withfds_async (skabus_rpc_t *, char const *, struc #define skabus_rpc_sendvpm_async(a, cname, v, vlen, limit, r) skabus_rpc_sendvpm_withfds_async(a, cname, v, vlen, 0, 0, unixmessage_bits_closenone, limit, r) extern uint64_t skabus_rpc_sendvpm_withfds (skabus_rpc_t *, char const *, struct iovec const *, unsigned int, int const *, unsigned int, unsigned char const *, tain_t const *, tain_t const *, tain_t *) ; -#define skabus_rpc_sendvpm_withfds_g(a, cname, v, vlen, fds, nfds, bits, limit, deadline) skabus_rpc_sendvpm_withfds(a, cname, v, vlen, fds, nfds, limit, (deadline), &STAMP) +#define skabus_rpc_sendvpm_withfds_g(a, cname, v, vlen, fds, nfds, bits, limit, deadline) skabus_rpc_sendvpm_withfds(a, cname, v, vlen, fds, nfds, bits, limit, (deadline), &STAMP) #define skabus_rpc_sendvpm(a, cname, v, vlen, limit, deadline, stamp) skabus_rpc_sendvpm_withfds(a, cname, v, vlen, 0, 0, unixmessage_bits_closenone, limit, deadline, stamp) #define skabus_rpc_sendvpm_g(a, cname, v, vlen, limit, deadline) skabus_rpc_sendvpm(a, cname, v, vlen, limit, (deadline), &STAMP) 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) ; } diff --git a/src/rpc/skabus-rpcd.c b/src/rpc/skabus-rpcd.c index 36c6a78..b9398cc 100644 --- a/src/rpc/skabus-rpcd.c +++ b/src/rpc/skabus-rpcd.c @@ -100,6 +100,11 @@ int parse_protocol_async (unixmessage_t const *m, void *p) unixmessage_drop(m) ; return 1 ; } + if (INTERFACE(QUERY(qq)->interface)->client != *(uint32_t *)p) + { + unixmessage_drop(m) ; + return 1 ; + } query_reply(qq, m->s[9], &mtosend) ; return 1 ; } |
