diff --git a/doc/tasks.md b/doc/tasks.md index 8a507b73..ded5f828 100644 --- a/doc/tasks.md +++ b/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`), при одиночном diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index 1ddb7cc6..da0329e8 100644 --- a/src/routing_layer/topo_group.c +++ b/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; } diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index 99b5d79c..d3cdf2a1 100644 --- a/src/routing_layer/topo_group.h +++ b/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; }; /** diff --git a/src/transport_layer/etcp_api.h b/src/transport_layer/etcp_api.h index 367e3859..9b42f7e9 100644 --- a/src/transport_layer/etcp_api.h +++ b/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); /** diff --git a/tests/test_chat_join_e2e.c b/tests/test_chat_join_e2e.c index c526a445..890699b9 100644 --- a/tests/test_chat_join_e2e.c +++ b/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); - } - { - 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)); +/* 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; } - 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;