diff --git a/src/etcp_connections.c b/src/etcp_connections.c index ad95909c..8606dd55 100644 --- a/src/etcp_connections.c +++ b/src/etcp_connections.c @@ -50,7 +50,7 @@ static void tcp_server_on_link(struct stcp_link *link, void *arg) { void etcp_connections_read_callback_socket(socket_t sock, void* arg); static void etcp_link_remove_from_connections(struct ETCP_SOCKET* conn, struct ETCP_LINK* link); -static void etcp_link_send_init(struct ETCP_LINK* link, uint8_t reset); +static void etcp_link_send_init(struct ETCP_LINK* link, uint8_t reset, uint8_t collision); //static int etcp_link_send_reset(struct ETCP_LINK* link); static void etcp_link_init_timer_cbk(void* arg); static void etcp_link_send_keepalive(struct ETCP_LINK* link); @@ -101,9 +101,9 @@ static void burst_resp_timeout_cb(void* arg) { #define INIT_TIMEOUT_INITIAL 500 #define INIT_TIMEOUT_MAX 50000 -static void etcp_link_send_init(struct ETCP_LINK* link, uint8_t reset) { +static void etcp_link_send_init(struct ETCP_LINK* link, uint8_t reset, uint8_t collision) { DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "link=%p, is_server=%d, reset=%d", link, link ? link->is_server : -1, reset); + DEBUG_INFO(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; struct ETCP_DGRAM* dgram = u_malloc(PACKET_DATA_SIZE); @@ -134,6 +134,7 @@ static void etcp_link_send_init(struct ETCP_LINK* link, uint8_t reset) { memset(req->src_ipv4, 0, 4); memset(req->src_port, 0, 2); } + req->collision = collision; memcpy(req->ed25519_pubkey, link->etcp->instance->my_ed25519_pubkey, SC_PUBKEY_SIZE); size_t offset = ETCP_INIT_REQ_SIZE; @@ -196,8 +197,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 == 1) etcp_link_send_init(link,1);// init (with etcp reset) - else etcp_link_send_init(link,0);// no etcp reset (reinit) + 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) } void etcp_link_restart_init_timer(struct ETCP_LINK* link) { @@ -212,7 +213,7 @@ void etcp_link_enter_init(struct ETCP_LINK* link) {// if (!link) return; link->link_state = 1; // handshake if (link->is_server != 0) return; - etcp_link_send_init(link,1);// init with reset + etcp_link_send_init(link,1,0);// init with reset etcp_link_restart_init_timer(link); } @@ -222,7 +223,7 @@ void etcp_link_enter_reinit(struct ETCP_LINK* link) { link->link_state = 2; // reconnect etcp_on_link_down(link->etcp); if (link->is_server != 0) return; - etcp_link_send_init(link,0);// init without reset + etcp_link_send_init(link,0,0);// init without reset if (link->keepalive_timer) {// keepalive заменяяется reinit запросами uasync_cancel_timeout(link->etcp->instance->ua, link->keepalive_timer); @@ -1743,19 +1744,48 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { link->remote_only_local = req->only_local; link->remote_type = req->type; - // For CHANNEL_INIT (0x04): if link already initialized - no reset, otherwise reset - // For INIT_REQUEST (0x02): always reset - // Check session_id: if same - no reinit, if different - client restarted, do reinit - if (code == ETCP_INIT_REQUEST || conn->session_id != session_id || !conn->got_initial_pkt) { - send_reset = 1; // Client explicitly requested reset, or new session, or server waiting for first packet - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] REINIT existing link: code=0x%02x sess=%08x→%08x got_init=%d initialized=%d links_up=%d reset_done=%d", - conn->log_name, code, conn->session_id, session_id, conn->got_initial_pkt, conn->initialized, conn->links_up, conn->reset_done); + // ── Collision handling ── + if (req->collision) { + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] INIT collision=1 from peer=%016llx — remote is master, becoming slave", + conn->log_name, (unsigned long long)peer_id); conn->session_id = session_id; - if (!conn->reset_done) etcp_conn_reinit(conn); - else DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] REINIT skipped — reset_done=1", conn->log_name); + etcp_conn_reinit(conn); + send_reset = 1; } else { - send_reset = 0; // Same session, no reinit needed - DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "same session_id, skip reinit"); + // 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; } + if (ml) { + if (conn->instance->node_id < peer_id) { + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] COLLISION: we=master node_id=%016llx < peer=%016llx — sending collision INIT", + conn->log_name, (unsigned long long)conn->instance->node_id, (unsigned long long)peer_id); + memory_pool_free(e_sock->instance->pkt_pool, pkt); + etcp_link_send_init(ml, 0, 1); + return; + } + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] COLLISION: peer smaller, yielding master, processing as slave", + conn->log_name); + } + + // Normal reinit check (skip if reset_done=1 unless session_id changed) + if (!conn->reset_done) { + if (code == ETCP_INIT_REQUEST || conn->session_id != session_id || !conn->got_initial_pkt) { + send_reset = 1; + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%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); + conn->session_id = session_id; + etcp_conn_reinit(conn); + } + } else if (conn->session_id != session_id) { + send_reset = 1; + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] REINIT existing link (session changed): sess=%08x→%08x", + conn->log_name, conn->session_id, session_id); + conn->session_id = session_id; + etcp_conn_reinit(conn); + } else { + send_reset = 0; + DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "same session_id, skip reinit"); + } } // Cancel existing timers @@ -1773,14 +1803,41 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { link->remote_only_local = req->only_local; link->remote_type = req->type; - // For new links: reset if client requested or session changed - if (code == ETCP_INIT_REQUEST || conn->session_id != session_id || !conn->got_initial_pkt) { - send_reset = 1; - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] REINIT new link: code=0x%02x sess=%08x→%08x got_init=%d initialized=%d links_up=%d reset_done=%d", - conn->log_name, code, conn->session_id, session_id, conn->got_initial_pkt, conn->initialized, conn->links_up, conn->reset_done); + // ── Collision handling ── + if (req->collision) { + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] INIT collision=1 from peer=%016llx — remote is master, becoming slave", + conn->log_name, (unsigned long long)peer_id); conn->session_id = session_id; - if (!conn->reset_done) etcp_conn_reinit(conn); - else DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] REINIT skipped — reset_done=1", conn->log_name); + etcp_conn_reinit(conn); + send_reset = 1; + } else { + // Check if WE have an outbound (master) link + struct ETCP_LINK* ml = conn->links; + while (ml) { if (ml->is_server == 0) break; ml = ml->next; } + if (ml && conn->instance->node_id < peer_id) { + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] COLLISION: we=master node_id=%016llx < peer=%016llx — sending collision INIT", + conn->log_name, (unsigned long long)conn->instance->node_id, (unsigned long long)peer_id); + memory_pool_free(e_sock->instance->pkt_pool, pkt); + etcp_link_send_init(ml, 0, 1); + return; + } + + // Normal reinit check (skip if reset_done=1 unless session_id changed) + if (!conn->reset_done) { + if (code == ETCP_INIT_REQUEST || conn->session_id != session_id || !conn->got_initial_pkt) { + send_reset = 1; + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] REINIT new 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); + conn->session_id = session_id; + etcp_conn_reinit(conn); + } + } else if (conn->session_id != session_id) { + send_reset = 1; + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] REINIT new link (session changed): sess=%08x→%08x", + conn->log_name, conn->session_id, session_id); + conn->session_id = session_id; + etcp_conn_reinit(conn); + } } link->keepalive_interval=(req->keepalive[0]<<8) | req->keepalive[1]; link->recovery_interval=((req->recovery[0]<<8) | req->recovery[1])*100;// timebase в link, timebase/100 в кодограмме diff --git a/src/etcp_connections.h b/src/etcp_connections.h index e3501779..85010d65 100644 --- a/src/etcp_connections.h +++ b/src/etcp_connections.h @@ -76,8 +76,8 @@ struct ETCP_INIT_REQUEST_PKT { // V2 fields (NAT_DIRECT detection): uint8_t src_ipv4[4]; // 23: client interface_addr IPv4 (big-endian, 0 if N/A) uint8_t src_port[2]; // 27: client interface_addr port (big-endian) - // V3 fields: - uint8_t ed25519_pubkey[SC_PUBKEY_SIZE]; // 29: client Ed25519 pubkey (32 bytes) + uint8_t collision; // 29: 1 = cross-connect, remote claims master + uint8_t ed25519_pubkey[SC_PUBKEY_SIZE]; // 30: client Ed25519 pubkey (32 bytes) } __attribute__((packed)); #define ETCP_INIT_REQ_SIZE sizeof(struct ETCP_INIT_REQUEST_PKT)