summaryrefslogtreecommitdiff
path: root/src/lib/dev.c
diff options
context:
space:
mode:
Diffstat (limited to 'src/lib/dev.c')
-rw-r--r--src/lib/dev.c135
1 files changed, 102 insertions, 33 deletions
diff --git a/src/lib/dev.c b/src/lib/dev.c
index 3fb8d831..eb706691 100644
--- a/src/lib/dev.c
+++ b/src/lib/dev.c
@@ -28,6 +28,9 @@
#include "config.h"
#include "ssm.h"
+#include "poa/poa.h"
+
+#define OUROBOROS_PREFIX "libouroboros"
#include <ouroboros/atomics.h>
#include <ouroboros/bitmap.h>
@@ -45,6 +48,7 @@
#include <ouroboros/ipcp-dev.h>
#include <ouroboros/list.h>
#include <ouroboros/local-dev.h>
+#include <ouroboros/logs.h>
#include <ouroboros/np1_flow.h>
#include <ouroboros/pthread.h>
#include <ouroboros/random.h>
@@ -63,8 +67,10 @@
#ifdef HAVE_LIBGCRYPT
#include <gcrypt.h>
#endif
+#include <arpa/inet.h>
#include <stdarg.h>
#include <stdbool.h>
+#include <stddef.h>
#include <inttypes.h>
#include <stdio.h>
#include <stdlib.h>
@@ -98,8 +104,8 @@ struct flow {
ssize_t part_idx;
struct crypt_ctx * crypt;
- int headsz; /* Selector */
- int tailsz; /* Tag + CRC */
+ int headsz; /* Selector */
+ int tailsz; /* Tag + CRC */
struct timespec rk_grace; /* TX-promote deadline */
struct timespec rk_attempt; /* Last re-key attempt */
@@ -116,6 +122,8 @@ struct flow {
struct timespec rcv_timeo;
struct frcti * frcti;
+
+ struct poa_flow * poa; /* NULL for shared memory flows */
};
struct flow_set {
@@ -696,6 +704,9 @@ static void do_flow_fini(int fd)
{
assert(fd >= 0 && fd < PROC_MAX_FLOWS);
+ if (proc.flows[fd].poa != NULL)
+ poa_flow_detach(proc.flows[fd].poa);
+
if (proc.flows[fd].frcti != NULL)
frcti_destroy(proc.flows[fd].frcti);
@@ -750,9 +761,11 @@ static __inline__ size_t flow_user_mtu(const struct flow * flow,
return raw > hdr ? raw - hdr : 0;
}
+/* A PoA flow transmits on its own socket; it has no tx ring. */
static int flow_init(struct flow_info * info,
struct crypt_sk * sk,
- time_t rtt_hint)
+ time_t rtt_hint,
+ struct poa_flow * pf)
{
struct timespec now;
struct flow * flow;
@@ -777,13 +790,15 @@ static int flow_init(struct flow_info * info,
if (flow->rx_rb == NULL)
goto fail_rx_rb;
- flow->tx_rb = ssm_rbuff_open(info->n_1_pid, info->id);
- if (flow->tx_rb == NULL)
- goto fail_tx_rb;
+ if (pf == NULL) {
+ flow->tx_rb = ssm_rbuff_open(info->n_1_pid, info->id);
+ if (flow->tx_rb == NULL)
+ goto fail_tx_rb;
- flow->set = ssm_flow_set_open(info->n_1_pid);
- if (flow->set == NULL)
- goto fail_set;
+ flow->set = ssm_flow_set_open(info->n_1_pid);
+ if (flow->set == NULL)
+ goto fail_set;
+ }
flow->oflags = FLOWFDEFAULT;
flow->part_idx = NO_PART;
@@ -792,11 +807,14 @@ static int flow_init(struct flow_info * info,
flow->crypt = NULL;
flow->headsz = 0;
flow->tailsz = 0;
+ flow->poa = pf;
if (IS_ENCRYPTED(sk)) {
flow->crypt = crypt_create_ctx(sk);
- if (flow->crypt == NULL)
+ if (flow->crypt == NULL) {
+ err = -ECRYPT;
goto fail_crypt;
+ }
flow->headsz = crypt_get_headsz(flow->crypt);
flow->tailsz = crypt_get_tagsz(flow->crypt);
}
@@ -815,6 +833,9 @@ static int flow_init(struct flow_info * info,
proc.id_to_fd[info->id].fd = fd;
+ if (pf != NULL)
+ poa_flow_attach(pf, info->id, flow->rx_rb);
+
flow_set_state(&proc.id_to_fd[info->id], FLOW_ALLOCATED);
pthread_rwlock_unlock(&proc.lock);
@@ -824,9 +845,11 @@ static int flow_init(struct flow_info * info,
fail_frcti:
crypt_destroy_ctx(flow->crypt);
fail_crypt:
- ssm_flow_set_close(flow->set);
+ if (flow->set != NULL)
+ ssm_flow_set_close(flow->set);
fail_set:
- ssm_rbuff_close(flow->tx_rb);
+ if (flow->tx_rb != NULL)
+ ssm_rbuff_close(flow->tx_rb);
fail_tx_rb:
ssm_rbuff_close(flow->rx_rb);
fail_rx_rb:
@@ -1077,14 +1100,20 @@ static void fini(void)
__attribute__((section(INIT_SECTION))) __typeof__(init) * __init = init;
__attribute__((section(FINI_SECTION))) __typeof__(fini) * __fini = fini;
+/*
+ * A PoA flow is announced by its peer before the accept; from the
+ * reply on, the peer may transmit, so the flow must be able to
+ * receive.
+ */
int flow_accept(qosspec_t * qs,
const struct timespec * timeo)
{
struct flow_info flow;
- struct crypt_sk crypt;
+ struct crypt_sk crypt;
uint8_t buf[SOCK_BUF_SIZE];
buffer_t msg = {SOCK_BUF_SIZE, buf};
uint8_t key[SYMMKEYSZ];
+ struct poa_flow * pf;
int fd;
int err;
@@ -1116,8 +1145,13 @@ int flow_accept(qosspec_t * qs,
if (err < 0)
return err;
- /* No RTT in accept; rtt_hint=0 bootstraps from first ACK. */
- fd = flow_init(&flow, &crypt, 0);
+ pf = poa_flow_take_pending(flow.id);
+
+ fd = flow_init(&flow, &crypt, 0, pf);
+ if (fd >= 0)
+ poa_flow_ready(pf);
+ else if (pf != NULL)
+ poa_flow_detach(pf);
crypt_secure_clear(key, SYMMKEYSZ);
@@ -1173,7 +1207,7 @@ int flow_alloc(const char * dst,
if (err < 0)
return err;
- fd = flow_init(&flow, &crypt, ts_diff_ns(&t1, &t0));
+ fd = flow_init(&flow, &crypt, ts_diff_ns(&t1, &t0), NULL);
crypt_secure_clear(key, SYMMKEYSZ);
@@ -1214,7 +1248,7 @@ int flow_join(const char * dst,
if (err < 0)
return err;
- fd = flow_init(&flow, &crypt, 0);
+ fd = flow_init(&flow, &crypt, 0, NULL);
crypt_secure_clear(key, SYMMKEYSZ);
@@ -1275,7 +1309,8 @@ int flow_dealloc(int fd)
pthread_cleanup_push(__cleanup_rwlock_unlock, &proc.lock);
- ssm_rbuff_fini(flow->tx_rb);
+ if (flow->tx_rb != NULL)
+ ssm_rbuff_fini(flow->tx_rb);
pthread_cleanup_pop(true);
@@ -1428,7 +1463,11 @@ int fccntl(int fd,
break;
case FLOWGTXQLEN:
qlen = va_arg(l, size_t *);
- *qlen = ssm_rbuff_queued(flow->tx_rb);
+
+ if (flow->poa != NULL)
+ *qlen = poa_flow_qpkts(flow->poa);
+ else
+ *qlen = ssm_rbuff_queued(flow->tx_rb);
break;
case FLOWGMTU:
maxp = va_arg(l, size_t *);
@@ -1456,16 +1495,18 @@ int fccntl(int fd,
if (flow->oflags & FLOWFDOWN) {
ssm_rbuff_set_bits(flow->rx_rb, RB_FLOWDOWN);
- ssm_rbuff_set_bits(flow->tx_rb, RB_FLOWDOWN);
- ssm_flow_set_notify(flow->set,
- flow->info.id,
- FLOW_DOWN);
+ if (flow->tx_rb != NULL)
+ ssm_rbuff_set_bits(flow->tx_rb, RB_FLOWDOWN);
+ if (flow->set != NULL)
+ ssm_flow_set_notify(flow->set, flow->info.id,
+ FLOW_DOWN);
} else {
ssm_rbuff_clr_bits(flow->rx_rb, RB_FLOWDOWN);
- ssm_rbuff_clr_bits(flow->tx_rb, RB_FLOWDOWN);
- ssm_flow_set_notify(flow->set,
- flow->info.id,
- FLOW_UP);
+ if (flow->tx_rb != NULL)
+ ssm_rbuff_clr_bits(flow->tx_rb, RB_FLOWDOWN);
+ if (flow->set != NULL)
+ ssm_flow_set_notify(flow->set, flow->info.id,
+ FLOW_UP);
}
break;
@@ -1599,6 +1640,9 @@ static int flow_tx_spb(struct flow * flow,
goto enomem;
}
+ if (flow->poa != NULL)
+ return poa_flow_tx(flow->poa, spb, block, abstime);
+
if (!block)
ret = ssm_rbuff_write(flow->tx_rb, idx);
else
@@ -2450,7 +2494,7 @@ int np1_flow_alloc(pid_t n_pid,
/* np1 flow: n_1_pid is the upper. */
flow.n_1_pid = n_pid;
- return flow_init(&flow, &crypt, 0);
+ return flow_init(&flow, &crypt, 0, NULL);
}
int np1_flow_dealloc(int flow_id,
@@ -2584,7 +2628,7 @@ int ipcp_flow_req_arr(const buffer_t * dst,
crypt.nid = NID_undef;
- return flow_init(&flow, &crypt, 0);
+ return flow_init(&flow, &crypt, 0, NULL);
}
int ipcp_flow_update_arr(int flow_id,
@@ -2853,11 +2897,13 @@ int ipcp_flow_fini(int fd)
}
ssm_rbuff_set_bits(proc.flows[fd].rx_rb, RB_FLOWDOWN);
- ssm_rbuff_set_bits(proc.flows[fd].tx_rb, RB_FLOWDOWN);
- ssm_flow_set_notify(proc.flows[fd].set,
- proc.flows[fd].info.id,
- FLOW_DEALLOC);
+ if (proc.flows[fd].tx_rb != NULL)
+ ssm_rbuff_set_bits(proc.flows[fd].tx_rb, RB_FLOWDOWN);
+
+ if (proc.flows[fd].set != NULL)
+ ssm_flow_set_notify(proc.flows[fd].set, proc.flows[fd].info.id,
+ FLOW_DEALLOC);
rx_rb = proc.flows[fd].rx_rb;
@@ -2892,9 +2938,30 @@ size_t ipcp_flow_queued(int fd)
assert(fd >= 0 && fd < PROC_MAX_FLOWS);
assert(proc.flows[fd].info.id >= 0);
+ if (proc.flows[fd].poa != NULL)
+ return poa_flow_qlen(proc.flows[fd].poa);
+
return ssm_rbuff_queued(proc.flows[fd].tx_rb);
}
+int ipcp_flow_queue_id(int fd)
+{
+ int qid;
+
+ assert(fd >= 0 && fd < PROC_MAX_FLOWS);
+ assert(proc.flows[fd].info.id >= 0);
+
+ if (proc.flows[fd].poa == NULL)
+ return fd;
+
+ /* An unidentified PoA answers for itself, never for an fd. */
+ qid = poa_flow_qid(proc.flows[fd].poa);
+ if (qid < 0 || qid >= POA_MAX_POAS)
+ return fd;
+
+ return PROC_MAX_FLOWS + qid;
+}
+
int local_flow_transfer(int src_fd,
int dst_fd,
struct ssm_pool * src_pool,
@@ -2960,3 +3027,5 @@ int local_flow_transfer(int src_fd,
return ret;
}
+
+#include "poa/poa.c"