diff options
Diffstat (limited to 'src/libskabus')
40 files changed, 809 insertions, 0 deletions
diff --git a/src/libskabus/deps-lib/skabus b/src/libskabus/deps-lib/skabus new file mode 100644 index 0000000..423b672 --- /dev/null +++ b/src/libskabus/deps-lib/skabus @@ -0,0 +1,39 @@ +skabus_rpc_cancel.o +skabus_rpc_cancel_async.o +skabus_rpc_end.o +skabus_rpc_get.o +skabus_rpc_idstr.o +skabus_rpc_idstr_async.o +skabus_rpc_interface_register.o +skabus_rpc_interface_register_async.o +skabus_rpc_interface_register_cb.o +skabus_rpc_interface_unregister.o +skabus_rpc_interface_unregister_async.o +skabus_rpc_interface_zero.o +skabus_rpc_qinfo_zero.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_rinfo_pack.o +skabus_rpc_rinfo_unpack.o +skabus_rpc_rinfo_zero.o +skabus_rpc_send.o +skabus_rpc_send_async.o +skabus_rpc_send_cb.o +skabus_rpc_sendpm.o +skabus_rpc_sendpm_async.o +skabus_rpc_sendq.o +skabus_rpc_sendq_async.o +skabus_rpc_sendv.o +skabus_rpc_sendv_async.o +skabus_rpc_sendvpm.o +skabus_rpc_sendvpm_async.o +skabus_rpc_sendvq.o +skabus_rpc_sendvq_async.o +skabus_rpc_start.o +skabus_rpc_start_async.o +skabus_rpc_update.o +skabus_rpc_zero.o +-lskarnet diff --git a/src/libskabus/skabus-rpc-internal.h b/src/libskabus/skabus-rpc-internal.h new file mode 100644 index 0000000..04f29e4 --- /dev/null +++ b/src/libskabus/skabus-rpc-internal.h @@ -0,0 +1,20 @@ +/* ISC license. */ + +#ifndef SKABUS_RPC_INTERNAL_H +#define SKABUS_RPC_INTERNAL_H + +#include <sys/uio.h> +#include <skalibs/uint64.h> +#include <skalibs/tai.h> +#include <skalibs/unixmessage.h> +#include <skabus/rpc.h> + +extern int skabus_rpc_sendq_withfds_async (skabus_rpc_t *, char const *, size_t, char const *, char const *, size_t, int const *, unsigned int, unsigned char const *, tain_t const *, skabus_rpc_send_result_t *) ; +extern uint64_t skabus_rpc_sendq_withfds (skabus_rpc_t *, char const *, size_t, char const *, char const *, size_t, int const *, unsigned int, unsigned char const *, tain_t const *, tain_t const *, tain_t *) ; +extern int skabus_rpc_sendvq_withfds_async (skabus_rpc_t *, char const *, size_t, char const *, struct iovec const *, unsigned int, int const *, unsigned int, unsigned char const *, tain_t const *, skabus_rpc_send_result_t *) ; +extern uint64_t skabus_rpc_sendvq_withfds (skabus_rpc_t *, char const *, size_t, char const *, struct iovec const *, unsigned int, int const *, unsigned int, unsigned char const *, tain_t const *, tain_t const *, tain_t *) ; + +extern unixmessage_handler_func_t skabus_rpc_send_cb ; +extern unixmessage_handler_func_t skabus_rpc_interface_register_cb ; + +#endif diff --git a/src/libskabus/skabus_rpc_cancel.c b/src/libskabus/skabus_rpc_cancel.c new file mode 100644 index 0000000..d731c6e --- /dev/null +++ b/src/libskabus/skabus_rpc_cancel.c @@ -0,0 +1,15 @@ +/* ISC license. */ + +#include <errno.h> +#include <skalibs/tai.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> + +int skabus_rpc_cancel (skabus_rpc_t *a, uint64_t serial, tain_t const *deadline, tain_t *stamp) +{ + unsigned char r ; + if (!skabus_rpc_cancel_async(a, serial, &r)) return 0 ; + if (!skaclient_syncify(&a->connection, deadline, stamp)) return 0 ; + if (r) return (errno = r, 0) ; + return 1 ; +} diff --git a/src/libskabus/skabus_rpc_cancel_async.c b/src/libskabus/skabus_rpc_cancel_async.c new file mode 100644 index 0000000..f836ed1 --- /dev/null +++ b/src/libskabus/skabus_rpc_cancel_async.c @@ -0,0 +1,25 @@ +/* ISC license. */ + +#include <stdint.h> +#include <errno.h> +#include <skalibs/uint64.h> +#include <skalibs/gensetdyn.h> +#include <skalibs/avltree.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> + +int skabus_rpc_cancel_async (skabus_rpc_t *a, uint64_t serial, unsigned char *err) +{ + uint32_t id ; + skabus_rpc_qinfo_t *p ; + char pack[9] = "C" ; + if (!avltree_search(&a->qmap, &serial, &id)) return 0 ; + p = GENSETDYN_P(skabus_rpc_qinfo_t, &a->q, id) ; + if (p->status != EAGAIN) return (errno = EINVAL, 0) ; + uint64_pack_big(pack+1, serial) ; + if (!skaclient_put(&a->connection, pack, 9, &skaclient_default_cb, err)) return 0 ; + p->status = EINVAL ; + avltree_delete(&a->qmap, &serial) ; + gensetdyn_delete(&a->q, id) ; + return 1 ; +} diff --git a/src/libskabus/skabus_rpc_end.c b/src/libskabus/skabus_rpc_end.c new file mode 100644 index 0000000..34c95fb --- /dev/null +++ b/src/libskabus/skabus_rpc_end.c @@ -0,0 +1,19 @@ +/* ISC license. */ + +#include <stdint.h> +#include <skalibs/uint64.h> +#include <skalibs/genalloc.h> +#include <skalibs/gensetdyn.h> +#include <skalibs/avltree.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> + +void skabus_rpc_end (skabus_rpc_t *a) +{ + skaclient_end(&a->connection) ; + genalloc_free(uint64_t, &a->qlist) ; + avltree_free(&a->qmap) ; + gensetdyn_free(&a->q) ; + gensetdyn_free(&a->r) ; + a->pmid = (uint32_t)-1 ; +} diff --git a/src/libskabus/skabus_rpc_get.c b/src/libskabus/skabus_rpc_get.c new file mode 100644 index 0000000..d780211 --- /dev/null +++ b/src/libskabus/skabus_rpc_get.c @@ -0,0 +1,19 @@ +/* ISC license. */ + +#include <stdint.h> +#include <errno.h> +#include <skalibs/gensetdyn.h> +#include <skalibs/avltree.h> +#include <skabus/rpc.h> + +int skabus_rpc_get (skabus_rpc_t *a, uint64_t serial, int *result, unixmessage_t *m) +{ + uint32_t id ; + skabus_rpc_qinfo_t *p ; + if (!avltree_search(&a->qmap, &serial, &id)) return 0 ; + p = GENSETDYN_P(skabus_rpc_qinfo_t, &a->q, id) ; + if (p->status) return (errno = p->status, 0) ; + *result = (int)(unsigned char)p->result ; + *m = p->message ; + return 1 ; +} diff --git a/src/libskabus/skabus_rpc_idstr.c b/src/libskabus/skabus_rpc_idstr.c new file mode 100644 index 0000000..887e100 --- /dev/null +++ b/src/libskabus/skabus_rpc_idstr.c @@ -0,0 +1,14 @@ +/* ISC license. */ + +#include <errno.h> +#include <skalibs/tai.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> + +int skabus_rpc_idstr (skabus_rpc_t *a, char const *idstr, skabus_rpc_interface_t const *ifbody, char const *re, tain_t const *deadline, tain_t *stamp) +{ + skabus_rpc_interface_result_t r ; + if (!skabus_rpc_idstr_async(a, idstr, ifbody, re, &r)) return 0 ; + if (!skaclient_syncify(&a->connection, deadline, stamp)) return 0 ; + return r.err ? (errno = r.err, 0) : 1 ; +} diff --git a/src/libskabus/skabus_rpc_idstr_async.c b/src/libskabus/skabus_rpc_idstr_async.c new file mode 100644 index 0000000..9737086 --- /dev/null +++ b/src/libskabus/skabus_rpc_idstr_async.c @@ -0,0 +1,31 @@ +/* ISC license. */ + +#include <string.h> +#include <sys/uio.h> +#include <errno.h> +#include <skalibs/uint32.h> +#include <skalibs/gensetdyn.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +int skabus_rpc_idstr_async (skabus_rpc_t *a, char const *idstr, skabus_rpc_interface_t const *ifbody, char const *re, skabus_rpc_interface_result_t *result) +{ + size_t idlen = strlen(idstr) ; + size_t relen = strlen(re) ; + skabus_rpc_ifnode_t *ifnode ; + char pack[10] = "S" ; + struct iovec v[3] = { { .iov_base = pack, .iov_len = 10 }, { .iov_base = (char *)idstr, .iov_len = idlen+1 }, { .iov_base = (char *)re, .iov_len = relen+1 } } ; + if (idlen > SKABUS_RPC_IDSTR_SIZE) return (errno = ENAMETOOLONG, 0) ; + if (relen > 0x6ffffffe) return (errno = ENAMETOOLONG, 0) ; + if (!gensetdyn_new(&a->r, &a->pmid)) return 0 ; + result->ifid = a->pmid ; + result->r = &a->r ; + ifnode = GENSETDYN_P(skabus_rpc_ifnode_t, &a->r, a->pmid) ; + ifnode->name[0] = 0 ; + ifnode->body = *ifbody ; + uint32_pack_big(pack+1, a->pmid) ; + pack[5] = (unsigned char)idlen ; + uint32_pack_big(pack+6, relen) ; + return skaclient_putv(&a->connection, v, 3, &skabus_rpc_interface_register_cb, result) ; +} diff --git a/src/libskabus/skabus_rpc_interface_register.c b/src/libskabus/skabus_rpc_interface_register.c new file mode 100644 index 0000000..5978a7b --- /dev/null +++ b/src/libskabus/skabus_rpc_interface_register.c @@ -0,0 +1,16 @@ +/* ISC license. */ + +#include <errno.h> +#include <skalibs/tai.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> + +int skabus_rpc_interface_register (skabus_rpc_t *a, uint32_t *ifid, char const *ifname, skabus_rpc_interface_t const *ifbody, char const *re, tain_t const *deadline, tain_t *stamp) +{ + skabus_rpc_interface_result_t r ; + if (!skabus_rpc_interface_register_async(a, ifname, ifbody, re, &r)) return 0 ; + if (!skaclient_syncify(&a->connection, deadline, stamp)) return 0 ; + if (r.err) return (errno = r.err, 0) ; + *ifid = r.ifid ; + return 1 ; +} diff --git a/src/libskabus/skabus_rpc_interface_register_async.c b/src/libskabus/skabus_rpc_interface_register_async.c new file mode 100644 index 0000000..986ae5f --- /dev/null +++ b/src/libskabus/skabus_rpc_interface_register_async.c @@ -0,0 +1,30 @@ +/* ISC license. */ + +#include <string.h> +#include <sys/uio.h> +#include <errno.h> +#include <skalibs/uint32.h> +#include <skalibs/error.h> +#include <skalibs/gensetdyn.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +int skabus_rpc_interface_register_async (skabus_rpc_t *a, char const *ifname, skabus_rpc_interface_t const *ifbody, char const *re, skabus_rpc_interface_result_t *result) +{ + size_t ifnamelen = strlen(ifname) ; + size_t relen = strlen(re) ; + skabus_rpc_ifnode_t *ifnode ; + char pack[10] = "I" ; + struct iovec v[3] = { { .iov_base = pack, .iov_len = 10 }, { .iov_base = (char *)ifname, .iov_len = ifnamelen + 1 }, { .iov_base = (char *)re, .iov_len = relen + 1 } } ; + if (ifnamelen > SKABUS_RPC_INTERFACE_MAXLEN) return (errno = ENAMETOOLONG, 0) ; + if (!gensetdyn_new(&a->r, &result->ifid)) return 0 ; + result->r = &a->r ; + ifnode = GENSETDYN_P(skabus_rpc_ifnode_t, &a->r, result->ifid) ; + memcpy(ifnode->name, ifname, ifnamelen + 1) ; + ifnode->body = *ifbody ; + uint32_pack_big(pack+1, result->ifid) ; + pack[5] = (unsigned char)ifnamelen ; + uint32_pack_big(pack+6, relen) ; + return skaclient_putv(&a->connection, v, 3, &skabus_rpc_interface_register_cb, result) ; +} diff --git a/src/libskabus/skabus_rpc_interface_register_cb.c b/src/libskabus/skabus_rpc_interface_register_cb.c new file mode 100644 index 0000000..88d27e7 --- /dev/null +++ b/src/libskabus/skabus_rpc_interface_register_cb.c @@ -0,0 +1,17 @@ + /* ISC license. */ + +#include <errno.h> +#include <skalibs/error.h> +#include <skalibs/gensetdyn.h> +#include <skalibs/unixmessage.h> +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +int skabus_rpc_interface_register_cb (unixmessage_t const *m, void *p) +{ + skabus_rpc_interface_result_t *r = p ; + if (m->len != 1 || m->nfds) return (errno = EPROTO, 0) ; + r->err = m->s[0] ; + if (r->err) gensetdyn_delete(r->r, r->ifid) ; + return 1 ; +} diff --git a/src/libskabus/skabus_rpc_interface_unregister.c b/src/libskabus/skabus_rpc_interface_unregister.c new file mode 100644 index 0000000..9263c0c --- /dev/null +++ b/src/libskabus/skabus_rpc_interface_unregister.c @@ -0,0 +1,15 @@ +/* ISC license. */ + +#include <errno.h> +#include <skalibs/tai.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> + +int skabus_rpc_interface_unregister (skabus_rpc_t *a, uint32_t ifid, tain_t const *deadline, tain_t *stamp) +{ + skabus_rpc_interface_result_t r ; + if (!skabus_rpc_interface_unregister_async(a, ifid, &r)) return 0 ; + if (!skaclient_syncify(&a->connection, deadline, stamp)) return 0 ; + if (r.err) return (errno = r.err, 0) ; + return 1 ; +} diff --git a/src/libskabus/skabus_rpc_interface_unregister_async.c b/src/libskabus/skabus_rpc_interface_unregister_async.c new file mode 100644 index 0000000..54a05ba --- /dev/null +++ b/src/libskabus/skabus_rpc_interface_unregister_async.c @@ -0,0 +1,29 @@ +/* ISC license. */ + +#include <string.h> +#include <sys/uio.h> +#include <errno.h> +#include <skalibs/error.h> +#include <skalibs/gensetdyn.h> +#include <skalibs/unixmessage.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> + +static int skabus_rpc_interface_unregister_cb (unixmessage_t const *m, void *p) +{ + skabus_rpc_interface_result_t *r = p ; + if (m->len != 1 || m->nfds) return (errno = EPROTO, 0) ; + r->err = m->s[0] ; + return r->err ? 1 : gensetdyn_delete(r->r, r->ifid) ; +} + +int skabus_rpc_interface_unregister_async (skabus_rpc_t *a, uint32_t ifid, skabus_rpc_interface_result_t *result) +{ + char *name = GENSETDYN_P(skabus_rpc_ifnode_t, &a->r, ifid)->name ; + size_t n = strlen(name) ; + char pack[2] = { 'i', (unsigned char)n } ; + struct iovec v[2] = { { .iov_base = pack, .iov_len = 2 }, { .iov_base = name, .iov_len = n+1 } } ; + result->r = &a->r ; + result->ifid = ifid ; + return skaclient_putv(&a->connection, v, 2, &skabus_rpc_interface_unregister_cb, result) ; +} diff --git a/src/libskabus/skabus_rpc_interface_zero.c b/src/libskabus/skabus_rpc_interface_zero.c new file mode 100644 index 0000000..df94b97 --- /dev/null +++ b/src/libskabus/skabus_rpc_interface_zero.c @@ -0,0 +1,5 @@ +/* ISC license. */ + +#include <skabus/rpc.h> + +skabus_rpc_interface_t const skabus_rpc_interface_zero = SKABUS_RPC_INTERFACE_ZERO ; diff --git a/src/libskabus/skabus_rpc_qinfo_zero.c b/src/libskabus/skabus_rpc_qinfo_zero.c new file mode 100644 index 0000000..257f34b --- /dev/null +++ b/src/libskabus/skabus_rpc_qinfo_zero.c @@ -0,0 +1,5 @@ +/* ISC license. */ + +#include <skabus/rpc.h> + +skabus_rpc_qinfo_t const skabus_rpc_qinfo_zero = SKABUS_RPC_QINFO_ZERO ; diff --git a/src/libskabus/skabus_rpc_r_notimpl.c b/src/libskabus/skabus_rpc_r_notimpl.c new file mode 100644 index 0000000..8ec3922 --- /dev/null +++ b/src/libskabus/skabus_rpc_r_notimpl.c @@ -0,0 +1,11 @@ +/* ISC license. */ + +#include <errno.h> +#include <skabus/rpc.h> + +int skabus_rpc_r_notimpl (skabus_rpc_t *a, skabus_rpc_rinfo_t const *info, unixmessage_t const *m, void *data) +{ + (void)m ; + (void)data ; + return skabus_rpc_reply(a, info->serial, ENOSYS, "", 0) ; +} diff --git a/src/libskabus/skabus_rpc_rcancel_ignore.c b/src/libskabus/skabus_rpc_rcancel_ignore.c new file mode 100644 index 0000000..9697201 --- /dev/null +++ b/src/libskabus/skabus_rpc_rcancel_ignore.c @@ -0,0 +1,11 @@ +/* ISC license. */ + +#include <skabus/rpc.h> + +int skabus_rpc_rcancel_ignore (uint64 serial, char reason, void *data) +{ + (void)serial ; + (void)reason ; + (void)data ; + return 1 ; +} diff --git a/src/libskabus/skabus_rpc_release.c b/src/libskabus/skabus_rpc_release.c new file mode 100644 index 0000000..d0a6875 --- /dev/null +++ b/src/libskabus/skabus_rpc_release.c @@ -0,0 +1,25 @@ +/* ISC license. */ + +#include <stdint.h> +#include <errno.h> +#include <skalibs/uint64.h> +#include <skalibs/alloc.h> +#include <skalibs/gensetdyn.h> +#include <skalibs/avltree.h> +#include <skabus/rpc.h> + +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 ; + p = GENSETDYN_P(skabus_rpc_qinfo_t, &a->q, id) ; + if (p->status) return (errno = p->status, 0) ; + alloc_free(p->message.s) ; + alloc_free(p->message.fds) ; + /* fds are purposefully left open for the client. */ + p->status = EINVAL ; + avltree_delete(&a->qmap, &serial) ; + gensetdyn_delete(&a->q, id) ; + return 1 ; +} diff --git a/src/libskabus/skabus_rpc_reply.c b/src/libskabus/skabus_rpc_reply.c new file mode 100644 index 0000000..491772b --- /dev/null +++ b/src/libskabus/skabus_rpc_reply.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 (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 new file mode 100644 index 0000000..ef727e9 --- /dev/null +++ b/src/libskabus/skabus_rpc_replyv.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 (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_rinfo_pack.c b/src/libskabus/skabus_rpc_rinfo_pack.c new file mode 100644 index 0000000..4e0d34c --- /dev/null +++ b/src/libskabus/skabus_rpc_rinfo_pack.c @@ -0,0 +1,18 @@ +/* ISC license. */ + +#include <string.h> +#include <skalibs/uint32.h> +#include <skalibs/uint64.h> +#include <skalibs/types.h> +#include <skalibs/tai.h> +#include <skabus/rpc.h> + +void skabus_rpc_rinfo_pack (char *s, skabus_rpc_rinfo_t const *ri) +{ + uint64_pack_big(s, ri->serial) ; s += 8 ; + tain_pack(s, &ri->limit) ; s += TAIN_PACK ; + tain_pack(s, &ri->timestamp) ; s += TAIN_PACK ; + uid_pack_big(s, ri->uid) ; s += UID_PACK ; + gid_pack_big(s, ri->gid) ; s += GID_PACK ; + memcpy(s, ri->idstr, SKABUS_RPC_IDSTR_SIZE + 1) ; +} diff --git a/src/libskabus/skabus_rpc_rinfo_unpack.c b/src/libskabus/skabus_rpc_rinfo_unpack.c new file mode 100644 index 0000000..ed04219 --- /dev/null +++ b/src/libskabus/skabus_rpc_rinfo_unpack.c @@ -0,0 +1,18 @@ +/* ISC license. */ + +#include <string.h> +#include <skalibs/uint32.h> +#include <skalibs/uint64.h> +#include <skalibs/types.h> +#include <skalibs/tai.h> +#include <skabus/rpc.h> + +void skabus_rpc_rinfo_unpack (char const *s, skabus_rpc_rinfo_t *ri) +{ + uint64_unpack_big(s, &ri->serial) ; s += 8 ; + tain_unpack(s, &ri->limit) ; s += TAIN_PACK ; + tain_unpack(s, &ri->timestamp) ; s += TAIN_PACK ; + uid_unpack_big(s, &ri->uid) ; s += UID_PACK ; + gid_unpack_big(s, &ri->gid) ; s += GID_PACK ; + memcpy(ri->idstr, s, SKABUS_RPC_IDSTR_SIZE + 1) ; +} diff --git a/src/libskabus/skabus_rpc_rinfo_zero.c b/src/libskabus/skabus_rpc_rinfo_zero.c new file mode 100644 index 0000000..57be272 --- /dev/null +++ b/src/libskabus/skabus_rpc_rinfo_zero.c @@ -0,0 +1,5 @@ +/* ISC license. */ + +#include <skabus/rpc.h> + +skabus_rpc_rinfo_t const skabus_rpc_rinfo_zero = SKABUS_RPC_RINFO_ZERO ; diff --git a/src/libskabus/skabus_rpc_send.c b/src/libskabus/skabus_rpc_send.c new file mode 100644 index 0000000..aa0131c --- /dev/null +++ b/src/libskabus/skabus_rpc_send.c @@ -0,0 +1,9 @@ +/* ISC license. */ + +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +uint64_t skabus_rpc_send_withfds (skabus_rpc_t *a, char const *ifname, char const *s, size_t len, int const *fds, unsigned int nfds, unsigned char const *bits, tain_t const *limit, tain_t const *deadline, tain_t *stamp) +{ + return skabus_rpc_sendq_withfds(a, 0, 0, ifname, s, len, fds, nfds, bits, limit, deadline, stamp) ; +} diff --git a/src/libskabus/skabus_rpc_send_async.c b/src/libskabus/skabus_rpc_send_async.c new file mode 100644 index 0000000..7fc29d2 --- /dev/null +++ b/src/libskabus/skabus_rpc_send_async.c @@ -0,0 +1,9 @@ +/* ISC license. */ + +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +int skabus_rpc_send_withfds_async (skabus_rpc_t *a, char const *ifname, char const *s, size_t len, int const *fds, unsigned int nfds, unsigned char const *bits, tain_t const *limit, skabus_rpc_send_result_t *r) +{ + return skabus_rpc_sendq_withfds_async(a, 0, 0, ifname, s, len, fds, nfds, bits, limit, r) ; +} diff --git a/src/libskabus/skabus_rpc_send_cb.c b/src/libskabus/skabus_rpc_send_cb.c new file mode 100644 index 0000000..7aa7fbc --- /dev/null +++ b/src/libskabus/skabus_rpc_send_cb.c @@ -0,0 +1,53 @@ +/* ISC license. */ + +#include <sys/uio.h> +#include <errno.h> +#include <skalibs/uint64.h> +#include <skalibs/error.h> +#include <skalibs/gensetdyn.h> +#include <skalibs/avltree.h> +#include <skalibs/unixmessage.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +static int autocancel_cb (unixmessage_t const *m, void *p) +{ + (void)m ; + (void)p ; + return 1 ; +} + +int skabus_rpc_send_cb (unixmessage_t const *m, void *p) +{ + skabus_rpc_send_result_t *r = p ; + skabus_rpc_qinfo_t *info = GENSETDYN_P(skabus_rpc_qinfo_t, &r->a->q, r->i) ; + if (!m->len) return (errno = EPROTO, 0) ; + if (m->s[0]) + { + r->err = m->s[0] ; + r->u = 0 ; + info->status = EINVAL ; + gensetdyn_delete(&r->a->q, r->i) ; + return 1 ; + } + if (m->len != 9) return (errno = EPROTO, 0) ; + uint64_unpack_big(m->s+1, &info->serial) ; + if (!avltree_insert(&r->a->qmap, r->i)) + { + /* the client can't store the info but the server is performing the query, + so we try to send a stealthy cancel */ + char what = 'C' ; + struct iovec v[2] = { { .iov_base = &what, .iov_len = 1 }, { .iov_base = m->s + 1, .iov_len = 8 } } ; + r->err = errno ; + r->u = 0 ; + info->status = EINVAL ; + gensetdyn_delete(&r->a->q, r->i) ; + if (skaclient_putv(&r->a->connection, v, 2, &autocancel_cb, 0)) + skaclient_flush(&r->a->connection) ; + return 1 ; + } + r->u = info->serial ; + info->status = EBUSY ; + return 1 ; +} diff --git a/src/libskabus/skabus_rpc_sendpm.c b/src/libskabus/skabus_rpc_sendpm.c new file mode 100644 index 0000000..5f175e4 --- /dev/null +++ b/src/libskabus/skabus_rpc_sendpm.c @@ -0,0 +1,9 @@ +/* ISC license. */ + +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +uint64_t skabus_rpc_sendpm_withfds (skabus_rpc_t *a, char const *cname, char const *s, size_t len, int const *fds, unsigned int nfds, unsigned char const *bits, tain_t const *limit, tain_t const *deadline, tain_t *stamp) +{ + return skabus_rpc_sendq_withfds(a, "\xff", 1, cname, s, len, fds, nfds, bits, limit, deadline, stamp) ; +} diff --git a/src/libskabus/skabus_rpc_sendpm_async.c b/src/libskabus/skabus_rpc_sendpm_async.c new file mode 100644 index 0000000..b1b09d4 --- /dev/null +++ b/src/libskabus/skabus_rpc_sendpm_async.c @@ -0,0 +1,9 @@ +/* ISC license. */ + +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +int skabus_rpc_sendpm_withfds_async (skabus_rpc_t *a, char const *cname, char const *s, size_t len, int const *fds, unsigned int nfds, unsigned char const *bits, tain_t const *limit, skabus_rpc_send_result_t *r) +{ + return skabus_rpc_sendq_withfds_async(a, "\xff", 1, cname, s, len, fds, nfds, bits, limit, r) ; +} diff --git a/src/libskabus/skabus_rpc_sendq.c b/src/libskabus/skabus_rpc_sendq.c new file mode 100644 index 0000000..06944fa --- /dev/null +++ b/src/libskabus/skabus_rpc_sendq.c @@ -0,0 +1,14 @@ +/* ISC license. */ + +#include <errno.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +uint64_t skabus_rpc_sendq_withfds (skabus_rpc_t *a, char const *prefix, size_t plen, char const *ifname, char const *s, size_t len, int const *fds, unsigned int nfds, unsigned char const *bits, tain_t const *limit, tain_t const *deadline, tain_t *stamp) +{ + skabus_rpc_send_result_t r ; + if (!skabus_rpc_sendq_withfds_async(a, prefix, plen, ifname, s, len, fds, nfds, bits, limit, &r)) return 0 ; + if (!skaclient_syncify(&a->connection, deadline, stamp)) return 0 ; + return r.err ? (errno = r.err, 0) : r.u ; +} diff --git a/src/libskabus/skabus_rpc_sendq_async.c b/src/libskabus/skabus_rpc_sendq_async.c new file mode 100644 index 0000000..b579291 --- /dev/null +++ b/src/libskabus/skabus_rpc_sendq_async.c @@ -0,0 +1,34 @@ +/* ISC license. */ + +#include <string.h> +#include <sys/uio.h> +#include <errno.h> +#include <skalibs/tai.h> +#include <skalibs/gensetdyn.h> +#include <skalibs/unixmessage.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +int skabus_rpc_sendq_withfds_async (skabus_rpc_t *a, char const *prefix, size_t plen, char const *ifname, char const *s, size_t len, int const *fds, unsigned int nfds, unsigned char const *bits, tain_t const *limit, skabus_rpc_send_result_t *r) +{ + size_t iflen = strlen(ifname) ; + char pack[2 + TAIN_PACK] = "Q" ; + struct iovec v[4] = { { .iov_base = pack, .iov_len = 2 + TAIN_PACK }, { .iov_base = (char *)prefix, .iov_len = plen }, { .iov_base = (char *)ifname, .iov_len = iflen + 1 }, { .iov_base = (char *)s, .iov_len = len } } ; + unixmessage_v_t m = { .v = v, .vlen = 4, .fds = (int *)fds, .nfds = nfds } ; + iflen += plen ; + if (iflen > SKABUS_RPC_INTERFACE_MAXLEN) return (errno = ENAMETOOLONG, 0) ; + if (!gensetdyn_new(&a->q, &r->i)) return 0 ; + r->a = a ; + if (limit) tain_pack(pack + 1, limit) ; else memset(pack + 1, 0, TAIN_PACK) ; + pack[1 + TAIN_PACK] = (unsigned char)iflen ; + if (!skaclient_putmsgv_and_close(&a->connection, &m, bits, &skabus_rpc_send_cb, r)) + { + int e = errno ; + gensetdyn_delete(&a->q, r->i) ; + errno = e ; + return 0 ; + } + GENSETDYN_P(skabus_rpc_qinfo_t, &a->q, r->i)->status = EINPROGRESS ; + return 1 ; +} diff --git a/src/libskabus/skabus_rpc_sendv.c b/src/libskabus/skabus_rpc_sendv.c new file mode 100644 index 0000000..24d598d --- /dev/null +++ b/src/libskabus/skabus_rpc_sendv.c @@ -0,0 +1,9 @@ +/* ISC license. */ + +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +uint64_t skabus_rpc_sendv_withfds (skabus_rpc_t *a, char const *ifname, struct iovec const *v, unsigned int vlen, int const *fds, unsigned int nfds, unsigned char const *bits, tain_t const *limit, tain_t const *deadline, tain_t *stamp) +{ + return skabus_rpc_sendvq_withfds(a, 0, 0, ifname, v, vlen, fds, nfds, bits, limit, deadline, stamp) ; +} diff --git a/src/libskabus/skabus_rpc_sendv_async.c b/src/libskabus/skabus_rpc_sendv_async.c new file mode 100644 index 0000000..c9b7910 --- /dev/null +++ b/src/libskabus/skabus_rpc_sendv_async.c @@ -0,0 +1,9 @@ +/* ISC license. */ + +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +int skabus_rpc_sendv_withfds_async (skabus_rpc_t *a, char const *ifname, struct iovec const *v, unsigned int vlen, int const *fds, unsigned int nfds, unsigned char const *bits, tain_t const *limit, skabus_rpc_send_result_t *r) +{ + return skabus_rpc_sendvq_withfds_async(a, 0, 0, ifname, v, vlen, fds, nfds, bits, limit, r) ; +} diff --git a/src/libskabus/skabus_rpc_sendvpm.c b/src/libskabus/skabus_rpc_sendvpm.c new file mode 100644 index 0000000..8ca74e6 --- /dev/null +++ b/src/libskabus/skabus_rpc_sendvpm.c @@ -0,0 +1,9 @@ +/* ISC license. */ + +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +uint64_t skabus_rpc_sendvpm_withfds (skabus_rpc_t *a, char const *cname, struct iovec const *v, unsigned int vlen, int const *fds, unsigned int nfds, unsigned char const *bits, tain_t const *limit, tain_t const *deadline, tain_t *stamp) +{ + return skabus_rpc_sendvq_withfds(a, "\xff", 1, cname, v, vlen, fds, nfds, bits, limit, deadline, stamp) ; +} diff --git a/src/libskabus/skabus_rpc_sendvpm_async.c b/src/libskabus/skabus_rpc_sendvpm_async.c new file mode 100644 index 0000000..090423c --- /dev/null +++ b/src/libskabus/skabus_rpc_sendvpm_async.c @@ -0,0 +1,9 @@ +/* ISC license. */ + +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +int skabus_rpc_sendvpm_withfds_async (skabus_rpc_t *a, char const *cname, struct iovec const *v, unsigned int vlen, int const *fds, unsigned int nfds, unsigned char const *bits, tain_t const *limit, skabus_rpc_send_result_t *r) +{ + return skabus_rpc_sendvq_withfds_async(a, "\xff", 1, cname, v, vlen, fds, nfds, bits, limit, r) ; +} diff --git a/src/libskabus/skabus_rpc_sendvq.c b/src/libskabus/skabus_rpc_sendvq.c new file mode 100644 index 0000000..3cd0e3a --- /dev/null +++ b/src/libskabus/skabus_rpc_sendvq.c @@ -0,0 +1,13 @@ +/* ISC license. */ + +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +uint64_t skabus_rpc_sendvq_withfds (skabus_rpc_t *a, char const *prefix, size_t plen, char const *ifname, struct iovec const *v, unsigned int vlen, int const *fds, unsigned int nfds, unsigned char const *bits, tain_t const *limit, tain_t const *deadline, tain_t *stamp) +{ + skabus_rpc_send_result_t r ; + if (!skabus_rpc_sendvq_withfds_async(a, prefix, plen, ifname, v, vlen, fds, nfds, bits, limit, &r)) return 0 ; + if (!skaclient_syncify(&a->connection, deadline, stamp)) return 0 ; + return r.err ? (errno = r.err, 0) : r.u ; +} diff --git a/src/libskabus/skabus_rpc_sendvq_async.c b/src/libskabus/skabus_rpc_sendvq_async.c new file mode 100644 index 0000000..97117ea --- /dev/null +++ b/src/libskabus/skabus_rpc_sendvq_async.c @@ -0,0 +1,38 @@ +/* ISC license. */ + +#include <string.h> +#include <sys/uio.h> +#include <errno.h> +#include <skalibs/tai.h> +#include <skalibs/gensetdyn.h> +#include <skalibs/unixmessage.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> +#include "skabus-rpc-internal.h" + +int skabus_rpc_sendvq_withfds_async (skabus_rpc_t *a, char const *prefix, size_t plen, char const *ifname, struct iovec const *v, unsigned int vlen, int const *fds, unsigned int nfds, unsigned char const *bits, tain_t const *limit, skabus_rpc_send_result_t *r) +{ + size_t iflen = strlen(ifname) ; + char pack[2 + TAIN_PACK] = "Q" ; + struct iovec vv[vlen + 3] ; + unixmessage_v_t m = { .v = vv, .vlen = vlen + 3, .fds = (int *)fds, .nfds = nfds } ; + iflen += plen ; + if (iflen > SKABUS_RPC_INTERFACE_MAXLEN) return (errno = ENAMETOOLONG, 0) ; + if (!gensetdyn_new(&a->q, &r->i)) return 0 ; + vv[0].iov_base = pack ; vv[0].iov_len = 2 + TAIN_PACK ; + vv[1].iov_base = (char *)prefix ; vv[1].iov_len = plen ; + vv[2].iov_base = (char *)ifname ; vv[2].iov_len = iflen + 1 ; + for (unsigned int i = 0 ; i < vlen ; i++) vv[3 + i] = v[i] ; + r->a = a ; + if (limit) tain_pack(pack + 1, limit) ; else memset(pack + 1, 0, TAIN_PACK) ; + pack[1 + TAIN_PACK] = (unsigned char)iflen ; + if (!skaclient_putmsgv_and_close(&a->connection, &m, bits, &skabus_rpc_send_cb, r)) + { + int e = errno ; + gensetdyn_delete(&a->q, r->i) ; + errno = e ; + return 0 ; + } + GENSETDYN_P(skabus_rpc_qinfo_t, &a->q, r->i)->status = EINPROGRESS ; + return 1 ; +} diff --git a/src/libskabus/skabus_rpc_start.c b/src/libskabus/skabus_rpc_start.c new file mode 100644 index 0000000..8e9c99b --- /dev/null +++ b/src/libskabus/skabus_rpc_start.c @@ -0,0 +1,19 @@ +/* ISC license. */ + +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> + +int skabus_rpc_start (skabus_rpc_t *a, char const *path, tain_t const *deadline, tain_t *stamp) +{ + return skaclient_start_b( + &a->connection, + &a->buffers, + path, + SKACLIENT_OPTION_ASYNC_ACCEPT_FDS, + SKABUS_RPC_BANNER1, + SKABUS_RPC_BANNER1_LEN, + SKABUS_RPC_BANNER2, + SKABUS_RPC_BANNER2_LEN, + deadline, + stamp) ; +} diff --git a/src/libskabus/skabus_rpc_start_async.c b/src/libskabus/skabus_rpc_start_async.c new file mode 100644 index 0000000..d5b3e8c --- /dev/null +++ b/src/libskabus/skabus_rpc_start_async.c @@ -0,0 +1,18 @@ +/* ISC license. */ + +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> + +int skabus_rpc_start_async (skabus_rpc_t *a, char const *path, skabus_rpc_start_result_t *data) +{ + return skaclient_start_async_b( + &a->connection, + &a->buffers, + path, + SKACLIENT_OPTION_ASYNC_ACCEPT_FDS, + SKABUS_RPC_BANNER1, + SKABUS_RPC_BANNER1_LEN, + SKABUS_RPC_BANNER2, + SKABUS_RPC_BANNER2_LEN, + &data->skaclient_cbdata) ; +} diff --git a/src/libskabus/skabus_rpc_update.c b/src/libskabus/skabus_rpc_update.c new file mode 100644 index 0000000..d727de1 --- /dev/null +++ b/src/libskabus/skabus_rpc_update.c @@ -0,0 +1,120 @@ +/* ISC license. */ + +#include <string.h> +#include <stdint.h> +#include <errno.h> +#include <skalibs/error.h> +#include <skalibs/uint32.h> +#include <skalibs/uint64.h> +#include <skalibs/alloc.h> +#include <skalibs/bytestr.h> +#include <skalibs/genalloc.h> +#include <skalibs/gensetdyn.h> +#include <skalibs/avltree.h> +#include <skalibs/unixmessage.h> +#include <skalibs/skaclient.h> +#include <skabus/rpc.h> + +typedef int localhandler_func_t (skabus_rpc_t *, unixmessage_t *) ; +typedef localhandler_func_t *localhandler_func_t_ref ; + +static int skabus_rpc_serve (skabus_rpc_t *a, unixmessage_t *m) +{ + skabus_rpc_rinfo_t rinfo ; + uint32_t ifid ; + skabus_rpc_interface_t *p ; + if (m->len < 4 + SKABUS_RPC_RINFO_PACK) return (errno = EPROTO, 0) ; + uint32_unpack_big(m->s, &ifid) ; m->s += 4 ; m->len -= 4 ; + skabus_rpc_rinfo_unpack(m->s, &rinfo) ; + m->s += SKABUS_RPC_RINFO_PACK ; m->len -= SKABUS_RPC_RINFO_PACK ; + p = GENSETDYN_P(skabus_rpc_interface_t, &a->r, ifid) ; + if (!p) return (errno = ESRCH, 0) ; + return (*p->f)(a, &rinfo, m, p->data) ; +} + +static int skabus_rpc_cancelr (skabus_rpc_t *a, unixmessage_t *m) +{ + uint64_t serial ; + uint32_t ifid ; + skabus_rpc_interface_t *p ; + if (m->len != 13 || m->nfds) return (errno = EPROTO, 0) ; + uint32_unpack_big(m->s, &ifid) ; + uint64_unpack_big(m->s+4, &serial) ; + p = GENSETDYN_P(skabus_rpc_interface_t, &a->r, ifid) ; + if (!p) return (errno = ESRCH, 0) ; + return (*p->cancelf)(serial, m->s[12], p->data) ; +} + +static int skabus_rpc_handle_reply (skabus_rpc_t *a, unixmessage_t *m) +{ + skabus_rpc_qinfo_t *p ; + uint64_t serial ; + uint32_t id ; + if (m->len < 9) return (errno = EPROTO, 0) ; + uint64_unpack_big(m->s, &serial) ; + if (!avltree_search(&a->qmap, &serial, &id)) + { + unixmessage_drop(m) ; + return 1 ; + } + p = GENSETDYN_P(skabus_rpc_qinfo_t, &a->q, id) ; + if (!genalloc_readyplus(uint64_t, &a->qlist, 1)) return 0 ; + if (!m->s[8]) + { + if (m->len < 10) return (errno = EPROTO, 0) ; + p->message.s = alloc(m->len - 10) ; + if (!p->message.s) return 0 ; + p->message.fds = (int *)alloc(m->nfds * sizeof(int)) ; + if (!p->message.fds) + { + alloc_free(p->message.s) ; + p->message.s = 0 ; + return 0 ; + } + p->result = m->s[9] ; + p->message.len = m->len - 10 ; + p->message.nfds = m->nfds ; + memcpy(p->message.s, m->s + 10, p->message.len) ; + memcpy(p->message.fds, m->fds, p->message.nfds * sizeof(int)) ; + } + p->status = m->s[8] ; + genalloc_append(uint64_t, &a->qlist, &serial) ; + return 1 ; +} + +static int skabus_rpc_error (skabus_rpc_t *a, unixmessage_t *m) +{ + (void)a ; + (void)m ; + return (errno = EPROTO, 0) ; +} + +static int handler (unixmessage_t const *m, void *x) +{ + skabus_rpc_t *a = x ; + unixmessage_t mm = { .s = m->s + 1, .len = m->len - 1, .fds = m->fds, .nfds = m->nfds } ; + static localhandler_func_t_ref const f[4] = + { + &skabus_rpc_serve, + &skabus_rpc_handle_reply, + &skabus_rpc_cancelr, + &skabus_rpc_error + } ; + if (!m->len) + { + unixmessage_drop(m) ; + return (errno = EPROTO, 0) ; + } + if (!(*f[byte_chr("QRC", 3, m->s[0])])(a, &mm)) + { + unixmessage_drop(m) ; + return 0 ; + } + return 1 ; +} + +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/libskabus/skabus_rpc_zero.c b/src/libskabus/skabus_rpc_zero.c new file mode 100644 index 0000000..b9db269 --- /dev/null +++ b/src/libskabus/skabus_rpc_zero.c @@ -0,0 +1,5 @@ +/* ISC license. */ + +#include <skabus/rpc.h> + +skabus_rpc_t const skabus_rpc_zero = SKABUS_RPC_ZERO ; |
