You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

439 lines
19 KiB

/*
* chat_status.c — сбор статуса (NTP + ETCP connections) и отправка
*
* Вынесено из chat_core.c для уменьшения размера модуля.
*/
#include "chat_core_priv.h"
#include "chat_event.h"
#include "../utun_instance.h"
#include "../ntp_time.h"
#include "../../lib/ll_queue.h"
#include "../../lib/mem.h"
#include "../../lib/platform_compat.h"
#include "../transport_layer/etcp.h"
#include "../transport_layer/etcp_connections.h"
#include "../transport_layer/stcp_link.h"
#include "../transport_layer/pkt_normalizer.h"
#include <string.h>
static const char* nat_type_str(uint8_t t) {
switch (t) { case 0: return "UNKNOWN"; case 1: return "EIM"; case 2: return "STRICT"; case 3: return "DIRECT"; default: return "?"; }
}
static void get_node_name(uint64_t node_id, char* out, size_t sz) {
out[0] = '\0';
if (!g_cc.db || node_id == 0) return;
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(g_cc.db, "SELECT name FROM nodes WHERE node_id=?", -1, &st, NULL) != SQLITE_OK) return;
sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id);
if (sqlite3_step(st) == SQLITE_ROW) {
const char* n = (const char*)sqlite3_column_text(st, 0);
if (n) snprintf(out, sz, "%s", n);
}
sqlite3_finalize(st);
}
static void chat_core_collect_status(void) {
char buf[8192]; int off = 0;
if (!g_cc.initialized || !g_cc.inst) {
off = snprintf(buf, sizeof(buf), "uTun not initialized\n");
chat_event_post(CHAT_EVT_STATUS_REFRESH, (const uint8_t*)buf, off);
return;
}
struct NTP_TIME* ntp = &g_cc.inst->ntp;
off += snprintf(buf + off, sizeof(buf) - off, "=== NTP ===\n");
off += snprintf(buf + off, sizeof(buf) - off, "Enabled: %s\n", ntp->enabled ? "yes" : "no");
off += snprintf(buf + off, sizeof(buf) - off, "Synced: %s\n", ntp->synced ? "yes" : "no");
off += snprintf(buf + off, sizeof(buf) - off, "Offset: %.1f s\n", ntp->offset_us / 1000000.0);
if (ntp->server_count > 0 && ntp->servers) {
off += snprintf(buf + off, sizeof(buf) - off, "Servers: %d\n", ntp->server_count);
for (int i = 0; i < ntp->server_count && ntp->servers[i]; i++)
off += snprintf(buf + off, sizeof(buf) - off, " [%d] %s%s\n", i, ntp->servers[i],
i == ntp->server_current ? " (current)" : "");
}
if (ntp->synced && ntp->last_sync_tb > 0) {
uint64_t now_tb = get_time_tb();
uint64_t ago_sec = (now_tb - ntp->last_sync_tb) / 10000;
off += snprintf(buf + off, sizeof(buf) - off, "Last sync: %llu sec ago\n", (unsigned long long)ago_sec);
}
int conn_count = 0;
struct ll_entry* entry = g_cc.inst->connections->head;
while (entry) { conn_count++; entry = entry->next; }
off += snprintf(buf + off, sizeof(buf) - off, "=== Connections (%d) ===\n", conn_count);
uint64_t my_id = g_cc.inst->node_id;
entry = g_cc.inst->connections->head;
while (entry) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data;
if (!ce || !ce->conn) { entry = entry->next; continue; }
uint64_t pid = ce->peer_node_id;
struct ETCP_CONN* conn = ce->conn;
int link_count = 0;
struct ETCP_LINK* tl = conn->links;
while (tl) { link_count++; tl = tl->next; }
const char* init_str = conn->initialized ? "" : "(!init)";
char myhex[5], peerhex[5];
snprintf(myhex, sizeof(myhex), "%04llX", (unsigned long long)(my_id & 0xFFFF));
snprintf(peerhex, sizeof(peerhex), "%04llX", (unsigned long long)(pid & 0xFFFF));
char peername[64];
get_node_name(pid, peername, sizeof(peername));
if (peername[0])
off += snprintf(buf + off, sizeof(buf) - off, "[%s]→[%s] %s ETCP:%s%s(%dL)\n",
myhex, peerhex, peername,
conn->links_up ? "UP" : "DOWN", init_str, link_count);
else
off += snprintf(buf + off, sizeof(buf) - off, "[%s]→[%s] ETCP:%s%s(%dL)\n",
myhex, peerhex,
conn->links_up ? "UP" : "DOWN", init_str, link_count);
int link_idx = 0;
struct ETCP_LINK* link = conn->links;
while (link) {
char rtt_str[32]; rtt_str[0] = '\0';
if (link->rtt_last > 0) snprintf(rtt_str, sizeof(rtt_str), " / rtt=%ums", link->rtt_last / 10);
off += snprintf(buf + off, sizeof(buf) - off, " LINK#%d: %s /NAT=%s%s\n",
link_idx,
link->link_status ? "UP" : "DOWN",
nat_type_str(link->nat_type),
rtt_str);
link = link->next; link_idx++;
}
entry = entry->next;
}
chat_event_post(CHAT_EVT_STATUS_REFRESH, (const uint8_t*)buf, off);
}
/* ── Бинарный список соединений: [count:2][entry:46B]* ── */
/* entry: peer_node_id:8 flags:1 link_count:1 rtt:2 inflight_kb:2 name:32 */
#define CONN_LIST_ENTRY_SIZE 46
static void collect_conn_list(void) {
if (!g_cc.initialized || !g_cc.inst || !g_cc.inst->connections) return;
uint16_t count = 0;
struct ll_entry* entry = g_cc.inst->connections->head;
while (entry) { count++; entry = entry->next; }
if (count > 250) count = 250;
size_t buf_sz = 2 + (size_t)count * CONN_LIST_ENTRY_SIZE;
uint8_t* buf = (uint8_t*)u_malloc(buf_sz);
if (!buf) return;
uint16_t* hdr = (uint16_t*)buf;
*hdr = count;
uint8_t* p = buf + 2;
entry = g_cc.inst->connections->head;
for (uint16_t i = 0; i < count && entry; i++, entry = entry->next) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data;
if (!ce || !ce->conn) { memset(p, 0, CONN_LIST_ENTRY_SIZE); p += CONN_LIST_ENTRY_SIZE; continue; }
struct ETCP_CONN* conn = ce->conn;
memcpy(p, &ce->peer_node_id, 8); p += 8;
uint8_t flags = 0;
if (conn->links_up) flags |= 1;
if (conn->initialized) flags |= 2;
if (conn->links) {
if (!conn->links->is_server) flags |= 4;
else flags |= 8;
}
*p++ = flags;
uint8_t lc = 0;
struct ETCP_LINK* tl = conn->links; while (tl) { lc++; tl = tl->next; }
*p++ = lc;
uint16_t rtt = 0;
if (conn->links) rtt = conn->links->rtt_last;
memcpy(p, &rtt, 2); p += 2;
uint16_t inflight_kb = (uint16_t)(conn->unacked_bytes / 1024);
memcpy(p, &inflight_kb, 2); p += 2;
char name[32]; memset(name, 0, 32);
get_node_name(ce->peer_node_id, name, 31);
if (!name[0]) snprintf(name, 32, "%016llX", (unsigned long long)ce->peer_node_id);
memcpy(p, name, 32); p += 32;
}
chat_event_post(CHAT_EVT_CONN_LIST, buf, (int)buf_sz);
u_free(buf);
}
/* ── Полный текстовый дамп одного соединения ── */
static void collect_conn_metrics(uint64_t peer_node_id) {
char buf[8192]; int off = 0;
if (!g_cc.initialized || !g_cc.inst || !g_cc.inst->connections) {
off = snprintf(buf, sizeof(buf), "uTun not initialized\n");
chat_event_post(CHAT_EVT_CONN_METRICS, (const uint8_t*)buf, off);
return;
}
struct ll_entry* e = queue_find_data_by_index(g_cc.inst->connections, (const uint8_t*)&peer_node_id);
if (!e) {
off = snprintf(buf, sizeof(buf), "Connection %016llX not found\n", (unsigned long long)peer_node_id);
chat_event_post(CHAT_EVT_CONN_METRICS, (const uint8_t*)buf, off);
return;
}
struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data;
if (!ce || !ce->conn) {
off = snprintf(buf, sizeof(buf), "Connection %016llX has null data\n", (unsigned long long)peer_node_id);
chat_event_post(CHAT_EVT_CONN_METRICS, (const uint8_t*)buf, off);
return;
}
struct ETCP_CONN* conn = ce->conn;
uint64_t my_id = g_cc.inst->node_id;
char peername[64]; get_node_name(peer_node_id, peername, sizeof(peername));
if (!peername[0]) snprintf(peername, sizeof(peername), "%016llX", (unsigned long long)peer_node_id);
int is_server = conn->links ? conn->links->is_server : -1;
const char* role_str = is_server < 0 ? "none" : (is_server ? "server" : "client");
/* Count links */
int link_count = 0, links_up_count = 0;
struct ETCP_LINK* tl = conn->links; while (tl) { if (tl->link_status) links_up_count++; link_count++; tl = tl->next; }
off += snprintf(buf + off, sizeof(buf) - off,
"=== PEER: %04llX\u2192%04llX [%s] ===\n"
"Role: %s\n"
"Status: %s Links: %d/%d Initialized: %s Session: 0x%X MTU: %d\n"
"Reinit: %u Reset: %u Tx_state: %d Routing_exchange: %d\n"
"\n--- ETCP ---\n"
"RTT last/avg10: %.1f/%.1f ms Jitter: %.1f ms\n"
"Bytes sent: %u Retrans: %u ACKs: %u\n"
"Unacked bytes: %u Max inflight: %u\n"
"RX dup: %u TX dup: %u\n"
"IDs: next_tx=%u last_rx=%u last_del=%u rx_ack_till=%u\n"
"\n--- Queues ---\n",
(unsigned long long)(my_id & 0xFFFF), (unsigned long long)(peer_node_id & 0xFFFF), peername,
role_str,
conn->links_up ? "UP" : "DOWN", links_up_count, link_count,
conn->initialized ? "yes" : "no", conn->session_id, conn->mtu,
conn->reinit_count, conn->reset_count, conn->tx_state, conn->routing_exchange_active,
conn->rtt_last / 10.0f, conn->rtt_avg_10 / 10.0f, conn->jitter / 10.0f,
conn->bytes_sent_total, conn->retransmissions_count, conn->ack_packets_count,
conn->unacked_bytes, conn->max_inflight,
conn->rx_dup_count, conn->tx_dup_count,
conn->next_tx_id, conn->last_rx_id, conn->last_delivered_id, conn->rx_ack_till);
if (conn->input_queue) off += snprintf(buf + off, sizeof(buf) - off, "input_q: %uB/%dp ", (unsigned)queue_total_bytes(conn->input_queue), queue_entry_count(conn->input_queue));
if (conn->input_send_q) off += snprintf(buf + off, sizeof(buf) - off, "send_q: %uB/%dp ", (unsigned)queue_total_bytes(conn->input_send_q), queue_entry_count(conn->input_send_q));
if (conn->input_wait_ack) off += snprintf(buf + off, sizeof(buf) - off, "wait_ack: %uB/%dp\n", (unsigned)queue_total_bytes(conn->input_wait_ack), queue_entry_count(conn->input_wait_ack));
if (conn->ack_q) off += snprintf(buf + off, sizeof(buf) - off, "ack_q: %uB/%dp ", (unsigned)queue_total_bytes(conn->ack_q), queue_entry_count(conn->ack_q));
if (conn->recv_q) off += snprintf(buf + off, sizeof(buf) - off, "recv_q: %uB/%dp ", (unsigned)queue_total_bytes(conn->recv_q), queue_entry_count(conn->recv_q));
if (conn->output_queue) off += snprintf(buf + off, sizeof(buf) - off, "out_q: %uB/%dp\n", (unsigned)queue_total_bytes(conn->output_queue), queue_entry_count(conn->output_queue));
off += snprintf(buf + off, sizeof(buf) - off, "\n--- ACK Debug ---\n"
"hit_inf=%u hit_sndq=%u miss=%u link_wait=%u\n",
conn->cnt_ack_hit_inf, conn->cnt_ack_hit_sndq, conn->cnt_ack_miss, conn->cnt_link_wait);
/* Normalizer */
if (conn->normalizer) {
struct PKTNORM* pn = conn->normalizer;
off += snprintf(buf + off, sizeof(buf) - off,
"\n--- Normalizer ---\n"
"in: %llu pkts / %llu B out: %llu pkts / %llu B\n"
"frag_size=%u data_ptr=%u data_size=%u alloc_err=%u logic_err=%u\n",
(unsigned long long)pn->in_total_pkts, (unsigned long long)pn->in_total_bytes,
(unsigned long long)pn->out_total_pkts, (unsigned long long)pn->out_total_bytes,
pn->frag_size, pn->data_ptr, pn->data_size, pn->alloc_errors, pn->logic_errors);
}
/* Links */
int link_idx = 0;
struct ETCP_LINK* link = conn->links;
while (link) {
const char* is_tcp_str = link->is_tcp ? "TCP" : "UDP";
char local_str[54], remote_str[54];
strncpy(local_str, link->conn ? sockaddr_storage_to_str(&link->conn->local_addr).str : "stcp", sizeof(local_str) - 1);
strncpy(remote_str, sockaddr_storage_to_str(&link->remote_addr).str, sizeof(remote_str) - 1);
off += snprintf(buf + off, sizeof(buf) - off,
"\n--- LINK#%d: %s %s NAT=%s ---\n"
"addr: %s -> %s\n"
"RTT: %.1f ms TT: %.1f ms\n"
"Keepalive: sent=%u recv=%u period=%ums\n"
"Inflight: %u / %u B (%u pkts)\n"
"Encrypt: sent=%zu err=%zu Decrypt: rcvd=%zu err=%zu\n",
link_idx, link->link_status ? "UP" : "DOWN", is_tcp_str, nat_type_str(link->nat_type),
local_str, remote_str,
link->rtt_last / 10.0f, link->tt_last / 10.0f,
link->keepalive_sent_count, link->keepalive_recv_count, link->keepalive_interval,
link->inflight_bytes, link->inflight_lim_bytes, link->inflight_packets,
link->total_encrypted, link->encrypt_errors, link->total_decrypted, link->decrypt_errors);
link = link->next; link_idx++;
}
chat_event_post(CHAT_EVT_CONN_METRICS, (const uint8_t*)buf, off);
}
void chat_core_collect_status_trampoline(void* arg) {
(void)arg;
chat_core_collect_status();
}
void chat_core_collect_conn_list_trampoline(void* arg) {
(void)arg;
collect_conn_list();
}
void chat_core_collect_conn_metrics_trampoline(void* arg) {
uint64_t peer_id;
memcpy(&peer_id, arg, sizeof(peer_id));
u_free(arg);
collect_conn_metrics(peer_id);
}
/* ── Live-снапшот для панели деталей участника ──
* packet: [node_id:8][flags:1][own_tcp_active:2][link_count:1](link:25B)*[sock_count:1](sock:48B)*
* flags: bit0=conn_present bit1=conn_up(links_up) bit2=conn_initialized
* link: [is_tcp:1][link_state:1][link_status:1][tcp_ready:1][is_server:1][family:1][addr_len:1][addr:16][port:2]
* sock: [sock_id:4][port:2][link_count:4][ifname:38]
* Собирается в worker-потоке — читает живые структуры UTUN_INSTANCE безопасно. */
#define MEMBER_DETAIL_MAX_LINKS 16
#define MEMBER_DETAIL_MAX_SOCKS 16
#define MEMBER_DETAIL_LINK_SIZE 25
#define MEMBER_DETAIL_SOCK_SIZE 48
/* Повторяет Qt: name.lastIndexOf("_v") → name.mid(3, idx - 3) */
static void sock_ifname(const char* name, char* out, size_t out_sz) {
out[0] = '\0';
if (!name || !name[0]) return;
const char* last = NULL;
for (const char* q = name; (q = strstr(q, "_v")) != NULL; q++) last = q;
if (!last) return;
int idx = (int)(last - name);
if (idx <= 3) return;
int len = idx - 3;
if (len >= (int)out_sz) len = (int)out_sz - 1;
memcpy(out, name + 3, len);
out[len] = '\0';
}
static void collect_member_detail(uint64_t node_id) {
if (!g_cc.initialized || !g_cc.inst) return;
struct UTUN_INSTANCE* inst = g_cc.inst;
uint8_t flags = 0;
struct ETCP_CONN* conn = instance_find_conn(inst, node_id);
if (conn) {
flags |= 1;
if (conn->links_up) flags |= 2;
if (conn->initialized) flags |= 4;
}
uint16_t own_tcp_active = 0;
if (node_id == inst->node_id) {
if (inst->connections) {
for (struct ll_entry* e = inst->connections->head; e; e = e->next) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data;
if (ce && ce->conn && ce->conn->initialized) {
for (struct ETCP_LINK* lk = ce->conn->links; lk; lk = lk->next)
if (lk->is_tcp) { own_tcp_active++; break; }
}
}
}
if (inst->tcp_connections) {
for (struct ll_entry* e = inst->tcp_connections->head; e; e = e->next) {
struct tcp_conn_entry* te = (struct tcp_conn_entry*)e->data;
if (te && te->etcp_conn && te->etcp_conn->initialized) own_tcp_active++;
}
}
}
uint8_t link_count = 0;
for (struct ETCP_LINK* lk = conn ? conn->links : NULL; lk && link_count < MEMBER_DETAIL_MAX_LINKS; lk = lk->next)
link_count++;
uint8_t sock_count = 0;
for (struct ETCP_SOCKET* s = inst->etcp_sockets; s && sock_count < MEMBER_DETAIL_MAX_SOCKS; s = s->next)
sock_count++;
for (struct TCP_SOCKET* ts = inst->tcp_sockets; ts && sock_count < MEMBER_DETAIL_MAX_SOCKS; ts = ts->next)
sock_count++;
size_t buf_sz = 12 + (size_t)link_count * MEMBER_DETAIL_LINK_SIZE
+ 1 + (size_t)sock_count * MEMBER_DETAIL_SOCK_SIZE;
uint8_t* buf = (uint8_t*)u_malloc(buf_sz);
if (!buf) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "collect_member_detail: buf alloc failed"); return; }
uint8_t* p = buf;
memcpy(p, &node_id, 8); p += 8;
*p++ = flags;
memcpy(p, &own_tcp_active, 2); p += 2;
*p++ = link_count;
for (struct ETCP_LINK* lk = conn ? conn->links : NULL; lk && link_count; lk = lk->next, link_count--) {
*p++ = lk->is_tcp ? 1 : 0;
*p++ = lk->link_state;
*p++ = lk->link_status;
*p++ = (lk->is_tcp && lk->tcp_link && stcp_link_is_ready(lk->tcp_link)) ? 1 : 0;
*p++ = lk->is_server;
uint8_t family = 0, addr_len = 0; uint8_t addr[16] = {0}; uint16_t port = 0;
const struct sockaddr_storage* sa = &lk->remote_addr;
if (sa->ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)sa;
family = 4; addr_len = 4; memcpy(addr, &sin->sin_addr, 4); port = ntohs(sin->sin_port);
} else if (sa->ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)sa;
family = 6; addr_len = 16; memcpy(addr, &sin6->sin6_addr, 16); port = ntohs(sin6->sin6_port);
}
*p++ = family;
*p++ = addr_len;
memcpy(p, addr, 16); p += 16;
memcpy(p, &port, 2); p += 2;
}
*p++ = sock_count;
for (struct ETCP_SOCKET* s = inst->etcp_sockets; s && sock_count; s = s->next, sock_count--) {
uint32_t sid = s->sock_id;
uint16_t sp = 0;
const struct sockaddr_storage* sa = s->interface_addr.ss_family ? &s->interface_addr : &s->local_addr;
if (sa->ss_family == AF_INET) sp = ntohs(((const struct sockaddr_in*)sa)->sin_port);
else if (sa->ss_family == AF_INET6) sp = ntohs(((const struct sockaddr_in6*)sa)->sin6_port);
uint32_t lc = s->links_queue ? (uint32_t)queue_entry_count(s->links_queue) : 0;
char ifname[38] = {0}; sock_ifname(s->name, ifname, sizeof(ifname));
memcpy(p, &sid, 4); p += 4;
memcpy(p, &sp, 2); p += 2;
memcpy(p, &lc, 4); p += 4;
memcpy(p, ifname, 38); p += 38;
}
for (struct TCP_SOCKET* ts = inst->tcp_sockets; ts && sock_count; ts = ts->next, sock_count--) {
uint32_t sid = ts->sock_id;
uint16_t sp = 0;
const struct sockaddr_storage* sa = ts->interface_addr.ss_family ? &ts->interface_addr : &ts->local_addr;
if (sa->ss_family == AF_INET) sp = ntohs(((const struct sockaddr_in*)sa)->sin_port);
else if (sa->ss_family == AF_INET6) sp = ntohs(((const struct sockaddr_in6*)sa)->sin6_port);
char ifname[38] = {0}; sock_ifname(ts->name, ifname, sizeof(ifname));
memcpy(p, &sid, 4); p += 4;
memcpy(p, &sp, 2); p += 2;
uint32_t lc = 0; memcpy(p, &lc, 4); p += 4;
memcpy(p, ifname, 38); p += 38;
}
chat_event_post(CHAT_EVT_MEMBER_DETAIL, buf, (int)(p - buf));
u_free(buf);
}
void chat_core_collect_member_detail_trampoline(void* arg) {
uint64_t node_id;
memcpy(&node_id, arg, sizeof(node_id));
u_free(arg);
collect_member_detail(node_id);
}