aboutsummaryrefslogtreecommitdiffstats
path: root/src/libskabus/skabus_rpc_update.c
blob: 84de1b0bb017b8f7f755c7b8bb3c55949f2d7357 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
/* ISC license. */

#include <string.h>
#include <stdint.h>
#include <errno.h>

#include <skalibs/posixishard.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)
{
  return skaclient_update(&a->connection, &handler, a) ;
}