diff --git a/lib/debug_config.c b/lib/debug_config.c index 95e88f1b..617ff3b4 100644 --- a/lib/debug_config.c +++ b/lib/debug_config.c @@ -376,6 +376,8 @@ void debug_output(debug_level_t level, debug_category_t category_idx, int offset = 0; size_t remaining = BUFFER_SIZE; +#define BUF_ADVANCE(n) do { offset += (int)(n); if (offset > (int)BUFFER_SIZE - 1) offset = (int)BUFFER_SIZE - 1; remaining = BUFFER_SIZE - (size_t)offset; } while(0) + va_list args; va_start(args, format); @@ -386,42 +388,35 @@ void debug_output(debug_level_t level, debug_category_t category_idx, struct tm* tm_info = localtime(&tv_sec); char time_str[32]; strftime(time_str, sizeof(time_str), "%H:%M:%S", tm_info); - offset += snprintf(buffer + offset, remaining, "[%s-%03ld.%03ld] ", time_str, tv.tv_usec / 1000, tv.tv_usec % 1000); - remaining = BUFFER_SIZE - offset; + BUF_ADVANCE(snprintf(buffer + offset, remaining, "[%s-%03ld.%03ld] ", time_str, tv.tv_usec / 1000, tv.tv_usec % 1000)); /* Add thread marker if set */ if (g_debug_config.thread_marker > 0) { - offset += snprintf(buffer + offset, remaining, "[%d] ", g_debug_config.thread_marker); - remaining = BUFFER_SIZE - offset; + BUF_ADVANCE(snprintf(buffer + offset, remaining, "[%d] ", g_debug_config.thread_marker)); } /* Add level */ const char* level_name = get_level_name(level); - offset += snprintf(buffer + offset, remaining, "[%s] ", level_name); - remaining = BUFFER_SIZE - offset; + BUF_ADVANCE(snprintf(buffer + offset, remaining, "[%s] ", level_name)); /* Add category name */ const char* cat_name = debug_get_category_name(category_idx); - offset += snprintf(buffer + offset, remaining, "[%s] ", cat_name); - remaining = BUFFER_SIZE - offset; + BUF_ADVANCE(snprintf(buffer + offset, remaining, "[%s] ", cat_name)); /* Add file:line if enabled */ if (g_debug_config.file_line_enabled && file) { const char* bn = strrchr(file, '/'); bn = bn ? bn + 1 : file; - offset += snprintf(buffer + offset, remaining, "(%s:%d) ", bn, line); - remaining = BUFFER_SIZE - offset; + BUF_ADVANCE(snprintf(buffer + offset, remaining, "(%s:%d) ", bn, line)); } /* Add function name if enabled */ if (g_debug_config.function_name_enabled && function) { - offset += snprintf(buffer + offset, remaining, "%s() ", function); - remaining = BUFFER_SIZE - offset; + BUF_ADVANCE(snprintf(buffer + offset, remaining, "%s() ", function)); } /* Add the actual message */ - offset += vsnprintf(buffer + offset, remaining, format, args); - remaining = BUFFER_SIZE - offset; + BUF_ADVANCE(vsnprintf(buffer + offset, remaining, format, args)); /* Add newline */ if (remaining > 1) { @@ -431,6 +426,7 @@ void debug_output(debug_level_t level, debug_category_t category_idx, buffer[BUFFER_SIZE - 2] = '\n'; buffer[BUFFER_SIZE - 1] = '\0'; } +#undef BUF_ADVANCE va_end(args); diff --git a/lib/libopus/Makefile.am b/lib/libopus/Makefile.am index ab7028b5..5b22c4a5 100644 --- a/lib/libopus/Makefile.am +++ b/lib/libopus/Makefile.am @@ -1,5 +1,8 @@ noinst_LIBRARIES = libopus_internal.a +# Override global CFLAGS to disable ASAN (makes .a incompatible with PIE linking) +CFLAGS = -g -Og -fno-omit-frame-pointer -fno-pie -fPIC + libopus_internal_a_SOURCES = \ celt/bands.c \ celt/celt.c \ diff --git a/lib/strbuf.c b/lib/strbuf.c new file mode 100644 index 00000000..1a066dba --- /dev/null +++ b/lib/strbuf.c @@ -0,0 +1,79 @@ +/* + * strbuf.c — safe growable printf buffer + */ +#include "strbuf.h" +#include "mem.h" +#include +#include +#include + +#define STRBUF_MIN_CAP 16 + +void strbuf_init(struct strbuf *sb, size_t initial_cap) { + if (initial_cap < STRBUF_MIN_CAP) initial_cap = STRBUF_MIN_CAP; + sb->buf = u_malloc(initial_cap); + if (sb->buf) { sb->buf[0] = '\0'; sb->cap = initial_cap; } + else { sb->cap = 0; } + sb->len = 0; +} + +void strbuf_free(struct strbuf *sb) { + if (sb->buf) { u_free(sb->buf); sb->buf = NULL; } + sb->cap = 0; sb->len = 0; +} + +char *strbuf_detach(struct strbuf *sb) { + char *ret = sb->buf; + sb->buf = NULL; sb->cap = 0; sb->len = 0; + return ret; +} + +void strbuf_reset(struct strbuf *sb) { + if (sb->buf) sb->buf[0] = '\0'; + sb->len = 0; +} + +static int strbuf_grow(struct strbuf *sb, size_t need) { + size_t nc = sb->cap ? sb->cap : STRBUF_MIN_CAP; + while (nc < need) { nc *= 2; if (nc < STRBUF_MIN_CAP) return -1; } + char *n = u_realloc(sb->buf, nc); + if (!n) return -1; + sb->buf = n; sb->cap = nc; + return 0; +} + +int strbuf_addf(struct strbuf *sb, const char *fmt, ...) { + va_list ap; + + /* measure */ + va_start(ap, fmt); + int need = vsnprintf(NULL, 0, fmt, ap); + va_end(ap); + if (need < 0) { return -1; } + + /* grow */ + size_t want = sb->len + (size_t)need + 1; /* +1 for '\0' */ + if (want > sb->cap && strbuf_grow(sb, want) < 0) return -1; + + /* write */ + va_start(ap, fmt); + int w = vsnprintf(sb->buf + sb->len, sb->cap - sb->len, fmt, ap); + va_end(ap); + if (w > 0) sb->len += (size_t)w; + return w; +} + +void strbuf_addc(struct strbuf *sb, char c) { + if (sb->len + 2 > sb->cap && strbuf_grow(sb, sb->len + 2) < 0) return; + sb->buf[sb->len++] = c; + sb->buf[sb->len] = '\0'; +} + +void strbuf_adds(struct strbuf *sb, const char *s) { + if (!s) return; + size_t slen = strlen(s); + if (sb->len + slen + 1 > sb->cap && strbuf_grow(sb, sb->len + slen + 1) < 0) return; + memcpy(sb->buf + sb->len, s, slen); + sb->len += slen; + sb->buf[sb->len] = '\0'; +} diff --git a/lib/strbuf.h b/lib/strbuf.h new file mode 100644 index 00000000..48946e99 --- /dev/null +++ b/lib/strbuf.h @@ -0,0 +1,50 @@ +/* + * strbuf.h — safe growable printf buffer + * + * Replaces snprintf(buf+off, sizeof(buf)-off, ...) chains that crash + * on overflow (snprintf returns un-truncated size, offset exceeds buffer, + * next call gets negative size → FORTIFY abort on Bionic/Android). + * + * Usage: + * struct strbuf sb = strbuf_new(); + * strbuf_addf(&sb, "{\"ver\":\"%d\"", ver); + * strbuf_addf(&sb, ",\"name\":\"%s\"", name); + * strbuf_addf(&sb, "}"); + * // use strbuf_str(&sb) or strbuf_detach(&sb) + * strbuf_free(&sb); + */ +#ifndef STRBUF_H +#define STRBUF_H + +#include +#include + +struct strbuf { + char *buf; /* allocated buffer, always null-terminated */ + size_t cap; /* total capacity (including '\0') */ + size_t len; /* current length (excluding '\0') */ +}; + +/* zero-init on stack: struct strbuf sb = {0}; */ +#define strbuf_new() ((struct strbuf){NULL, 0, 0}) + +void strbuf_init(struct strbuf *sb, size_t initial_cap); +void strbuf_free(struct strbuf *sb); + +/* returns buf and detaches — caller must u_free(). sb is reset to empty */ +char *strbuf_detach(struct strbuf *sb); + +/* reset len=0, keep allocation */ +void strbuf_reset(struct strbuf *sb); + +/* printf-like append. auto-grows. returns written bytes or -1 on OOM */ +int strbuf_addf(struct strbuf *sb, const char *fmt, ...) + __attribute__((format(printf, 2, 3))); + +void strbuf_addc(struct strbuf *sb, char c); +void strbuf_adds(struct strbuf *sb, const char *s); + +/* safe accessor — never returns NULL */ +#define strbuf_str(sb) ((sb)->buf ? (sb)->buf : "") + +#endif /* STRBUF_H */ diff --git a/src/chat/chat_core.c b/src/chat/chat_core.c index a80b7fea..673604b9 100644 --- a/src/chat/chat_core.c +++ b/src/chat/chat_core.c @@ -593,9 +593,10 @@ int chat_member_tags_commit(struct chat_member_tags* t) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: tags_commit — invalid state", CC_ID); return -1; } char json[4096]; int off = 0, ver = 1; +#define TAGS_ADVANCE(n) do { if (off < (int)sizeof(json) - 10) off += (int)(n); else off = (int)sizeof(json); } while(0) int vi = tags_find(t, "ver"); if (vi >= 0) ver = atoi(t->keys[vi].val) + 1; - off += snprintf(json + off, sizeof(json) - off, "{\"ver\":\"%d\"", ver); + TAGS_ADVANCE(snprintf(json, sizeof(json), "{\"ver\":\"%d\"", ver)); for (int i = 0; i < t->key_count; i++) { struct chat_member_tags_key* k = &t->keys[i]; @@ -621,9 +622,10 @@ int chat_member_tags_commit(struct chat_member_tags* t) { } else { out_val = k->val; } - off += snprintf(json + off, sizeof(json) - off, ",\"%s\":\"%s\"", k->key, out_val); + TAGS_ADVANCE(snprintf(json + off, sizeof(json) - (size_t)off, ",\"%s\":\"%s\"", k->key, out_val)); } - off += snprintf(json + off, sizeof(json) - off, "}"); + TAGS_ADVANCE(snprintf(json + off, sizeof(json) - (size_t)off, "}")); +#undef TAGS_ADVANCE uint8_t ch_ed_priv[32]; if (topo_node_sqlite_channel_get_priv(g_cc.db, t->ch_id, ch_ed_priv) != 0) { diff --git a/src/media_delivery/media_delivery.c b/src/media_delivery/media_delivery.c index 75eb5c10..b149e163 100644 --- a/src/media_delivery/media_delivery.c +++ b/src/media_delivery/media_delivery.c @@ -712,43 +712,50 @@ static void stream_send_chunk_cb(struct ll_queue* q, void* arg) { struct media_delivery_ctx* md = sc->md; if (!md || sc->remaining == 0) { u_free(sc); return; } - /* read up to MD_CHUNK_SIZE bytes */ - size_t to_read = sc->remaining < MD_CHUNK_SIZE ? (size_t)sc->remaining : MD_CHUNK_SIZE; - uint8_t buf[MD_CHUNK_SIZE]; - fseeko(sc->file, (off_t)sc->offset, SEEK_SET); - size_t rd = fread(buf, 1, to_read, sc->file); - if (rd == 0) { /* EOF or error */ - u_free(sc); return; + while (sc->remaining > 0) { + size_t to_read = sc->remaining < MD_CHUNK_SIZE ? (size_t)sc->remaining : MD_CHUNK_SIZE; + uint8_t buf[MD_CHUNK_SIZE]; + fseeko(sc->file, (off_t)sc->offset, SEEK_SET); + size_t rd = fread(buf, 1, to_read, sc->file); + if (rd == 0) { u_free(sc); return; } + + uint8_t pkt[MEDIA_BLOCK_CHUNK_HDR_SIZE + MD_CHUNK_SIZE]; + struct media_pkt_block_chunk* ch = (struct media_pkt_block_chunk*)pkt; + memset(ch, 0, MEDIA_BLOCK_CHUNK_HDR_SIZE); + ch->subcmd = MEDIA_SUBCMD_BLOCK_CHUNK; + memcpy(ch->media_id, sc->media_id, 16); + memcpy(ch->block_id, sc->block_id, 16); + ch->chunk = sc->chunk; + ch->offset = (uint32_t)(sc->offset - sc->block_start); + ch->data_len = (uint16_t)rd; + memcpy(pkt + MEDIA_BLOCK_CHUNK_HDR_SIZE, buf, rd); + + md_send(md->inst, sc->group_id, sc->dst_node_id, pkt, MEDIA_BLOCK_CHUNK_HDR_SIZE + rd); + + sc->offset += rd; + sc->remaining -= rd; + + if (sc->remaining == 0) break; + + if (!etcp_router_send_q_has_room(md->inst, sc->group_id, sc->dst_node_id, ETCP_RT_ID_MEDIA_DELIVERY)) { + etcp_router_on_send_ready(md->inst, sc->group_id, + sc->dst_node_id, ETCP_RT_ID_MEDIA_DELIVERY, + &sc->waiter, stream_send_chunk_cb, sc); + return; + } } - /* build CHUNK packet */ - uint8_t pkt[MEDIA_BLOCK_CHUNK_HDR_SIZE + MD_CHUNK_SIZE]; - struct media_pkt_block_chunk* ch = (struct media_pkt_block_chunk*)pkt; - memset(ch, 0, MEDIA_BLOCK_CHUNK_HDR_SIZE); - ch->subcmd = MEDIA_SUBCMD_BLOCK_CHUNK; - memcpy(ch->media_id, sc->media_id, 16); - memcpy(ch->block_id, sc->block_id, 16); - ch->chunk = sc->chunk; - ch->offset = (uint32_t)(sc->offset - sc->block_start); - ch->data_len = (uint16_t)rd; - memcpy(pkt + MEDIA_BLOCK_CHUNK_HDR_SIZE, buf, rd); - - md_send(md->inst, sc->group_id, sc->dst_node_id, pkt, MEDIA_BLOCK_CHUNK_HDR_SIZE + rd); - - sc->offset += rd; - sc->remaining -= rd; - - if (sc->remaining == 0) { - /* all data sent — send BLOCK_DONE */ + /* all data sent — send BLOCK_DONE */ + { + uint32_t chunk = sc->chunk; uint8_t done_pkt[MEDIA_BLOCK_DONE_SIZE]; struct media_pkt_block_done* bd = (struct media_pkt_block_done*)done_pkt; memset(bd, 0, sizeof(*bd)); bd->subcmd = MEDIA_SUBCMD_BLOCK_DONE; memcpy(bd->media_id, sc->media_id, 16); memcpy(bd->block_id, sc->block_id, 16); - bd->chunk = sc->chunk; + bd->chunk = chunk; bd->total_size = (uint32_t)sc->block_data_len; - /* block_sig is zero — receiver will verify vs block_sigs from media_index_result */ md_send(md->inst, sc->group_id, sc->dst_node_id, done_pkt, MEDIA_BLOCK_DONE_SIZE); md_file_load_dec(md, sc->media_id, sc->dst_node_id); @@ -756,12 +763,7 @@ static void stream_send_chunk_cb(struct ll_queue* q, void* arg) { fclose(sc->file); u_free(sc); DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_DONE sent chunk=%d total=%u", - MD_ID, sc->chunk, bd->total_size); - } else { - /* continue streaming — register waiter for backpressure */ - etcp_router_on_send_ready(md->inst, sc->group_id, - sc->dst_node_id, ETCP_RT_ID_MEDIA_DELIVERY, - &sc->waiter, stream_send_chunk_cb, sc); + MD_ID, chunk, bd->total_size); } } diff --git a/src/routing_layer/etcp_router.c b/src/routing_layer/etcp_router.c index 176be211..9436341d 100644 --- a/src/routing_layer/etcp_router.c +++ b/src/routing_layer/etcp_router.c @@ -1404,6 +1404,14 @@ void etcp_router_cancel_send_ready(struct UTUN_INSTANCE* inst, uint64_t group_id queue_waiter_cancel(rconn->send_q, h); } +int etcp_router_send_q_has_room(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t node_id, uint8_t svc_id) { + if (!inst) return 0; + struct ETCP_ROUTER_CONN* rconn = router_conn_find(inst, group_id, node_id, svc_id); + if (!rconn || !rconn->send_q) return 0; + return rconn->send_q->count <= rconn->send_q->threshold_max_packets + && (rconn->send_q->threshold_max_bytes == 0 || rconn->send_q->total_bytes <= rconn->send_q->threshold_max_bytes); +} + // ==================================================================== // Seq-connection API // ==================================================================== diff --git a/src/routing_layer/etcp_router.h b/src/routing_layer/etcp_router.h index bb1de883..84f9a744 100644 --- a/src/routing_layer/etcp_router.h +++ b/src/routing_layer/etcp_router.h @@ -228,6 +228,9 @@ void etcp_router_on_send_ready(struct UTUN_INSTANCE* inst, uint64_t group_id, ui void etcp_router_cancel_send_ready(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t node_id, uint8_t svc_id, struct queue_waiter_handle* h); +// Backpressure: проверить есть ли место в send_q (без регистрации waiter) +int etcp_router_send_q_has_room(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t node_id, uint8_t svc_id); + // Сбросить состояние роутера для конкретного peer+svc в группе (перезапуск удалённой стороны). void etcp_router_conn_restart(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t remote_node_id, uint8_t svc_id); diff --git a/src/routing_layer/topo_node.c b/src/routing_layer/topo_node.c index 226dc62c..891286ca 100644 --- a/src/routing_layer/topo_node.c +++ b/src/routing_layer/topo_node.c @@ -654,6 +654,17 @@ int topo_node_format_all(struct TOPO_GROUP* group, char* buf, size_t buf_size) { int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group) { if (!instance || !group) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args"); return -1; } + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "TOPO_UPD_SELF: ENTER inst=%p nid=0x%016llx grp=%p grp_inst=%p grp_type=%d", + (void*)instance, (unsigned long long)instance->node_id, (void*)group, (void*)group->instance, group->group_type); + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "TOPO_UPD_SELF: topo_groups=%p etcp=%p tcp=%p reg=%p", + (void*)instance->topo_groups, (void*)instance->etcp_sockets, (void*)instance->tcp_sockets, + instance->topo_groups ? (void*)instance->topo_groups->node_registry : NULL); + if (instance->topo_groups) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "TOPO_UPD_SELF: pools sm4=%p a4=%p sm6=%p a6=%p sub4=%p sub6=%p", + (void*)instance->topo_groups->v4_sock_meta_pool, (void*)instance->topo_groups->v4_addr_pool, + (void*)instance->topo_groups->v6_sock_meta_pool, (void*)instance->topo_groups->v6_addr_pool, + (void*)instance->topo_groups->v4_subnet_pool, (void*)instance->topo_groups->v6_subnet_pool); + } size_t name_len = 0; if (instance->name[0]) { name_len = strlen(instance->name); if (name_len > 63) name_len = 63; } int vc = 0, vc6 = 0; @@ -750,7 +761,11 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR lq->last_ver = saved_ver; } e_sock = instance->etcp_sockets; + int etcp_iter_cnt = 0; while (e_sock) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "TOPO_UPD_SELF: etcp_iter[%d] e_sock=%p type=%d fam=%d next=%p", + etcp_iter_cnt, (void*)e_sock, e_sock->type, e_sock->local_addr.ss_family, (void*)e_sock->next); + etcp_iter_cnt++; if (e_sock->type == CFG_SERVER_TYPE_PRIVATE || e_sock->type == CFG_SERVER_TYPE_LOCAL) { e_sock = e_sock->next; continue; } if (e_sock->local_addr.ss_family == AF_INET) { { struct TOPO_SOCKMETA4* sm = memory_pool_alloc(group->instance->topo_groups->v4_sock_meta_pool); @@ -786,7 +801,10 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR e_sock = e_sock->next; } { struct TCP_SOCKET* ts = instance->tcp_sockets; - while (ts) { if (ts->type == CFG_SERVER_TYPE_PRIVATE || ts->type == CFG_SERVER_TYPE_LOCAL) { ts = ts->next; continue; } + int tcp_iter_cnt = 0; + while (ts) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "TOPO_UPD_SELF: tcp_iter[%d] ts=%p type=%d fam=%d next=%p", + tcp_iter_cnt, (void*)ts, ts->type, ts->local_addr.ss_family, (void*)ts->next); + tcp_iter_cnt++; if (ts->type == CFG_SERVER_TYPE_PRIVATE || ts->type == CFG_SERVER_TYPE_LOCAL) { ts = ts->next; continue; } struct sockaddr_storage* addr = ts->interface_addr.ss_family ? &ts->interface_addr : &ts->local_addr; if (!addr || !addr->ss_family) { ts = ts->next; continue; } if (addr->ss_family == AF_INET) { diff --git a/src/transport_layer/etcp.c b/src/transport_layer/etcp.c index 82dd9769..64d741e9 100644 --- a/src/transport_layer/etcp.c +++ b/src/transport_layer/etcp.c @@ -19,6 +19,7 @@ #include // For UINT16_MAX #include // For strftime in metrics snapshot #include "../lib/mem.h" +#include "../lib/strbuf.h" #include "../lib/memory_pool.h" // Enable comprehensive debug output for ETCP module @@ -309,34 +310,36 @@ void etcp_cbk_fire(struct ETCP_CONN* conn, int event) { static void etcp_on_up(struct ETCP_CONN* etcp) { int total_links = 0; for (struct ETCP_LINK* l = etcp->links; l; l = l->next) total_links++; uint64_t elapsed_ms = (get_time_tb() - etcp->setup_start_tb) / 10; - char links_str[256] = {0}; int pp = 0; + struct strbuf sb = strbuf_new(); int any = 0; for (struct ETCP_LINK* l = etcp->links; l; l = l->next) { if (l->link_status) { - if (pp) pp += snprintf(links_str + pp, sizeof(links_str) - pp, ", "); - pp += snprintf(links_str + pp, sizeof(links_str) - pp, "%s", sockaddr_storage_to_str(&l->remote_addr).str); + if (any) strbuf_adds(&sb, ", "); any = 1; + strbuf_adds(&sb, sockaddr_storage_to_str(&l->remote_addr).str); } } DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection UP (%d/%d links, mtu=%d, reinit=%u)%s%s", etcp->log_name, etcp->links_up, total_links, etcp->mtu, etcp->reinit_count, - pp > 0 ? " — up: " : "", links_str); + any > 0 ? " — up: " : "", strbuf_str(&sb)); + strbuf_free(&sb); (void)elapsed_ms; etcp_cbk_fire(etcp, ETCP_CBK_EVENT_UP); etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_UP); } static void etcp_on_down(struct ETCP_CONN* etcp, struct ETCP_LINK* down_link) { - char links_str[256] = {0}; int pp = 0; + struct strbuf sb = strbuf_new(); int any = 0; for (struct ETCP_LINK* l = etcp->links; l; l = l->next) { if (l == down_link) continue; - if (pp) pp += snprintf(links_str + pp, sizeof(links_str) - pp, ", "); - pp += snprintf(links_str + pp, sizeof(links_str) - pp, "%s", sockaddr_storage_to_str(&l->remote_addr).str); + if (any) strbuf_adds(&sb, ", "); any = 1; + strbuf_adds(&sb, sockaddr_storage_to_str(&l->remote_addr).str); } if (down_link) DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection DOWN — link %s lost%s%s", etcp->log_name, sockaddr_storage_to_str(&down_link->remote_addr).str, - pp > 0 ? ", remaining: " : "", links_str); + any > 0 ? ", remaining: " : "", strbuf_str(&sb)); else - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection DOWN — closing, links: %s", etcp->log_name, pp ? links_str : "(none)"); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection DOWN — closing, links: %s", etcp->log_name, any ? strbuf_str(&sb) : "(none)"); + strbuf_free(&sb); etcp_cbk_fire(etcp, ETCP_CBK_EVENT_DOWN); etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_DOWN); } @@ -1037,13 +1040,14 @@ static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp) {// вызыв DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); struct ETCP_DGRAM* dgram; if (etcp->tx_state!=ETCP_TX_STATE_DATA_WAIT) { - char l_status[256]={0}; + struct strbuf sb = strbuf_new(); struct ETCP_LINK* link = etcp->links; while (link) { - snprintf (l_status+strlen(l_status), 256-strlen(l_status), "L%d%d%s ", link->recv_keepalive, link->remote_keepalive, (link->shaper_timer)?"wait":"rdy"); + strbuf_addf(&sb, "L%d%d%s ", link->recv_keepalive, link->remote_keepalive, (link->shaper_timer)?"wait":"rdy"); link = link->next; } - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] TX state: %d (link not ready, skip send) %s", etcp->log_name, etcp->tx_state, l_status); + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] TX state: %d (link not ready, skip send) %s", etcp->log_name, etcp->tx_state, strbuf_str(&sb)); + strbuf_free(&sb); return; } dgram = etcp_request_pkt(etcp); diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index ca992f26..a0a58bc0 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -54,6 +54,28 @@ void tcp_server_on_link(struct stcp_link *link, void *arg) { struct ETCP_CONN *conn = instance_find_conn(inst, node_id); DEBUG_INFO(DEBUG_CATEGORY_ETCP, "tcp_server_on_link: %s conn=%p peer=0x%016llx pubkey=%016llx", conn ? "existing" : "NEW", (void*)conn, (unsigned long long)node_id, *(const uint64_t*)pubkey); + + /* Синхронизация got_initial_pkt и session_id: если грязный/другой → сбросить */ + { + uint8_t client_gop = stcp_link_get_peer_got_initial_pkt(link); + uint32_t client_sid = stcp_link_get_peer_session_id(link); + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "tcp_server_on_link: server_gop=%d server_sid=%08x client_gop=%d client_sid=%08x conn=%p", + conn ? conn->got_initial_pkt : 0, conn ? conn->session_id : 0, + client_gop, client_sid, (void*)conn); + if (conn) { + if (conn->got_initial_pkt == 1 && client_gop == 0) { + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "tcp_server_on_link: server dirty (gop=1) client clean (gop=0) → reinit"); + etcp_conn_reinit(conn, "stcp client clean"); + } + if (conn->session_id != client_sid) { + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "tcp_server_on_link: session_id changed my=%08x client=%08x → reinit", + conn->session_id, client_sid); + etcp_conn_reinit(conn, "stcp session changed"); + } + conn->session_id = client_sid; + } + } + if (!conn) { conn = etcp_connection_create(inst, NULL); if (!conn) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "tcp_server_on_link: etcp_connection_create failed"); return; } @@ -138,6 +160,7 @@ static void etcp_link_send_init(struct ETCP_LINK* link, uint8_t reset, uint8_t c DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "link=%p, is_server=%d, reset=%d, collision=%d", link, link ? link->is_server : -1, reset, collision); if (!link || !link->etcp || !link->etcp->instance) return; + if (link->is_tcp) return; // TCP uses STCP handshake instead of ETCP INIT struct ETCP_DGRAM* dgram = u_malloc(PACKET_DATA_SIZE); if (!dgram) { @@ -197,8 +220,11 @@ static void etcp_link_send_init(struct ETCP_LINK* link, uint8_t reset, uint8_t c uint8_t salt[SC_PUBKEY_ENC_SALT_SIZE]; random_bytes(salt, sizeof(salt)); +#pragma GCC diagnostic push +#pragma GCC diagnostic ignored "-Wstringop-overflow" memcpy(dgram->data + offset, salt, SC_PUBKEY_ENC_SALT_SIZE); offset += SC_PUBKEY_ENC_SALT_SIZE; +#pragma GCC diagnostic pop uint8_t obfuscated_pubkey[SC_PUBKEY_SIZE]; sc_obfuscate_pubkey(salt, link->etcp->crypto_ctx.peer_public_key, @@ -211,9 +237,12 @@ static void etcp_link_send_init(struct ETCP_LINK* link, uint8_t reset, uint8_t c link->etcp->instance->my_keys.public_key[0], link->etcp->instance->my_keys.public_key[1], link->etcp->instance->my_keys.public_key[2], link->etcp->instance->my_keys.public_key[3], link->etcp->crypto_ctx.session_key[0], link->etcp->crypto_ctx.session_key[1], - link->etcp->crypto_ctx.session_key[2], link->etcp->crypto_ctx.session_key[3]); + link->etcp->crypto_ctx.session_key[2], link->etcp->crypto_ctx.session_key[3]); +#pragma GCC diagnostic push +#pragma GCC diagnostic ignored "-Wstringop-overflow" memcpy(dgram->data + offset, obfuscated_pubkey, SC_PUBKEY_SIZE); offset += SC_PUBKEY_SIZE; +#pragma GCC diagnostic pop dgram->data_len = offset; @@ -243,6 +272,7 @@ static void etcp_link_init_timer_cbk(void* arg) { if (link->link_state == 3 && link->initialized) { DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] init_timer: SUPPRESSED state=%d init=%d link_status=%d recv_ka=%d remote_ka=%d links_up=%d", link->etcp->log_name, link->link_state, link->initialized, link->link_status, link->recv_keepalive, link->remote_keepalive, link->etcp->links_up); return; } if (link->etcp->links_up > 0) { DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] init_timer: SUPPRESSED (links_up=%d > 0)", link->etcp->log_name, link->etcp->links_up); return; } if (link->etcp->fin_wait) { DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] init_timer: SUPPRESSED (fin_wait)", link->etcp->log_name); return; } + if (link->is_tcp) { etcp_tcp_link_start_reconnect(link); return; } // TCP uses STCP handshake, not ETCP INIT if (link->etcp->got_initial_pkt == 0) etcp_link_send_init(link,1,0); else etcp_link_send_init(link,0,0); } @@ -352,6 +382,11 @@ static void keepalive_timer_cb(void* arg) { // Check if all links are down and start recovery if needed (client only) if (link->is_server == 0 && etcp_all_links_down(link->etcp)) { + if (link->is_tcp) { + if (link->tcp_link) { stcp_link_close(link->tcp_link); link->tcp_link = NULL; } + etcp_tcp_link_start_reconnect(link); + return; + } DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] All links are down, starting recovery", link->etcp->log_name); etcp_link_enter_reinit(link);// keepalive timr после reinit не нужен return; @@ -376,6 +411,7 @@ static void keepalive_timer_cb(void* arg) { link->link_status = 0; etcp_fire_link_status_cbk(link, link->link_state, old_link_status); etcp_on_link_down(link->etcp, link); + if (link->is_tcp && link->tcp_link) { stcp_link_close(link->tcp_link); link->tcp_link = NULL; } if (old_link_status) { DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link %d down: addr=%s ka=%d remote_ka=%d state=%d init=%d tmo: %llu>%llu els=%llums", link->etcp->log_name, link->local_link_id, sockaddr_storage_to_str(&link->remote_addr).str, link->recv_keepalive, link->remote_keepalive, link->link_state, link->initialized, (unsigned long long)timeout_units, (unsigned long long)elapsed, (unsigned long long)(elapsed/10)); } @@ -1087,9 +1123,29 @@ void etcp_link_close(struct ETCP_LINK* link) { void etcp_link_enter_ready_tcp(struct ETCP_LINK *link) { if (!link || !link->etcp) return; struct ETCP_CONN *etcp = link->etcp; + + /* Синхронизация got_initial_pkt: если я грязный (1) а пир чистый (0) → сбросить */ + if (link->tcp_link) { + uint8_t peer_gop = stcp_link_get_peer_got_initial_pkt(link->tcp_link); + uint32_t peer_sid = stcp_link_get_peer_session_id(link->tcp_link); + if (etcp->got_initial_pkt == 1 && peer_gop == 0) { + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] TCP ready: I'm dirty (gop=1) but peer clean (gop=0) → reinit", + etcp->log_name); + etcp_conn_reinit(etcp, "tcp peer clean"); + } + if (etcp->session_id != peer_sid) { + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] TCP ready: session_id mismatch my=%08x peer=%08x → reinit", + etcp->log_name, etcp->session_id, peer_sid); + etcp_conn_reinit(etcp, "tcp session changed"); + } + etcp->session_id = peer_sid; + } + link->initialized = 1; link->link_state = 3; link->link_status = 1; link->recv_keepalive = 1; if (!link->mtu_remote) link->mtu_remote = link->mtu; + etcp->got_initial_pkt = 1; + etcp->reset_done = 1; etcp->initialized = 1; etcp->links_up = 1; etcp->tcp_link_count++; if (etcp->tx_state == 0) etcp->tx_state = ETCP_TX_STATE_DATA_WAIT; if (link->tcp_link) { const uint8_t* ed = stcp_link_get_peer_ed25519_pubkey(link->tcp_link); if (ed) memcpy(etcp->peer_ed25519_pubkey, ed, SC_PUBKEY_SIZE); } @@ -1110,7 +1166,7 @@ static void tcp_link_reconnect_cb(void *arg) { if (link->tcp_reconnect_delay_ms == 0) link->tcp_reconnect_delay_ms = 1000; DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] TCP link %d reconnect attempt (delay=%ums)", link->etcp->log_name, link->local_link_id, link->tcp_reconnect_delay_ms); uint16_t port = ntohs(((struct sockaddr_in *)&link->remote_addr)->sin_port); - struct stcp_link_config tcp_cfg = {.ua = link->etcp->instance->ua, .my_keys = &link->etcp->instance->my_keys, .inst = link->etcp->instance, .peer_pubkey = link->etcp->crypto_ctx.peer_public_key, .peer_pubkey_mode = 0, .remote_addr = &link->remote_addr, .remote_port = port}; + struct stcp_link_config tcp_cfg = {.ua = link->etcp->instance->ua, .my_keys = &link->etcp->instance->my_keys, .inst = link->etcp->instance, .peer_pubkey = link->etcp->crypto_ctx.peer_public_key, .peer_pubkey_mode = 0, .remote_addr = &link->remote_addr, .remote_port = port, .got_initial_pkt = link->etcp->got_initial_pkt, .session_id = link->etcp->session_id}; struct stcp_link *sl = stcp_link_connect(&tcp_cfg); if (!sl) { link->tcp_reconnect_delay_ms *= 2; if (link->tcp_reconnect_delay_ms > 30000) link->tcp_reconnect_delay_ms = 30000; link->tcp_reconnect_timer = uasync_set_timeout(link->etcp->instance->ua, (int)(link->tcp_reconnect_delay_ms * 10), link, tcp_link_reconnect_cb, "tcp_rct"); return; } link->tcp_link = sl; stcp_link_set_etcp_conn(sl, link->etcp); stcp_link_set_etcp_link(sl, link); @@ -1120,7 +1176,7 @@ static void tcp_link_reconnect_cb(void *arg) { void etcp_tcp_link_start_connect(struct ETCP_LINK *link, struct sockaddr_storage *addr, uint16_t port) { if (!link || !link->etcp || !addr) return; memcpy(&link->remote_addr, addr, sizeof(*addr)); - struct stcp_link_config tcp_cfg = {.ua = link->etcp->instance->ua, .my_keys = &link->etcp->instance->my_keys, .inst = link->etcp->instance, .peer_pubkey = link->etcp->crypto_ctx.peer_public_key, .peer_pubkey_mode = 0, .remote_addr = addr, .remote_port = port}; + struct stcp_link_config tcp_cfg = {.ua = link->etcp->instance->ua, .my_keys = &link->etcp->instance->my_keys, .inst = link->etcp->instance, .peer_pubkey = link->etcp->crypto_ctx.peer_public_key, .peer_pubkey_mode = 0, .remote_addr = addr, .remote_port = port, .got_initial_pkt = link->etcp->got_initial_pkt, .session_id = link->etcp->session_id}; struct stcp_link *sl = stcp_link_connect(&tcp_cfg); if (!sl) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_tcp_link_start_connect: stcp_link_connect failed"); return; } link->tcp_link = sl; stcp_link_set_etcp_conn(sl, link->etcp); stcp_link_set_etcp_link(sl, link); diff --git a/src/transport_layer/node_conn_direct.c b/src/transport_layer/node_conn_direct.c index 6aefa23a..6dd80f79 100644 --- a/src/transport_layer/node_conn_direct.c +++ b/src/transport_layer/node_conn_direct.c @@ -442,8 +442,13 @@ static void ncd_recv_control_handler(struct ETCP_CONN* conn, struct ll_entry* e) etcp_send(conn, qe); } } } else { - DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] deferred close conn for CLOSE from 0x%016llx", (unsigned long long)sender_id); - if (conn) uasync_call_soon(conn->instance->ua, conn, ncd_deferred_close_conn); + DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] deferred close entry for CLOSE from 0x%016llx", (unsigned long long)sender_id); + if (entry->conn && entry->conn->conn_queue && entry->conn->conn_queue_entry) { + queue_remove_data(entry->conn->conn_queue, entry->conn->conn_queue_entry); + queue_entry_free(entry->conn->conn_queue_entry); + entry->conn->conn_queue_entry = NULL; entry->conn->conn_queue = NULL; + } + uasync_call_soon(entry->ua, entry, ncd_deferred_close); } break; case NCD_SUBCMD_KEEP_ALIVE: @@ -682,7 +687,7 @@ int node_conn_direct_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id, /* 2. Ищем conn через instance_find_conn */ { struct ETCP_CONN* conn = instance_find_conn(inst, node_id); - if (conn) { + if (conn && conn->state != 2) { if (conn->fin_wait) { DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open_node REUSED clearing fin_wait (existing conn) node=0x%016llx", (unsigned long long)node_id); conn->fin_wait = 0; conn->fin_wait_clear_cb = NULL; conn->fin_wait_clear_arg = NULL; @@ -907,6 +912,11 @@ void node_conn_direct_force_close(struct NODE_CONN_DIRECT* h) { etcp_conn_remove_cbk(conn, ncd_up_cb, entry); DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] force_close: removing ncd_down_cb"); etcp_conn_remove_cbk(conn, ncd_down_cb, entry); + if (conn->conn_queue && conn->conn_queue_entry) { + queue_remove_data(conn->conn_queue, conn->conn_queue_entry); + queue_entry_free(conn->conn_queue_entry); + conn->conn_queue_entry = NULL; conn->conn_queue = NULL; + } DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] force_close: scheduling deferred close"); if (conn->state != 2) uasync_call_soon(entry->ua, conn, ncd_deferred_close_conn); DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] force_close: removing from registry"); diff --git a/src/transport_layer/stcp.h b/src/transport_layer/stcp.h index 41cd6c24..54c770c2 100644 --- a/src/transport_layer/stcp.h +++ b/src/transport_layer/stcp.h @@ -15,16 +15,18 @@ extern "C" { #include "../lib/debug_config.h" #include +struct UTUN_INSTANCE; + #define STCP_MAX_MSG_SIZE 65535 #define STCP_RECV_BUF_INIT 8192 #define STCP_RECV_BUF_MAX 131072 #define STCP_HS_TIMEOUT 50000 // 5s in 0.1ms timebase units #define STCP_CONNECT_TIMEOUT 100000 // 10s in 0.1ms timebase units -#define STCP_HS_CLIENT_MIN (SC_PUBKEY_ENC_SIZE + STCP_HS_ENC_CLIENT) // 72+38=110 -#define STCP_HS_SERVER_MIN (SC_PUBKEY_ENC_SIZE + STCP_HS_ENC_SERVER) // 72+39=111 -#define STCP_HS_ENC_CLIENT 38 // ed25519_pubkey(32) + padding_size(2) + CRC32(4) -#define STCP_HS_ENC_SERVER 39 // ed25519_pubkey(32) + status(1) + padding_size(2) + CRC32(4) +#define STCP_HS_CLIENT_MIN (SC_PUBKEY_ENC_SIZE + STCP_HS_ENC_CLIENT) // 72+43=115 +#define STCP_HS_SERVER_MIN (SC_PUBKEY_ENC_SIZE + STCP_HS_ENC_SERVER) // 72+43=115 +#define STCP_HS_ENC_CLIENT 43 // ed25519_pubkey(32) + got_initial_pkt(1) + session_id(4) + padding_size(2) + CRC32(4) +#define STCP_HS_ENC_SERVER 43 // ed25519_pubkey(32) + got_initial_pkt(1) + session_id(4) + padding_size(2) + CRC32(4) #define STCP_STREAM_CLIENT_SEND 0 #define STCP_STREAM_SERVER_SEND 1 @@ -52,6 +54,8 @@ struct stcp_conn { struct UASYNC *ua; void *socket_id; + struct UTUN_INSTANCE *inst; // для lookup ETCP_CONN при хендшейке (сервер) + enum stcp_state state; uint8_t is_server; uint8_t allocated; // 1 = allocated by u_calloc, free in stcp_conn_free @@ -68,6 +72,11 @@ struct stcp_conn { uint8_t peer_ed25519_set; uint8_t my_ed25519_pubkey[SC_PUBKEY_SIZE]; + uint8_t got_initial_pkt; // local, sent to peer during handshake + uint8_t peer_got_initial_pkt; // remote, received from peer during handshake + uint32_t session_id; // local, sent to peer during handshake + uint32_t peer_session_id; // remote, received from peer during handshake + struct ll_queue *rx_queue; struct ll_queue *tx_queue; void (*tx_cb)(struct ll_queue *q, void *arg); diff --git a/src/transport_layer/stcp_client.c b/src/transport_layer/stcp_client.c index afae2544..d0c6ff90 100644 --- a/src/transport_layer/stcp_client.c +++ b/src/transport_layer/stcp_client.c @@ -54,11 +54,14 @@ static void client_send_handshake(struct stcp_conn *c, const uint8_t *server_pub memcpy(hs, salt, SC_PUBKEY_ENC_SALT_SIZE); sc_obfuscate_pubkey(salt, server_pubkey, c->my_keys.public_key, hs + SC_PUBKEY_ENC_SALT_SIZE); - uint8_t plain[34]; memcpy(plain, my_ed25519, 32); plain[32] = (uint8_t)padding; plain[33] = (uint8_t)(padding >> 8); - uint32_t crc = crc32_calc(plain, 34); + uint8_t plain[39]; memcpy(plain, my_ed25519, 32); + plain[32] = c->got_initial_pkt; + memcpy(plain + 33, &c->session_id, 4); + plain[37] = (uint8_t)padding; plain[38] = (uint8_t)(padding >> 8); + uint32_t crc = crc32_calc(plain, 39); uint8_t *enc_dst = hs + SC_PUBKEY_ENC_SIZE; - memcpy(enc_dst, plain, 34); - enc_dst[34] = (uint8_t)(crc >> 0); enc_dst[35] = (uint8_t)(crc >> 8); enc_dst[36] = (uint8_t)(crc >> 16); enc_dst[37] = (uint8_t)(crc >> 24); + memcpy(enc_dst, plain, 39); + enc_dst[39] = (uint8_t)(crc >> 0); enc_dst[40] = (uint8_t)(crc >> 8); enc_dst[41] = (uint8_t)(crc >> 16); enc_dst[42] = (uint8_t)(crc >> 24); if (sc_stream_xor(&c->stream_send, enc_dst, STCP_HS_ENC_CLIENT) != SC_OK) { u_free(hs); stcp_conn_do_close(c, 2); return; } for (int i = 0; i < padding; i++) hs[SC_PUBKEY_ENC_SIZE + STCP_HS_ENC_CLIENT + i] = (uint8_t)(salt[0] ^ i); @@ -88,10 +91,11 @@ static void client_hs_cb(struct stcp_conn *c, uint8_t *data, size_t len) { if (stcp_frame_decrypt(enc_hs, STCP_HS_ENC_SERVER, &c->stream_recv, &hs_data_len)) { stcp_conn_do_close(c, 3); return; } log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO, "stcp_client enc_hs AFTER xor", enc_hs, STCP_HS_ENC_SERVER); memcpy(c->peer_ed25519_pubkey, enc_hs, SC_PUBKEY_SIZE); c->peer_ed25519_set = 1; - uint8_t status = enc_hs[32]; - uint16_t padding_size = (uint16_t)enc_hs[33] | ((uint16_t)enc_hs[34] << 8); - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_client: server response OK status=%d padding=%u", status, padding_size); - if (status != 0) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "server handshake status=%d", status); stcp_conn_do_close(c, 4); return; } + c->peer_got_initial_pkt = enc_hs[32]; + memcpy(&c->peer_session_id, enc_hs + 33, 4); + uint16_t padding_size = (uint16_t)enc_hs[37] | ((uint16_t)enc_hs[38] << 8); + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_client: server response OK gop=%d sid=%08x padding=%u", + c->peer_got_initial_pkt, c->peer_session_id, padding_size); stcp_recv_set(c, padding_size, 0, client_hs_padding_cb); } @@ -140,6 +144,7 @@ static void client_connect_write_cb(socket_t sock, void *arg) { struct stcp_client *stcp_client_connect(struct UASYNC *ua, const char *addr, uint16_t port, struct SC_MYKEYS *keys, const uint8_t *peer_pubkey, const uint8_t *my_ed25519_pubkey, + uint8_t got_initial_pkt, uint32_t session_id, stcp_ready_cb ready_cb, void *arg, stcp_close_cb close_cb, void *close_arg) { if (!ua || !addr || !keys || !peer_pubkey || !ready_cb) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "invalid args"); return NULL; } @@ -154,6 +159,8 @@ struct stcp_client *stcp_client_connect(struct UASYNC *ua, const char *addr, uin c->on_close = close_cb; c->close_arg = close_arg; c->on_write_error = stcp_conn_do_close; c->tx_cb = stcp_tx_queue_cb; + c->got_initial_pkt = got_initial_pkt; + c->session_id = session_id; struct addrinfo hints = {0}; hints.ai_family = AF_UNSPEC; diff --git a/src/transport_layer/stcp_client.h b/src/transport_layer/stcp_client.h index 70d46a88..dea9d797 100644 --- a/src/transport_layer/stcp_client.h +++ b/src/transport_layer/stcp_client.h @@ -15,6 +15,7 @@ typedef void (*stcp_close_cb)(struct stcp_conn *conn, int err, void *arg); struct stcp_client *stcp_client_connect(struct UASYNC *ua, const char *addr, uint16_t port, struct SC_MYKEYS *keys, const uint8_t *peer_pubkey, const uint8_t *my_ed25519_pubkey, + uint8_t got_initial_pkt, uint32_t session_id, stcp_ready_cb ready_cb, void *arg, stcp_close_cb close_cb, void *close_arg); void stcp_client_destroy(struct stcp_client *cli); diff --git a/src/transport_layer/stcp_link.c b/src/transport_layer/stcp_link.c index 1dff31a5..c4d888dc 100644 --- a/src/transport_layer/stcp_link.c +++ b/src/transport_layer/stcp_link.c @@ -43,6 +43,9 @@ struct stcp_link { struct ll_queue *saved_rx_queue; // saved rx_queue for pre-closed conn cleanup uint8_t conn_pre_closed; // 1 = conn already CLOSED before stcp_link_close + + uint8_t peer_got_initial_pkt; // received from peer during handshake + uint32_t peer_session_id; // received from peer during handshake }; // ====== rx dispatch ====== @@ -110,6 +113,8 @@ static void server_accept_cb(struct stcp_conn *conn, void *arg) { link->cfg = ss->cfg; link->ready = 1; link->conn = conn; + link->peer_got_initial_pkt = conn->peer_got_initial_pkt; + link->peer_session_id = conn->peer_session_id; stcp_conn_set_on_close(conn, stcp_link_on_stcp_close, link); @@ -133,6 +138,8 @@ static void client_ready_cb(struct stcp_conn *conn, void *arg) { if (!conn) return; link->ready = 1; link->conn = conn; + link->peer_got_initial_pkt = conn->peer_got_initial_pkt; + link->peer_session_id = conn->peer_session_id; stcp_conn_set_on_close(conn, stcp_link_on_stcp_close, link); @@ -161,7 +168,7 @@ struct stcp_server *stcp_server_listen(struct stcp_link_config *cfg, uint16_t po ss->cfg = *cfg; ss->on_link = on_link; ss->on_link_arg = arg; - ss->srv = stcp_server_create(cfg->ua, port, cfg->my_keys, cfg->inst ? cfg->inst->my_ed25519_pubkey : NULL, server_accept_cb, ss, NULL, NULL, cfg->listen_family); + ss->srv = stcp_server_create(cfg->ua, port, cfg->my_keys, cfg->inst ? cfg->inst->my_ed25519_pubkey : NULL, cfg->inst, server_accept_cb, ss, NULL, NULL, cfg->listen_family); if (!ss->srv) { u_free(ss); return NULL; } DEBUG_INFO(DEBUG_CATEGORY_ETCP, "port=%u", port); return ss; @@ -212,6 +219,7 @@ struct stcp_link *stcp_link_connect(struct stcp_link_config *cfg) { link->cli = stcp_client_connect(cfg->ua, addr_str, port, cfg->my_keys, pubkey, cfg->inst ? cfg->inst->my_ed25519_pubkey : NULL, + cfg->got_initial_pkt, cfg->session_id, client_ready_cb, link, NULL, NULL); if (!link->cli) { u_free(link); return NULL; } @@ -317,6 +325,14 @@ const uint8_t *stcp_link_get_peer_ed25519_pubkey(struct stcp_link *link) { return link && link->conn && link->conn->peer_ed25519_set ? link->conn->peer_ed25519_pubkey : NULL; } +uint8_t stcp_link_get_peer_got_initial_pkt(struct stcp_link *link) { + return link ? link->peer_got_initial_pkt : 0; +} + +uint32_t stcp_link_get_peer_session_id(struct stcp_link *link) { + return link ? link->peer_session_id : 0; +} + void stcp_server_list_add(struct UTUN_INSTANCE *inst, struct stcp_server *srv) { if (!inst || !srv) return; srv->next = inst->stcp_servers; diff --git a/src/transport_layer/stcp_link.h b/src/transport_layer/stcp_link.h index 64f3cb70..f2de1dad 100644 --- a/src/transport_layer/stcp_link.h +++ b/src/transport_layer/stcp_link.h @@ -29,6 +29,8 @@ struct stcp_link_config { const struct sockaddr_storage *remote_addr; // адрес пира (клиент) uint16_t remote_port; // порт пира (клиент) int listen_family; // AF_INET или AF_INET6 для сервера (0 = v4) + uint8_t got_initial_pkt; // моё значение, отправляется пиру при handshake + uint32_t session_id; // мой ETCP session_id, отправляется пиру при handshake }; // ====== TCP server ====== @@ -62,6 +64,8 @@ void stcp_link_set_on_ready(struct stcp_link *link, stcp_link_cb cb, void *arg); void stcp_link_set_on_close(struct stcp_link *link, void (*cb)(struct stcp_link *link, int err, void *arg), void *arg); const uint8_t *stcp_link_get_peer_pubkey(struct stcp_link *link); const uint8_t *stcp_link_get_peer_ed25519_pubkey(struct stcp_link *link); +uint8_t stcp_link_get_peer_got_initial_pkt(struct stcp_link *link); +uint32_t stcp_link_get_peer_session_id(struct stcp_link *link); #ifdef __cplusplus diff --git a/src/transport_layer/stcp_server.c b/src/transport_layer/stcp_server.c index 5d9caea8..f6bfc52a 100644 --- a/src/transport_layer/stcp_server.c +++ b/src/transport_layer/stcp_server.c @@ -8,6 +8,8 @@ #include "../lib/mem.h" #include "../lib/debug_config.h" #include "../lib/platform_compat.h" +#include "../utun_instance.h" +#include "etcp.h" #include #include #include @@ -23,6 +25,7 @@ struct stcp_server { void *listen_id; struct SC_MYKEYS my_keys; uint8_t my_ed25519_pubkey[SC_PUBKEY_SIZE]; + struct UTUN_INSTANCE *inst; stcp_connect_cb connect_cb; void *cb_arg; stcp_close_cb close_cb; @@ -72,8 +75,11 @@ static void server_hs_phase1_cb(struct stcp_conn *c, uint8_t *data, size_t len) stcp_conn_do_close(c, 2); return; } memcpy(c->peer_ed25519_pubkey, enc_hs, SC_PUBKEY_SIZE); c->peer_ed25519_set = 1; - uint16_t padding_size = (uint16_t)enc_hs[32] | ((uint16_t)enc_hs[33] << 8); - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_server: client handshake key processed, expecting %u bytes padding", padding_size); + c->peer_got_initial_pkt = enc_hs[32]; + memcpy(&c->peer_session_id, enc_hs + 33, 4); + uint16_t padding_size = (uint16_t)enc_hs[37] | ((uint16_t)enc_hs[38] << 8); + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_server: client handshake OK gop=%d sid=%08x padding=%u", + c->peer_got_initial_pkt, c->peer_session_id, padding_size); stcp_recv_set(c, padding_size, 0, server_hs_phase2_cb); } @@ -82,6 +88,18 @@ static void server_hs_phase2_cb(struct stcp_conn *c, uint8_t *data, size_t len) if (c->hs_timer) { uasync_cancel_timeout(c->ua, c->hs_timer); c->hs_timer = NULL; } + uint8_t server_gop = 0; + uint32_t server_sid = 0; + if (c->inst && c->peer_pubkey_set) { + uint64_t node_id = sc_derive_node_id_from_pubkey(c->peer_pubkey); + struct ETCP_CONN *conn = instance_find_conn(c->inst, node_id); + server_gop = conn ? conn->got_initial_pkt : 0; + server_sid = conn ? conn->session_id : 0; + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_server: node=0x%016llx conn=%p server_gop=%d server_sid=%08x client_gop=%d client_sid=%08x", + (unsigned long long)node_id, (void*)conn, server_gop, server_sid, + c->peer_got_initial_pkt, c->peer_session_id); + } + uint8_t salt2[SC_PUBKEY_ENC_SALT_SIZE]; if (random_bytes(salt2, SC_PUBKEY_ENC_SALT_SIZE) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "random_bytes failed"); @@ -94,11 +112,14 @@ static void server_hs_phase2_cb(struct stcp_conn *c, uint8_t *data, size_t len) memcpy(resp, salt2, SC_PUBKEY_ENC_SALT_SIZE); sc_obfuscate_pubkey(salt2, c->peer_pubkey, c->my_keys.public_key, resp + SC_PUBKEY_ENC_SALT_SIZE); - uint8_t plain_hs[35]; memcpy(plain_hs, c->my_ed25519_pubkey, SC_PUBKEY_SIZE); plain_hs[32] = 0; plain_hs[33] = (uint8_t)padding; plain_hs[34] = (uint8_t)(padding >> 8); - uint32_t crc = crc32_calc(plain_hs, 35); + uint8_t plain_hs[39]; memcpy(plain_hs, c->my_ed25519_pubkey, SC_PUBKEY_SIZE); + plain_hs[32] = server_gop; + memcpy(plain_hs + 33, &server_sid, 4); + plain_hs[37] = (uint8_t)padding; plain_hs[38] = (uint8_t)(padding >> 8); + uint32_t crc = crc32_calc(plain_hs, 39); uint8_t *enc_dst = resp + SC_PUBKEY_ENC_SIZE; - memcpy(enc_dst, plain_hs, 35); - enc_dst[35] = (uint8_t)(crc >> 0); enc_dst[36] = (uint8_t)(crc >> 8); enc_dst[37] = (uint8_t)(crc >> 16); enc_dst[38] = (uint8_t)(crc >> 24); + memcpy(enc_dst, plain_hs, 39); + enc_dst[39] = (uint8_t)(crc >> 0); enc_dst[40] = (uint8_t)(crc >> 8); enc_dst[41] = (uint8_t)(crc >> 16); enc_dst[42] = (uint8_t)(crc >> 24); log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO, "stcp_server hs_resp BEFORE xor", enc_dst, STCP_HS_ENC_SERVER); if (sc_stream_xor(&c->stream_send, enc_dst, STCP_HS_ENC_SERVER) != SC_OK) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "encrypt failed"); @@ -157,6 +178,7 @@ static void server_accept_cb(socket_t listen_sock, void *arg) { c->on_close = srv->close_cb; c->close_arg = srv->close_arg; c->socket_id = uasync_add_socket_t(srv->ua, cli_sock, server_conn_read_cb, stcp_write_cb, NULL, c); + c->inst = srv->inst; stcp_recv_set(c, SC_PUBKEY_ENC_SIZE + STCP_HS_ENC_CLIENT, 0, server_hs_phase1_cb); c->hs_timer = uasync_set_timeout(c->ua, STCP_HS_TIMEOUT, c, hs_timeout_cb, "stcp_hs"); DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_server: accepted connection fd=%d", (int)cli_sock); @@ -165,6 +187,7 @@ static void server_accept_cb(socket_t listen_sock, void *arg) { struct stcp_server *stcp_server_create(struct UASYNC *ua, uint16_t port, struct SC_MYKEYS *keys, const uint8_t *my_ed25519_pubkey, + struct UTUN_INSTANCE *inst, stcp_connect_cb connect_cb, void *arg, stcp_close_cb close_cb, void *close_arg, int family) { @@ -174,6 +197,7 @@ struct stcp_server *stcp_server_create(struct UASYNC *ua, uint16_t port, srv->ua = ua; srv->my_keys = *keys; if (my_ed25519_pubkey) memcpy(srv->my_ed25519_pubkey, my_ed25519_pubkey, SC_PUBKEY_SIZE); + srv->inst = inst; srv->connect_cb = connect_cb; srv->cb_arg = arg; srv->close_cb = close_cb; diff --git a/src/transport_layer/stcp_server.h b/src/transport_layer/stcp_server.h index 0f2e4eeb..f2085b5d 100644 --- a/src/transport_layer/stcp_server.h +++ b/src/transport_layer/stcp_server.h @@ -9,12 +9,15 @@ extern "C" { #include "stcp.h" +struct UTUN_INSTANCE; + typedef void (*stcp_connect_cb)(struct stcp_conn *conn, void *arg); typedef void (*stcp_close_cb)(struct stcp_conn *conn, int err, void *arg); struct stcp_server *stcp_server_create(struct UASYNC *ua, uint16_t port, struct SC_MYKEYS *keys, const uint8_t *my_ed25519_pubkey, + struct UTUN_INSTANCE *inst, stcp_connect_cb connect_cb, void *arg, stcp_close_cb close_cb, void *close_arg, int family); diff --git a/tests/Makefile.am b/tests/Makefile.am index a573a1c9..bce44f6e 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -50,6 +50,7 @@ check_PROGRAMS = \ test_broadcast \ test_conn_mgr \ test_conn_mgr_already_connected \ + test_invite_group_create \ test_etcp_connect \ test_node_conn_direct \ test_db_sync \ @@ -294,6 +295,10 @@ test_conn_mgr_already_connected_SOURCES = test_conn_mgr_already_connected.c test_conn_mgr_already_connected_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_conn_mgr_already_connected_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_invite_group_create_SOURCES = test_invite_group_create.c +test_invite_group_create_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib +test_invite_group_create_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_etcp_connect_SOURCES = test_etcp_connect.c test_etcp_connect_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_etcp_connect_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_invite_group_create.c b/tests/test_invite_group_create.c new file mode 100644 index 00000000..4c4e75a7 --- /dev/null +++ b/tests/test_invite_group_create.c @@ -0,0 +1,166 @@ +/** + * @file test_invite_group_create.c + * @brief Full invite flow: conn_mgr_open_invite → B is member → CONN_EVENT_JOIN + */ +#include +#include +#include +#include +#include "../lib/platform_compat.h" +#include "../lib/mem.h" +#include "../lib/debug_config.h" +#include "test_utils.h" +#ifndef _WIN32 +#include +#endif + +#include "../src/transport_layer/etcp.h" +#include "../src/transport_layer/etcp_connections.h" +#include "../src/config_parser.h" +#include "../src/config_updater.h" +#include "../src/utun_instance.h" +#include "../src/routing_layer/topo_group.h" +#include "../src/routing_layer/topo_node.h" +#include "../src/routing_layer/topo_node_sqlite.h" +#include "../src/routing_layer/conn_mgr.h" + +#define TIMEOUT_TB 50000 +#define POLL_MS 20 + +static void fail(const char* msg) { fprintf(stderr, "FAIL: %s\n", msg); fflush(stderr); exit(1); } + +static void wf(const char* path, const char* fmt, ...) { + va_list ap; va_start(ap, fmt); + FILE* f = fopen(path, "w"); + if (!f) { fail("fopen"); return; } + vfprintf(f, fmt, ap); fclose(f); + va_end(ap); +} + +static char* ls(const char* p, const char* k) { + char b[1024]; FILE* f = fopen(p, "r"); if (!f) return NULL; + size_t n = fread(b, 1, sizeof(b) - 1, f); fclose(f); b[n] = 0; + char* x = strstr(b, k); if (!x) return NULL; + x += strlen(k) + 1; while (*x == ' ' || *x == '\t') x++; + char* r = u_strdup(x); char* e = r; while (*e && *e != '\n' && *e != '\r') e++; *e = 0; + return r; +} + +static int result = 0; + +static void ccb(struct CONN_MGR_HANDLE* h, uint64_t nid, uint64_t gid, + enum conn_mgr_event ev, void* arg) { + (void)h; (void)gid; (void)arg; + fprintf(stderr, "CB: ev=%d node=0x%llx\n", (int)ev, (unsigned long long)nid); fflush(stderr); + if (ev == CONN_EVENT_JOIN) result = 1; + if (ev == CONN_EVENT_TIMEOUT) result = 2; +} + +static void to_cb(void* arg) { (void)arg; fprintf(stderr, "TIMEOUT\n"); fflush(stderr); result = 2; } + +int main(void) { + debug_config_init(); + debug_set_level(DEBUG_LEVEL_TRACE); + debug_enable_file_output("/tmp/test_inv_crash.log", 1); + utun_instance_set_tun_init_enabled(0); + + char tdir[] = "/tmp/utin_XXXXXX"; + if (test_mkdtemp(tdir) != 0) { fail("mkdtemp"); return 1; } + char ca[256], cb[256]; + snprintf(ca, sizeof(ca), "%s/a.conf", tdir); + snprintf(cb, sizeof(cb), "%s/b.conf", tdir); + int porta = 52000 + (getpid() % 10000), portb = porta + 1; + + wf(ca, "[global]\ntun_ip=10.94.0.1/24\ntun_ifname=tun89\ndb_path=%s\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", tdir, porta); + wf(cb, "[global]\ntun_ip=10.94.0.2/24\ntun_ifname=tun88\ndb_path=%s\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", tdir, portb); + config_ensure_keys_and_node_id(ca); config_ensure_keys_and_node_id(cb); + + char* rva = ls(ca, "priv"); char* pua = ls(ca, "pub"); char* pub = ls(cb, "pub"); + wf(ca, "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.94.0.1/24\ntun_ifname=tun89\ndb_path=%s\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[client: to_b]\nkeepalive=1\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n[allowed_keys]\nallow_all=1\n", rva, pua, tdir, porta, pub, portb); + char* rvb = ls(cb, "priv"); + wf(cb, "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.94.0.2/24\ntun_ifname=tun88\ndb_path=%s\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", rvb, pub, tdir, portb); + u_free(rva); u_free(pua); u_free(pub); u_free(rvb); + + struct UASYNC* ua = uasync_create(); + struct UTUN_INSTANCE* a = utun_instance_create(ua, ca); + if (!a) { fail("create A"); goto clean; } + utun_instance_init(a); + + /* ── Add a TCP socket (simulating Android STCP server) ── */ + { + struct TCP_SOCKET* ts = u_calloc(1, sizeof(*ts)); + ts->instance = a; + ts->type = CFG_SERVER_TYPE_PUBLIC; + ts->sock_id = (uint8_t)a->next_socket_id++; + struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); + sin.sin_family = AF_INET; sin.sin_addr.s_addr = htonl(0x7f000001); /* 127.0.0.1 */ + sin.sin_port = htons((uint16_t)55555); + memcpy(&ts->local_addr, &sin, sizeof(sin)); + ts->interface_addr = ts->local_addr; + ts->next = a->tcp_sockets; a->tcp_sockets = ts; + fprintf(stderr, "Added TCP socket ts=%p next=%p\n", (void*)ts, (void*)ts->next); fflush(stderr); + } + fprintf(stderr, "tcp_sockets=%p etcp_sockets=%p\n", (void*)a->tcp_sockets, (void*)a->etcp_sockets); fflush(stderr); + + /* Get B's node_id and pubkey from its instance */ + struct UTUN_INSTANCE* b = utun_instance_create(ua, cb); + if (!b) { fail("create B"); a->running = 0; utun_instance_destroy(a); uasync_destroy(ua, 0); goto clean; } + utun_instance_init(b); + + /* ── Set up B as member of group 0x7777777700000001 ── */ + { + uint64_t gid = 0x7777777700000001ULL; + char ch_str[32]; snprintf(ch_str, sizeof(ch_str), "%llu", (unsigned long long)gid); + topo_groups_create_group(b->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, ch_str); + const char* ch_name = "TestCh"; + uint8_t sig_msg[256]; size_t slen = 0; + /* msg: ch_id\0 || name\0 || owner(8 LE) || x25519_pub(32) || ed_pub(32) */ + size_t ch_id_len = strlen(ch_str) + 1; memcpy(sig_msg + slen, ch_str, ch_id_len); slen += ch_id_len; + size_t name_len = strlen(ch_name) + 1; memcpy(sig_msg + slen, ch_name, name_len); slen += name_len; + uint64_t owner = b->node_id; memcpy(sig_msg + slen, &owner, 8); slen += 8; + memcpy(sig_msg + slen, b->my_keys.public_key, 32); slen += 32; + memcpy(sig_msg + slen, b->my_ed25519_pubkey, 32); slen += 32; + uint8_t ch_sig[64]; + sc_ed25519_sign(b->my_ed25519_privkey, sig_msg, slen, ch_sig); + topo_node_sqlite_channel_put(b->topo_sqlite_db, ch_str, ch_name, owner, + b->my_keys.public_key, NULL, + b->my_ed25519_pubkey, NULL, ch_sig); + fprintf(stderr, "B: channel setup gid=0x%016llX ch=%s name=%s\n", + (unsigned long long)gid, ch_str, ch_name); fflush(stderr); + } + + uint64_t nid_b = b->node_id; + uint8_t pk_b[32]; memcpy(pk_b, b->my_keys.public_key, 32); + + /* Build TOPO_NODE for invite */ + struct TOPO_NODE* ni = u_calloc(1, sizeof(*ni)); + ni->node_id = nid_b; memcpy(ni->public_key, pk_b, 32); + struct TOPO_ADDR4* a4 = u_calloc(1, sizeof(*a4)); + uint8_t ip[4] = {127, 0, 0, 1}; memcpy(a4->addr, ip, 4); + a4->port = (uint16_t)portb; a4->type = TOPO_ADDR_NAT; a4->protocol = 1; + ni->v4_addrs = a4; + + uint64_t gid = 0x7777777700000001ULL; + fprintf(stderr, "=== conn_mgr_open_invite gid=0x%016llX nid=0x%016llX ===\n", + (unsigned long long)gid, (unsigned long long)nid_b); fflush(stderr); + + int r = conn_mgr_open_invite(a, gid, ni, nid_b, ccb, NULL, NULL); + fprintf(stderr, "conn_mgr_open_invite => %d\n", r); fflush(stderr); + + void* tt = uasync_set_timeout(ua, TIMEOUT_TB, NULL, to_cb, "to"); + int el = 0; + while (!result && el < TIMEOUT_TB + 5000) { uasync_poll(ua, POLL_MS); el += POLL_MS; } + if (tt) uasync_cancel_timeout(ua, tt); + + a->running = 0; utun_instance_destroy(a); + b->running = 0; utun_instance_destroy(b); + uasync_destroy(ua, 0); + u_free(a4); u_free(ni); + + fprintf(stderr, "=== DONE result=%d ===\n", result); fflush(stderr); + +clean: + test_unlink(ca); test_unlink(cb); test_rmdir(tdir); + debug_disable_file_output(); + return (result == 1) ? 0 : 1; +} diff --git a/tests/test_stcp.c b/tests/test_stcp.c index d06024e9..911ebeea 100644 --- a/tests/test_stcp.c +++ b/tests/test_stcp.c @@ -103,8 +103,8 @@ static int test1_sizes(void) { struct test_peer srv = {0}, cli = {0}; uint16_t port = BASE_PORT + 1; - struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, server_connect_cb, &srv, peer_close_cb, &srv, AF_INET); TASSERT(ss); - struct stcp_client *sc = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, s_keys.public_key, NULL, client_ready_cb, &cli, peer_close_cb, &cli); TASSERT(sc); + struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, NULL, server_connect_cb, &srv, peer_close_cb, &srv, AF_INET); TASSERT(ss); + struct stcp_client *sc = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, s_keys.public_key, NULL, 0, 0, client_ready_cb, &cli, peer_close_cb, &cli); TASSERT(sc); size_t sizes[] = {0, 1, 16, 17, 255, 256, 1000, 65535}; int n_sizes = 8; @@ -143,8 +143,8 @@ static int test2_many(void) { struct test_peer srv = {0}, cli = {0}; uint16_t port = BASE_PORT + 2; - struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, server_connect_cb, &srv, peer_close_cb, &srv, AF_INET); TASSERT(ss); - struct stcp_client *sc = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, s_keys.public_key, NULL, client_ready_cb, &cli, peer_close_cb, &cli); TASSERT(sc); + struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, NULL, server_connect_cb, &srv, peer_close_cb, &srv, AF_INET); TASSERT(ss); + struct stcp_client *sc = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, s_keys.public_key, NULL, 0, 0, client_ready_cb, &cli, peer_close_cb, &cli); TASSERT(sc); int sent = 0, ticks = 0; while (srv.msg_count < 200 && ticks < 200) { @@ -181,11 +181,11 @@ static int test3_wrong_key(void) { struct test_peer srv = {0}, cli = {0}; uint16_t port = BASE_PORT + 3; - struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, server_connect_cb, &srv, peer_close_cb, &srv, AF_INET); TASSERT(ss); + struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, NULL, server_connect_cb, &srv, peer_close_cb, &srv, AF_INET); TASSERT(ss); struct SC_MYKEYS rogue; TASSERT(sc_generate_keypair(&rogue) == SC_OK); - struct stcp_client *sc = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, rogue.public_key, NULL, client_ready_cb, &cli, peer_close_cb, &cli); TASSERT(sc); + struct stcp_client *sc = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, rogue.public_key, NULL, 0, 0, client_ready_cb, &cli, peer_close_cb, &cli); TASSERT(sc); int ticks = 0; while (ticks < 200) { @@ -209,8 +209,8 @@ static int test4_close(void) { struct test_peer srv = {0}, cli = {0}; uint16_t port = BASE_PORT + 4; - struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, server_connect_cb, &srv, peer_close_cb, &srv, AF_INET); TASSERT(ss); - struct stcp_client *sc = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, s_keys.public_key, NULL, client_ready_cb, &cli, peer_close_cb, &cli); TASSERT(sc); + struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, NULL, server_connect_cb, &srv, peer_close_cb, &srv, AF_INET); TASSERT(ss); + struct stcp_client *sc = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, s_keys.public_key, NULL, 0, 0, client_ready_cb, &cli, peer_close_cb, &cli); TASSERT(sc); int closed = 0, ticks = 0; while (!srv.closed && ticks < 200) { @@ -256,11 +256,11 @@ static int test5_multi(void) { memset(srvp, 0, sizeof(srvp)); memset(clip, 0, sizeof(clip)); g_multi_peers = srvp; g_multi_idx = 0; g_multi_max = NCLI; - struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, multi_connect_cb, NULL, NULL, NULL, AF_INET); TASSERT(ss); + struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, NULL, multi_connect_cb, NULL, NULL, NULL, AF_INET); TASSERT(ss); struct stcp_client *clients[NCLI] = {0}; for (int i = 0; i < NCLI; i++) { - clients[i] = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, s_keys.public_key, NULL, client_ready_cb, &clip[i], peer_close_cb, &clip[i]); + clients[i] = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, s_keys.public_key, NULL, 0, 0, client_ready_cb, &clip[i], peer_close_cb, &clip[i]); TASSERT(clients[i]); } @@ -308,8 +308,8 @@ static int test6_interleaved(void) { struct test_peer srv = {0}, cli = {0}; uint16_t port = BASE_PORT + 6; - struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, server_connect_cb, &srv, peer_close_cb, &srv, AF_INET); TASSERT(ss); - struct stcp_client *sc = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, s_keys.public_key, NULL, client_ready_cb, &cli, peer_close_cb, &cli); TASSERT(sc); + struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, NULL, server_connect_cb, &srv, peer_close_cb, &srv, AF_INET); TASSERT(ss); + struct stcp_client *sc = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, s_keys.public_key, NULL, 0, 0, client_ready_cb, &cli, peer_close_cb, &cli); TASSERT(sc); int round = 0, ticks = 0; while (srv.msg_count < 50 || cli.msg_count < 50) { @@ -342,8 +342,8 @@ static int test7_bulk_4mb(void) { struct test_peer srv = {0}, cli = {0}; uint16_t port = BASE_PORT + 7; - struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, server_connect_cb, &srv, peer_close_cb, &srv, AF_INET); TASSERT(ss); - struct stcp_client *sc = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, s_keys.public_key, NULL, client_ready_cb, &cli, peer_close_cb, &cli); TASSERT(sc); + struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, NULL, server_connect_cb, &srv, peer_close_cb, &srv, AF_INET); TASSERT(ss); + struct stcp_client *sc = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, s_keys.public_key, NULL, 0, 0, client_ready_cb, &cli, peer_close_cb, &cli); TASSERT(sc); #define N_BULK 64 #define SZ_BULK 65535 @@ -381,8 +381,8 @@ static int test8_srv_recv_close(void) { struct test_peer srv = {0}, cli = {0}; uint16_t port = BASE_PORT + 8; - struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, server_connect_cb, &srv, peer_close_cb, &srv, AF_INET); TASSERT(ss); - struct stcp_client *sc = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, s_keys.public_key, NULL, client_ready_cb, &cli, peer_close_cb, &cli); TASSERT(sc); + struct stcp_server *ss = stcp_server_create(ua, port, &s_keys, NULL, NULL, server_connect_cb, &srv, peer_close_cb, &srv, AF_INET); TASSERT(ss); + struct stcp_client *sc = stcp_client_connect(ua, "127.0.0.1", port, &c_keys, s_keys.public_key, NULL, 0, 0, client_ready_cb, &cli, peer_close_cb, &cli); TASSERT(sc); int ticks = 0; while ((!srv.ready || !cli.ready) && ticks < 200) { uasync_poll(ua, 10); ticks++; } diff --git a/tools/chatgui-android/app/src/main/AndroidManifest.xml b/tools/chatgui-android/app/src/main/AndroidManifest.xml index c3dd3b73..0047a460 100644 --- a/tools/chatgui-android/app/src/main/AndroidManifest.xml +++ b/tools/chatgui-android/app/src/main/AndroidManifest.xml @@ -5,7 +5,8 @@ - + + @@ -30,6 +31,10 @@ + android:foregroundServiceType="specialUse"> + + diff --git a/tools/chatgui-android/jni_bridge/android_jni_bridge.c b/tools/chatgui-android/jni_bridge/android_jni_bridge.c index 400b73a2..2e3a6aac 100644 --- a/tools/chatgui-android/jni_bridge/android_jni_bridge.c +++ b/tools/chatgui-android/jni_bridge/android_jni_bridge.c @@ -62,20 +62,17 @@ static char g_db_path[512]; static void bridge_log(int level, const char* fmt, ...) { va_list ap; + char fmt_buf[1024]; va_start(ap, fmt); + vsnprintf(fmt_buf, sizeof(fmt_buf), fmt, ap); + va_end(ap); #ifdef __ANDROID__ - __android_log_vprint(blevel_to_android(level), LOG_TAG, fmt, ap); + __android_log_print(blevel_to_android(level), LOG_TAG, "%s", fmt_buf); #else - vfprintf(stderr, fmt, ap); fputc('\n', stderr); + fprintf(stderr, "%s\n", fmt_buf); #endif - va_end(ap); if (g_on_log) { - va_list ap2; - char buf[512]; - va_start(ap2, fmt); - vsnprintf(buf, sizeof(buf), fmt, ap2); - va_end(ap2); const char* cat; switch (level) { case BLEV_TRACE: cat = "TRACE"; break; @@ -85,7 +82,7 @@ static void bridge_log(int level, const char* fmt, ...) { case BLEV_ERROR: cat = "ERROR"; break; default: cat = "?"; break; } - g_on_log(level, cat, buf); + g_on_log(level, cat, fmt_buf); } } @@ -904,13 +901,14 @@ void utun_bridge_set_member_flags(const char* channel_id, uint64_t node_id, int new_ver = cur_ver + 1; /* build JSON — supernode/storage: direct; admin/moder: append-only from existing raw */ - char json[384]; int off = 0; - off += snprintf(json + off, sizeof(json) - (size_t)off, "{\"ver\":\"%d\"", new_ver); + char json[4096]; int off = 0; +#define SETMF_ADVANCE(n) do { if (off < (int)sizeof(json) - 10) off += (int)(n); else off = (int)sizeof(json); } while(0) + SETMF_ADVANCE(snprintf(json, sizeof(json), "{\"ver\":\"%d\"", new_ver)); if (json_flat_get(cur_tags, "supernode", buf, sizeof(buf)) == 0) snprintf(buf, sizeof(buf), "%s", buf); else buf[0] = '\0'; const char* sn_val = supernode ? "yes" : (buf[0] ? buf : NULL); - if (sn_val) off += snprintf(json + off, sizeof(json) - (size_t)off, ",\"supernode\":\"%s\"", sn_val); + if (sn_val) SETMF_ADVANCE(snprintf(json + off, sizeof(json) - (size_t)off, ",\"supernode\":\"%s\"", sn_val)); /* admin: append-only */ { @@ -920,9 +918,9 @@ void utun_bridge_set_member_flags(const char* channel_id, uint64_t node_id, if ((admin && last != 'e') || (!admin && last == 'e')) { uint64_t ts = (uint64_t)ntp_time_get_seconds(chat_core_get_inst()); size_t ol = strlen(raw); - snprintf(raw + ol, sizeof(raw) - ol, "%c%llu", admin ? 'e' : 'd', (unsigned long long)ts); + if (ol < sizeof(raw) - 30) snprintf(raw + ol, sizeof(raw) - ol, "%c%llu", admin ? 'e' : 'd', (unsigned long long)ts); } - if (raw[0]) off += snprintf(json + off, sizeof(json) - (size_t)off, ",\"admin\":\"%s\"", raw); + if (raw[0]) SETMF_ADVANCE(snprintf(json + off, sizeof(json) - (size_t)off, ",\"admin\":\"%s\"", raw)); } /* moder: append-only */ @@ -933,18 +931,19 @@ void utun_bridge_set_member_flags(const char* channel_id, uint64_t node_id, if ((moder && last != 'e') || (!moder && last == 'e')) { uint64_t ts = (uint64_t)ntp_time_get_seconds(chat_core_get_inst()); size_t ol = strlen(raw); - snprintf(raw + ol, sizeof(raw) - ol, "%c%llu", moder ? 'e' : 'd', (unsigned long long)ts); + if (ol < sizeof(raw) - 30) snprintf(raw + ol, sizeof(raw) - ol, "%c%llu", moder ? 'e' : 'd', (unsigned long long)ts); } - if (raw[0]) off += snprintf(json + off, sizeof(json) - (size_t)off, ",\"moder\":\"%s\"", raw); + if (raw[0]) SETMF_ADVANCE(snprintf(json + off, sizeof(json) - (size_t)off, ",\"moder\":\"%s\"", raw)); } /* storage */ if (json_flat_get(cur_tags, "storage", buf, sizeof(buf)) == 0) snprintf(buf, sizeof(buf), "%s", buf); else buf[0] = '\0'; const char* st_val = storage_flag ? "yes" : (buf[0] ? buf : NULL); - if (st_val) off += snprintf(json + off, sizeof(json) - (size_t)off, ",\"storage\":\"%s\"", st_val); + if (st_val) SETMF_ADVANCE(snprintf(json + off, sizeof(json) - (size_t)off, ",\"storage\":\"%s\"", st_val)); - off += snprintf(json + off, sizeof(json) - (size_t)off, "}"); + SETMF_ADVANCE(snprintf(json + off, sizeof(json) - (size_t)off, "}")); +#undef SETMF_ADVANCE /* get channel ed25519 private key for signing */ uint8_t ch_ed_priv[32] = {0}; @@ -1058,7 +1057,8 @@ char* utun_bridge_get_member_links_json(uint64_t node_id) { "%s{\"fm\":%d,\"ad\":\"%s\",\"po\":%d,\"st\":%d,\"ls\":%d}", sep, fm, ip, (int)po, ready ? 3 : 1, ready ? 1 : 0); } - pos += snprintf(buf + pos, sizeof(buf) - (size_t)pos, "]}"); + if (pos < (int)sizeof(buf)) + pos += snprintf(buf + pos, sizeof(buf) - (size_t)pos, "]}"); bridge_log(BLEV_INFO, "getMemberLinks node=0x%016llx → %s", (unsigned long long)node_id, buf); return u_strdup(buf); }