Browse Source

topo: remember REQUEST_TABLE for not-yet-created groups, fulfill on create

При REQUEST_TABLE для ещё не созданной группы (CHAT-канал не загружен) пир
дропал запрос, и проситель навсегда не узнавал таблицу — односторонняя
видимость (звонок android→linux падал «peer offline»). Теперь запрос
запоминается в pending_table_reqs[32] и выполняется при создании группы
(topo_group_fulfill_table_reqs → handle_request_table).

Попутно: test_chat_join_e2e — критерий успеха через CHAT_EVT_CONNECT_RESULT
(JOIN_READY) вместо senders_list; doxygen-комментарии в etcp_api.h.
proxy
evgeny 2 weeks ago
parent
commit
6d3ad965ee
  1. 14
      doc/tasks.md
  2. 39
      src/routing_layer/topo_group.c
  3. 11
      src/routing_layer/topo_group.h
  4. 141
      src/transport_layer/etcp_api.h
  5. 33
      tests/test_chat_join_e2e.c

14
doc/tasks.md

@ -186,9 +186,17 @@
`router_no_route`), не влияет на результат — стоит разобрать отдельно.
## Открытые флаки
[ ] **test_chat_join_e2e** — стабильно падает (2/3 сценария): «J not signed by A»
`get_sign_rc=-1` (signed_by-верификация invite). Не связан с socket-классификацией
(обнаружен при прогоне после фикса Q8/LAN). Скорее всего регресс join-протокола.
[+] **test_chat_join_e2e** — сделано. Причина: ложный критерий успеха джойнера J.
`wait_group_started()` проверял `senders_list` непуст в предположении «senders_list
непуст ⟺ J получил JOIN_READY». После коммита b483335c `tgi_to_channel_cb`
(topo_group_invite.c) при `NCD_EVENT_UP` стал вызывать `topo_group_new_conn` →
conn попадает в `senders_list` сразу после установки TCP/UDP-соединения, ДО
завершения хендшейка (JOIN_INFO_REQ→RESP→REQUEST→READY). J выходил «OK» раньше,
чем A получал JOIN_REQUEST → A не подписывал J → `wait_signed` на A/C падал
(`get_sign_rc=-1`). Фикс: критерий J — событие `CHAT_EVT_CONNECT_RESULT(result=0)`
(JOIN_READY), ловится через `chat_event_set_handler`; `wait_group_started` удалён.
Это также убрало гонку «C вышел раньше, чем J синкает свой рекорд» (C выходит сразу
после `wait_signed`, не дожидаясь merkle-синка J). 5 прогонов стабильно 3/3.
[ ] **test_etcp_reconnect** — флаки под параллельной нагрузкой `make check -j4`:
phase 4 (reconnect после server restart) таймаутит (`sent=480 recv=0`), при одиночном

39
src/routing_layer/topo_group.c

@ -102,6 +102,42 @@ static void topo_group_handle_join_group(struct TOPO_GROUP* group, struct ETCP_C
static void topo_group_send_resync(struct ETCP_CONN* conn);
static void topo_group_handle_resync(struct UTUN_INSTANCE* instance, struct ETCP_CONN* conn);
/* Запоминает REQUEST_TABLE для группы, которой у нас ещё нет (CHAT-канал не загружен). */
static void topo_group_remember_table_req(struct TOPO_GROUPS* g, uint64_t node_id, uint64_t group_id) {
if (!g) return;
for (int i = 0; i < g->pending_table_req_count; i++)
if (g->pending_table_reqs[i].node_id == node_id && g->pending_table_reqs[i].group_id == group_id) return;
if (g->pending_table_req_count >= TOPO_MAX_PENDING_TABLE_REQS) {
memmove(g->pending_table_reqs, g->pending_table_reqs + 1,
(TOPO_MAX_PENDING_TABLE_REQS - 1) * sizeof(g->pending_table_reqs[0]));
g->pending_table_req_count = TOPO_MAX_PENDING_TABLE_REQS - 1;
}
g->pending_table_reqs[g->pending_table_req_count].node_id = node_id;
g->pending_table_reqs[g->pending_table_req_count].group_id = group_id;
g->pending_table_req_count++;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "remembered table request node=%016llx grp=%016llx (group not created yet)",
(unsigned long long)node_id, (unsigned long long)group_id);
}
/* Когда группа создана — отвечаем тем, кто просил её таблицу раньше. */
static void topo_group_fulfill_table_reqs(struct TOPO_GROUPS* g, uint64_t group_id) {
if (!g) return;
struct TOPO_GROUP* group = topo_groups_find(g, group_id);
if (!group) return;
int w = 0;
for (int i = 0; i < g->pending_table_req_count; i++) {
struct topo_pending_table_req* req = &g->pending_table_reqs[i];
if (req->group_id != group_id) { g->pending_table_reqs[w++] = *req; continue; }
struct ETCP_CONN* conn = instance_find_conn(g->instance, req->node_id);
if (conn && conn->links_up) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "fulfilling remembered table request node=%016llx grp=%016llx",
(unsigned long long)req->node_id, (unsigned long long)group_id);
topo_group_handle_request_table(group, conn);
}
}
g->pending_table_req_count = w;
}
/* Краткий DEBUG-дамп принятого NODEINFO (id/ver/счётчики/hops). */
static void nodeinfo_dump_log(const uint8_t* data, size_t len) {
if (!data || len < sizeof(struct TOPOMSG_NODEINFO_PKT)) return;
@ -219,6 +255,8 @@ static void topo_group_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry*
struct TOPO_GROUP* group = topo_groups_find(instance->topo_groups, pkt_group_id);
if (!group) {
if (subcmd == TOPO_SUBCMD_REQUEST_TABLE)
topo_group_remember_table_req(instance->topo_groups, from_conn->peer_node_id, pkt_group_id);
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "BGP recv %s from %s: group %016llx not found, dropping", group_subcmd_name(subcmd), from_conn->log_name, (unsigned long long)pkt_group_id);
queue_dgram_free(entry); queue_entry_free(entry); return;
}
@ -580,6 +618,7 @@ struct TOPO_GROUP* topo_groups_create_group(struct TOPO_GROUPS* g, uint64_t grou
DEBUG_INFO(DEBUG_CATEGORY_BGP, "Group created: group_id=%016llx type=%u ch_id=%s", (unsigned long long)group_id, group_type, group->channel_id);
if (group_type == TOPO_GROUP_TYPE_CHAT && chat_setting_get_int(group->instance, "group_autoconnect", 1)) topo_group_connect_init(group);
topo_group_fulfill_table_reqs(g, group_id);
return group;
}

11
src/routing_layer/topo_group.h

@ -195,6 +195,15 @@ struct TOPO_GROUP {
struct topo_node_cbk_entry* node_cbks; // цепочка подписчиков на события узлов
};
/* Отложенный REQUEST_TABLE: пир запросил таблицу группы, которой у нас ещё нет.
* Запоминаем (node_id, group_id) и отвечаем, когда группа будет создана. */
#define TOPO_MAX_PENDING_TABLE_REQS 32
struct topo_pending_table_req {
uint64_t node_id; /* кто запросил таблицу */
uint64_t group_id; /* какую группу */
};
/**
* @brief Контейнер всех групп топологии экземпляра
*/
@ -209,6 +218,8 @@ struct TOPO_GROUPS {
struct memory_pool* v4_subnet_pool;
struct memory_pool* v6_subnet_pool;
topo_node_updated_fn node_updated_cb; /* optional — set by chatgui's member_sync */
struct topo_pending_table_req pending_table_reqs[TOPO_MAX_PENDING_TABLE_REQS];
int pending_table_req_count;
};
/**

141
src/transport_layer/etcp_api.h

@ -108,6 +108,15 @@ struct etcp_cbk_entry {
struct etcp_cbk_entry* next;
};
/**
* @brief Рассылка per-connection события
*
* @param conn ETCP соединение
* @param event Битовая маска события (ETCP_CBK_EVENT_*)
*
* @note Вызываются все коллбэки соединения, у которых event_mask
* пересекается с переданным event
*/
void etcp_cbk_fire(struct ETCP_CONN* conn, int event);
// Instance-level callback types (separate from per-connection etcp_cbk_fn in etcp.h)
@ -119,12 +128,52 @@ struct etcp_inst_cbk_entry {
};
// ---- Per-connection unified callback API ----
/**
* @brief Зарегистрировать per-connection коллбэк с маской событий
*
* @param conn ETCP соединение
* @param fn Коллбэк (сигнатура etcp_cbk_fn)
* @param arg Пользовательский аргумент, передаваемый в коллбэк
* @param event_mask Битовая маска интересующих событий (ETCP_CBK_EVENT_*)
*/
void etcp_conn_add_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg, int event_mask);
/**
* @brief Удалить per-connection коллбэк
*
* @param conn ETCP соединение
* @param fn Коллбэк
* @param arg Пользовательский аргумент (для идентификации записи)
*/
void etcp_conn_remove_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg);
/**
* @brief Обновить маску событий у уже зарегистрированного коллбэка
*
* @param conn ETCP соединение
* @param fn Коллбэк
* @param arg Пользовательский аргумент
* @param event_mask Новая битовая маска событий (ETCP_CBK_EVENT_*)
*/
void etcp_conn_update_cbk_mask(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg, int event_mask);
// ---- Instance-level callbacks ----
/**
* @brief Зарегистрировать instance-level коллбэк «создано новое соединение»
*
* @param inst UTUN instance
* @param fn Коллбэк (сигнатура etcp_inst_cbk_fn)
* @param arg Пользовательский аргумент
*/
void etcp_add_new_conn_cbk(struct UTUN_INSTANCE* inst, etcp_inst_cbk_fn fn, void* arg);
/**
* @brief Удалить instance-level коллбэк «создано новое соединение»
*
* @param inst UTUN instance
* @param fn Коллбэк
* @param arg Пользовательский аргумент
*/
void etcp_remove_new_conn_cbk(struct UTUN_INSTANCE* inst, etcp_inst_cbk_fn fn, void* arg);
// ---- Link status change callback (instance-level) ----
@ -137,8 +186,34 @@ struct etcp_link_status_cbk_entry {
struct etcp_link_status_cbk_entry* next;
};
/**
* @brief Зарегистрировать instance-level коллбэк смены состояния линка
*
* @param inst UTUN instance
* @param fn Коллбэк (сигнатура etcp_link_status_cbk_fn)
* @param arg Пользовательский аргумент
*/
void etcp_add_link_status_cbk(struct UTUN_INSTANCE* inst, etcp_link_status_cbk_fn fn, void* arg);
/**
* @brief Удалить instance-level коллбэк смены состояния линка
*
* @param inst UTUN instance
* @param fn Коллбэк
* @param arg Пользовательский аргумент
*/
void etcp_remove_link_status_cbk(struct UTUN_INSTANCE* inst, etcp_link_status_cbk_fn fn, void* arg);
/**
* @brief Внутренняя функция: рассылка события смены состояния линка
*
* @param link Линк, состояние которого изменилось
* @param old_state Предыдущее link_state (0-init, 1-handshake, 2-try_reconnect, 3-connected)
* @param old_status Предыдущее link_status (1-up, 0-down)
*
* @note Вызывается автоматически из etcp_connections.c при изменении
* link_state/link_status; не предназначен для вызова извне
*/
void etcp_fire_link_status_cbk(struct ETCP_LINK* link, int old_state, int old_status);
// ---- Background connection initialization ----
@ -148,6 +223,21 @@ void etcp_fire_link_status_cbk(struct ETCP_LINK* link, int old_state, int old_st
typedef void (*etcp_connect_callback_t)(void* arg, struct ETCP_CONN* conn, int type);
/**
* @brief Фоновая инициализация соединения к узлу
*
* @param instance UTUN instance
* @param node Узел из node registry (TOPO_GROUP_NODE)
* @param cb Коллбэк завершения (сигнатура etcp_connect_callback_t)
* @param arg Пользовательский аргумент
* @param flags Битовая маска этапов (ETCP_CONNECT_EARLY / LATE / BGP_READY)
* @return 0 при успехе, -1 при ошибке
*
* @note Коллбэк вызывается на каждом запрошенном этапе с параметром
* type = ETCP_CONNECT_EARLY / ETCP_CONNECT_LATE / ETCP_CONNECT_BGP_READY
* (или 0 при таймауте/неудаче). При уже готовом соединении EARLY+LATE
* доставляются немедленно.
*/
int etcp_connect(struct UTUN_INSTANCE* instance, struct TOPO_GROUP_NODE* node,
etcp_connect_callback_t cb, void* arg, uint8_t flags);
@ -195,7 +285,24 @@ int etcp_bind(struct UTUN_INSTANCE* inst, uint8_t id, etcp_recv_fn callback);
*/
int etcp_unbind(struct UTUN_INSTANCE* inst, uint8_t id);
/**
* @brief Зарегистрировать instance-level коллбэк статуса соединения
*
* @param inst UTUN instance
* @param fn Коллбэк (сигнатура etcp_conn_status_fn)
* @param arg Пользовательский аргумент
*
* @note status = ETCP_CONN_STATUS_NEW / UP / DOWN / DELETE
*/
void etcp_add_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg);
/**
* @brief Удалить instance-level коллбэк статуса соединения
*
* @param inst UTUN instance
* @param fn Коллбэк
* @param arg Пользовательский аргумент
*/
void etcp_remove_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg);
/* Socket property change events (instance-level) */
@ -210,10 +317,44 @@ struct etcp_socket_cbk_entry {
struct etcp_socket_cbk_entry* next;
};
/**
* @brief Зарегистрировать instance-level коллбэк событий сокета
*
* @param inst UTUN instance
* @param fn Коллбэк (сигнатура etcp_socket_cbk_fn)
* @param arg Пользовательский аргумент
* @param event_mask Битовая маска событий (ETCP_SOCKET_EVENT_ADDR_CHANGED / STATUS_CHANGED)
*/
void etcp_add_socket_cbk(struct UTUN_INSTANCE* inst, etcp_socket_cbk_fn fn, void* arg, int event_mask);
/**
* @brief Удалить instance-level коллбэк событий сокета
*
* @param inst UTUN instance
* @param fn Коллбэк
* @param arg Пользовательский аргумент
*/
void etcp_remove_socket_cbk(struct UTUN_INSTANCE* inst, etcp_socket_cbk_fn fn, void* arg);
/**
* @brief Внутренняя функция: рассылка события сокета
*
* @param sock Сокет, у которого произошло событие
* @param event Битовая маска события (ETCP_SOCKET_EVENT_*)
*
* @note Вызываются коллбэки, у которых event_mask пересекается с event
*/
void etcp_socket_cbk_fire(struct ETCP_SOCKET* sock, int event);
/**
* @brief Установить состояние обмена маршрут-таблицами (BGP)
*
* @param conn ETCP соединение
* @param new_state 0-не активен, 1-надо инициировать (клиент), 2-обмен идёт,
* 3-завершён, 4-пропущен (нет BGP)
*
* @note При new_state >= 3 и наличии bgp_ready_cbk вызывается этот коллбэк
*/
void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state);
/**

33
tests/test_chat_join_e2e.c

@ -17,6 +17,7 @@
#include "invite_link.h"
#include "invite_build.h"
#include "member_sync.h"
#include "chat_event.h"
#include "../routing_layer/topo_node_sqlite.h"
#include "../routing_layer/topo_group.h"
#include "../routing_layer/topo_node.h"
@ -235,21 +236,22 @@ static int wait_file(struct UTUN_INSTANCE* inst, const char* path, int max_iter)
return 0;
}
/* J получил JOIN_READY: cs_handle_join_ready вызвал topo_group_new_conn → senders_list непуст */
static int wait_group_started(struct UTUN_INSTANCE* inst, int max_iter) {
uint64_t gid = strtoull(g_sh.ch_id, NULL, 10);
for (int a = 0; a < max_iter; a++) {
struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, gid);
if (g && g->senders_list && g->senders_list->head) return 1;
uasync_poll(inst->ua, POLL_MS);
/* JOIN_READY (успех join) наблюдается через CHAT_EVT_CONNECT_RESULT(result=0):
* сигнал точный (в отличие от senders_list, заполняемого и ncd-подключением)
* и не зависит от того, успел ли connection-узел остаться в живых до синка мембера. */
static int g_j_join_ok = 0;
static void join_event_handler(struct UTUN_INSTANCE* inst, int type, const uint8_t* data, int len) {
(void)inst;
if (type == CHAT_EVT_CONNECT_RESULT && len >= 12) {
int r = 0; memcpy(&r, data + 8, 4);
if (r == 0) g_j_join_ok = 1;
}
{
struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, gid);
fprintf(stderr, " [join-ready-fail] self=0x%016llx group=%p senders_head=%p\n",
(unsigned long long)inst->node_id, (void*)g,
(void*)(g && g->senders_list ? g->senders_list->head : NULL));
}
return 0;
static int wait_join_result(struct UTUN_INSTANCE* inst, int max_iter) {
for (int a = 0; a < max_iter && !g_j_join_ok; a++) uasync_poll(inst->ua, POLL_MS);
return g_j_join_ok;
}
/* ── настройка канала и мемберов ── */
@ -291,6 +293,7 @@ static int child_main(const char* role, const char* dir, int invalid_key) {
struct UTUN_INSTANCE* inst = utun_instance_create(ua, cfg);
if (!inst) { fprintf(stderr, "%s: create failed\n", role); uasync_destroy(ua, 0); return 1; }
utun_instance_init(inst);
chat_event_set_handler(inst, join_event_handler);
int rc = 1;
uint64_t ch_num = strtoull(g_sh.ch_id, NULL, 10);
@ -397,8 +400,8 @@ static int child_main(const char* role, const char* dir, int invalid_key) {
rc = wait_not_member(inst, g_sh.nid[IDX_J], 200) ? 0 : 1;
if (rc) fprintf(stderr, "J: unexpectedly became member with wrong key\n");
} else {
/* критерий J по ТЗ: получил JOIN_READY (topo_group_new_conn → senders_list) */
if (!wait_group_started(inst, 8000)) {
/* критерий J по ТЗ: получил JOIN_READY (CHAT_EVT_CONNECT_RESULT result=0) */
if (!wait_join_result(inst, 8000)) {
fprintf(stderr, "J: no JOIN_READY\n"); goto out;
}
rc = 0;

Loading…
Cancel
Save