From a3b00984ed3b5de82ded14f954e6bef322e24133 Mon Sep 17 00:00:00 2001 From: evgeny Date: Thu, 30 Jul 2026 01:21:40 +0300 Subject: [PATCH] fix: stable ETCP reconnect (no REINIT loop), own voice playback, peersOnline via conn state, UDP log diagnostics --- src/chat/chat_msg.c | 11 ++++- src/transport_layer/etcp_connections.c | 27 ++++++------ .../jni_bridge/android_jni_bridge.c | 13 +++++- .../jni_bridge/android_udp_log.c | 42 +++++++++++++++---- 4 files changed, 67 insertions(+), 26 deletions(-) diff --git a/src/chat/chat_msg.c b/src/chat/chat_msg.c index 0bab0904..82e5548a 100644 --- a/src/chat/chat_msg.c +++ b/src/chat/chat_msg.c @@ -82,6 +82,7 @@ void chat_core_submit_trampoline(void* arg) { struct media_submit_ctx { char channel_id[64]; char content_type[32]; + char media_dest[512]; char media_base[512]; uint8_t msg_data[4096]; uint32_t msg_data_len; @@ -153,11 +154,16 @@ static void on_media_registered(void* arg, int err, const struct media_index_res } int ret = db_sync_insert_signed(si, json, strlen(json), json_sig, 64, ts, NULL); - if (ret != 0) + if (ret != 0) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: media insert failed ret=%d", CC_ID, ret); - else + } else { DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: media message inserted ch=%s ts=%llu", CC_ID, mctx->channel_id, (unsigned long long)ts); + const char* base = strrchr(mctx->media_dest, '/'); + const char* filename = base ? base + 1 : mctx->media_dest; + char attrs[320]; snprintf(attrs, sizeof(attrs), "{\"st\":\"fl\",\"fp\":\"%s\"}", filename); + chat_core_update_local_attrs(mctx->channel_id, ts, json_sig, attrs); + } u_free(sig_msg); u_free(json); u_free(full_data); u_free(mctx); } @@ -190,6 +196,7 @@ static void chat_core_submit_media_message(struct chat_msg_submit* req) { if (!mctx) return; snprintf(mctx->channel_id, sizeof(mctx->channel_id), "%s", req->channel_id); snprintf(mctx->content_type, sizeof(mctx->content_type), "%s", req->content_type); + snprintf(mctx->media_dest, sizeof(mctx->media_dest), "%s", req->media_dest); snprintf(mctx->media_base, sizeof(mctx->media_base), "%s", media_base); mctx->timestamp = req->timestamp; if (req->data && req->data_len > 0) { diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index dabbad9a..b6feca47 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -195,6 +195,8 @@ static void etcp_link_init_timer_cbk(void* arg) { } link->init_timer = uasync_set_timeout(link->etcp->instance->ua, link->init_timeout, link, etcp_link_init_timer_cbk, "link_init"); + if (link->link_state == 3 && link->initialized) return; + if (link->etcp->links_up > 0) return; /* another link is already UP on this connection */ if (link->link_state == 1) etcp_link_send_init(link,1,0);// init (with etcp reset) else etcp_link_send_init(link,0,0);// no etcp reset (reinit) } @@ -1533,11 +1535,7 @@ static int handle_init_response_client(struct ETCP_SOCKET* e_sock, struct ETCP_D link->init_timer = NULL; } - if (pkt_code == ETCP_INIT_RESPONSE && link->etcp->got_initial_pkt) { - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] REINIT from client: INIT_RESPONSE(0x03) received, reinit conn=%p", - link->etcp->log_name, link->etcp); - etcp_conn_reinit(link->etcp); - } + if (link->etcp->initialized == 0) { etcp_conn_ready(link->etcp); @@ -1895,6 +1893,7 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { // Check if WE have an outbound (master) link on this conn struct ETCP_LINK* ml = conn->links; while (ml) { if (ml->is_server == 0) break; ml = ml->next; } + int yielding = 0; if (ml) { if (conn->instance->node_id < peer_id) { DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] COLLISION: we=master node_id=%016llx < peer=%016llx — sending collision INIT", @@ -1905,17 +1904,17 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { } DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] COLLISION: peer smaller, yielding master, processing as slave", conn->log_name); + yielding = 1; } - // Always sync session_id; reset if new session or data was flowing - { int sess_changed = (conn->session_id != session_id); - conn->session_id = session_id; - if (sess_changed || conn->got_initial_pkt) { - send_reset = 1; - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] REINIT existing link: code=0x%02x sess=%08x→%08x got_init=%d initialized=%d links_up=%d", - conn->log_name, code, conn->session_id, session_id, conn->got_initial_pkt, conn->initialized, conn->links_up); - etcp_conn_reinit(conn); - } } + /* Sync session_id; reset only if new session (not when yielding to master) */ + conn->session_id = session_id; + if (!yielding && conn->got_initial_pkt) { + send_reset = 1; + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] REINIT existing link: code=0x%02x sess=%08x→%08x got_init=%d initialized=%d links_up=%d", + conn->log_name, code, conn->session_id, session_id, conn->got_initial_pkt, conn->initialized, conn->links_up); + etcp_conn_reinit(conn); + } } // Cancel existing timers diff --git a/tools/chatgui-android/jni_bridge/android_jni_bridge.c b/tools/chatgui-android/jni_bridge/android_jni_bridge.c index c6993df5..eb2162d2 100644 --- a/tools/chatgui-android/jni_bridge/android_jni_bridge.c +++ b/tools/chatgui-android/jni_bridge/android_jni_bridge.c @@ -19,6 +19,8 @@ #include "../../../lib/sqlite3.h" #include "../../../src/chat/chat_core.h" #include "../../../src/chat/chat_setting.h" +#include "../../../src/utun_instance.h" +#include "../../../src/transport_layer/etcp.h" #include #include #include @@ -323,10 +325,17 @@ char* utun_bridge_get_channels_json(void) { } int online = 0; + uint64_t my_id = g_cc.my_node_id; + struct UTUN_INSTANCE* inst = chat_core_get_inst(); sqlite3_stmt* pc = NULL; - snprintf(sql, sizeof(sql), "SELECT COUNT(*) FROM \"%s\"", tbl_peers); + snprintf(sql, sizeof(sql), "SELECT DISTINCT node_id FROM \"%s\"", tbl_peers); if (sqlite3_prepare_v2(db, sql, -1, &pc, NULL) == SQLITE_OK) { - if (sqlite3_step(pc) == SQLITE_ROW) online = sqlite3_column_int(pc, 0); + while (sqlite3_step(pc) == SQLITE_ROW) { + uint64_t nid = (uint64_t)sqlite3_column_int64(pc, 0); + if (nid == my_id) continue; + if (inst) { struct ETCP_CONN* c = instance_find_conn(inst, nid); if (c && c->initialized && c->links_up) online++; } + else online++; + } sqlite3_finalize(pc); } diff --git a/tools/chatgui-android/jni_bridge/android_udp_log.c b/tools/chatgui-android/jni_bridge/android_udp_log.c index 268ae664..6665ed11 100644 --- a/tools/chatgui-android/jni_bridge/android_udp_log.c +++ b/tools/chatgui-android/jni_bridge/android_udp_log.c @@ -7,18 +7,23 @@ #include "android_udp_log.h" #include "../libutun_lite/instance_lite.h" #include "../../../lib/u_async.h" +#include "../../../lib/platform_compat.h" +#include "../../../lib/socket_compat.h" #include #include #include -#include -#include -#include -#include #define UDP_BUF_SIZE 65536 #define UDP_FLUSH_BYTES 16384 #define UDP_FLUSH_TB 2000 /* 200ms in uasync timebase (0.1ms units) */ +#ifdef __ANDROID__ +#include +#define ULOG(fmt, ...) __android_log_print(ANDROID_LOG_DEBUG, "utun-udplog", fmt, ##__VA_ARGS__) +#else +#define ULOG(fmt, ...) fprintf(stderr, "utun-udplog: " fmt "\n", ##__VA_ARGS__) +#endif + static int g_udp_fd = -1; static struct sockaddr_in g_udp_addr; static int g_udp_has_addr = 0; @@ -26,6 +31,9 @@ static int g_udp_stopping = 0; static char g_udp_buf[UDP_BUF_SIZE]; static int g_udp_buf_len = 0; static void* g_udp_timer = NULL; +static int g_udp_sent_first = 0; +static int64_t g_udp_total_sent = 0; +static int g_udp_send_errors = 0; static void udp_reschedule_timer(void) { if (g_udp_stopping) return; @@ -35,19 +43,30 @@ static void udp_reschedule_timer(void) { static void udp_open_socket(void) { if (g_udp_fd >= 0 || !g_udp_has_addr) return; - g_udp_fd = socket(AF_INET, SOCK_DGRAM | SOCK_NONBLOCK, 0); - if (g_udp_fd < 0) return; + g_udp_fd = (int)socket(AF_INET, SOCK_DGRAM | SOCK_NONBLOCK, 0); + if (g_udp_fd < 0) { ULOG("socket() failed: %s (errno=%d)", strerror(errno), errno); return; } if (connect(g_udp_fd, (struct sockaddr*)&g_udp_addr, sizeof(g_udp_addr)) < 0) { + ULOG("connect() to %s:%d failed: %s (errno=%d)", + inet_ntoa(g_udp_addr.sin_addr), (int)ntohs(g_udp_addr.sin_port), strerror(errno), errno); close(g_udp_fd); g_udp_fd = -1; return; } + ULOG("socket opened fd=%d -> %s:%d", g_udp_fd, + inet_ntoa(g_udp_addr.sin_addr), (int)ntohs(g_udp_addr.sin_port)); udp_reschedule_timer(); } static void udp_do_send(const char* data, int len) { if (g_udp_fd < 0) { udp_open_socket(); if (g_udp_fd < 0) return; } ssize_t sent = send(g_udp_fd, data, len, MSG_DONTWAIT); - (void)sent; - if (sent < 0 && (errno == ECONNREFUSED || errno == ENETUNREACH || errno == EHOSTUNREACH)) { + if (sent >= 0) { + g_udp_total_sent += sent; + if (!g_udp_sent_first) { ULOG("first send OK: %zd bytes, total=%lld", sent, (long long)g_udp_total_sent); g_udp_sent_first = 1; } + return; + } + g_udp_send_errors++; + if (g_udp_send_errors <= 3 || g_udp_send_errors % 100 == 0) + ULOG("send() error #%d: %s (errno=%d)", g_udp_send_errors, strerror(errno), errno); + if (errno == ECONNREFUSED || errno == ENETUNREACH || errno == EHOSTUNREACH) { close(g_udp_fd); g_udp_fd = -1; } @@ -66,13 +85,19 @@ void udp_trampoline(void* arg) { void udp_log_set_target(const char* ip, int port) { g_udp_buf_len = 0; g_udp_stopping = 0; + g_udp_sent_first = 0; + g_udp_total_sent = 0; + g_udp_send_errors = 0; + if (g_udp_fd >= 0) { close(g_udp_fd); g_udp_fd = -1; } if (ip && port > 0) { memset(&g_udp_addr, 0, sizeof(g_udp_addr)); g_udp_addr.sin_family = AF_INET; g_udp_addr.sin_port = htons((unsigned short)port); g_udp_has_addr = (inet_pton(AF_INET, ip, &g_udp_addr.sin_addr) == 1); + ULOG("set target %s:%d -> has_addr=%d", ip, port, g_udp_has_addr); } else { g_udp_has_addr = 0; + ULOG("set target disabled (ip=%s port=%d)", ip ? ip : "null", port); } } @@ -103,6 +128,7 @@ void udp_log_cancel_timer(void) { void udp_log_stop(void) { g_udp_stopping = 1; if (g_udp_buf_len > 0) { udp_do_send(g_udp_buf, g_udp_buf_len); g_udp_buf_len = 0; } + ULOG("stop: total_sent=%lld bytes, errors=%d", (long long)g_udp_total_sent, g_udp_send_errors); if (g_udp_fd >= 0) { close(g_udp_fd); g_udp_fd = -1; } g_udp_has_addr = 0; }