Browse Source

fix: cm_invite_pending — deferred free + cancelled guard

Корень краша: cm_handle_invite_info_resp вызывал cm_invite_cleanup(inv)
без unlink из invite_list — висячий указатель при conn_mgr_destroy → SIGSEGV.

Механизм защиты:
- cancelled:1 в cm_invite_pending — блокирует все коллбэки после cancel
- cm_invite_cancel: cancelled=1 первой строкой, deferred free через uasync_post
- Все коллбэки (NCD, TCP ready/close, overall timer): guard + DEBUG_WARN
- Success path: temp_nq=NULL (передача владения), вызов cm_invite_cancel
- INV_CANCEL логи: skip (NULL) / done — видно что реально сделано
topo_upd
evgeny 2 months ago
parent
commit
dc8c9f7875
  1. 67
      close_bug.txt
  2. 107
      src/routing_layer/conn_mgr_core.c
  3. 2
      src/routing_layer/conn_mgr_priv.h
  4. 30
      src/routing_layer/topo_group.c
  5. 2
      src/transport_layer/etcp_connections.c
  6. 228
      src/utun_instance.c
  7. 1
      tools/chatgui-android/app/src/main/java/com/utun/chat/viewmodel/ChatViewModel.kt
  8. 9
      tools/chatgui-android/jni_bridge/android_jni_bridge.c
  9. 209
      tools/chatgui-android/libutun_lite/instance_lite.c

67
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 (приемлемо).

107
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;
}

2
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;
};

30
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) {

2
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;

228
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);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] N done — router/nat_det/fw/stcp_srv");
// Cleanup networks queue
/* 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;
}
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] O done — networks");
// Cleanup config
if (instance->config) {
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Freeing configuration");
free_config(instance->config);
instance->config = NULL;
}
/* 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;
}
}

1
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

9
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) {

209
tools/chatgui-android/libutun_lite/instance_lite.c

@ -31,6 +31,8 @@
#include <stdlib.h>
#include <unistd.h>
#include <signal.h>
#include <time.h>
#include <errno.h>
#include <openssl/evp.h>
#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);
}

Loading…
Cancel
Save