Browse Source

Expose group peer readiness transitions to subscribers

master
evgeny 5 days ago
parent
commit
865657b3bd
  1. 35
      src/routing_layer/topo_group.c
  2. 13
      src/routing_layer/topo_group.h

35
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));

13
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 начат; новые попытки запрещены
};

Loading…
Cancel
Save