aboutsummaryrefslogtreecommitdiffstats
path: root/src/libskabus
diff options
context:
space:
mode:
Diffstat (limited to 'src/libskabus')
-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
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) ;
}