Browse Source

etcp: instance-level conn status callbacks (NEW/UP/DOWN/DELETE), fix member online in chatgui

- Add etcp_add_conn_status_cbk/remove/set — unified callback for all connections
- 4 statuses: NEW, UP, DOWN, DELETE — one subscription covers all connections
- Fire points in etcp.c: create→NEW, on_up→UP, on_down→DOWN, close→DELETE
- chat_sync.c: replace cs_on_new_conn+per-conn subs with cs_on_conn_status
- db_sync.c: same migration to instance-level callback
- cs_is_peer_online: use conn->links_up instead of link->link_status (was 0 during INIT handshake)
topo_upd
Evgeny 3 months ago
parent
commit
aa1ed006d5
  1. 38
      src/db_sync.c
  2. 12
      src/etcp.c
  3. 28
      src/etcp_api.c
  4. 17
      src/etcp_api.h
  5. 1
      src/utun_instance.h
  6. 36
      tools/chatgui/transport/chat_sync.c

38
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);

12
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);

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

17
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);
/**

1
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;// для входных-выходных данных пакета

36
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);

Loading…
Cancel
Save