diff --git a/src/db_sync.c b/src/db_sync.c index 873cd4e1..236dce04 100644 --- a/src/db_sync.c +++ b/src/db_sync.c @@ -21,7 +21,7 @@ struct DB_SYNC; struct DB_SYNC_INSTANCE; static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry); -static void db_sync_on_new_conn(struct ETCP_CONN* conn, void* arg); +static void db_sync_on_conn_status(struct ETCP_CONN* conn, int status, void* arg); static void db_sync_on_conn_up(struct ETCP_CONN* conn, void* arg); static void db_sync_on_conn_down(struct ETCP_CONN* conn, void* arg); static void db_sync_peer_check_cb(void* arg); @@ -1094,6 +1094,15 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) // Connection callbacks // ============================================================ +static void db_sync_on_conn_status(struct ETCP_CONN* conn, int status, void* arg) { + switch (status) { + case ETCP_CONN_STATUS_UP: db_sync_on_conn_up(conn, arg); break; + case ETCP_CONN_STATUS_DOWN: db_sync_on_conn_down(conn, arg); break; + default: break; + } +} + +#if 0 /* replaced by db_sync_on_conn_status */ static void db_sync_on_new_conn(struct ETCP_CONN* conn, void* arg) { (void)arg; @@ -1102,6 +1111,7 @@ static void db_sync_on_new_conn(struct ETCP_CONN* conn, void* arg) etcp_conn_add_up_cbk(conn, db_sync_on_conn_up, NULL); etcp_conn_add_down_cbk(conn, db_sync_on_conn_down, NULL); } +#endif static void db_sync_on_conn_up(struct ETCP_CONN* conn, void* arg) { @@ -1355,19 +1365,7 @@ int db_sync_init(struct UTUN_INSTANCE* inst) } etcp_bind(inst, ETCP_RT_ID_DB_SYNC, db_sync_recv_cb); - etcp_set_new_conn_cbk(inst, db_sync_on_new_conn, NULL); - - { - struct ll_entry* entry = inst->connections->head; - while (entry) { - struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; - if (ce && ce->conn) { - etcp_conn_add_up_cbk(ce->conn, db_sync_on_conn_up, NULL); - etcp_conn_add_down_cbk(ce->conn, db_sync_on_conn_down, NULL); - } - entry = entry->next; - } - } + etcp_add_conn_status_cbk(inst, db_sync_on_conn_status, NULL); if (db->enabled) db->peer_check_timer = uasync_set_timeout(inst->ua, DB_SYNC_PEER_CHECK_INTERVAL * 10000u, @@ -1394,17 +1392,7 @@ void db_sync_destroy(struct UTUN_INSTANCE* inst) if (si->peers) u_free(si->peers); } - { - struct ll_entry* entry = inst->connections->head; - while (entry) { - struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; - if (ce && ce->conn) { - etcp_conn_remove_up_cbk(ce->conn, db_sync_on_conn_up, NULL); - etcp_conn_remove_down_cbk(ce->conn, db_sync_on_conn_down, NULL); - } - entry = entry->next; - } - } + etcp_remove_conn_status_cbk(inst, db_sync_on_conn_status, NULL); db_sqlite_close(db); if (db->instances) u_free(db->instances); diff --git a/src/etcp.c b/src/etcp.c index 46efddcc..424b307f 100644 --- a/src/etcp.c +++ b/src/etcp.c @@ -46,6 +46,7 @@ static void wait_ack_cb(struct ll_queue* q, void* arg); static void send_ack_req_cb(struct ll_queue* q, void* arg); static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp); static void etcp_connection_free_deferred(void* arg); +static void etcp_fire_conn_status(struct ETCP_CONN* etcp, int status); struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp); static void clear_queue(struct ll_queue* q) { @@ -279,16 +280,24 @@ struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* n struct etcp_cbk_entry* cbe = instance->new_conn_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; cbe->fn(etcp, cbe->arg); cbe = n; } } + etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_NEW); return etcp; } +static void etcp_fire_conn_status(struct ETCP_CONN* etcp, int status) { + if (!etcp || !etcp->instance) return; + struct etcp_status_cbk_entry* cbe = etcp->instance->conn_status_cbks; + while (cbe) { struct etcp_status_cbk_entry* n = cbe->next; cbe->fn(etcp, status, cbe->arg); cbe = n; } +} + static void etcp_on_up(struct ETCP_CONN* etcp) { DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] Connection UP (links_up=%d initialized=%d)", etcp->log_name, etcp->links_up, etcp->initialized); etcp->callbacks_running = 1; struct etcp_cbk_entry* cbe = etcp->up_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; cbe->fn(etcp, cbe->arg); cbe = n; } + etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_UP); etcp->callbacks_running = 0; } @@ -297,6 +306,7 @@ static void etcp_on_down(struct ETCP_CONN* etcp) { etcp->callbacks_running = 1; struct etcp_cbk_entry* cbe = etcp->down_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; cbe->fn(etcp, cbe->arg); cbe = n; } + etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_DOWN); etcp->callbacks_running = 0; } @@ -382,6 +392,8 @@ void etcp_connection_close(struct ETCP_CONN* etcp) { etcp->state = 2; // deleted + etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_DELETE); + // === PHASE 2: deferred resource cleanup (only if no outstanding refs) === if (etcp->ref_count == 0) uasync_call_soon(etcp->instance->ua, etcp, etcp_connection_free_deferred); diff --git a/src/etcp_api.c b/src/etcp_api.c index 8aa75ab0..54fd3f43 100644 --- a/src/etcp_api.c +++ b/src/etcp_api.c @@ -69,6 +69,34 @@ void etcp_conn_remove_down_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg void etcp_add_new_conn_cbk(struct UTUN_INSTANCE* inst, etcp_cbk_fn fn, void* arg) { if (inst) etcp_cbk_add_to_chain(&inst->new_conn_cbks, fn, arg); } void etcp_remove_new_conn_cbk(struct UTUN_INSTANCE* inst, etcp_cbk_fn fn, void* arg) { if (inst) etcp_cbk_remove_from_chain(&inst->new_conn_cbks, fn, arg); } +static void etcp_status_cbk_add_chain(struct etcp_status_cbk_entry** head, etcp_conn_status_fn fn, void* arg) { + if (!head || !fn) return; + struct etcp_status_cbk_entry* e = u_malloc(sizeof(struct etcp_status_cbk_entry)); + if (!e) return; + e->fn = fn; e->arg = arg; e->next = *head; + *head = e; +} +static void etcp_status_cbk_remove_chain(struct etcp_status_cbk_entry** head, etcp_conn_status_fn fn, void* arg) { + if (!head || !fn) return; + struct etcp_status_cbk_entry** p = head; + while (*p) { + if ((*p)->fn == fn && (*p)->arg == arg) { + struct etcp_status_cbk_entry* rm = *p; + *p = rm->next; u_free(rm); return; + } + p = &(*p)->next; + } +} +void etcp_add_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg) { if (inst) etcp_status_cbk_add_chain(&inst->conn_status_cbks, fn, arg); } +void etcp_remove_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg) { if (inst) etcp_status_cbk_remove_chain(&inst->conn_status_cbks, fn, arg); } +void etcp_set_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg) { + if (!inst) return; + struct etcp_status_cbk_entry* entry = inst->conn_status_cbks; + while (entry) { struct etcp_status_cbk_entry* next = entry->next; u_free(entry); entry = next; } + inst->conn_status_cbks = NULL; + if (fn) etcp_status_cbk_add_chain(&inst->conn_status_cbks, fn, arg); +} + void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state) { if (!conn) return; conn->routing_exchange_active = new_state; diff --git a/src/etcp_api.h b/src/etcp_api.h index 216dde06..d057050b 100644 --- a/src/etcp_api.h +++ b/src/etcp_api.h @@ -44,10 +44,23 @@ extern "C" { #define ETCP_RT_ID_CONN_MGR 0x11 // Connection Manager — management connections #define ETCP_RT_ID_NTP_TIME 0x12 // NTP time sync между узлами +// Connection status events (instance-level callback) +#define ETCP_CONN_STATUS_NEW 0 // соединение создано +#define ETCP_CONN_STATUS_UP 1 // соединение поднялось +#define ETCP_CONN_STATUS_DOWN 2 // соединение упало +#define ETCP_CONN_STATUS_DELETE 3 // соединение удалено + // Forward declarations struct ETCP_CONN; struct UTUN_INSTANCE; +typedef void (*etcp_conn_status_fn)(struct ETCP_CONN* conn, int status, void* arg); +struct etcp_status_cbk_entry { + etcp_conn_status_fn fn; + void* arg; + struct etcp_status_cbk_entry* next; +}; + /** * @brief Тип коллбэка для приёма метрик от пира * @param user_ptr пользовательский указатель @@ -162,6 +175,10 @@ void etcp_conn_remove_down_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg void etcp_add_new_conn_cbk(struct UTUN_INSTANCE* inst, etcp_cbk_fn fn, void* arg); void etcp_remove_new_conn_cbk(struct UTUN_INSTANCE* inst, etcp_cbk_fn fn, void* arg); +void etcp_add_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg); +void etcp_remove_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg); +void etcp_set_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg); + void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state); /** diff --git a/src/utun_instance.h b/src/utun_instance.h index 5c91115e..fe6d0a8d 100644 --- a/src/utun_instance.h +++ b/src/utun_instance.h @@ -102,6 +102,7 @@ struct UTUN_INSTANCE { // Callback chain for new ETCP connections struct etcp_cbk_entry* new_conn_cbks; + struct etcp_status_cbk_entry* conn_status_cbks; // instance-level: NEW/UP/DOWN/DELETE void* test_user_ptr; // Generic user pointer (used by tests) struct memory_pool* data_pool;// для входных-выходных данных пакета diff --git a/tools/chatgui/transport/chat_sync.c b/tools/chatgui/transport/chat_sync.c index 5f146b76..3ac53304 100644 --- a/tools/chatgui/transport/chat_sync.c +++ b/tools/chatgui/transport/chat_sync.c @@ -458,8 +458,7 @@ static int cs_is_peer_online(struct UTUN_INSTANCE* inst, uint64_t peer_id) { struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&peer_id); if (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; - struct ETCP_LINK* l = ce->conn->links; - while (l) { if (l->initialized && l->link_status) return 1; l = l->next; } + return ce->conn->initialized && ce->conn->links_up; } return 0; } @@ -557,6 +556,15 @@ static void cs_on_conn_down(struct ETCP_CONN* conn, void* arg) { } } +static void cs_on_conn_status(struct ETCP_CONN* conn, int status, void* arg) { + switch (status) { + case ETCP_CONN_STATUS_UP: cs_on_conn_up(conn, arg); break; + case ETCP_CONN_STATUS_DOWN: cs_on_conn_down(conn, arg); break; + default: break; + } +} + +#if 0 /* replaced by cs_on_conn_status */ static void cs_on_new_conn(struct ETCP_CONN* conn, void* arg) { (void)arg; if (!conn) return; @@ -564,6 +572,7 @@ static void cs_on_new_conn(struct ETCP_CONN* conn, void* arg) { etcp_conn_add_up_cbk(conn, cs_on_conn_up, NULL); etcp_conn_add_down_cbk(conn, cs_on_conn_down, NULL); } +#endif /* ── Periodic refresh from DB ── */ @@ -633,17 +642,7 @@ int chat_sync_init(struct UTUN_INSTANCE* inst, g_cs = cs; etcp_bind(inst, ETCP_RT_ID_CHAT_SYNC, chat_sync_recv_cb); - etcp_add_new_conn_cbk(inst, cs_on_new_conn, NULL); - - { - struct ll_entry* entry = inst->connections->head; - while (entry) { - struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; - etcp_conn_add_up_cbk(ce->conn, cs_on_conn_up, NULL); - etcp_conn_add_down_cbk(ce->conn, cs_on_conn_down, NULL); - entry = entry->next; - } - } + etcp_add_conn_status_cbk(inst, cs_on_conn_status, NULL); cs_refresh_channels(cs); @@ -667,19 +666,10 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) { cs->initialized = 0; g_cs = NULL; etcp_unbind(inst, ETCP_RT_ID_CHAT_SYNC); + etcp_remove_conn_status_cbk(inst, cs_on_conn_status, NULL); if (cs->refresh_timer) { uasync_cancel_timeout(inst->ua, cs->refresh_timer); cs->refresh_timer = NULL; } cs_cancel_proto_timers(cs); - { - struct ll_entry* entry = inst->connections->head; - while (entry) { - struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; - etcp_conn_remove_up_cbk(ce->conn, cs_on_conn_up, NULL); - etcp_conn_remove_down_cbk(ce->conn, cs_on_conn_down, NULL); - entry = entry->next; - } - } - for (int i = 0; i < cs->channel_count; i++) if (cs->channels[i].peer_ids) u_free(cs->channels[i].peer_ids); if (cs->channels) u_free(cs->channels);