aboutsummaryrefslogtreecommitdiffstats
path: root/src
diff options
context:
space:
mode:
Diffstat (limited to 'src')
-rw-r--r--src/include/skabus/rpc.h31
-rw-r--r--src/libskabus/deps-lib/skabus4
-rw-r--r--src/libskabus/skabus_rpc_get.c10
-rw-r--r--src/libskabus/skabus_rpc_qlist.c12
-rw-r--r--src/libskabus/skabus_rpc_qlist_ack.c15
-rw-r--r--src/libskabus/skabus_rpc_r_notimpl.c2
-rw-r--r--src/libskabus/skabus_rpc_release.c6
-rw-r--r--src/libskabus/skabus_rpc_reply.c13
-rw-r--r--src/libskabus/skabus_rpc_reply_async.c17
-rw-r--r--src/libskabus/skabus_rpc_replyv.c15
-rw-r--r--src/libskabus/skabus_rpc_replyv_async.c19
-rw-r--r--src/libskabus/skabus_rpc_update.c1
-rw-r--r--src/rpc/skabus-rpcd.c5
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 ;
}