diff --git a/src/chat/chat_sync.c b/src/chat/chat_sync.c index 3ac8b785..dc156ab5 100644 --- a/src/chat/chat_sync.c +++ b/src/chat/chat_sync.c @@ -150,6 +150,7 @@ static int cs_send(struct chat_sync* cs, const char* ch_id, uint64_t dst, struct ETCP_CONN* conn = cs_find_conn_for_node(cs->inst, dst); if (!conn) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: no conn for node %016llx", CS_ID, (unsigned long long)dst); u_free(buf); queue_entry_free(entry); return -1; } int rc = etcp_send(conn, entry); + if (rc != 0) { queue_dgram_free(entry); queue_entry_free(entry); } DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: etcp_send rc=%d id=0x%02x to=%016llx conn=%s", CS_ID, rc, payload[0], (unsigned long long)dst, conn->log_name); return rc; } diff --git a/src/chat/db_sync.c b/src/chat/db_sync.c index d204bb5c..ac630d92 100644 --- a/src/chat/db_sync.c +++ b/src/chat/db_sync.c @@ -700,10 +700,11 @@ static int db_sync_send_gid(struct DB_SYNC* db, uint64_t node_id, uint64_t group DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync: send DROP — no conn to %04llX (group=%016llx, type=%02x)", (unsigned long long)(node_id >> 16), group_id, plen > 0 ? payload[0] : 0, (void*)conn, conn ? conn->links_up : -1); - queue_entry_free(entry); + queue_dgram_free(entry); queue_entry_free(entry); return -1; } int ret = etcp_send(conn, entry); + if (ret != 0) { queue_dgram_free(entry); queue_entry_free(entry); } DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "sync: send OK → etcp_send ret=%d conn=%s l_up=%d init=%d q=%p", ret, conn->log_name, conn->links_up, conn->initialized, (void*)conn->send_input_q); return ret; diff --git a/src/routing_layer/route_ping.c b/src/routing_layer/route_ping.c index d40a2619..94c0beaf 100644 --- a/src/routing_layer/route_ping.c +++ b/src/routing_layer/route_ping.c @@ -286,14 +286,17 @@ static void route_ping_series_finish(struct route_ping_series_ctx* ctx) { if (e) { e->dgram = (uint8_t*)resp; e->len = sizeof(struct NATDET_PING_RESP); - etcp_send(ctx->reply_conn, e); - DEBUG_INFO(DEBUG_CATEGORY_BGP, - "PING series done %s request_id=%08x sent=%u ok=%u avg_rtt=%u", - ctx->reply_conn->log_name, - (unsigned)ctx->request_id, - (unsigned)ctx->count_sent, - (unsigned)ctx->count_ok, - (unsigned)avg_rtt); + if (etcp_send(ctx->reply_conn, e) != 0) { + queue_dgram_free(e); queue_entry_free(e); + } else { + DEBUG_INFO(DEBUG_CATEGORY_BGP, + "PING series done %s request_id=%08x sent=%u ok=%u avg_rtt=%u", + ctx->reply_conn->log_name, + (unsigned)ctx->request_id, + (unsigned)ctx->count_sent, + (unsigned)ctx->count_ok, + (unsigned)avg_rtt); + } } else { u_free(resp); } diff --git a/src/routing_layer/topo_group_invite.c b/src/routing_layer/topo_group_invite.c index b4d318c9..8b693231 100644 --- a/src/routing_layer/topo_group_invite.c +++ b/src/routing_layer/topo_group_invite.c @@ -140,7 +140,8 @@ static void tgi_handle_info_req(struct ETCP_CONN* conn, const uint8_t* data, siz if (sg) memcpy(resp->ed25519_pubkey, sg->ed25519_public_key, SC_PUBKEY_SIZE); resp->node_name_len = (uint8_t)nl; if (nl) memcpy(resp->node_name, nm, nl); struct ll_entry* qe = queue_entry_new(0); - if (qe) { qe->dgram = rb; qe->len = (uint16_t)rs; etcp_send(conn, qe); } else u_free(rb); + if (qe) { qe->dgram = rb; qe->len = (uint16_t)rs; + if (etcp_send(conn, qe) != 0) { queue_dgram_free(qe); queue_entry_free(qe); } } else u_free(rb); DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "invite: group=%016llx NOT in DB — responding is_member=NO to 0x%016llx", (unsigned long long)req->group_id, (unsigned long long)conn->peer_node_id); return; @@ -226,7 +227,8 @@ static void tgi_handle_info_req(struct ETCP_CONN* conn, const uint8_t* data, siz (int)ch_nl, ch_name, (unsigned long long)owner, (unsigned)inv_flags, rs); { struct ll_entry* qe = queue_entry_new(0); - if (qe) { qe->dgram = rb; qe->len = (uint16_t)rs; etcp_send(conn, qe); } else u_free(rb); + if (qe) { qe->dgram = rb; qe->len = (uint16_t)rs; + if (etcp_send(conn, qe) != 0) { queue_dgram_free(qe); queue_entry_free(qe); } } else u_free(rb); } } diff --git a/src/transport_layer/etcp_api.c b/src/transport_layer/etcp_api.c index 59bd86b7..75edea78 100644 --- a/src/transport_layer/etcp_api.c +++ b/src/transport_layer/etcp_api.c @@ -167,6 +167,8 @@ int etcp_unbind(struct UTUN_INSTANCE* inst, uint8_t id) { } int etcp_send(struct ETCP_CONN* conn, struct ll_entry* entry) { + /* Владение entry: при успехе entry уходит в очередь (может быть освобождён + * синхронно); при ошибке (-1) entry НЕ освобождается — вызывающий владеет. */ if (!conn || !entry || !conn->send_input_q || conn->state == 2) return -1; return queue_data_put(conn->send_input_q, entry); } diff --git a/src/transport_layer/etcp_api.h b/src/transport_layer/etcp_api.h index 08dfa520..448ce426 100644 --- a/src/transport_layer/etcp_api.h +++ b/src/transport_layer/etcp_api.h @@ -289,11 +289,14 @@ struct ETCP_BINDINGS { * @param entry Элемент очереди с данными для отправки * @return 0 при успехе, -1 при ошибке * - * @note Функция забирает ownership entry — вызывающий код не должен - * ни освобождать, ни читать entry после вызова: если очередь - * normalizer'а пуста, entry освобождается СИНХРОННО внутри этого - * вызова (queue_data_put → callback → queue_entry_free). Любые - * нужные поля entry надо сохранять до etcp_send(). + * @note Владение entry: + * - при успехе (0) entry передаётся в очередь normalizer'а; вызывающий + * НЕ должен ни освобождать, ни читать entry после вызова — если очередь + * пуста, entry обрабатывается и освобождается СИНХРОННО внутри вызова + * (queue_data_put → callback → queue_entry_free). Нужные поля entry надо + * сохранять до etcp_send(). + * - при ошибке (-1) entry НЕ освобождается: вызывающий владеет entry и + * должен сам освободить его (queue_dgram_free(entry); queue_entry_free(entry)). */ int etcp_send(struct ETCP_CONN* conn, struct ll_entry* entry); diff --git a/src/transport_layer/node_conn_direct.c b/src/transport_layer/node_conn_direct.c index 0713d88d..cc1fbcd3 100644 --- a/src/transport_layer/node_conn_direct.c +++ b/src/transport_layer/node_conn_direct.c @@ -470,8 +470,10 @@ static void ncd_recv_control_handler(struct ETCP_CONN* conn, struct ll_entry* e) { struct ncd_control_msg resp; resp.cmd = ETCP_RT_ID_NCD_CONTROL; resp.subcmd = NCD_SUBCMD_KEEP_ALIVE; resp.node_id = conn->instance->node_id; struct ll_entry* qe = queue_entry_new(0); - if (qe) { qe->dgram = u_malloc(sizeof(resp)); memcpy(qe->dgram, &resp, sizeof(resp)); qe->len = sizeof(resp); - etcp_send(conn, qe); } + if (qe) { qe->dgram = u_malloc(sizeof(resp)); + if (qe->dgram) { memcpy(qe->dgram, &resp, sizeof(resp)); qe->len = sizeof(resp); + if (etcp_send(conn, qe) != 0) { queue_dgram_free(qe); queue_entry_free(qe); } } + else queue_entry_free(qe); } } } else { DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] deferred close entry for CLOSE from 0x%016llx", (unsigned long long)sender_id); @@ -915,9 +917,11 @@ void node_conn_direct_close(struct NODE_CONN_DIRECT* h) { { struct ncd_control_msg msg; msg.cmd = ETCP_RT_ID_NCD_CONTROL; msg.subcmd = NCD_SUBCMD_CLOSE; msg.node_id = conn->instance->node_id; struct ll_entry* qe = queue_entry_new(0); - if (qe) { qe->dgram = u_malloc(sizeof(msg)); memcpy(qe->dgram, &msg, sizeof(msg)); qe->len = sizeof(msg); - etcp_send(conn, qe); - DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] sent CLOSE to 0x%016llx", (unsigned long long)node_id); } + if (qe) { qe->dgram = u_malloc(sizeof(msg)); + if (qe->dgram) { memcpy(qe->dgram, &msg, sizeof(msg)); qe->len = sizeof(msg); + if (etcp_send(conn, qe) != 0) { queue_dgram_free(qe); queue_entry_free(qe); } + else DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] sent CLOSE to 0x%016llx", (unsigned long long)node_id); } + else queue_entry_free(qe); } } entry->fin_wait_timer = uasync_set_timeout(entry->ua, NCD_FIN_WAIT_TIMEOUT_TB, entry, ncd_fin_wait_timeout_cb, "ncd_fin_wait"); DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] fin_wait started node=0x%016llx", (unsigned long long)node_id);