diff --git a/src/etcp_router.c b/src/etcp_router.c index 569a2224..985a0e59 100644 --- a/src/etcp_router.c +++ b/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,24 +556,34 @@ 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); - 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); - } - 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); + 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); } + queue_free(rconn->recv_q); + } + 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); + } + 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); + } } - 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); } diff --git a/src/etcp_router.h b/src/etcp_router.h index 91bb8406..d65410cf 100644 --- a/src/etcp_router.h +++ b/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);