Browse Source

etcp_router: CLOSE/RST signaling flags in seq, close notification to services

- ROUTER_SEQ_CLOSE_FLAG (0x80000000): normal conn close
- ROUTER_SEQ_RST_FLAG   (0x40000000): conn not found
- router_send_to(): SVC_ROUTE with flags, no conn needed
- router_send_close_to_service(): notify local service [svc_id],len=1
- router_close_and_notify(): CLOSE to remote + close to local + cleanup
- etcp_router_conn_close(): sends CLOSE to remote before cleanup
- etcp_router_conn_close_all_for_node(): close all conns for dead peer
- recv_cb: flags checked only when pl_len==0 (avoid seq collision)
- etcp_router_destroy: close notification to each service
congestion
Evgeny 4 months ago
parent
commit
53d3f26e87
  1. 117
      src/etcp_router.c
  2. 7
      src/etcp_router.h

117
src/etcp_router.c

@ -22,8 +22,12 @@ static void router_schedule_ack(struct ETCP_ROUTER_CONN* rconn);
static void router_try_assembly(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* conn);
static void router_deliver(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* conn, const uint8_t* payload, size_t payload_len);
static int router_send_one(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len);
static int router_send_one_flags(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len, uint32_t seq_flags);
static void router_send_to(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t svc_id, uint32_t seq_flags);
static void router_drain_send_q(struct ETCP_ROUTER_CONN* rconn);
static void router_send_resume_cb(void* arg);
static void router_send_close_to_service(struct ETCP_ROUTER_CONN* rconn);
static void router_close_and_notify(struct ETCP_ROUTER_CONN* rconn);
// ====================================================================
// Управление ROUTER_CONN
@ -43,11 +47,17 @@ static struct ETCP_ROUTER_CONN* router_conn_find(struct UTUN_INSTANCE* inst,
// ====================================================================
static int router_send_one(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len) {
return router_send_one_flags(rconn, payload, pl_len, 0);
}
static int router_send_one_flags(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len, uint32_t seq_flags) {
struct UTUN_INSTANCE* inst = rconn->inst;
uint32_t seq = rconn->tx_seq++;
uint32_t seq = (seq_flags & (ROUTER_SEQ_CLOSE_FLAG | ROUTER_SEQ_RST_FLAG))
? 0 : rconn->tx_seq++;
seq |= seq_flags;
size_t total_len = SVC_ROUTE_HDR_SIZE + pl_len;
uint8_t* dgram = u_malloc(total_len);
if (!dgram) { rconn->tx_seq--; return -1; }
if (!dgram) { if (!seq_flags) rconn->tx_seq--; return -1; }
struct SVC_ROUTE_HDR* hdr = (struct SVC_ROUTE_HDR*)dgram;
hdr->cmd = ETCP_ID_SVC_ROUTE;
hdr->dst_node_id = rconn->remote_node_id;
@ -57,7 +67,7 @@ static int router_send_one(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payloa
if (pl_len > 0) memcpy(dgram + SVC_ROUTE_HDR_SIZE, payload, pl_len);
struct ll_entry* entry = queue_entry_new(0);
if (!entry) { u_free(dgram); rconn->tx_seq--; return -1; }
if (!entry) { u_free(dgram); if (!seq_flags) rconn->tx_seq--; return -1; }
entry->dgram = dgram;
entry->len = total_len;
@ -65,7 +75,7 @@ static int router_send_one(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payloa
if (!conn) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "router_send_one: no route to %016llx svc_id=%u",
(unsigned long long)rconn->remote_node_id, rconn->svc_id);
queue_entry_free(entry); queue_dgram_free(entry); rconn->tx_seq--; return -1;
queue_entry_free(entry); queue_dgram_free(entry); if (!seq_flags) rconn->tx_seq--; return -1;
}
rconn->last_dgram_ts = get_current_timestamp();
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "router_send: seq=%u → %016llx svc_id=%u len=%zu inflight=%u",
@ -74,6 +84,70 @@ static int router_send_one(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payloa
return etcp_send(conn, entry);
}
// Отправка без conn (RST для неизвестного src)
static void router_send_to(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t svc_id, uint32_t seq_flags) {
struct ETCP_CONN* conn = route_bgp_find_conn_for_node(inst->bgp, dst);
if (!conn) return;
struct SVC_ROUTE_HDR* hdr = (struct SVC_ROUTE_HDR*)u_malloc(SVC_ROUTE_HDR_SIZE);
if (!hdr) return;
hdr->cmd = ETCP_ID_SVC_ROUTE;
hdr->dst_node_id = dst;
hdr->src_node_id = inst->node_id;
hdr->seq = seq_flags;
hdr->svc_id = svc_id;
struct ll_entry* entry = queue_entry_new(0);
if (!entry) { u_free(hdr); return; }
entry->dgram = (uint8_t*)hdr;
entry->len = SVC_ROUTE_HDR_SIZE;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "router_send_to: seq_flags=0x%08x → %016llx svc_id=%u",
seq_flags, (unsigned long long)dst, svc_id);
etcp_send(conn, entry);
}
static void router_send_close_to_service(struct ETCP_ROUTER_CONN* rconn) {
struct UTUN_INSTANCE* inst = rconn->inst;
etcp_recv_fn cb = inst->router_bindings.callbacks[rconn->svc_id];
if (!cb) return;
struct ll_entry* e = queue_entry_new(0);
if (!e) return;
e->dgram = u_malloc(1);
if (!e->dgram) { queue_entry_free(e); return; }
e->dgram[0] = rconn->svc_id;
e->len = 1;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "router_close_svc: svc_id=%u remote=%016llx",
rconn->svc_id, (unsigned long long)rconn->remote_node_id);
cb(NULL, e);
}
// Закрыть conn: отправить CLOSE удалённой стороне, уведомить локальный сервис, очистить
static void router_close_and_notify(struct ETCP_ROUTER_CONN* rconn) {
if (!rconn) return;
struct UTUN_INSTANCE* inst = rconn->inst;
router_send_close_to_service(rconn);
// Отправляем CLOSE удалённой стороне
struct ETCP_CONN* conn = route_bgp_find_conn_for_node(inst->bgp, rconn->remote_node_id);
if (conn) router_send_one_flags(rconn, NULL, 0, ROUTER_SEQ_CLOSE_FLAG);
// Таймеры
if (rconn->ack_timer) { uasync_cancel_timeout(inst->ua, rconn->ack_timer); rconn->ack_timer = NULL; }
if (rconn->idle_ack_timer) { uasync_cancel_timeout(inst->ua, rconn->idle_ack_timer); rconn->idle_ack_timer = NULL; }
if (rconn->send_resume_timer) { uasync_cancel_timeout(inst->ua, rconn->send_resume_timer); rconn->send_resume_timer = NULL; }
if (rconn->recv_q) {
struct ll_entry* f;
while ((f = queue_data_get(rconn->recv_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); }
queue_free(rconn->recv_q); rconn->recv_q = NULL;
}
if (rconn->send_q) {
struct ll_entry* f;
while ((f = queue_data_get(rconn->send_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); }
queue_free(rconn->send_q); rconn->send_q = NULL;
}
if (inst->router_conns)
queue_remove_data(inst->router_conns, &rconn->ll);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "router_conn: closed remote=%016llx svc_id=%u",
(unsigned long long)rconn->remote_node_id, rconn->svc_id);
queue_entry_free(&rconn->ll);
}
static void router_drain_send_q(struct ETCP_ROUTER_CONN* rconn) {
while (1) {
if (rconn->tx_seq - rconn->rx_acked >= ROUTER_MAX_INFLIGHT) break;
@ -236,9 +310,15 @@ static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry)
if (hdr->dst_node_id == inst->node_id) {
// ========== Мы — целевая нода ==========
// Авто-создаём conn при первом пакете
struct ETCP_ROUTER_CONN* rconn = router_conn_find(inst, hdr->src_node_id, hdr->svc_id);
// CLOSE или RST — pl_len==0 и флаг в seq (закрываем conn)
if (pl_len == 0 && (hdr->seq & (ROUTER_SEQ_CLOSE_FLAG | ROUTER_SEQ_RST_FLAG))) {
if (rconn) router_close_and_notify(rconn);
queue_dgram_free(entry); queue_entry_free(entry);
return;
}
if (pl_len == 0) {
// ACK-пакет: обновляем rx_acked, drain send_q, выходим
if (rconn) {
@ -343,6 +423,7 @@ void etcp_router_destroy(struct UTUN_INSTANCE* inst) {
struct ll_entry* entry;
while ((entry = queue_data_get(inst->router_conns)) != NULL) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)entry;
router_send_close_to_service(rconn);
if (rconn->ack_timer) uasync_cancel_timeout(inst->ua, rconn->ack_timer);
if (rconn->idle_ack_timer) uasync_cancel_timeout(inst->ua, rconn->idle_ack_timer);
if (rconn->send_resume_timer) uasync_cancel_timeout(inst->ua, rconn->send_resume_timer);
@ -475,11 +556,19 @@ int etcp_router_conn_send(struct ETCP_ROUTER_CONN* rconn,
}
void etcp_router_conn_close(struct ETCP_ROUTER_CONN* rconn) {
if (!rconn) return;
struct UTUN_INSTANCE* inst = rconn->inst;
if (rconn->ack_timer) uasync_cancel_timeout(inst->ua, rconn->ack_timer);
if (rconn->idle_ack_timer) uasync_cancel_timeout(inst->ua, rconn->idle_ack_timer);
if (rconn->send_resume_timer) uasync_cancel_timeout(inst->ua, rconn->send_resume_timer);
router_close_and_notify(rconn);
}
void etcp_router_conn_close_all_for_node(struct UTUN_INSTANCE* inst, uint64_t remote_node_id) {
if (!inst || !inst->router_conns) return;
struct ll_entry* entry;
while ((entry = queue_data_get(inst->router_conns)) != NULL) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)entry;
if (rconn->remote_node_id == remote_node_id) {
router_send_close_to_service(rconn);
if (rconn->ack_timer) { uasync_cancel_timeout(inst->ua, rconn->ack_timer); rconn->ack_timer = NULL; }
if (rconn->idle_ack_timer) { uasync_cancel_timeout(inst->ua, rconn->idle_ack_timer); rconn->idle_ack_timer = NULL; }
if (rconn->send_resume_timer) { uasync_cancel_timeout(inst->ua, rconn->send_resume_timer); rconn->send_resume_timer = NULL; }
if (rconn->recv_q) {
struct ll_entry* f;
while ((f = queue_data_get(rconn->recv_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); }
@ -490,9 +579,11 @@ void etcp_router_conn_close(struct ETCP_ROUTER_CONN* rconn) {
while ((f = queue_data_get(rconn->send_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); }
queue_free(rconn->send_q);
}
if (inst->router_conns)
queue_remove_data(inst->router_conns, &rconn->ll);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "router_conn: closed remote=%016llx svc_id=%u",
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "router_conn: closed by node remote=%016llx svc_id=%u",
(unsigned long long)rconn->remote_node_id, rconn->svc_id);
queue_entry_free(&rconn->ll);
} else {
queue_data_put(inst->router_conns, entry);
}
}
}

7
src/etcp_router.h

@ -51,6 +51,10 @@ struct ETCP_ROUTER_CONN {
#define ROUTER_ACK_IDLE_TB 5000 // idle таймаут: 500ms
#define ROUTER_SEND_RESUME_TB 500 // retry интервал send_q: 50ms
// Флаги в старших битах seq (взаимоисключающие)
#define ROUTER_SEQ_CLOSE_FLAG 0x80000000 // нормальное закрытие conn
#define ROUTER_SEQ_RST_FLAG 0x40000000 // conn не найден (reset)
// Bindings для сервисов внутри etcp_router (аналогично ETCP_BINDINGS)
struct ETCP_ROUTER_BINDINGS {
etcp_recv_fn callbacks[SVC_ROUTE_MAX_BINDINGS];
@ -81,6 +85,9 @@ int etcp_router_conn_send(struct ETCP_ROUTER_CONN* rconn,
// Закрыть seq-подключение
void etcp_router_conn_close(struct ETCP_ROUTER_CONN* rconn);
// Закрыть все router_conn для указанного remote_node_id (peer умер)
void etcp_router_conn_close_all_for_node(struct UTUN_INSTANCE* inst, uint64_t remote_node_id);
// Возвращает количество пакетов в очереди normalizer->input для узла node_id
int etcp_router_input_q_count(struct UTUN_INSTANCE* inst, uint64_t node_id);

Loading…
Cancel
Save