Browse Source

etcp: add node_changed callback chain + chatgui subscription

- etcp.h: add node_changed_cbks field to ETCP_CONN
- etcp_api.h/c: add etcp_conn_node_changed() + add/remove_node_changed_cbk
- etcp_connections.c: fire node_changed callbacks on pubkey mismatch
  (with DEBUG_WARN(GENERAL) logging source addr, peer_id, old/new pubkey)
- etcp.c: cleanup node_changed_cbks on connection close
- chat_conn_mgr.c: subscribe to node_changed — clear DB addresses,
  reset in-memory cache, close connections (deferred for current conn),
  post GUI_EVT_NODE_CHANGED
- gui_bridge.h: GUI_EVT_NODE_CHANGED=11
- gui_bridge_impl.cpp: handle GUI_EVT_NODE_CHANGED
topo_upd
Evgeny 2 months ago
parent
commit
7a5df3e196
  1. 1
      src/etcp.c
  2. 1
      src/etcp.h
  3. 3
      src/etcp_api.c
  4. 3
      src/etcp_api.h
  5. 13
      src/etcp_connections.c
  6. 63
      tools/chatgui/transport/chat_conn_mgr.c
  7. 1
      tools/chatgui/transport/gui_bridge.h
  8. 8
      tools/chatgui/transport/gui_bridge_impl.cpp

1
src/etcp.c

@ -338,6 +338,7 @@ static void etcp_connection_free_resources(struct ETCP_CONN* etcp) {
{ struct etcp_cbk_entry* cbe = etcp->init_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->init_cbks = NULL; } { struct etcp_cbk_entry* cbe = etcp->init_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->init_cbks = NULL; }
{ struct etcp_cbk_entry* cbe = etcp->up_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->up_cbks = NULL; } { struct etcp_cbk_entry* cbe = etcp->up_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->up_cbks = NULL; }
{ struct etcp_cbk_entry* cbe = etcp->down_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->down_cbks = NULL; } { struct etcp_cbk_entry* cbe = etcp->down_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->down_cbks = NULL; }
{ struct etcp_cbk_entry* cbe = etcp->node_changed_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->node_changed_cbks = NULL; }
if (etcp->inflight_pool) { memory_pool_destroy(etcp->inflight_pool); etcp->inflight_pool = NULL; } if (etcp->inflight_pool) { memory_pool_destroy(etcp->inflight_pool); etcp->inflight_pool = NULL; }
if (etcp->io_pool) { memory_pool_destroy(etcp->io_pool); etcp->io_pool = NULL; } if (etcp->io_pool) { memory_pool_destroy(etcp->io_pool); etcp->io_pool = NULL; }

1
src/etcp.h

@ -240,6 +240,7 @@ struct ETCP_CONN {
struct etcp_cbk_entry* init_cbks; // цепочка callback'ов при инициализации соединения struct etcp_cbk_entry* init_cbks; // цепочка callback'ов при инициализации соединения
struct etcp_cbk_entry* up_cbks; // цепочка callback'ов при поднятии канала struct etcp_cbk_entry* up_cbks; // цепочка callback'ов при поднятии канала
struct etcp_cbk_entry* down_cbks; // цепочка callback'ов при падении канала struct etcp_cbk_entry* down_cbks; // цепочка callback'ов при падении канала
struct etcp_cbk_entry* node_changed_cbks; // цепочка callback'ов при смене pubkey узла (peer_id совпал, ключ изменился)
void (*bgp_ready_cbk)(struct ETCP_CONN* conn); // вызывается когда BGP готов (завершён или пропущен) void (*bgp_ready_cbk)(struct ETCP_CONN* conn); // вызывается когда BGP готов (завершён или пропущен)
uint32_t cnt_ack_hit_inf; // счетчик удлений из inflight uint32_t cnt_ack_hit_inf; // счетчик удлений из inflight

3
src/etcp_api.c

@ -37,6 +37,9 @@ void etcp_conn_add_up_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { i
void etcp_conn_remove_up_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn) etcp_cbk_remove_from_chain(&conn->up_cbks, fn, arg); } void etcp_conn_remove_up_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn) etcp_cbk_remove_from_chain(&conn->up_cbks, fn, arg); }
void etcp_conn_add_down_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn) etcp_cbk_add_to_chain(&conn->down_cbks, fn, arg); } void etcp_conn_add_down_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn) etcp_cbk_add_to_chain(&conn->down_cbks, fn, arg); }
void etcp_conn_remove_down_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn) etcp_cbk_remove_from_chain(&conn->down_cbks, fn, arg); } void etcp_conn_remove_down_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn) etcp_cbk_remove_from_chain(&conn->down_cbks, fn, arg); }
void etcp_conn_node_changed(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn && fn) fn(conn, arg); }
void etcp_conn_add_node_changed_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn) etcp_cbk_add_to_chain(&conn->node_changed_cbks, fn, arg); }
void etcp_conn_remove_node_changed_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn) etcp_cbk_remove_from_chain(&conn->node_changed_cbks, fn, 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_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); } 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); }

3
src/etcp_api.h

@ -148,6 +148,9 @@ void etcp_conn_add_up_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg);
void etcp_conn_remove_up_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg); void etcp_conn_remove_up_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg);
void etcp_conn_add_down_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg); void etcp_conn_add_down_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg);
void etcp_conn_remove_down_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg); void etcp_conn_remove_down_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg);
void etcp_conn_node_changed(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg);
void etcp_conn_add_node_changed_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg);
void etcp_conn_remove_node_changed_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_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_remove_new_conn_cbk(struct UTUN_INSTANCE* inst, etcp_cbk_fn fn, void* arg);

13
src/etcp_connections.c

@ -1711,7 +1711,18 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) {
queue_entry_count(e_sock->instance->connections)); queue_entry_count(e_sock->instance->connections));
} }
else {// check keys если существующее подключение else {// check keys если существующее подключение
if (memcmp(conn->crypto_ctx.peer_public_key, sc.peer_public_key, SC_PUBKEY_SIZE)) { errorcode=5; DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "peer key mismatch for node %016llx", (unsigned long long)peer_id); goto ec_fr; }// коллизия - peer id совпал а ключи разные. if (memcmp(conn->crypto_ctx.peer_public_key, sc.peer_public_key, SC_PUBKEY_SIZE)) {
errorcode=5;
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "node_changed: different pubkey from %s claiming peer=0x%016llx — conn=%s old_pub=%016llx new_pub=%016llx",
sockaddr_storage_to_str(&addr).str, (unsigned long long)peer_id, conn->log_name,
*(uint64_t*)conn->crypto_ctx.peer_public_key, *(uint64_t*)sc.peer_public_key);
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "peer key mismatch for node %016llx, firing node_changed callbacks", (unsigned long long)peer_id);
conn->callbacks_running = 1;
struct etcp_cbk_entry* cbe = conn->node_changed_cbks;
while (cbe) { struct etcp_cbk_entry* n = cbe->next; cbe->fn(conn, cbe->arg); cbe = n; }
conn->callbacks_running = 0;
goto ec_fr;
}// коллизия - peer id совпал а ключи разные.
} }
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "INIT conn=%s new_conn=%d peer=0x%016llx state=%d links_up=%d links=%p", conn->log_name, new_conn, (unsigned long long)peer_id, conn->state, conn->links_up, (void*)conn->links); DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "INIT conn=%s new_conn=%d peer=0x%016llx state=%d links_up=%d links=%p", conn->log_name, new_conn, (unsigned long long)peer_id, conn->state, conn->links_up, (void*)conn->links);

63
tools/chatgui/transport/chat_conn_mgr.c

@ -73,6 +73,11 @@ struct cm_slot_ctx {
int slot_idx; int slot_idx;
}; };
struct cm_node_changed_ctx {
struct chat_conn_mgr* mgr;
uint64_t node_id;
};
struct cm_node { struct cm_node {
uint64_t node_id; uint64_t node_id;
uint8_t pubkey[SC_PUBKEY_SIZE]; uint8_t pubkey[SC_PUBKEY_SIZE];
@ -114,6 +119,8 @@ static void cm_ping_timer_cb(void* arg);
static void cm_ping_result_cb(int success, uint16_t rtt, void* arg, uint64_t nonce, const uint8_t* resp_data, size_t resp_data_len); static void cm_ping_result_cb(int success, uint16_t rtt, void* arg, uint64_t nonce, const uint8_t* resp_data, size_t resp_data_len);
static struct ETCP_SOCKET* cm_best_socket(struct chat_conn_mgr* mgr); static struct ETCP_SOCKET* cm_best_socket(struct chat_conn_mgr* mgr);
static void cm_invite_result_cb(int result, uint64_t node_id, void* arg); static void cm_invite_result_cb(int result, uint64_t node_id, void* arg);
static void cm_node_changed_cb(struct ETCP_CONN* conn, void* arg);
static void cm_node_changed_deferred(void* arg);
/* ─── init / destroy ─── */ /* ─── init / destroy ─── */
@ -309,6 +316,7 @@ static void cm_rotation(struct cm_node* node) {
if (!ctx) { etcp_connection_close(conn); slot->state = CM_SLOT_FREE; return; } if (!ctx) { etcp_connection_close(conn); slot->state = CM_SLOT_FREE; return; }
ctx->node = node; ctx->slot_idx = (int)(slot - node->slots); ctx->node = node; ctx->slot_idx = (int)(slot - node->slots);
etcp_conn_add_init_cbk(conn, cm_init_cb, ctx); etcp_conn_add_init_cbk(conn, cm_init_cb, ctx);
etcp_conn_add_node_changed_cbk(conn, cm_node_changed_cb, mgr);
if (!etcp_link_new(conn, sock, &node->addrs[idx].sa, 0)) { if (!etcp_link_new(conn, sock, &node->addrs[idx].sa, 0)) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: link_new failed node=0x%016llx idx=%d", CM_ID, (unsigned long long)node->node_id, idx); DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: link_new failed node=0x%016llx idx=%d", CM_ID, (unsigned long long)node->node_id, idx);
etcp_connection_close(conn); u_free(ctx); slot->state = CM_SLOT_DEAD; etcp_connection_close(conn); u_free(ctx); slot->state = CM_SLOT_DEAD;
@ -332,6 +340,7 @@ static void cm_rotation(struct cm_node* node) {
if (!ctx) { etcp_connection_close(conn); slot->state = CM_SLOT_FREE; return; } if (!ctx) { etcp_connection_close(conn); slot->state = CM_SLOT_FREE; return; }
ctx->node = node; ctx->slot_idx = (int)(slot - node->slots); ctx->node = node; ctx->slot_idx = (int)(slot - node->slots);
etcp_conn_add_init_cbk(conn, cm_init_cb, ctx); etcp_conn_add_init_cbk(conn, cm_init_cb, ctx);
etcp_conn_add_node_changed_cbk(conn, cm_node_changed_cb, mgr);
if (!etcp_link_new(conn, sock, &node->addrs[slot->addr_idx].sa, 0)) { if (!etcp_link_new(conn, sock, &node->addrs[slot->addr_idx].sa, 0)) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: link_new(retry) failed node=0x%016llx idx=%d", CM_ID, (unsigned long long)node->node_id, slot->addr_idx); DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: link_new(retry) failed node=0x%016llx idx=%d", CM_ID, (unsigned long long)node->node_id, slot->addr_idx);
etcp_connection_close(conn); u_free(ctx); slot->state = CM_SLOT_DEAD; etcp_connection_close(conn); u_free(ctx); slot->state = CM_SLOT_DEAD;
@ -438,6 +447,60 @@ static void cm_retry_cb(void* arg) {
node->retry_timer = uasync_set_timeout(node->mgr->inst->ua, CM_RETRY_TB, node, cm_retry_cb, "cm_retry"); node->retry_timer = uasync_set_timeout(node->mgr->inst->ua, CM_RETRY_TB, node, cm_retry_cb, "cm_retry");
} }
/* ─── node_changed ─── */
static void cm_node_changed_deferred(void* arg) {
struct cm_node_changed_ctx* ctx = (struct cm_node_changed_ctx*)arg;
struct cm_node* node = cm_find_node(ctx->mgr, ctx->node_id);
if (node) {
for (int i = 0; i < CM_MAX_PARALLEL; i++) {
struct cm_slot* slot = &node->slots[i];
if (slot->conn) { etcp_connection_close(slot->conn); slot->conn = NULL; }
slot->state = CM_SLOT_FREE; slot->addr_idx = 255;
}
cm_deliver_all(node, CC_ERR_INTERNAL);
node->delivered = 0;
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: NODE_CHANGED deferred close done node=0x%016llx",
CM_ID, (unsigned long long)ctx->node_id);
}
u_free(ctx);
}
static void cm_node_changed_cb(struct ETCP_CONN* conn, void* arg) {
struct chat_conn_mgr* mgr = (struct chat_conn_mgr*)arg;
if (!mgr || !conn) return;
uint64_t node_id = conn->peer_node_id;
struct cm_node* node = cm_find_node(mgr, node_id);
if (!node) return;
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: NODE_CHANGED node=0x%016llx — clearing addresses and closing connections",
CM_ID, (unsigned long long)node_id);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(mgr->db, "DELETE FROM node_addresses WHERE node_id=?", -1, &st, NULL) == SQLITE_OK) {
sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id);
sqlite3_step(st); sqlite3_finalize(st);
}
node->addr_loaded = 0;
node->addr_count = 0;
if (node->retry_timer) { uasync_cancel_timeout(mgr->inst->ua, node->retry_timer); node->retry_timer = NULL; }
for (int i = 0; i < CM_MAX_PARALLEL; i++) {
struct cm_slot* slot = &node->slots[i];
if (slot->init_timer) { uasync_cancel_timeout(mgr->inst->ua, slot->init_timer); slot->init_timer = NULL; }
if (slot->conn && slot->conn != conn) { etcp_connection_close(slot->conn); slot->conn = NULL; }
if (slot->conn != conn) { slot->state = CM_SLOT_FREE; slot->addr_idx = 255; }
}
struct cm_node_changed_ctx* ctx = u_calloc(1, sizeof(*ctx));
if (ctx) { ctx->mgr = mgr; ctx->node_id = node_id; uasync_post(mgr->inst->ua, cm_node_changed_deferred, ctx); }
uint8_t evt[8]; memcpy(evt, &node_id, 8);
gui_bridge_post(GUI_EVT_NODE_CHANGED, evt, 8);
}
/* ─── фоновый пинг ─── */ /* ─── фоновый пинг ─── */
static void cm_ping_timer_cb(void* arg) { static void cm_ping_timer_cb(void* arg) {

1
tools/chatgui/transport/gui_bridge.h

@ -22,6 +22,7 @@ struct UASYNC;
#define GUI_EVT_CHANNEL_PEERS_ONLINE 8 /* data: [ch_id_len:1][ch_id:var][online_count:2] */ #define GUI_EVT_CHANNEL_PEERS_ONLINE 8 /* data: [ch_id_len:1][ch_id:var][online_count:2] */
#define GUI_EVT_DB_READY 9 /* data: none — DB is open, tables created */ #define GUI_EVT_DB_READY 9 /* data: none — DB is open, tables created */
#define GUI_EVT_STATUS_REFRESH 10 /* data: status text (null-terminated string) */ #define GUI_EVT_STATUS_REFRESH 10 /* data: status text (null-terminated string) */
#define GUI_EVT_NODE_CHANGED 11 /* data: [node_id:8] */
/* ── API ── */ /* ── API ── */

8
tools/chatgui/transport/gui_bridge_impl.cpp

@ -130,6 +130,14 @@ void GuiBridgeReceiver::processPost(int eventType, QByteArray data) {
case GUI_EVT_STATUS_REFRESH: case GUI_EVT_STATUS_REFRESH:
if (g_status_refresh_cb) g_status_refresh_cb((const char*)d, dlen); if (g_status_refresh_cb) g_status_refresh_cb((const char*)d, dlen);
break; break;
case GUI_EVT_NODE_CHANGED:
if (dlen >= 8) {
uint64_t nodeId; memcpy(&nodeId, d, 8);
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "gui_bridge: NODE_CHANGED node=0x%016llx — pubkey changed, connection closed", (unsigned long long)nodeId);
} else {
DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "gui_bridge: NODE_CHANGED data too short %d", dlen);
}
break;
default: default:
DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "gui_bridge: unknown event type %d", eventType); DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "gui_bridge: unknown event type %d", eventType);
break; break;

Loading…
Cancel
Save