diff --git a/close_bug.txt b/close_bug.txt new file mode 100644 index 00000000..2ace0830 --- /dev/null +++ b/close_bug.txt @@ -0,0 +1,67 @@ +# cm_invite_pending — утечка указателей и corruption при destroy + +## Корень проблемы + +У `cm_invite_pending` три внешних держателя указателя в разных подсистемах: + + uasync heap (overall_timer) → node->arg = inv + NCD (node_conn_direct) → arg = inv в cm_invite_ncd_callback + STCP (stcp_link) → arg = inv в cm_tcp_ready_cb + +Когда cm_invite_cancel делает u_free(inv), коллбэки от уже запланированных +событий приходят после освобождения → пишут в переиспользованную память → +corruption (float-значения в полях mgr/node_id). + +## Решение + +1. Поле uint8_t cancelled:1 в struct cm_invite_pending + (рядом с tcp_ready:1 в conn_mgr_priv.h) + +2. cm_invite_cancel: + - inv->cancelled = 1 первой же строкой + - отмена таймера, закрытие ncd/tcp handles + - удаление из invite_list + - uasync_post(ua, cm_invite_free_deferred, inv) вместо немедленного u_free + +3. Все коллбэки — guard первой строкой: + void cm_invite_ncd_callback(..., void* arg) { + struct cm_invite_pending* inv = arg; + if (inv->cancelled) return; + ... + } + +4. cm_invite_fail: + - if (inv->cancelled) return + - сохранить cb/handle/node_id/group_id ДО вызова cm_invite_cancel + - cm_invite_cancel(inv) (cancelled=1, deferred free) + - cb(handle, ...) (коллбэк после cancel, inv ещё жив) + - u_free(handle) + +## Почему корректно + + cm_invite_cancel: + inv->cancelled = 1 ← блокировка коллбэков (однопоточно, немедленно) + uasync_cancel_timeout ← таймер удалён из heap + node_conn_direct_close ← handle закрыт + uasync_post(free, inv) ← освобождение в конец очереди событий + + [уже запланированные коллбэки]: + if (inv->cancelled) return ← выход без действий + + [deferred free]: + u_free(inv) ← все коллбэки уже отработали + +## Краевые случаи + + conn_mgr_destroy → cm_invite_cancel: + Отмена таймера, закрытие handles, deferred free. Ни один коллбэк + не дёрнет освобождённую память. + + NCD/STCP коллбэк после cancel: + cancelled=1 → мгновенный return. + + Таймаут приглашения после cancel: + cancelled=1 → return, не вызывает cm_invite_fail повторно. + + uasync_destroy до deferred free: + inv утекает при shutdown (приемлемо). diff --git a/src/routing_layer/conn_mgr_core.c b/src/routing_layer/conn_mgr_core.c index 8b144835..4a58da91 100644 --- a/src/routing_layer/conn_mgr_core.c +++ b/src/routing_layer/conn_mgr_core.c @@ -29,6 +29,7 @@ void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry); static void cm_handle_invite_info_req(struct CONN_MGR* mgr, struct ETCP_CONN* conn, const uint8_t* data, size_t len); static void cm_handle_invite_info_resp(struct CONN_MGR* mgr, struct ETCP_CONN* conn, const uint8_t* data, size_t len); +static void cm_invite_cancel(struct cm_invite_pending* inv); /* ═══════ утилиты ═══════ */ @@ -299,6 +300,7 @@ void cm_ncd_callback(struct NODE_CONN_DIRECT* ncd_h, enum ncd_event ncd_ev, void * TIMEOUT/DOWN (фейлим invite с CONN_EVENT_TIMEOUT). */ void cm_invite_ncd_callback(struct NODE_CONN_DIRECT* ncd_h, enum ncd_event ncd_ev, void* arg) { struct cm_invite_pending* inv = (struct cm_invite_pending*)arg; + if (inv->cancelled) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "invite: NCD callback IGNORED — already cancelled node=0x%016llx", (unsigned long long)inv->node_id); return; } switch (ncd_ev) { case NCD_EVENT_UP: if (inv->state != CM_INVITE_CONNECTING) return; @@ -348,20 +350,31 @@ struct CONN_MGR* conn_mgr_init(struct TOPO_GROUP* group) { void conn_mgr_destroy(struct CONN_MGR* mgr) { if (!mgr) return; + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 0 enter mgr=%p entries=%p invite=%p exch=%p rev=%p", + mgr, mgr->entries, mgr->invite_list, mgr->exchange_pending, mgr->reverse_pending); if (mgr->bg_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->bg_ping_timer); mgr->bg_ping_timer = NULL; } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 1 bg_ping done"); if (mgr->candidate_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->candidate_ping_timer); mgr->candidate_ping_timer = NULL; } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 2 cand_ping done"); { size_t ec = queue_entry_count(mgr->entries); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 3 entries count=%zu", ec); struct ll_entry* e = mgr->entries->head; - while (e) { struct ll_entry* next = e->next; struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)e; cm_entry_cleanup(entry); e = next; } - while (mgr->invite_list) cm_invite_fail(mgr->invite_list); + int clean = 0; + while (e) { struct ll_entry* next = e->next; struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)e; cm_entry_cleanup(entry); clean++; e = next; } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 3 entries cleaned (%d)", clean); + int inv_clean = 0; + while (mgr->invite_list) { cm_invite_cancel(mgr->invite_list); inv_clean++; } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 4 invite_list cleaned (%d)", inv_clean); while (mgr->exchange_pending) { struct cm_exchange_pending* ep = mgr->exchange_pending; mgr->exchange_pending = ep->next; if (ep->timer) uasync_cancel_timeout(mgr->instance->ua, ep->timer); u_free(ep); } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 5 exchange_pending done"); while (mgr->reverse_pending) { struct cm_reverse_pending* rp = mgr->reverse_pending; mgr->reverse_pending = rp->next; u_free(rp); } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 6 reverse_pending done"); queue_free(mgr->entries); mgr->initialized = 0; - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: destroyed, entries=%zu", ec); } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 7 queue_free done (was %zu)", ec); } u_free(mgr); } @@ -807,49 +820,102 @@ void cm_handle_invite_info_resp(struct CONN_MGR* mgr, struct ETCP_CONN* conn, co } inv->ncd_handle = NULL; } - cm_invite_cleanup(inv); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_SUCCESS] ncd transferred, temp_nq kept in group — calling cancel"); + inv->temp_nq = NULL; + cm_invite_cancel(inv); } void cm_invite_overall_timeout(void* arg) { - struct cm_invite_pending* inv = (struct cm_invite_pending*)arg; inv->overall_timer = NULL; + struct cm_invite_pending* inv = (struct cm_invite_pending*)arg; + if (inv->cancelled) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "invite: overall timer IGNORED — already cancelled node=0x%016llx", (unsigned long long)inv->node_id); return; } + inv->overall_timer = NULL; DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "invite: overall TIMEOUT to 0x%016llx — no response within %ds", (unsigned long long)inv->node_id, CM_INVITE_DEFAULT_TIMEOUT_MS/1000); cm_invite_fail(inv); } -/* Очистка invite при ошибке/таймауте: отменяет overall_timer, закрывает TCP-link - * и NCD-handle, удаляет временный TOPO_GROUP_NODE, доставляет CONN_EVENT_TIMEOUT - * через handle и освобождает его. Вызывается также из conn_mgr_destroy. */ -void cm_invite_fail(struct cm_invite_pending* inv) { +/* Cleanup invite resources WITHOUT calling external callbacks. + * Safe to call during teardown/destroy. */ +static void cm_invite_free_deferred(void* arg) { + struct cm_invite_pending* inv = (struct cm_invite_pending*)arg; + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_FREE] freeing inv=%p node=0x%016llx", inv, (unsigned long long)inv->node_id); + u_free(inv); +} + +static void cm_invite_cancel(struct cm_invite_pending* inv) { if (!inv) return; - DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "cm_invite_fail: node=0x%016llx state=%d ncd=%p tcp=%p", - (unsigned long long)inv->node_id, (int)inv->state, - (void*)inv->ncd_handle, (void*)inv->tcp_link); - if (inv->overall_timer) { uasync_cancel_timeout(inv->mgr->instance->ua, inv->overall_timer); inv->overall_timer = NULL; } + inv->cancelled = 1; + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] 0 inv=%p node=0x%016llx mgr=%p", + inv, (unsigned long long)inv->node_id, inv->mgr); + if (!inv->mgr) { DEBUG_ERROR(DEBUG_CATEGORY_SYS, "[INV_CANCEL] inv->mgr is NULL!"); goto inv_cancel_free; } + if (!inv->mgr->instance) { DEBUG_ERROR(DEBUG_CATEGORY_SYS, "[INV_CANCEL] inv->mgr->instance is NULL!"); goto inv_cancel_free; } + if (inv->overall_timer) { + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] 1 overall_timer=%p — cancelling", inv->overall_timer); + uasync_cancel_timeout(inv->mgr->instance->ua, inv->overall_timer); inv->overall_timer = NULL; + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] 1 overall_timer done"); + } else + { DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] 1 overall_timer skip (NULL)"); } if (inv->tcp_link) { + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] 2 tcp_link=%p — closing", inv->tcp_link); struct ll_entry* e = queue_find_data_by_index(inv->mgr->instance->tcp_connections, (const uint8_t*)&inv->node_id); if (e) queue_remove_data(inv->mgr->instance->tcp_connections, e); stcp_link_close(inv->tcp_link); inv->tcp_link = NULL; - } - if (inv->ncd_handle) { node_conn_direct_close(inv->ncd_handle); inv->ncd_handle = NULL; } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] 2 tcp_link done"); + } else + { DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] 2 tcp_link skip (NULL)"); } + if (inv->ncd_handle) { + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] 3 ncd_handle=%p — closing", inv->ncd_handle); + node_conn_direct_close(inv->ncd_handle); inv->ncd_handle = NULL; + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] 3 ncd done"); + } else + { DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] 3 ncd skip (NULL)"); } if (inv->temp_nq) { + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] 4 temp_nq=%p group=%p — removing from nodes", + inv->temp_nq, inv->mgr->group); queue_remove_data(inv->mgr->group->nodes, &inv->temp_nq->ll); topo_nodeq_free_group_fields(inv->mgr->instance->topo_groups, inv->temp_nq); queue_entry_free(&inv->temp_nq->ll); inv->temp_nq = NULL; - } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] 4 temp_nq done"); + } else + { DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] 4 temp_nq skip (NULL)"); } { struct cm_invite_pending** pp = &inv->mgr->invite_list; while (*pp) { if (*pp == inv) { *pp = inv->next; break; } pp = &(*pp)->next; } } - if (inv->cb) inv->cb(inv->handle, inv->node_id, inv->mgr->group->group_id, CONN_EVENT_TIMEOUT, inv->cb_arg); - if (inv->handle) u_free(inv->handle); - cm_invite_cleanup(inv); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] 5 unlink done"); +inv_cancel_free: + if (inv->mgr && inv->mgr->instance && inv->mgr->instance->ua) + { DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INV_CANCEL] posting deferred free inv=%p", inv); + uasync_post(inv->mgr->instance->ua, cm_invite_free_deferred, inv); } + else + { DEBUG_ERROR(DEBUG_CATEGORY_SYS, "[INV_CANCEL] no ua — immediate free inv=%p", inv); + u_free(inv); } +} + +/* Cancel invite AND notify caller. For normal operation (timeout, error, invite failure). + * Not safe during teardown — use cm_invite_cancel for destroy paths. */ +void cm_invite_fail(struct cm_invite_pending* inv) { + if (!inv) return; + if (inv->cancelled) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "cm_invite_fail: IGNORED — already cancelled node=0x%016llx", (unsigned long long)inv->node_id); return; } + DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "cm_invite_fail: node=0x%016llx state=%d ncd=%p tcp=%p", + (unsigned long long)inv->node_id, (int)inv->state, + (void*)inv->ncd_handle, (void*)inv->tcp_link); + conn_mgr_cb_t cb = inv->cb; + struct CONN_MGR_HANDLE* handle = inv->handle; + uint64_t node_id = inv->node_id; + uint64_t group_id = inv->mgr->group->group_id; + void* cb_arg = inv->cb_arg; + cm_invite_cancel(inv); + if (cb) cb(handle, node_id, group_id, CONN_EVENT_TIMEOUT, cb_arg); + if (handle) u_free(handle); } void cm_invite_cleanup(struct cm_invite_pending* inv) { if (inv) u_free(inv); } static void cm_tcp_ready_cb(struct stcp_link* link, void* arg) { struct cm_invite_pending* inv = (struct cm_invite_pending*)arg; - if (!inv || inv->tcp_ready) return; inv->tcp_ready = 1; + if (!inv) return; + if (inv->cancelled) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "invite: TCP ready IGNORED — already cancelled node=0x%016llx", (unsigned long long)inv->node_id); return; } + if (inv->tcp_ready) return; inv->tcp_ready = 1; struct ETCP_CONN* etcp = stcp_link_get_etcp_conn(link); etcp->peer_node_id = inv->node_id; struct ll_entry* qe = queue_entry_new(sizeof(struct tcp_conn_entry)); if (!qe) { cm_invite_fail(inv); return; } @@ -870,6 +936,7 @@ static void cm_tcp_ready_cb(struct stcp_link* link, void* arg) { static void cm_tcp_close_cb(struct stcp_link* link, int err, void* arg) { struct cm_invite_pending* inv = (struct cm_invite_pending*)arg; if (!inv) return; + if (inv->cancelled) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "invite: TCP close IGNORED — already cancelled node=0x%016llx err=%d", (unsigned long long)inv->node_id, err); return; } DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "invite: TCP link CLOSED to 0x%016llx err=%d", (unsigned long long)inv->node_id, err); inv->tcp_link = NULL; } diff --git a/src/routing_layer/conn_mgr_priv.h b/src/routing_layer/conn_mgr_priv.h index a7d9e914..2be1ab47 100644 --- a/src/routing_layer/conn_mgr_priv.h +++ b/src/routing_layer/conn_mgr_priv.h @@ -110,7 +110,7 @@ struct cm_invite_pending { struct cm_invite_pending* next; struct CONN_MGR* mgr; uint64_t node_id; uint8_t state; struct NODE_CONN_DIRECT* ncd_handle; struct stcp_link* tcp_link; - uint8_t tcp_ready:1; void* overall_timer; + uint8_t tcp_ready:1, cancelled:1; void* overall_timer; struct CONN_MGR_HANDLE* handle; struct TOPO_GROUP_NODE* temp_nq; conn_mgr_cb_t cb; void* cb_arg; }; diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index a34acc14..caf4a652 100644 --- a/src/routing_layer/topo_group.c +++ b/src/routing_layer/topo_group.c @@ -240,18 +240,24 @@ static struct TOPO_GROUP* topo_group_create(struct UTUN_INSTANCE* instance, uint static void topo_group_destroy(struct TOPO_GROUP* group) { if (!group) return; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "group_id=%016llx", (unsigned long long)group->group_id); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 0 enter grp=%p id=0x%016llx rec_list=%p connect=%p", + group, (unsigned long long)group->group_id, group->recovery_list, group->connect); topo_recovery_cancel_all(group); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 1 recovery_cancel done grp=%p", group); topo_group_connect_destroy(group); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 2 connect_destroy done"); struct ll_entry* e; while ((e = queue_data_get(group->senders_list)) != NULL) queue_entry_free(e); queue_free(group->senders_list); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 3 senders_list done"); if (group->conn_mgr) { conn_mgr_destroy(group->conn_mgr); group->conn_mgr = NULL; } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 4 conn_mgr done"); if (group->nodes) queue_free(group->nodes); if (group->local_node) { topo_nodeq_free_group_fields(group->instance->topo_groups, group->local_node); u_free(group->local_node); } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 5 done"); } struct TOPO_GROUPS* topo_groups_init(struct UTUN_INSTANCE* instance) { @@ -323,27 +329,38 @@ struct TOPO_GROUPS* topo_groups_init(struct UTUN_INSTANCE* instance) { void topo_groups_destroy(struct UTUN_INSTANCE* instance) { if (!instance) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "instance is NULL"); return; } if (!instance->topo_groups) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "topo_groups is NULL"); return; } - DEBUG_INFO(DEBUG_CATEGORY_BGP, "node_id=%016llx", (unsigned long long)instance->node_id); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "node_id=%016llx tgs=%p", (unsigned long long)instance->node_id, instance->topo_groups); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[TOPO_DESTROY] 0 enter"); etcp_unbind(instance, ETCP_ID_TOPO_ENTRY); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[TOPO_DESTROY] 1 etcp_unbind done"); etcp_remove_conn_status_cbk(instance, topo_group_conn_status, instance->topo_groups); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[TOPO_DESTROY] 2 conn_cbk_remove done"); etcp_remove_socket_cbk(instance, topo_node_on_socket_changed, NULL); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[TOPO_DESTROY] 3 socket_cbk_remove done"); route_connectivity_cancel_all(instance); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[TOPO_DESTROY] 4 connectivity_cancel done"); struct TOPO_GROUPS* g = instance->topo_groups; struct ll_entry* ge; - while ((ge = queue_data_get(g->group_list)) != NULL) { struct TOPO_GROUP* grp = (struct TOPO_GROUP*)ge; topo_group_destroy(grp); queue_entry_free(ge); } + int gcount = 0; + while ((ge = queue_data_get(g->group_list)) != NULL) { struct TOPO_GROUP* grp = (struct TOPO_GROUP*)ge; topo_group_destroy(grp); queue_entry_free(ge); gcount++; } queue_free(g->group_list); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[TOPO_DESTROY] 5 groups destroyed (%d)", gcount); if (g->node_registry) { + int ncount = 0; struct ll_entry* re; while ((re = queue_data_get(g->node_registry)) != NULL) { struct TOPO_NODE* node; memcpy(&node, re->data + 8, sizeof(node)); topo_node_registry_unref(g, node->node_id); - queue_entry_free(re); + queue_entry_free(re); ncount++; } queue_free(g->node_registry); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[TOPO_DESTROY] 6 node_registry done (%d)", ncount); + } else { + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[TOPO_DESTROY] 6 node_registry skip (NULL)"); } memory_pool_destroy(g->v4_sock_meta_pool); @@ -352,10 +369,13 @@ void topo_groups_destroy(struct UTUN_INSTANCE* instance) { memory_pool_destroy(g->v6_addr_pool); memory_pool_destroy(g->v4_subnet_pool); memory_pool_destroy(g->v6_subnet_pool); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[TOPO_DESTROY] 7 pools destroyed"); if (instance->topo_sqlite_db) { sqlite3_close(instance->topo_sqlite_db); instance->topo_sqlite_db = NULL; } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[TOPO_DESTROY] 8 sqlite closed"); - u_free(g); instance->topo_groups = NULL; + u_free(g); instance->topo_groups = NULL; instance->conn_mgr = NULL; + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[TOPO_DESTROY] 9 done"); } struct TOPO_GROUP* topo_groups_get_default(struct TOPO_GROUPS* g) { diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index cc9b9a85..2355dc27 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -755,7 +755,7 @@ void etcp_socket_remove(struct ETCP_SOCKET* conn) { DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP] Removing socket %p, socket_id=%p", conn, conn->socket_id); // Remove from uasync if registered - if (conn->socket_id) { + if (conn->socket_id && conn->instance) { DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP] Removing socket from uasync, instance=%p, ua=%p", conn->instance, conn->instance->ua); uasync_remove_socket_t(conn->instance->ua, conn->fd); conn->socket_id = NULL; diff --git a/src/utun_instance.c b/src/utun_instance.c index 703f9269..6460ec07 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -403,187 +403,172 @@ struct UTUN_INSTANCE* utun_instance_create_from_str(struct UASYNC* ua, const cha void utun_instance_destroy(struct UTUN_INSTANCE *instance) { if (!instance) return; - DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Starting cleanup for instance %p", instance); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Starting cleanup for instance %p node=0x%016llx ua=%p", + instance, (unsigned long long)instance->node_id, instance->ua); - // Диагностика ресурсов ДО cleanup + /* Phase A: diagnose */ utun_instance_diagnose_leaks(instance, "BEFORE_CLEANUP"); if (instance->ua) uasync_print_resources(instance->ua, "INSTANCE_DESTROY_BEFORE"); - // Stop running if not already instance->running = 0; + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] A done — diagnose complete"); - // Cancel NTP timer + /* Phase B: NTP */ ntp_time_destroy(instance); ntp_node_time_destroy(instance); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] B done — NTP"); - // Shutdown control server first + /* Phase C: control server */ if (instance->control_srv) { - DEBUG_INFO(DEBUG_CATEGORY_CONTROL, "Shutting down control server"); control_server_shutdown(instance->control_srv); u_free(instance->control_srv); instance->control_srv = NULL; } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] C done — control_srv"); - // Close config-based connection handles (sends CLOSE before sockets die) - struct CONFIG_CONN_HANDLE* ch = instance->config_conn_handles; - while (ch) { - struct CONFIG_CONN_HANDLE* next = ch->next; - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[INSTANCE_DESTROY] closing config handle for node=0x%016llx", (unsigned long long)ch->node_id); - node_conn_direct_close(ch->handle); - u_free(ch); - ch = next; + /* Phase D: config conn handles */ + { + struct CONFIG_CONN_HANDLE* ch = instance->config_conn_handles; + int ch_count = 0; + while (ch) { + struct CONFIG_CONN_HANDLE* next = ch->next; + node_conn_direct_close(ch->handle); + u_free(ch); ch_count++; + ch = next; + } + if (ch_count) DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] D: closed %d config handles", ch_count); } instance->config_conn_handles = NULL; + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] D done — config_conn"); + /* Phase E: socket_monitor + auto_socket */ socket_monitor_destroy(instance); auto_socket_destroy(instance); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] E done — socket_monitor+auto_socket"); - // Cleanup BGP module BEFORE sockets (needs live conn_mgr for recovery cleanup) + /* Phase F: BGP */ if (instance->topo_groups) { - DEBUG_INFO(DEBUG_CATEGORY_BGP, "Destroying BGP module"); etcp_router_unbind(instance, ETCP_RT_ID_CONN_MGR); etcp_unbind(instance, ETCP_RT_ID_CONN_MGR); topo_groups_destroy(instance); } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] F done — BGP"); - // Cleanup conn_mgr BEFORE sockets (unbinds from etcp_router, closes NCD handles while conns alive) - - // Cleanup ETCP sockets and connections FIRST (before destroying uasync) - DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Cleaning up ETCP sockets and connections"); - struct ETCP_SOCKET* sock = instance->etcp_sockets; - while (sock) { - struct ETCP_SOCKET* next = sock->next; - DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Removing socket %p, fd=%d", sock, sock->fd); - etcp_socket_remove(sock); // Полный cleanup сокета - sock = next; + /* Phase G: ETCP sockets */ + { + struct ETCP_SOCKET* sock = instance->etcp_sockets; + int sc = 0; + while (sock) { + struct ETCP_SOCKET* next = sock->next; + if (sock->instance && sock->fd > 0) DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] G: removing socket %p fd=%d inst=%p", sock, sock->fd, sock->instance); + else if (sock->fd > 0) { DEBUG_ERROR(DEBUG_CATEGORY_SYS, "[DESTROY] G: socket %p fd=%d has NULL instance! skip remove", sock, sock->fd); sock = next; continue; } + etcp_socket_remove(sock); + sc++; + sock = next; + } + if (sc) DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] G: removed %d ETCP sockets", sc); } instance->etcp_sockets = NULL; - DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] ETCP sockets cleanup complete"); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] G done — ETCP sockets"); - // Cleanup chat before db_sync (chat releases db_sync instances) + /* Phase H: chat */ chat_sync_destroy(instance); chat_core_destroy(instance); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] H done — chat"); - // Cleanup db_sync before connections cleanup (db_sync_destroy iterates connections) + /* Phase I: db_sync */ db_sync_destroy(instance); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] I done — db_sync"); - // Cleanup ETCP connections (phase 1 detach + deferred phase 2 via call_soon) + /* Phase J: ETCP connections */ { struct ll_entry* entry = instance->connections->head; + int cc = 0; while (entry) { struct ll_entry* next = entry->next; struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; - if (ce && ce->conn) etcp_connection_close(ce->conn); + if (ce && ce->conn) { etcp_connection_close(ce->conn); cc++; } entry = next; } + if (cc) DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] J: closed %d ETCP connections", cc); } queue_free(instance->connections); instance->connections = NULL; - DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] ETCP connections cleanup complete"); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] J done — connections"); - // Cleanup TCP connections (stcp_links) + /* Phase K: TCP connections + sockets */ if (instance->tcp_connections) { struct ll_entry* entry = instance->tcp_connections->head; + int tc = 0; while (entry) { struct ll_entry* next = entry->next; struct tcp_conn_entry* te = (struct tcp_conn_entry*)entry->data; - if (te && te->link) stcp_link_close(te->link); + if (te && te->link) { stcp_link_close(te->link); tc++; } entry = next; } + if (tc) DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] K: closed %d TCP links", tc); } queue_free(instance->tcp_connections); instance->tcp_connections = NULL; - DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] TCP connections cleanup complete"); - - // Cleanup TCP sockets while (instance->tcp_sockets) tcp_socket_remove(instance->tcp_sockets); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] K done — TCP"); - // Wait for all deferred callbacks (etcp_connection_free_deferred) - // to finish before destroying pools - while (instance->ua && instance->ua->immediate_queue_head) - uasync_poll(instance->ua, 0); - - struct PING_CONTEXT* p = instance->pending_pings; - while (p) { - struct PING_CONTEXT* next = p->next; - if (p->timeout_timer) uasync_cancel_timeout(instance->ua, p->timeout_timer); - u_free(p); - p = next; + /* Phase L: deferred callbacks + pings */ + { + int deferred_loops = 0; + while (instance->ua && instance->ua->immediate_queue_head) { + uasync_poll(instance->ua, 0); + if (++deferred_loops > 1000) { DEBUG_ERROR(DEBUG_CATEGORY_SYS, "[DESTROY] L: deferred callback loop stuck"); break; } + } + if (deferred_loops) DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] L: drained %d deferred callbacks", deferred_loops); } - instance->pending_pings = NULL; - - // Cleanup TUN - if (instance->tun) { - DEBUG_INFO(DEBUG_CATEGORY_TUN, "Closing TUN interface: %s", instance->tun->ifname); - - // Delete system routes added at startup - if (instance->tun->ifindex && instance->route_subnets) { - int deleted = tun_route_del_all(instance->tun->ifindex, instance->tun->ifname, instance->route_subnets); - DEBUG_INFO(DEBUG_CATEGORY_TUN, "Deleted %d system routes for TUN interface", deleted); + { + struct PING_CONTEXT* p = instance->pending_pings; + int pc = 0; + while (p) { + struct PING_CONTEXT* next = p->next; + if (p->timeout_timer) uasync_cancel_timeout(instance->ua, p->timeout_timer); + u_free(p); pc++; + p = next; } - + instance->pending_pings = NULL; + if (pc) DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] L: freed %d pending pings", pc); + } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] L done — deferred+pings"); + + /* Phase M: TUN + proxy + routing + NAT */ + if (instance->tun) { + if (instance->tun->ifindex && instance->route_subnets) + tun_route_del_all(instance->tun->ifindex, instance->tun->ifname, instance->route_subnets); tun_close(instance->tun); instance->tun = NULL; } - - // Cleanup TCP proxy client module #ifndef _WIN32 - if (instance->tcp_proxy_client) { - DEBUG_INFO(DEBUG_CATEGORY_TUN, "Destroying TCP proxy client module"); - tcp_proxy_client_destroy(instance->tcp_proxy_client); - instance->tcp_proxy_client = NULL; - } - - // Cleanup TCP proxy server + if (instance->tcp_proxy_client) { tcp_proxy_client_destroy(instance->tcp_proxy_client); instance->tcp_proxy_client = NULL; } tcp_proxy_server_destroy(instance); #endif - - // Cleanup routing module (unbinds from etcp_router before etcp_router_destroy) routing_destroy(instance); + if (instance->md.initialized) media_delivery_destroy(instance); + media_async_destroy(instance->media_async); instance->media_async = NULL; + if (instance->nat_tr.initialized) nat_transport_destroy(instance); + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] M done — TUN/proxy/routing/NAT"); - // Cleanup media delivery (unbinds from etcp_router before etcp_router_destroy) - if (instance->md.initialized) { - media_delivery_destroy(instance); - } - - // Cleanup media async engine - media_async_destroy(instance->media_async); - instance->media_async = NULL; - - // Cleanup NAT (unbinds from etcp_router before etcp_router_destroy) - if (instance->nat_tr.initialized) { - nat_transport_destroy(instance); - } - - // Cleanup etcp_router + /* Phase N: etcp_router + nat_det + firewall + TCP servers */ etcp_router_destroy(instance); - - // Cleanup NAT detection - if (instance->nat_det) { - nat_detection_destroy(instance->nat_det); - instance->nat_det = NULL; - } - - // Cleanup firewall + if (instance->nat_det) { nat_detection_destroy(instance->nat_det); instance->nat_det = NULL; } fw_free(&instance->fw); - - // Cleanup TCP servers stcp_server_list_destroy_all(instance); - - // Cleanup networks queue + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] N done — router/nat_det/fw/stcp_srv"); + + /* Phase O: networks queue */ if (instance->networks) { struct ll_entry *entry; - while ((entry = queue_data_get(instance->networks)) != NULL) { - queue_entry_free(entry); - } + while ((entry = queue_data_get(instance->networks)) != NULL) queue_entry_free(entry); queue_free(instance->networks); instance->networks = NULL; } - - // Cleanup config - if (instance->config) { - DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Freeing configuration"); - free_config(instance->config); - instance->config = NULL; - } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] O done — networks"); + + /* Phase P: config + pools */ + if (instance->config) { free_config(instance->config); instance->config = NULL; } // Cleanup packet pool (ensure no leak if stop wasn't called) if (instance->pkt_pool) { @@ -755,22 +740,43 @@ void utun_instance_diagnose_leaks(struct UTUN_INSTANCE *instance, const char *ph int etcp_links_count; } report = {0}; - // Подсчёт ETCP сокетов + /* Check etcp_sockets linked list for corruption */ struct ETCP_SOCKET *sock = instance->etcp_sockets; + int sock_visited = 0; while (sock) { + /* Validate pointer looks like heap (not node_id or stack garbage) */ + if ((uintptr_t)sock < 0x1000) { + DEBUG_ERROR(DEBUG_CATEGORY_SYS, "[DIAGNOSE] CORRUPT socket ptr=%p (low addr) at index %d", sock, sock_visited); + break; + } + if (sock_visited++ > 10000) { + DEBUG_ERROR(DEBUG_CATEGORY_SYS, "[DIAGNOSE] socket list cycle or too long (%d)", sock_visited); + break; + } report.etcp_sockets_count++; - report.etcp_links_count += queue_entry_count(sock->links_queue); + if (sock->links_queue) report.etcp_links_count += queue_entry_count(sock->links_queue); sock = sock->next; } - // Подсчёт ETCP соединений - { + /* Check connections linked list for corruption */ + if (instance->connections) { struct ll_entry* entry = instance->connections->head; + int conn_visited = 0; while (entry) { + if ((uintptr_t)entry < 0x1000) { + DEBUG_ERROR(DEBUG_CATEGORY_SYS, "[DIAGNOSE] CORRUPT conn entry ptr=%p at index %d", entry, conn_visited); + break; + } + if (conn_visited++ > 10000) { + DEBUG_ERROR(DEBUG_CATEGORY_SYS, "[DIAGNOSE] conn list cycle or too long (%d)", conn_visited); + break; + } struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; - report.etcp_connections_count++; - struct ETCP_LINK *link = ce->conn->links; - while (link) { report.etcp_links_count++; link = link->next; } + if (ce && ce->conn) { + report.etcp_connections_count++; + struct ETCP_LINK *link = ce->conn->links; + while (link) { report.etcp_links_count++; link = link->next; } + } entry = entry->next; } } diff --git a/tools/chatgui-android/app/src/main/java/com/utun/chat/viewmodel/ChatViewModel.kt b/tools/chatgui-android/app/src/main/java/com/utun/chat/viewmodel/ChatViewModel.kt index e2fd22e2..1a857f14 100644 --- a/tools/chatgui-android/app/src/main/java/com/utun/chat/viewmodel/ChatViewModel.kt +++ b/tools/chatgui-android/app/src/main/java/com/utun/chat/viewmodel/ChatViewModel.kt @@ -485,6 +485,7 @@ class ChatViewModel : ViewModel() { fun clearState() { repo?.close() repo = null + dbReady = false _channels.value = emptyList() _messages.value = emptyList() _currentChannel.value = null diff --git a/tools/chatgui-android/jni_bridge/android_jni_bridge.c b/tools/chatgui-android/jni_bridge/android_jni_bridge.c index 509dab17..22941a85 100644 --- a/tools/chatgui-android/jni_bridge/android_jni_bridge.c +++ b/tools/chatgui-android/jni_bridge/android_jni_bridge.c @@ -157,7 +157,6 @@ void utun_bridge_destroy(void) { voice_recorder_deinit(); instance_lite_stop(); g_on_log = NULL; - g_on_event = NULL; g_cfg_string = NULL; g_cfg_int64 = NULL; g_cfg_int = NULL; @@ -1028,6 +1027,14 @@ void utun_bridge_restart(const char* config_text) { } #endif instance_lite_restart(config_text); + + /* Wait for instance to be ready before returning to Kotlin. + Normal restart (live): is_running() true immediately → no delay. + Dead restart (after stop): polls until g_running=1 → blocks until init done. */ + for (int i = 0; i < 500 && !instance_lite_is_running(); i++) + usleep(10000); + if (!instance_lite_is_running()) + bridge_log(BLEV_ERROR, "bridge restart: instance not ready after 5s"); } void utun_bridge_ping(void) { diff --git a/tools/chatgui-android/libutun_lite/instance_lite.c b/tools/chatgui-android/libutun_lite/instance_lite.c index 49c7c854..0c8dabc1 100644 --- a/tools/chatgui-android/libutun_lite/instance_lite.c +++ b/tools/chatgui-android/libutun_lite/instance_lite.c @@ -31,6 +31,8 @@ #include #include #include +#include +#include #include #ifdef __ANDROID__ @@ -44,19 +46,33 @@ #define IL_LOGE(fmt, ...) fprintf(stderr, fmt "\n", ##__VA_ARGS__) #endif -/* ── Static state ── */ +/* ── Static state ── + * + * Thread model: one Kotlin thread, one C worker thread. + * Each variable has exactly ONE writer — no data races by construction. + * + * Kotlin-writes (signals/commands to C): + * g_stop, g_do_restart, g_restart_config, g_event_handler, g_db_path + * C-writes (state to Kotlin): + * g_ua, g_inst, g_running, g_thread_exited, g_thread_running + */ -static struct UASYNC* g_ua = NULL; -static struct UTUN_INSTANCE* g_inst = NULL; +static struct UASYNC* g_ua = NULL; /* C-write, Kotlin-read via __atomic */ +static struct UTUN_INSTANCE* g_inst = NULL; /* C-write, Kotlin-read via __atomic */ static pthread_t g_thread; -static volatile int g_running = 0; -static volatile int g_stop = 0; -static volatile int g_do_restart = 0; -static char* g_restart_config = NULL; -static char g_db_path[512]; -static instance_lite_event_fn g_event_handler = NULL; +static volatile int g_stop = 0; /* Kotlin-write, C-read (poll loop) */ +static volatile int g_do_restart = 0; /* Kotlin-write, C atomic-xchg→0 */ +static volatile char* g_restart_config = NULL; /* Kotlin-write, C atomic-xchg→NULL (ownership transfer) */ +static volatile int g_running = 0; /* C-write, Kotlin-read */ +static volatile int g_thread_running = 0; /* C-write + Kotlin CAS-guard for start() */ +static volatile int g_thread_exited = 0; /* C-write (under mutex), Kotlin-read */ +static volatile int g_generation = 0; /* Kotlin-write, C-read */ +static char g_db_path[512]; /* Kotlin-write (start), C-read (init) */ +static instance_lite_event_fn g_event_handler = NULL; /* Kotlin-write, C-read+call */ static char g_generated_pub[65]; static char g_generated_priv[65]; +static pthread_mutex_t g_stop_mutex = PTHREAD_MUTEX_INITIALIZER; +static pthread_cond_t g_stop_cond = PTHREAD_COND_INITIALIZER; /* ── Generate X25519 key pair, hex-encode to out buffers ── */ @@ -200,15 +216,16 @@ static void heartbeat_cb(void* arg) { static void* instance_thread(void* arg) { char* config_text = (char*)arg; + int my_gen = g_generation; struct utun_config* config = parse_config_from_buf(config_text, strlen(config_text), "android"); u_free(config_text); - if (!config) { IL_LOGE("parse_config_from_buf failed"); return NULL; } + if (!config) { IL_LOGE("parse_config_from_buf failed"); __atomic_store_n(&g_thread_running, 0, __ATOMIC_RELEASE); pthread_detach(pthread_self()); return NULL; } install_crash_handlers(); g_ua = uasync_create(); - if (!g_ua) { IL_LOGE("uasync_create failed"); free_config(config); return NULL; } + if (!g_ua) { IL_LOGE("uasync_create failed"); free_config(config); __atomic_store_n(&g_thread_running, 0, __ATOMIC_RELEASE); pthread_detach(pthread_self()); return NULL; } if (config->global.log_udp_ip[0] && config->global.log_udp_port > 0) { udp_log_set_target(config->global.log_udp_ip, config->global.log_udp_port); @@ -225,6 +242,8 @@ static void* instance_thread(void* arg) { uasync_destroy(g_ua, 0); g_ua = NULL; u_report_unfreed_blocks(); + __atomic_store_n(&g_thread_running, 0, __ATOMIC_RELEASE); + pthread_detach(pthread_self()); return NULL; } @@ -257,6 +276,8 @@ static void* instance_thread(void* arg) { uasync_destroy(g_ua, 0); g_ua = NULL; u_report_unfreed_blocks(); + __atomic_store_n(&g_thread_running, 0, __ATOMIC_RELEASE); + pthread_detach(pthread_self()); return NULL; } @@ -282,19 +303,26 @@ static void* instance_thread(void* arg) { (const uint8_t*)config->global.my_public_key_hex, 64); } - g_running = 1; + __atomic_store_n(&g_running, 1, __ATOMIC_RELEASE); uasync_mark_running(g_ua); chat_event_post(CHAT_EVT_SERVICE_STARTED, NULL, 0); uasync_set_timeout(g_ua, 100000, NULL, heartbeat_cb, "hb"); - while (!g_stop) { + while (!__atomic_load_n(&g_stop, __ATOMIC_ACQUIRE)) { uasync_poll(g_ua, 100); } - while (g_do_restart && g_restart_config) { + while (__atomic_load_n(&g_do_restart, __ATOMIC_ACQUIRE)) { IL_LOGI("poll exit: restart requested, destroying old instance"); + __atomic_store_n(&g_do_restart, 0, __ATOMIC_RELEASE); + IL_LOGI("pre-destroy: g_inst=%p node_id=0x%016llx ua=%p running=%d sockets=%p conns=%p", + g_inst, g_inst ? (unsigned long long)g_inst->node_id : 0ULL, + g_ua, g_running, + g_inst ? g_inst->etcp_sockets : NULL, + g_inst ? g_inst->connections : NULL); uasync_mark_stopped(g_ua); + __atomic_store_n(&g_running, 0, __ATOMIC_RELEASE); utun_instance_destroy(g_inst); g_inst = NULL; uasync_print_resources(g_ua, "AFTER_DESTROY"); @@ -305,10 +333,9 @@ static void* instance_thread(void* arg) { chat_event_set_handler(chat_event_forward); - char* cfg = g_restart_config; - g_restart_config = NULL; - g_do_restart = 0; - g_stop = 0; + char* cfg = (char*)__atomic_exchange_n(&g_restart_config, NULL, __ATOMIC_ACQUIRE); + if (!cfg) { IL_LOGE("poll exit: restart_config is NULL"); break; } + __atomic_store_n(&g_stop, 0, __ATOMIC_RELEASE); IL_LOGI("poll exit: creating new instance from restart config"); struct utun_config* config = parse_config_from_buf(cfg, strlen(cfg), "android"); @@ -338,30 +365,40 @@ static void* instance_thread(void* arg) { chat_event_post(CHAT_EVT_KEYS_GENERATED, (const uint8_t*)config->global.my_public_key_hex, 64); uasync_set_timeout(g_ua, 100000, NULL, heartbeat_cb, "hb"); - g_running = 1; + __atomic_store_n(&g_running, 1, __ATOMIC_RELEASE); uasync_mark_running(g_ua); chat_event_post(CHAT_EVT_SERVICE_STARTED, NULL, 0); IL_LOGI("poll exit: restart done, entering new poll loop"); - while (!g_stop) uasync_poll(g_ua, 100); + while (!__atomic_load_n(&g_stop, __ATOMIC_ACQUIRE)) uasync_poll(g_ua, 100); } IL_LOGI("poll exit: final cleanup"); - if (g_inst) { - utun_instance_destroy(g_inst); - g_inst = NULL; - } - if (g_ua) { - uasync_destroy(g_ua, 0); - g_ua = NULL; + if (my_gen == g_generation) { + __atomic_store_n(&g_running, 0, __ATOMIC_RELEASE); + if (g_inst) { + utun_instance_destroy(g_inst); + g_inst = NULL; + } + if (g_ua) { + uasync_destroy(g_ua, 0); + g_ua = NULL; + } + char* stale_cfg = (char*)__atomic_exchange_n(&g_restart_config, NULL, __ATOMIC_ACQUIRE); + u_free(stale_cfg); + u_report_unfreed_blocks(); + chat_event_post(CHAT_EVT_SERVICE_STOPPED, NULL, 0); + IL_LOGI("cleanup complete (gen=%d)", my_gen); + } else { + IL_LOGI("cleanup skipped — thread gen=%d but current gen=%d (stale thread)", my_gen, g_generation); } - if (g_restart_config) { u_free(g_restart_config); g_restart_config = NULL; } - u_report_unfreed_blocks(); - chat_event_post(CHAT_EVT_SERVICE_STOPPED, NULL, 0); - g_event_handler = NULL; - pthread_detach(pthread_self()); - IL_LOGI("cleanup complete"); + __atomic_store_n(&g_thread_running, 0, __ATOMIC_RELEASE); + pthread_mutex_lock(&g_stop_mutex); + g_thread_exited = 1; + pthread_cond_signal(&g_stop_cond); + pthread_mutex_unlock(&g_stop_mutex); + IL_LOGI("thread exit signalled"); return NULL; } @@ -369,7 +406,17 @@ static void* instance_thread(void* arg) { int instance_lite_start(const char* config_text) { if (!config_text) return -1; - if (g_inst) return 0; + + /* CAS: prevent double-start race between multiple Kotlin restart calls */ + int expected = 0; + if (!__atomic_compare_exchange_n(&g_thread_running, &expected, 1, 0, + __ATOMIC_ACQUIRE, __ATOMIC_RELAXED)) + return 0; + + if (__atomic_load_n(&g_inst, __ATOMIC_ACQUIRE)) { + __atomic_store_n(&g_thread_running, 0, __ATOMIC_RELEASE); + return 0; + } debug_config_init(); debug_set_level(DEBUG_LEVEL_INFO); @@ -380,12 +427,15 @@ int instance_lite_start(const char* config_text) { cfg_get_val(config_text, "db_path", g_db_path, sizeof(g_db_path)); char* config_copy = u_strdup(config_text); - if (!config_copy) return -1; + if (!config_copy) { __atomic_store_n(&g_thread_running, 0, __ATOMIC_RELEASE); return -1; } - g_stop = 0; + g_generation++; + g_thread_exited = 0; + __atomic_store_n(&g_stop, 0, __ATOMIC_RELEASE); int rc = pthread_create(&g_thread, NULL, instance_thread, config_copy); if (rc != 0) { u_free(config_copy); + __atomic_store_n(&g_thread_running, 0, __ATOMIC_RELEASE); IL_LOGE("pthread_create failed rc=%d", rc); return -1; } @@ -395,38 +445,62 @@ int instance_lite_start(const char* config_text) { } void instance_lite_stop(void) { - if (!g_inst) return; + struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)__atomic_load_n(&g_inst, __ATOMIC_ACQUIRE); + if (!inst) return; IL_LOGI("stopping..."); - uasync_mark_stopped(g_ua); - g_stop = 1; - if (g_ua) uasync_wakeup(g_ua); - for (int i = 0; i < 500 && g_inst != NULL; i++) usleep(10000); - if (g_inst) { - IL_LOGE("stop timeout, force cleanup"); - g_inst = NULL; - g_ua = NULL; + + __atomic_store_n(&g_stop, 1, __ATOMIC_RELEASE); + struct UASYNC* ua = (struct UASYNC*)__atomic_load_n(&g_ua, __ATOMIC_ACQUIRE); + if (ua) uasync_wakeup(ua); + + struct timespec ts; + clock_gettime(CLOCK_REALTIME, &ts); + ts.tv_sec += 2; + + int timedout = 0; + pthread_mutex_lock(&g_stop_mutex); + while (!g_thread_exited) { + int rc = pthread_cond_timedwait(&g_stop_cond, &g_stop_mutex, &ts); + if (rc == ETIMEDOUT) { + IL_LOGE("stop timeout (2s), detaching thread and bumping generation to prevent stale writes"); + pthread_detach(g_thread); + g_generation++; + timedout = 1; + break; + } } - g_running = 0; + pthread_mutex_unlock(&g_stop_mutex); + + if (!timedout) { + int join_rc = pthread_join(g_thread, NULL); + if (join_rc != 0) IL_LOGE("pthread_join failed rc=%d errno=%d", join_rc, errno); + IL_LOGI("thread joined"); + } + + /* C-thread cleanup already set g_inst=NULL, g_ua=NULL, g_running=0. + Reset thread sync state only. */ + g_thread_exited = 0; } void instance_lite_restart(const char* new_config_text) { if (!new_config_text) return; - instance_lite_event_fn saved_handler = g_event_handler; - if (!g_ua || !g_inst || !g_running) { - IL_LOGI("restart: instance dead, doing hard restart via stop+start (saved_handler=%p)", (void*)saved_handler); - if (g_inst) instance_lite_stop(); - g_event_handler = saved_handler; + struct UASYNC* ua = (struct UASYNC*)__atomic_load_n(&g_ua, __ATOMIC_ACQUIRE); + struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)__atomic_load_n(&g_inst, __ATOMIC_ACQUIRE); + int running = __atomic_load_n(&g_running, __ATOMIC_ACQUIRE); + if (!ua || !inst || !running) { + IL_LOGI("restart: instance dead, doing hard restart via stop+start"); + if (inst) instance_lite_stop(); instance_lite_start(new_config_text); return; } IL_LOGI("restart: signaling poll loop to exit"); char* copy = u_strdup(new_config_text); if (!copy) return; - if (g_restart_config) u_free(g_restart_config); - g_restart_config = copy; - g_do_restart = 1; - g_stop = 1; - uasync_wakeup(g_ua); + char* old_cfg = (char*)__atomic_exchange_n(&g_restart_config, copy, __ATOMIC_RELEASE); + u_free(old_cfg); + __atomic_store_n(&g_do_restart, 1, __ATOMIC_RELEASE); + __atomic_store_n(&g_stop, 1, __ATOMIC_RELEASE); + uasync_wakeup(ua); } /* ── Health check ping ── */ @@ -440,23 +514,28 @@ static void ping_trampoline(void* arg) { } void instance_lite_ping(void) { - if (!g_ua || !g_inst) return; + struct UASYNC* ua = (struct UASYNC*)__atomic_load_n(&g_ua, __ATOMIC_ACQUIRE); + struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)__atomic_load_n(&g_inst, __ATOMIC_ACQUIRE); + if (!ua || !inst) return; g_ping_id++; - uasync_post(g_ua, ping_trampoline, NULL); + uasync_post(ua, ping_trampoline, NULL); } int instance_lite_is_responsive(void) { - if (!g_ua || !g_inst || !g_running) return 0; + struct UASYNC* ua = (struct UASYNC*)__atomic_load_n(&g_ua, __ATOMIC_ACQUIRE); + struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)__atomic_load_n(&g_inst, __ATOMIC_ACQUIRE); + if (!ua || !inst || !__atomic_load_n(&g_running, __ATOMIC_ACQUIRE)) return 0; if (g_ping_id == 0) return 0; return (g_pong_id == g_ping_id) ? 1 : 0; } int instance_lite_is_running(void) { - return (g_inst != NULL && g_running) ? 1 : 0; + struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)__atomic_load_n(&g_inst, __ATOMIC_ACQUIRE); + return (inst != NULL && __atomic_load_n(&g_running, __ATOMIC_ACQUIRE)) ? 1 : 0; } struct UASYNC* instance_lite_get_uasync(void) { - return g_ua; + return (struct UASYNC*)__atomic_load_n(&g_ua, __ATOMIC_ACQUIRE); } static void regenerate_trampoline(void* arg) { @@ -478,12 +557,14 @@ static void regenerate_trampoline(void* arg) { } void instance_lite_regenerate_keys(void) { - if (!g_ua || !g_inst) return; - uasync_post(g_ua, regenerate_trampoline, NULL); + struct UASYNC* ua = (struct UASYNC*)__atomic_load_n(&g_ua, __ATOMIC_ACQUIRE); + struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)__atomic_load_n(&g_inst, __ATOMIC_ACQUIRE); + if (!ua || !inst) return; + uasync_post(ua, regenerate_trampoline, NULL); } void instance_lite_set_event_handler(instance_lite_event_fn handler) { IL_LOGI("[EVT_DIAG] instance_lite_set_event_handler handler=%p g_running=%d", (void*)handler, g_running); g_event_handler = handler; - if (g_running) chat_event_set_handler(chat_event_forward); + if (__atomic_load_n(&g_running, __ATOMIC_ACQUIRE)) chat_event_set_handler(chat_event_forward); }