summaryrefslogtreecommitdiff
path: root/src/ipcpd/unicast/dt.c
diff options
context:
space:
mode:
Diffstat (limited to 'src/ipcpd/unicast/dt.c')
-rw-r--r--src/ipcpd/unicast/dt.c547
1 files changed, 283 insertions, 264 deletions
diff --git a/src/ipcpd/unicast/dt.c b/src/ipcpd/unicast/dt.c
index 0f504daa..381ed815 100644
--- a/src/ipcpd/unicast/dt.c
+++ b/src/ipcpd/unicast/dt.c
@@ -1,5 +1,5 @@
/*
- * Ouroboros - Copyright (C) 2016 - 2021
+ * Ouroboros - Copyright (C) 2016 - 2026
*
* Data Transfer Component
*
@@ -31,19 +31,23 @@
#define DT "dt"
#define OUROBOROS_PREFIX DT
+#include <ouroboros/atomics.h>
#include <ouroboros/bitmap.h>
#include <ouroboros/errno.h>
#include <ouroboros/logs.h>
#include <ouroboros/dev.h>
+#include <ouroboros/ipcp-dev.h>
#include <ouroboros/notifier.h>
#include <ouroboros/rib.h>
#ifdef IPCP_FLOW_STATS
#include <ouroboros/fccntl.h>
#endif
+#include "addr-auth.h"
#include "common/comp.h"
#include "common/connmgr.h"
#include "ca.h"
+#include "cap.h"
#include "ipcp.h"
#include "dt.h"
#include "pff.h"
@@ -59,7 +63,7 @@
#include <assert.h>
#define QOS_BLOCK_LEN 672
-#define RIB_FILE_STRLEN (189 + QOS_BLOCK_LEN * QOS_CUBE_MAX)
+#define RIB_FILE_STRLEN (169 + RIB_TM_STRLEN + QOS_BLOCK_LEN * QOS_CUBE_MAX)
#define RIB_NAME_STRLEN 256
#ifndef CLOCK_REALTIME_COARSE
@@ -67,7 +71,7 @@
#endif
struct comp_info {
- void (* post_packet)(void * comp, struct shm_du_buff * sdb);
+ void (* post_packet)(void * comp, struct ssm_pk_buff * spb);
void * comp;
char * name;
};
@@ -76,12 +80,14 @@ struct comp_info {
#define TTL_LEN 1
#define QOS_LEN 1
#define ECN_LEN 1
+#define CAP_LEN 1
struct dt_pci {
uint64_t dst_addr;
qoscube_t qc;
uint8_t ttl;
uint8_t ecn;
+ uint8_t cap;
uint64_t eid;
};
@@ -94,6 +100,7 @@ struct {
size_t qc_o;
size_t ttl_o;
size_t ecn_o;
+ size_t cap_o;
size_t eid_o;
/* Initial TTL value */
@@ -113,6 +120,7 @@ static void dt_pci_ser(uint8_t * head,
memcpy(head + dt_pci_info.qc_o, &dt_pci->qc, QOS_LEN);
memcpy(head + dt_pci_info.ttl_o, &ttl, TTL_LEN);
memcpy(head + dt_pci_info.ecn_o, &dt_pci->ecn, ECN_LEN);
+ memcpy(head + dt_pci_info.cap_o, &dt_pci->cap, CAP_LEN);
memcpy(head + dt_pci_info.eid_o, &dt_pci->eid, dt_pci_info.eid_size);
}
@@ -131,22 +139,26 @@ static void dt_pci_des(uint8_t * head,
memcpy(&dt_pci->qc, head + dt_pci_info.qc_o, QOS_LEN);
memcpy(&dt_pci->ttl, head + dt_pci_info.ttl_o, TTL_LEN);
memcpy(&dt_pci->ecn, head + dt_pci_info.ecn_o, ECN_LEN);
+ memcpy(&dt_pci->cap, head + dt_pci_info.cap_o, CAP_LEN);
memcpy(&dt_pci->eid, head + dt_pci_info.eid_o, dt_pci_info.eid_size);
}
-static void dt_pci_shrink(struct shm_du_buff * sdb)
+static void dt_pci_shrink(struct ssm_pk_buff * spb)
{
- assert(sdb);
+ assert(spb);
- shm_du_buff_head_release(sdb, dt_pci_info.head_size);
+ ssm_pk_buff_pop(spb, dt_pci_info.head_size);
}
struct {
struct psched * psched;
+ uint64_t addr;
+
struct pff * pff[QOS_CUBE_MAX];
struct routing_i * routing[QOS_CUBE_MAX];
#ifdef IPCP_FLOW_STATS
+ /* Flow stats use lock-free atomics; stamp is the validity flag. */
struct {
time_t stamp;
uint64_t addr;
@@ -164,31 +176,39 @@ struct {
size_t w_drp_bytes[QOS_CUBE_MAX];
size_t f_nhp_pkt[QOS_CUBE_MAX];
size_t f_nhp_bytes[QOS_CUBE_MAX];
- pthread_mutex_t lock;
- } stat[PROG_MAX_FLOWS];
+ } stat[PROC_MAX_FLOWS];
size_t n_flows;
#endif
struct bmp * res_fds;
- struct comp_info comps[PROG_RES_FDS];
+ struct comp_info comps[PROC_RES_FDS];
pthread_rwlock_t lock;
pthread_t listener;
} dt;
+#ifdef IPCP_FLOW_STATS
+#define dt_stat_inc(idx, name, qc, len) \
+ do { \
+ FETCH_ADD_RELAXED(&dt.stat[idx].name ## _pkt[qc], 1); \
+ FETCH_ADD_RELAXED(&dt.stat[idx].name ## _bytes[qc], (len)); \
+ } while (0)
+#define dt_stat_load(idx, field, qc) LOAD_RELAXED(&dt.stat[idx].field[qc])
+
static int dt_rib_read(const char * path,
char * buf,
size_t len)
{
-#ifdef IPCP_FLOW_STATS
int fd;
int i;
char str[QOS_BLOCK_LEN + 1];
char addrstr[20];
char * entry;
- char tmstr[20];
+ char tmstr[RIB_TM_STRLEN];
size_t rxqlen = 0;
size_t txqlen = 0;
+ time_t stamp;
+ uint64_t addr;
struct tm * tm;
/* NOTE: we may need stronger checks. */
@@ -202,32 +222,31 @@ static int dt_rib_read(const char * path,
buf[0] = '\0';
- pthread_mutex_lock(&dt.stat[fd].lock);
-
- if (dt.stat[fd].stamp == 0) {
- pthread_mutex_unlock(&dt.stat[fd].lock);
+ stamp = LOAD_ACQUIRE(&dt.stat[fd].stamp);
+ if (stamp == 0)
return 0;
- }
- if (dt.stat[fd].addr == ipcpi.dt_addr)
+ addr = LOAD_RELAXED(&dt.stat[fd].addr);
+
+ if (addr == dt.addr)
sprintf(addrstr, "%s", dt.comps[fd].name);
else
- sprintf(addrstr, "%" PRIu64, dt.stat[fd].addr);
+ sprintf(addrstr, ADDR_FMT32, ADDR_VAL32(&addr));
- tm = localtime(&dt.stat[fd].stamp);
- strftime(tmstr, sizeof(tmstr), "%F %T", tm);
+ tm = gmtime(&stamp);
+ strftime(tmstr, sizeof(tmstr), RIB_TM_FORMAT, tm);
- if (fd >= PROG_RES_FDS) {
+ if (fd >= PROC_RES_FDS) {
fccntl(fd, FLOWGRXQLEN, &rxqlen);
fccntl(fd, FLOWGTXQLEN, &txqlen);
}
sprintf(buf,
- "Flow established at: %20s\n"
+ "Flow established at: %.*s\n"
"Endpoint address: %20s\n"
"Queued packets (rx): %20zu\n"
"Queued packets (tx): %20zu\n\n",
- tmstr, addrstr, rxqlen, txqlen);
+ RIB_TM_STRLEN - 1, tmstr, addrstr, rxqlen, txqlen);
for (i = 0; i < QOS_CUBE_MAX; ++i) {
sprintf(str,
"Qos cube %3d:\n"
@@ -246,38 +265,29 @@ static int dt_rib_read(const char * path,
" failed nhop (packets): %20zu\n"
" failed nhop (bytes): %20zu\n",
i,
- dt.stat[fd].snd_pkt[i],
- dt.stat[fd].snd_bytes[i],
- dt.stat[fd].rcv_pkt[i],
- dt.stat[fd].rcv_bytes[i],
- dt.stat[fd].lcl_w_pkt[i],
- dt.stat[fd].lcl_w_bytes[i],
- dt.stat[fd].lcl_r_pkt[i],
- dt.stat[fd].lcl_r_bytes[i],
- dt.stat[fd].r_drp_pkt[i],
- dt.stat[fd].r_drp_bytes[i],
- dt.stat[fd].w_drp_pkt[i],
- dt.stat[fd].w_drp_bytes[i],
- dt.stat[fd].f_nhp_pkt[i],
- dt.stat[fd].f_nhp_bytes[i]
+ dt_stat_load(fd, snd_pkt, i),
+ dt_stat_load(fd, snd_bytes, i),
+ dt_stat_load(fd, rcv_pkt, i),
+ dt_stat_load(fd, rcv_bytes, i),
+ dt_stat_load(fd, lcl_w_pkt, i),
+ dt_stat_load(fd, lcl_w_bytes, i),
+ dt_stat_load(fd, lcl_r_pkt, i),
+ dt_stat_load(fd, lcl_r_bytes, i),
+ dt_stat_load(fd, r_drp_pkt, i),
+ dt_stat_load(fd, r_drp_bytes, i),
+ dt_stat_load(fd, w_drp_pkt, i),
+ dt_stat_load(fd, w_drp_bytes, i),
+ dt_stat_load(fd, f_nhp_pkt, i),
+ dt_stat_load(fd, f_nhp_bytes, i)
);
strcat(buf, str);
}
- pthread_mutex_unlock(&dt.stat[fd].lock);
-
return RIB_FILE_STRLEN;
-#else
- (void) path;
- (void) buf;
- (void) len;
- return 0;
-#endif
}
static int dt_rib_readdir(char *** buf)
{
-#ifdef IPCP_FLOW_STATS
char entry[RIB_PATH_LEN + 1];
size_t i;
int idx = 0;
@@ -285,81 +295,64 @@ static int dt_rib_readdir(char *** buf)
pthread_rwlock_rdlock(&dt.lock);
if (dt.n_flows < 1) {
- pthread_rwlock_unlock(&dt.lock);
- return 0;
+ *buf = NULL;
+ goto no_flows;
}
*buf = malloc(sizeof(**buf) * dt.n_flows);
- if (*buf == NULL) {
- pthread_rwlock_unlock(&dt.lock);
- return -ENOMEM;
- }
+ if (*buf == NULL)
+ goto fail_entries;
- for (i = 0; i < PROG_MAX_FLOWS; ++i) {
- pthread_mutex_lock(&dt.stat[i].lock);
-
- if (dt.stat[i].stamp == 0) {
- pthread_mutex_unlock(&dt.stat[i].lock);
- /* Optimization: skip unused res_fds. */
- if (i < PROG_RES_FDS)
- i = PROG_RES_FDS;
- continue;
- }
+ for (i = 0; i < PROC_MAX_FLOWS && idx < (int) dt.n_flows; ++i) {
+ if (LOAD_RELAXED(&dt.stat[i].stamp) == 0)
+ continue; /* n-1 flows start at PROC_RES_FDS */
sprintf(entry, "%zu", i);
(*buf)[idx] = malloc(strlen(entry) + 1);
- if ((*buf)[idx] == NULL) {
- while (idx-- > 0)
- free((*buf)[idx]);
- free(buf);
- pthread_mutex_unlock(&dt.stat[i].lock);
- pthread_rwlock_unlock(&dt.lock);
- return -ENOMEM;
- }
+ if ((*buf)[idx] == NULL)
+ goto fail_entry;
strcpy((*buf)[idx++], entry);
- pthread_mutex_unlock(&dt.stat[i].lock);
}
- assert((size_t) idx == dt.n_flows);
-
+ no_flows:
pthread_rwlock_unlock(&dt.lock);
return idx;
-#else
- (void) buf;
- return 0;
-#endif
+
+ fail_entry:
+ while (idx-- > 0)
+ free((*buf)[idx]);
+
+ free(*buf);
+ fail_entries:
+ pthread_rwlock_unlock(&dt.lock);
+ return -ENOMEM;
}
static int dt_rib_getattr(const char * path,
struct rib_attr * attr)
{
-#ifdef IPCP_FLOW_STATS
int fd;
char * entry;
+ time_t stamp;
entry = strstr(path, RIB_SEPARATOR) + 1;
assert(entry);
fd = atoi(entry);
- pthread_mutex_lock(&dt.stat[fd].lock);
+ stamp = LOAD_ACQUIRE(&dt.stat[fd].stamp);
- if (dt.stat[fd].stamp != -1) {
+ if (stamp != -1) {
attr->size = RIB_FILE_STRLEN;
- attr->mtime = dt.stat[fd].stamp;
+ attr->mtime = stamp;
} else {
attr->size = 0;
attr->mtime = 0;
}
- pthread_mutex_unlock(&dt.stat[fd].lock);
-#else
- (void) path;
- (void) attr;
-#endif
return 0;
}
@@ -369,29 +362,49 @@ static struct rib_ops r_ops = {
.getattr = dt_rib_getattr
};
-#ifdef IPCP_FLOW_STATS
static void stat_used(int fd,
uint64_t addr)
{
struct timespec now;
+ int i;
clock_gettime(CLOCK_REALTIME_COARSE, &now);
- pthread_mutex_lock(&dt.stat[fd].lock);
+ pthread_rwlock_wrlock(&dt.lock);
- memset(&dt.stat[fd], 0, sizeof(dt.stat[fd]));
+ STORE_RELEASE(&dt.stat[fd].stamp, 0);
- dt.stat[fd].stamp = (addr != INVALID_ADDR) ? now.tv_sec : 0;
- dt.stat[fd].addr = addr;
+ /* Don't memset: incremented without locks in fast path. */
+ for (i = 0; i < QOS_CUBE_MAX; ++i) {
+ STORE_RELAXED(&dt.stat[fd].snd_pkt[i], 0);
+ STORE_RELAXED(&dt.stat[fd].snd_bytes[i], 0);
+ STORE_RELAXED(&dt.stat[fd].rcv_pkt[i], 0);
+ STORE_RELAXED(&dt.stat[fd].rcv_bytes[i], 0);
+ STORE_RELAXED(&dt.stat[fd].lcl_r_pkt[i], 0);
+ STORE_RELAXED(&dt.stat[fd].lcl_r_bytes[i], 0);
+ STORE_RELAXED(&dt.stat[fd].lcl_w_pkt[i], 0);
+ STORE_RELAXED(&dt.stat[fd].lcl_w_bytes[i], 0);
+ STORE_RELAXED(&dt.stat[fd].r_drp_pkt[i], 0);
+ STORE_RELAXED(&dt.stat[fd].r_drp_bytes[i], 0);
+ STORE_RELAXED(&dt.stat[fd].w_drp_pkt[i], 0);
+ STORE_RELAXED(&dt.stat[fd].w_drp_bytes[i], 0);
+ STORE_RELAXED(&dt.stat[fd].f_nhp_pkt[i], 0);
+ STORE_RELAXED(&dt.stat[fd].f_nhp_bytes[i], 0);
+ }
- pthread_mutex_unlock(&dt.stat[fd].lock);
+ STORE_RELAXED(&dt.stat[fd].addr, addr);
- pthread_rwlock_wrlock(&dt.lock);
-
- (addr != INVALID_ADDR) ? ++dt.n_flows : --dt.n_flows;
+ if (addr != INVALID_ADDR) {
+ STORE_RELEASE(&dt.stat[fd].stamp, now.tv_sec);
+ ++dt.n_flows;
+ } else {
+ --dt.n_flows;
+ }
pthread_rwlock_unlock(&dt.lock);
}
+#else
+#define dt_stat_inc(idx, name, qc, len) ((void) 0)
#endif
static void handle_event(void * self,
@@ -399,6 +412,7 @@ static void handle_event(void * self,
const void * o)
{
struct conn * c;
+ int fd;
(void) self;
@@ -406,65 +420,57 @@ static void handle_event(void * self,
switch (event) {
case NOTIFY_DT_CONN_ADD:
+ fd = c->flow_info.fd;
#ifdef IPCP_FLOW_STATS
- stat_used(c->flow_info.fd, c->conn_info.addr);
+ stat_used(fd, c->conn_info.addr);
#endif
- psched_add(dt.psched, c->flow_info.fd);
- log_dbg("Added fd %d to packet scheduler.", c->flow_info.fd);
+ cap_reset(fd);
+ psched_add(dt.psched, fd);
+ log_dbg("Added fd %d to packet scheduler.", fd);
break;
case NOTIFY_DT_CONN_DEL:
+ fd = c->flow_info.fd;
#ifdef IPCP_FLOW_STATS
- stat_used(c->flow_info.fd, INVALID_ADDR);
+ stat_used(fd, INVALID_ADDR);
#endif
- psched_del(dt.psched, c->flow_info.fd);
- log_dbg("Removed fd %d from "
- "packet scheduler.", c->flow_info.fd);
+ psched_del(dt.psched, fd);
+ log_dbg("Removed fd %d from packet scheduler.", fd);
break;
default:
break;
}
}
-static void packet_handler(int fd,
- qoscube_t qc,
- struct shm_du_buff * sdb)
+static time_t packet_handler(int fd,
+ qoscube_t qc,
+ struct ssm_pk_buff * spb)
{
struct dt_pci dt_pci;
int ret;
int ofd;
uint8_t * head;
size_t len;
+ size_t qlen;
+ bool marks;
- len = shm_du_buff_tail(sdb) - shm_du_buff_head(sdb);
+ len = ssm_pk_buff_len(spb);
#ifndef IPCP_FLOW_STATS
- (void) fd;
-#else
- pthread_mutex_lock(&dt.stat[fd].lock);
-
- ++dt.stat[fd].rcv_pkt[qc];
- dt.stat[fd].rcv_bytes[qc] += len;
-
- pthread_mutex_unlock(&dt.stat[fd].lock);
+ (void) fd;
#endif
+ dt_stat_inc(fd, rcv, qc, len);
+
memset(&dt_pci, 0, sizeof(dt_pci));
- head = shm_du_buff_head(sdb);
+ head = ssm_pk_buff_head(spb);
dt_pci_des(head, &dt_pci);
- if (dt_pci.dst_addr != ipcpi.dt_addr) {
+ if (dt_pci.dst_addr != dt.addr) {
if (dt_pci.ttl == 0) {
log_dbg("TTL was zero.");
- ipcp_sdb_release(sdb);
-#ifdef IPCP_FLOW_STATS
- pthread_mutex_lock(&dt.stat[fd].lock);
-
- ++dt.stat[fd].r_drp_pkt[qc];
- dt.stat[fd].r_drp_bytes[qc] += len;
-
- pthread_mutex_unlock(&dt.stat[fd].lock);
-#endif
- return;
+ ipcp_spb_release(spb);
+ dt_stat_inc(fd, r_drp, qc, len);
+ return 0;
}
/* FIXME: Use qoscube from PCI instead of incoming flow. */
@@ -472,75 +478,56 @@ static void packet_handler(int fd,
if (ofd < 0) {
log_dbg("No next hop for %" PRIu64 ".",
dt_pci.dst_addr);
- ipcp_sdb_release(sdb);
-#ifdef IPCP_FLOW_STATS
- pthread_mutex_lock(&dt.stat[fd].lock);
+ ipcp_spb_release(spb);
+ dt_stat_inc(fd, f_nhp, qc, len);
+ return 0;
+ }
- ++dt.stat[fd].f_nhp_pkt[qc];
- dt.stat[fd].f_nhp_bytes[qc] += len;
+ marks = ca_marks_ecn();
+ qlen = marks ? ipcp_flow_queued(ofd) : 0;
- pthread_mutex_unlock(&dt.stat[fd].lock);
-#endif
- return;
- }
+ (void) ca_calc_ecn(qlen, head + dt_pci_info.ecn_o, qc, len);
- (void) ca_calc_ecn(ofd, head + dt_pci_info.ecn_o, qc, len);
+ if (marks)
+ cap_stamp(head + dt_pci_info.cap_o, cap_get(ofd));
- ret = ipcp_flow_write(ofd, sdb);
+ ret = ipcp_flow_write(ofd, spb);
if (ret < 0) {
log_dbg("Failed to write packet to fd %d.", ofd);
if (ret == -EFLOWDOWN)
notifier_event(NOTIFY_DT_FLOW_DOWN, &ofd);
- ipcp_sdb_release(sdb);
-#ifdef IPCP_FLOW_STATS
- pthread_mutex_lock(&dt.stat[ofd].lock);
-
- ++dt.stat[ofd].w_drp_pkt[qc];
- dt.stat[ofd].w_drp_bytes[qc] += len;
-
- pthread_mutex_unlock(&dt.stat[ofd].lock);
-#endif
- return;
+ ipcp_spb_release(spb);
+ dt_stat_inc(ofd, w_drp, qc, len);
+ return 0;
}
-#ifdef IPCP_FLOW_STATS
- pthread_mutex_lock(&dt.stat[ofd].lock);
- ++dt.stat[ofd].snd_pkt[qc];
- dt.stat[ofd].snd_bytes[qc] += len;
+ dt_stat_inc(ofd, snd, qc, len);
- pthread_mutex_unlock(&dt.stat[ofd].lock);
-#endif
+ if (marks)
+ cap_update(ofd, qlen, len);
} else {
- dt_pci_shrink(sdb);
- if (dt_pci.eid >= PROG_RES_FDS) {
+ dt_pci_shrink(spb);
+ if (dt_pci.eid >= PROC_RES_FDS) {
uint8_t ecn = *(head + dt_pci_info.ecn_o);
- fa_np1_rcv(dt_pci.eid, ecn, sdb);
- return;
+ uint8_t cap = *(head + dt_pci_info.cap_o);
+ fa_np1_rcv(dt_pci.eid, ecn, cap, spb);
+ return 0;
}
if (dt.comps[dt_pci.eid].post_packet == NULL) {
log_err("No registered component on eid %" PRIu64 ".",
dt_pci.eid);
- ipcp_sdb_release(sdb);
- return;
+ ipcp_spb_release(spb);
+ return 0;
}
-#ifdef IPCP_FLOW_STATS
- pthread_mutex_lock(&dt.stat[fd].lock);
-
- ++dt.stat[fd].lcl_r_pkt[qc];
- dt.stat[fd].lcl_r_bytes[qc] += len;
-
- pthread_mutex_unlock(&dt.stat[fd].lock);
- pthread_mutex_lock(&dt.stat[dt_pci.eid].lock);
+ dt_stat_inc(fd, lcl_r, qc, len);
+ dt_stat_inc(dt_pci.eid, snd, qc, len);
- ++dt.stat[dt_pci.eid].snd_pkt[qc];
- dt.stat[dt_pci.eid].snd_bytes[qc] += len;
-
- pthread_mutex_unlock(&dt.stat[dt_pci.eid].lock);
-#endif
dt.comps[dt_pci.eid].post_packet(dt.comps[dt_pci.eid].comp,
- sdb);
+ spb);
}
+
+ return 0;
}
static void * dt_conn_handle(void * o)
@@ -563,43 +550,49 @@ static void * dt_conn_handle(void * o)
return 0;
}
-int dt_init(enum pol_routing pr,
- uint8_t addr_size,
- uint8_t eid_size,
- uint8_t max_ttl)
+int dt_init(struct dt_config cfg)
{
int i;
int j;
+#ifdef IPCP_FLOW_STATS
char dtstr[RIB_NAME_STRLEN + 1];
- int pp;
+#endif
+ enum pol_pff pp;
struct conn_info info;
memset(&info, 0, sizeof(info));
+ dt.addr = addr_auth_address();
+ if (dt.addr == INVALID_ADDR) {
+ log_err("Failed to get address");
+ return -1;
+ }
+
strcpy(info.comp_name, DT_COMP);
strcpy(info.protocol, DT_PROTO);
info.pref_version = 1;
info.pref_syntax = PROTO_FIXED;
- info.addr = ipcpi.dt_addr;
+ info.addr = dt.addr;
- if (eid_size != 8) { /* only support 64 bits from now */
+ if (cfg.eid_size != 8) { /* only support 64 bits from now */
log_warn("Invalid EID size. Only 64 bit is supported.");
- eid_size = 8;
+ cfg.eid_size = 8;
}
- dt_pci_info.addr_size = addr_size;
- dt_pci_info.eid_size = eid_size;
- dt_pci_info.max_ttl = max_ttl;
+ dt_pci_info.addr_size = cfg.addr_size;
+ dt_pci_info.eid_size = cfg.eid_size;
+ dt_pci_info.max_ttl = cfg.max_ttl;
dt_pci_info.qc_o = dt_pci_info.addr_size;
dt_pci_info.ttl_o = dt_pci_info.qc_o + QOS_LEN;
dt_pci_info.ecn_o = dt_pci_info.ttl_o + TTL_LEN;
- dt_pci_info.eid_o = dt_pci_info.ecn_o + ECN_LEN;
+ dt_pci_info.cap_o = dt_pci_info.ecn_o + ECN_LEN;
+ dt_pci_info.eid_o = dt_pci_info.cap_o + CAP_LEN;
dt_pci_info.head_size = dt_pci_info.eid_o + dt_pci_info.eid_size;
- if (notifier_reg(handle_event, NULL)) {
- log_err("Failed to register with notifier.");
- goto fail_notifier_reg;
+ if (cap_init() < 0) {
+ log_err("Failed to init capacity estimator.");
+ goto fail_cap;
}
if (connmgr_comp_init(COMPID_DT, &info)) {
@@ -607,8 +600,7 @@ int dt_init(enum pol_routing pr,
goto fail_connmgr_comp_init;
}
- pp = routing_init(pr);
- if (pp < 0) {
+ if (routing_init(&cfg.routing, &pp) < 0) {
log_err("Failed to init routing.");
goto fail_routing;
}
@@ -637,34 +629,27 @@ int dt_init(enum pol_routing pr,
goto fail_rwlock_init;
}
- dt.res_fds = bmp_create(PROG_RES_FDS, 0);
+ dt.res_fds = bmp_create(PROC_RES_FDS, 0);
if (dt.res_fds == NULL)
goto fail_res_fds;
#ifdef IPCP_FLOW_STATS
memset(dt.stat, 0, sizeof(dt.stat));
- for (i = 0; i < PROG_MAX_FLOWS; ++i)
- if (pthread_mutex_init(&dt.stat[i].lock, NULL)) {
- for (j = 0; j < i; ++j)
- pthread_mutex_destroy(&dt.stat[j].lock);
- goto fail_stat_lock;
- }
-
dt.n_flows = 0;
-#endif
- sprintf(dtstr, "%s.%" PRIu64, DT, ipcpi.dt_addr);
- if (rib_reg(dtstr, &r_ops))
+
+ sprintf(dtstr, "%s." ADDR_FMT32, DT, ADDR_VAL32(&dt.addr));
+ if (rib_reg(dtstr, &r_ops)) {
+ log_err("Failed to register RIB.");
goto fail_rib_reg;
+ }
+#endif
return 0;
- fail_rib_reg:
#ifdef IPCP_FLOW_STATS
- for (i = 0; i < PROG_MAX_FLOWS; ++i)
- pthread_mutex_destroy(&dt.stat[i].lock);
- fail_stat_lock:
-#endif
+ fail_rib_reg:
bmp_destroy(dt.res_fds);
+#endif
fail_res_fds:
pthread_rwlock_destroy(&dt.lock);
fail_rwlock_init:
@@ -678,21 +663,21 @@ int dt_init(enum pol_routing pr,
fail_routing:
connmgr_comp_fini(COMPID_DT);
fail_connmgr_comp_init:
- notifier_unreg(&handle_event);
- fail_notifier_reg:
+ cap_fini();
+ fail_cap:
return -1;
}
void dt_fini(void)
{
+#ifdef IPCP_FLOW_STATS
char dtstr[RIB_NAME_STRLEN + 1];
+#endif
int i;
- sprintf(dtstr, "%s.%" PRIu64, DT, ipcpi.dt_addr);
- rib_unreg(dtstr);
#ifdef IPCP_FLOW_STATS
- for (i = 0; i < PROG_MAX_FLOWS; ++i)
- pthread_mutex_destroy(&dt.stat[i].lock);
+ sprintf(dtstr, "%s.%" PRIu64, DT, dt.addr);
+ rib_unreg(dtstr);
#endif
bmp_destroy(dt.res_fds);
@@ -708,46 +693,70 @@ void dt_fini(void)
connmgr_comp_fini(COMPID_DT);
- notifier_unreg(&handle_event);
+ cap_fini();
}
int dt_start(void)
{
- dt.psched = psched_create(packet_handler);
+ dt.psched = psched_create(packet_handler, ipcp_flow_read);
if (dt.psched == NULL) {
log_err("Failed to create N-1 packet scheduler.");
- return -1;
+ goto fail_psched;
+ }
+
+ if (notifier_reg(handle_event, NULL)) {
+ log_err("Failed to register with notifier.");
+ goto fail_notifier_reg;
}
if (pthread_create(&dt.listener, NULL, dt_conn_handle, NULL)) {
log_err("Failed to create listener thread.");
- psched_destroy(dt.psched);
- return -1;
+ goto fail_listener;
+ }
+
+ if (routing_start() < 0) {
+ log_err("Failed to start routing.");
+ goto fail_routing;
}
return 0;
+
+ fail_routing:
+ pthread_cancel(dt.listener);
+ pthread_join(dt.listener, NULL);
+ fail_listener:
+ notifier_unreg(&handle_event);
+ fail_notifier_reg:
+ psched_destroy(dt.psched);
+ fail_psched:
+ return -1;
}
void dt_stop(void)
{
+ routing_stop();
+
pthread_cancel(dt.listener);
pthread_join(dt.listener, NULL);
+
+ notifier_unreg(&handle_event);
+
psched_destroy(dt.psched);
}
int dt_reg_comp(void * comp,
- void (* func)(void * func, struct shm_du_buff *),
+ void (* func)(void * func, struct ssm_pk_buff *),
char * name)
{
int eid;
- assert(func);
+ assert(func != NULL);
pthread_rwlock_wrlock(&dt.lock);
eid = bmp_allocate(dt.res_fds);
if (!bmp_is_id_valid(dt.res_fds, eid)) {
- log_warn("Reserved EIDs depleted.");
+ log_err("Cannot allocate EID.");
pthread_rwlock_unlock(&dt.lock);
return -EBADF;
}
@@ -762,71 +771,90 @@ int dt_reg_comp(void * comp,
pthread_rwlock_unlock(&dt.lock);
#ifdef IPCP_FLOW_STATS
- stat_used(eid, ipcpi.dt_addr);
+ stat_used(eid, dt.addr);
#endif
return eid;
}
+void dt_unreg_comp(int eid)
+{
+ assert(eid >= 0 && eid < PROC_RES_FDS);
+
+ pthread_rwlock_wrlock(&dt.lock);
+
+ assert(dt.comps[eid].post_packet != NULL);
+
+ dt.comps[eid].post_packet = NULL;
+ dt.comps[eid].comp = NULL;
+ dt.comps[eid].name = NULL;
+
+ pthread_rwlock_unlock(&dt.lock);
+
+ return;
+}
+
int dt_write_packet(uint64_t dst_addr,
qoscube_t qc,
uint64_t eid,
- struct shm_du_buff * sdb)
+ struct ssm_pk_buff * spb,
+ uint8_t * ecn)
{
struct dt_pci dt_pci;
int fd;
int ret;
uint8_t * head;
size_t len;
+ size_t qlen;
+ bool marks;
- assert(sdb);
- assert(dst_addr != ipcpi.dt_addr);
-
- len = shm_du_buff_tail(sdb) - shm_du_buff_head(sdb);
+ assert(spb);
+ assert(dst_addr != dt.addr);
#ifdef IPCP_FLOW_STATS
- if (eid < PROG_RES_FDS) {
- pthread_mutex_lock(&dt.stat[eid].lock);
-
- ++dt.stat[eid].lcl_r_pkt[qc];
- dt.stat[eid].lcl_r_bytes[qc] += len;
+ len = ssm_pk_buff_len(spb);
- pthread_mutex_unlock(&dt.stat[eid].lock);
- }
+ if (eid < PROC_RES_FDS)
+ dt_stat_inc(eid, lcl_r, qc, len);
#endif
fd = pff_nhop(dt.pff[qc], dst_addr);
if (fd < 0) {
- log_dbg("Could not get nhop for addr %" PRIu64 ".", dst_addr);
+ log_dbg("Could not get nhop for " ADDR_FMT32 ".",
+ ADDR_VAL32(&dst_addr));
#ifdef IPCP_FLOW_STATS
- if (eid < PROG_RES_FDS) {
- pthread_mutex_lock(&dt.stat[eid].lock);
-
- ++dt.stat[eid].lcl_r_pkt[qc];
- dt.stat[eid].lcl_r_bytes[qc] += len;
-
- pthread_mutex_unlock(&dt.stat[eid].lock);
- }
+ if (eid < PROC_RES_FDS)
+ dt_stat_inc(eid, lcl_r, qc, len);
#endif
return -EPERM;
}
- head = shm_du_buff_head_alloc(sdb, dt_pci_info.head_size);
+ head = ssm_pk_buff_push(spb, dt_pci_info.head_size);
if (head == NULL) {
log_dbg("Failed to allocate DT header.");
goto fail_write;
}
- len = shm_du_buff_tail(sdb) - shm_du_buff_head(sdb);
+ len = ssm_pk_buff_len(spb);
dt_pci.dst_addr = dst_addr;
dt_pci.qc = qc;
dt_pci.eid = eid;
dt_pci.ecn = 0;
+ dt_pci.cap = 0;
- (void) ca_calc_ecn(fd, &dt_pci.ecn, qc, len);
+ marks = ca_marks_ecn();
+ qlen = marks ? ipcp_flow_queued(fd) : 0;
+
+ (void) ca_calc_ecn(qlen, &dt_pci.ecn, qc, len);
+
+ if (marks)
+ dt_pci.cap = cap_get(fd);
+
+ if (ecn != NULL)
+ *ecn = dt_pci.ecn;
dt_pci_ser(head, &dt_pci);
- ret = ipcp_flow_write(fd, sdb);
+ ret = ipcp_flow_write(fd, spb);
if (ret < 0) {
log_dbg("Failed to write packet to fd %d.", fd);
if (ret == -EFLOWDOWN)
@@ -834,31 +862,22 @@ int dt_write_packet(uint64_t dst_addr,
goto fail_write;
}
#ifdef IPCP_FLOW_STATS
- pthread_mutex_lock(&dt.stat[fd].lock);
-
- if (dt_pci.eid < PROG_RES_FDS) {
- ++dt.stat[fd].lcl_w_pkt[qc];
- dt.stat[fd].lcl_w_bytes[qc] += len;
- }
- ++dt.stat[fd].snd_pkt[qc];
- dt.stat[fd].snd_bytes[qc] += len;
+ if (dt_pci.eid < PROC_RES_FDS)
+ dt_stat_inc(fd, lcl_w, qc, len);
- pthread_mutex_unlock(&dt.stat[fd].lock);
+ dt_stat_inc(fd, snd, qc, len);
#endif
+ if (marks)
+ cap_update(fd, qlen, len);
+
return 0;
fail_write:
#ifdef IPCP_FLOW_STATS
- pthread_mutex_lock(&dt.stat[fd].lock);
-
- if (eid < PROG_RES_FDS) {
- ++dt.stat[fd].lcl_w_pkt[qc];
- dt.stat[fd].lcl_w_bytes[qc] += len;
- }
- ++dt.stat[fd].w_drp_pkt[qc];
- dt.stat[fd].w_drp_bytes[qc] += len;
+ if (eid < PROC_RES_FDS)
+ dt_stat_inc(fd, lcl_w, qc, len);
- pthread_mutex_unlock(&dt.stat[fd].lock);
+ dt_stat_inc(fd, w_drp, qc, len);
#endif
return -1;
}