diff --git a/src/routing_layer/conn_mgr.c b/src/routing_layer/conn_mgr.c index d5c4020a..0dfe4763 100644 --- a/src/routing_layer/conn_mgr.c +++ b/src/routing_layer/conn_mgr.c @@ -273,9 +273,8 @@ int conn_mgr_connect_node(struct CONN_MGR* mgr, uint64_t node_id, uint32_t idle_ struct ll_entry* ce = queue_find_data_by_index(mgr->instance->connections, (const uint8_t*)&node_id); while (ce) { struct conn_queue_entry* cqe = (struct conn_queue_entry*)ce->data; - if (cqe->conn && cqe->conn->peer_node_id == node_id && cqe->conn->links_up) { + if (cqe->conn && cqe->conn->peer_node_id == node_id) { existing = cqe->conn; - topo_group_new_conn(group, existing); break; } ce = queue_find_next_by_index(mgr->instance->connections, (const uint8_t*)&node_id, ce); @@ -297,6 +296,13 @@ int conn_mgr_connect_node(struct CONN_MGR* mgr, uint64_t node_id, uint32_t idle_ cm_deliver_result(entry, CONN_MGR_OK); return CONN_MGR_OK; } + entry = cm_ensure_entry(mgr, node_id); + if (!entry) return CONN_MGR_ERR_INTERNAL; + entry->state = CONN_MGR_STATE_CONNECTING; + { struct cm_cb_node* cn = u_calloc(1, sizeof(struct cm_cb_node)); + if (cn) { cn->cb = cb; cn->arg = cb_arg; cn->next = entry->cb_list; entry->cb_list = cn; } } + topo_group_new_conn(group, existing); + return CONN_MGR_OK; } entry = cm_ensure_entry(mgr, node_id); if (!entry) return CONN_MGR_ERR_INTERNAL; diff --git a/src/transport_layer/etcp.c b/src/transport_layer/etcp.c index 7416e46f..161f4769 100644 --- a/src/transport_layer/etcp.c +++ b/src/transport_layer/etcp.c @@ -661,6 +661,22 @@ void etcp_conn_set_peer_node_id(struct ETCP_CONN* conn, uint64_t peer_node_id) { if (conn->peer_node_id == peer_node_id) return; + if (peer_node_id != 0 && conn->instance && conn->instance->connections) { + struct ll_entry* clash = queue_find_data_by_index(conn->instance->connections, (const uint8_t*)&peer_node_id); + while (clash) { + struct conn_queue_entry* cqe = (struct conn_queue_entry*)clash->data; + if (cqe && cqe->conn && cqe->conn != conn && cqe->peer_node_id == peer_node_id) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, + "!!!!!!!!!!!! PEER NODE ID COLLISION !!!!!!!!!!!! " + "conn=%s trying peer=0x%016llx already_owned_by=%s(peer=0x%016llx state=%d)", + conn->log_name, (unsigned long long)peer_node_id, + cqe->conn->log_name, (unsigned long long)cqe->conn->peer_node_id, cqe->conn->state); + return; + } + clash = queue_find_next_by_index(conn->instance->connections, (const uint8_t*)&peer_node_id, clash); + } + } + if (conn->state == 0 && conn->peer_node_id == 0 && peer_node_id != 0) { /* первый раз ставим реальный peer_id: переиндексировать из key=0 в key=peer_node_id */ DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] peer_node_id 0 -> 0x%llx, reindexing pending conn", diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index 4c70a87a..daa348fa 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -371,6 +371,15 @@ static int insert_link_queue(struct ETCP_SOCKET* e_sock, struct ETCP_LINK* link) struct ll_entry* dup_qe = queue_find_data_by_index(e_sock->links_queue, key); if (dup_qe) { struct link_queue_entry* dup_lqe = (struct link_queue_entry*)dup_qe->data; + if (dup_lqe->link && dup_lqe->link->etcp != link->etcp) { + DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, + "!!!!!!!!!!!! LINK ADDR COLLISION !!!!!!!!!!!! " + "addr=%s new_conn=%s(peer=0x%016llx) already_used_by=%s(peer=0x%016llx)", + sockaddr_storage_to_str(&link->remote_addr).str, + link->etcp->log_name, (unsigned long long)link->etcp->peer_node_id, + dup_lqe->link->etcp->log_name, (unsigned long long)dup_lqe->link->etcp->peer_node_id); + return -1; + } DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "insert_link_queue: replacing stale DUP addr in [%s] old_link=%p", e_sock->name, dup_lqe->link); if (dup_lqe->link) dup_lqe->link->link_queue_entry = NULL; queue_remove_data(e_sock->links_queue, dup_qe);