From 865657b3bd920650cbbc5a506d14cc796345f656 Mon Sep 17 00:00:00 2001 From: evgeny Date: Sun, 4 Oct 2026 16:39:19 +0200 Subject: [PATCH] Expose group peer readiness transitions to subscribers --- src/routing_layer/topo_group.c | 35 ++++++++++++++++++++++++++++++++++ src/routing_layer/topo_group.h | 13 +++++++++++++ 2 files changed, 48 insertions(+) diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index 811c898e..3f423a2a 100644 --- a/src/routing_layer/topo_group.c +++ b/src/routing_layer/topo_group.c @@ -97,7 +97,18 @@ int topo_group_peer_ready(const struct TOPO_GROUP* group, uint64_t peer_id) { peer->accepted && peer->table_received && peer->table_sent; } +static void topo_group_notify_ready(struct TOPO_GROUP_CONN_ITEM* peer, int ready) { + if (peer->ready_notified == ready) return; + peer->ready_notified = ready; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group peer readiness: group=%016llx peer=%016llx ready=%d epoch=%016llx", + (unsigned long long)peer->group->group_id, (unsigned long long)peer->node_id, + ready, (unsigned long long)peer->local_epoch); + for (struct topo_peer_ready_cbk* cb = peer->group->peer_ready_cbks; cb; cb = cb->next) + cb->fn(peer->group, peer->node_id, ready, cb->arg); +} + static void topo_group_peer_free(struct TOPO_GROUP_CONN_ITEM* peer) { + topo_group_notify_ready(peer, 0); topo_exchange_cancel(peer); topo_tx_clear(peer); topo_adverts_clear(peer); @@ -212,6 +223,7 @@ static void topo_group_log_exchange(struct TOPO_GROUP* group, struct TOPO_GROUP_ topo_exchange_cancel(peer); peer->retained = 1; topo_group_connect_on_up(group, peer->conn); } + topo_group_notify_ready(peer, topo_group_peer_ready(group, peer->node_id)); topo_group_peer_progressed(group, peer->node_id); DEBUG_INFO(DEBUG_CATEGORY_BGP, "group exchange: group=%016llx peer=%016llx exchange=%016llx sent=%u received=%u ready=%d", (unsigned long long)group->group_id, (unsigned long long)peer->conn->peer_node_id, @@ -251,6 +263,7 @@ static int topo_group_reset_exchange(struct TOPO_GROUP_CONN_ITEM* peer) { topo_tx_clear(peer); peer->tx_failed = 0; peer->transport_down = 0; topo_adverts_clear(peer); peer->accepted = 0; peer->table_received = 0; peer->table_sent = 0; peer->peer_epoch = 0; + topo_group_notify_ready(peer, 0); uint64_t epoch; if (random_bytes((uint8_t*)&epoch, sizeof(epoch)) != 0 || !epoch) { peer->local_epoch = 0; @@ -668,6 +681,10 @@ static void topo_group_destroy(struct TOPO_GROUP* group) { struct topo_node_cbk_entry* c = group->node_cbks; group->node_cbks = c->next; u_free(c); } + while (group->peer_ready_cbks) { + struct topo_peer_ready_cbk* cb = group->peer_ready_cbks; + group->peer_ready_cbks = cb->next; u_free(cb); + } DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 5 done"); } @@ -1012,6 +1029,7 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, en /* Больше не публикуем в закрываемую сессию, даже если общий транспорт ещё UP. */ pending->transport_down = 1; + topo_group_notify_ready(pending, 0); int nodes_removed = 0; int cascaded = 0; struct ll_entry* node_entry = group->nodes ? group->nodes->head : NULL; @@ -1520,6 +1538,7 @@ static void topo_tx_fail(struct TOPO_GROUP_CONN_ITEM* peer) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group sender failed: group=%016llx peer=%016llx epoch=%016llx", (unsigned long long)peer->group->group_id, (unsigned long long)peer->node_id, (unsigned long long)peer->local_epoch); topo_tx_clear(peer); peer->tx_failed = 1; peer->table_sent = 0; + topo_group_notify_ready(peer, 0); topo_exchange_arm(peer); topo_recovery_changed(peer->group); topo_group_connect_changed(peer->group); @@ -1851,6 +1870,22 @@ static void topo_group_handle_resync(struct UTUN_INSTANCE* instance, struct ETCP /* ── BGP node event callbacks ── */ /* Регистрирует подписчика на события узлов группы (NEW/UPDATE/REMOVE). */ +int topo_group_add_peer_ready_cbk(struct TOPO_GROUP* group, topo_peer_ready_fn fn, void* arg) { + if (!group || !fn) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "peer readiness subscription: invalid arguments"); return -1; } + struct topo_peer_ready_cbk* cb = u_malloc(sizeof(*cb)); + if (!cb) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "peer readiness subscription allocation failed"); return -1; } + cb->fn = fn; cb->arg = arg; cb->next = group->peer_ready_cbks; group->peer_ready_cbks = cb; + return 0; +} + +void topo_group_remove_peer_ready_cbk(struct TOPO_GROUP* group, topo_peer_ready_fn fn, void* arg) { + for (struct topo_peer_ready_cbk** p = group ? &group->peer_ready_cbks : NULL; p && *p; p = &(*p)->next) { + if ((*p)->fn == fn && (*p)->arg == arg) { + struct topo_peer_ready_cbk* cb = *p; *p = cb->next; u_free(cb); return; + } + } +} + void topo_group_add_node_cbk(struct TOPO_GROUP* group, topo_node_event_fn fn, void* arg) { if (!group || !fn) return; struct topo_node_cbk_entry* e = u_malloc(sizeof(*e)); diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index 08b5ee21..37fb4970 100644 --- a/src/routing_layer/topo_group.h +++ b/src/routing_layer/topo_group.h @@ -82,6 +82,17 @@ struct topo_node_cbk_entry { struct topo_node_cbk_entry* next; }; +/* Переход готовности конкретной групповой сессии, независимо от общего транспорта. + * Callback выполняется в uasync; не должен удалять группу или её подписчиков. */ +typedef void (*topo_peer_ready_fn)(struct TOPO_GROUP* group, uint64_t node_id, int ready, void* arg); +struct topo_peer_ready_cbk { + topo_peer_ready_fn fn; + void* arg; + struct topo_peer_ready_cbk* next; +}; +int topo_group_add_peer_ready_cbk(struct TOPO_GROUP* group, topo_peer_ready_fn fn, void* arg); +void topo_group_remove_peer_ready_cbk(struct TOPO_GROUP* group, topo_peer_ready_fn fn, void* arg); + /* ── Global nodeinfo callbacks (one per instance, any node change) ── */ typedef void (*nodeinfo_cbk_fn)(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, void* arg); @@ -198,6 +209,7 @@ struct TOPO_GROUP_CONN_ITEM { struct ll_queue* advertised; // node_id -> хеш отправленного NODEINFO без RTT; владеет сессия uint8_t table_received; // TABLE_COMPLETE текущей сессии принят uint8_t table_sent; // ответ на запрос пира целиком передан транспорту + uint8_t ready_notified; // последнее опубликованное состояние READY }; #define TOPO_GROUP_SYNC_TIMEOUT_MS 5000 /* отсутствие прогресса: удалить пути пира и начать новый JOIN */ @@ -230,6 +242,7 @@ struct TOPO_GROUP { struct radio_ctx* radio; // radio (PTT walkie-talkie) per-group context uint8_t radio_active; // 1 = я слушаю рацию этого канала (TOPO_FLAG_RADIO в NODEINFO) struct topo_node_cbk_entry* node_cbks; // цепочка подписчиков на события узлов + struct topo_peer_ready_cbk* peer_ready_cbks; uint8_t stopping; // teardown начат; новые попытки запрещены };