diff --git a/src/etcp_connections.c b/src/etcp_connections.c index 9e588594..3414177c 100644 --- a/src/etcp_connections.c +++ b/src/etcp_connections.c @@ -42,22 +42,23 @@ static void keepalive_timer_cb(void* arg); #define INIT_TIMEOUT_INITIAL 500 #define INIT_TIMEOUT_MAX 50000 -static void etcp_link_send_init(struct ETCP_LINK* link) { +static void etcp_link_send_init_internal(struct ETCP_LINK* link, uint8_t reset) { DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "etcp_link_send_init link=%p, is_server=%d", link, link ? link->is_server : -1); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "etcp_link_send_init link=%p, is_server=%d, reset=%d", link, link ? link->is_server : -1, reset); if (!link || !link->etcp || !link->etcp->instance) return; - + struct ETCP_DGRAM* dgram = malloc(sizeof(struct ETCP_DGRAM) + 100); if (!dgram) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "etcp_link_send_init: malloc failed"); return; } - + dgram->link = link; dgram->noencrypt_len = SC_PUBKEY_SIZE; size_t offset = 0; - - dgram->data[offset++] = ETCP_INIT_REQUEST; + + // reset=1: ETCP_INIT_REQUEST (0x02), reset=0: ETCP_INIT_REQUEST_NOINIT (0x04) + dgram->data[offset++] = reset ? ETCP_INIT_REQUEST : ETCP_INIT_REQUEST_NOINIT; uint64_t node_id = link->etcp->instance->node_id; dgram->data[offset++] = (node_id >> 56) & 0xFF; @@ -94,9 +95,9 @@ static void etcp_link_send_init(struct ETCP_LINK* link) { etcp_encrypt_send(dgram); free(dgram); - + link->init_retry_count++; - + if (!link->init_timer && link->is_server == 0) { link->init_timeout = INIT_TIMEOUT_INITIAL; link->init_timer = uasync_set_timeout(link->etcp->instance->ua, link->init_timeout, link, etcp_link_init_timer_cbk); @@ -110,6 +111,16 @@ static void etcp_link_send_init(struct ETCP_LINK* link) { } } +// Wrapper for backward compatibility - sends init WITH reset (0x02) +static void etcp_link_send_init(struct ETCP_LINK* link) { + etcp_link_send_init_internal(link, 1); +} + +// Send init WITHOUT reset (0x04) - for link recovery +static void etcp_link_send_channel_init(struct ETCP_LINK* link) { + etcp_link_send_init_internal(link, 0); +} + static void etcp_link_init_timer_cbk(void* arg) { DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); struct ETCP_LINK* link = (struct ETCP_LINK*)arg; @@ -142,6 +153,58 @@ static void etcp_link_send_keepalive(struct ETCP_LINK* link) { free(dgram); } +// Check if all links for an ETCP_CONN are down +// Returns 1 if all links are down or no links exist, 0 otherwise +static int etcp_all_links_down(struct ETCP_CONN* etcp) { + if (!etcp || !etcp->links) return 1; + + struct ETCP_LINK* l = etcp->links; + while (l) { + if (l->link_status == 1) { + return 0; // At least one link is up + } + l = l->next; + } + return 1; // All links are down +} + +// Start link recovery process - send CHANNEL_INIT (0x04) on all links +static void etcp_start_link_recovery(struct ETCP_CONN* etcp) { + if (!etcp) return; + + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Starting link recovery - all links are down", + etcp->log_name ? etcp->log_name : "????→????"); + + struct ETCP_LINK* link = etcp->links; + while (link) { + if (link->is_server == 0) { // Only client links + // Reset link state for recovery + link->initialized = 0; + // Send CHANNEL_INIT (0x04) without reset + etcp_link_send_channel_init(link); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Sent CHANNEL_INIT on link %p for recovery", + etcp->log_name, link); + } + link = link->next; + } +} + +// Cancel init_timer for all links of an ETCP_CONN +static void etcp_cancel_all_init_timers(struct ETCP_CONN* etcp) { + if (!etcp) return; + + struct ETCP_LINK* link = etcp->links; + while (link) { + if (link->init_timer) { + uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer); + link->init_timer = NULL; + DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "[%s] Cancelled init_timer on link %p", + etcp->log_name, link); + } + link = link->next; + } +} + // Keepalive timer callback static void keepalive_timer_cb(void* arg) { DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); @@ -150,6 +213,15 @@ static void keepalive_timer_cb(void* arg) { link->keepalive_timer = NULL; + // Check if all links are down and start recovery if needed (client only) + if (link->is_server == 0 && etcp_all_links_down(link->etcp)) { + DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "[%s] All links are down, starting recovery", + link->etcp->log_name); + etcp_start_link_recovery(link->etcp); + // Don't restart keepalive timer during recovery + return; + } + // Skip if link is not initialized if (!link->initialized) { DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "[%s] Keepalive skipped - link not initialized", @@ -575,6 +647,46 @@ void etcp_link_close(struct ETCP_LINK* link) { free(link); } +// Reset link state (for INIT_RESPONSE with reset) +static void etcp_link_reset(struct ETCP_LINK* link) { + DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); + if (!link) return; + + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Resetting link %p (local_id=%d)", + link->etcp->log_name, link, link->local_link_id); + + // Cancel all timers + if (link->init_timer) { + uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer); + link->init_timer = NULL; + } + if (link->shaper_timer) { + uasync_cancel_timeout(link->etcp->instance->ua, link->shaper_timer); + link->shaper_timer = NULL; + } + + // Reset state + link->initialized = 0; + link->link_status = 0; + link->recv_keepalive = 0; + link->remote_keepalive = 0; + link->init_retry_count = 0; + link->init_timeout = 0; + + // Reset shaper state + link->shaper_load_time_tb = 0; + link->shaper_sub_nanotime = 0; + link->shaper_state = 0; + + // Reset counters + link->encrypt_errors = 0; + link->decrypt_errors = 0; + link->send_errors = 0; + link->recv_errors = 0; + link->total_encrypted = 0; + link->total_decrypted = 0; +} + int etcp_encrypt_send(struct ETCP_DGRAM* dgram) { DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); // DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_encrypt_send called, link=%p", dgram ? dgram->link : NULL); @@ -747,7 +859,7 @@ static void etcp_connections_read_callback_socket(socket_t sock, void* arg) { } *ack_hdr=(void*)&pkt->data[0]; uint64_t peer_id; memcpy(&peer_id, &ack_hdr->id[0], 8); - if (ack_hdr->code!=ETCP_INIT_REQUEST && ack_hdr->code!=ETCP_CHANNEL_INIT) { + if (ack_hdr->code!=ETCP_INIT_REQUEST && ack_hdr->code!=ETCP_INIT_REQUEST_NOINIT) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_connections_read_callback: not an init packet, code=%02x", ack_hdr->code); errorcode=4; goto ec_fr; @@ -779,10 +891,40 @@ static void etcp_connections_read_callback_socket(socket_t sock, void* arg) { else {// check keys если существующее подключение if (memcmp(conn->crypto_ctx.peer_public_key, sc.peer_public_key, SC_PUBKEY_SIZE)) { errorcode=5; DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "etcp_connections_read_callback: peer key mismatch for node %llu", (unsigned long long)peer_id); goto ec_fr; }// коллизия - peer id совпал а ключи разные. } - link = etcp_link_new(conn, e_sock, &addr, 1); - if (!link) { if (new_conn) etcp_connection_close(conn); errorcode=66; DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "etcp_connections_read_callback: failed to create link for connection"); goto ec_fr; }// облом - link->remote_link_id = ack_hdr->link_id; - if (ack_hdr->code==0x02) etcp_conn_reset(conn); + // Check if link already exists (for CHANNEL_INIT recovery) + struct ETCP_LINK* existing_link = etcp_link_find_by_addr(e_sock, &addr); + uint8_t send_reset = 0; + + if (existing_link && existing_link->etcp == conn) { + // Link exists - reuse it for recovery + link = existing_link; + link->remote_link_id = ack_hdr->link_id; + + // For CHANNEL_INIT (0x04): if link already initialized - no reset, otherwise reset + // For INIT_REQUEST (0x02): always reset + if (ack_hdr->code == ETCP_INIT_REQUEST_NOINIT && link->initialized) { + send_reset = 0; // Link is up, respond without reset + } else { + send_reset = 1; // INIT_REQUEST (0x02) or uninitialized link - send reset + etcp_conn_reset(conn); + } + + // Cancel existing timers + if (link->init_timer) { + uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer); + link->init_timer = NULL; + } + } else { + // Create new link + link = etcp_link_new(conn, e_sock, &addr, 1); + if (!link) { if (new_conn) etcp_connection_close(conn); errorcode=66; DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "etcp_connections_read_callback: failed to create link for connection"); goto ec_fr; }// облом + link->remote_link_id = ack_hdr->link_id; + + // For new links: INIT_REQUEST (0x02) causes reset, CHANNEL_INIT (0x04) does not + if (ack_hdr->code == ETCP_INIT_REQUEST) { + etcp_conn_reset(conn); + } + } struct { uint8_t code; @@ -792,7 +934,13 @@ static void etcp_connections_read_callback_socket(socket_t sock, void* arg) { uint8_t peer_ipv4[4]; uint8_t peer_port[2]; } *ack_repl_hdr=(void*)&pkt->data[0]; - ack_repl_hdr->code+=1; + + // Set response code: 0x03 (with reset) or 0x05 (without reset) + if (send_reset || ack_hdr->code == ETCP_INIT_REQUEST) { + ack_repl_hdr->code = ETCP_INIT_RESPONSE; // 0x03 - with reset + } else { + ack_repl_hdr->code = ETCP_INIT_RESPONSE_NOINIT; // 0x05 - without reset + } memcpy(ack_repl_hdr->id, &e_sock->instance->node_id, 8); int mtu=e_sock->instance->config->global.mtu; ack_repl_hdr->mtu[0]=mtu>>8; @@ -866,11 +1014,16 @@ process_decrypted: link->etcp->log_name, link, link->local_link_id); } + // Cancel all init timers - link is alive, no need for recovery + etcp_cancel_all_init_timers(link->etcp); + size_t offset = 0; uint8_t code = pkt->data[offset++]; - if (code == ETCP_INIT_RESPONSE || code == ETCP_CHANNEL_RESPONSE) { + if (code == ETCP_INIT_RESPONSE || code == ETCP_INIT_RESPONSE_NOINIT) { // Parse response + // ETCP_INIT_RESPONSE (0x03) - reset entire ETCP_CONN + // ETCP_INIT_RESPONSE_NOINIT (0x05) - no reset if (code == ETCP_INIT_RESPONSE) etcp_conn_reset(link->etcp); uint64_t server_node_id = 0; for (int i = 0; i < 8; i++) { diff --git a/src/etcp_connections.h b/src/etcp_connections.h index bf6e1d8d..72bd9cdd 100644 --- a/src/etcp_connections.h +++ b/src/etcp_connections.h @@ -13,8 +13,8 @@ // Типы кодограмм протокола #define ETCP_INIT_REQUEST 0x02 #define ETCP_INIT_RESPONSE 0x03 -#define ETCP_CHANNEL_INIT 0x04 -#define ETCP_CHANNEL_RESPONSE 0x05 +#define ETCP_INIT_REQUEST_NOINIT 0x04 +#define ETCP_INIT_RESPONSE_NOINIT 0x05 #pragma pack(push, 1) diff --git a/src/etcp_loadbalancer.c b/src/etcp_loadbalancer.c index 20545958..953d5545 100644 --- a/src/etcp_loadbalancer.c +++ b/src/etcp_loadbalancer.c @@ -152,9 +152,9 @@ void loadbalancer_link_ready(struct ETCP_LINK* link) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "loadbalancer_link_ready: invalid link (%p)", link); return; } - + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "loadbalancer_link_ready: link=%p now ready, notifying ETCP_CONN", link); - + // Call ETCP_CONN resume (assumes link_ready_for_send_fn in ETCP_CONN; add to etcp.h: void (*link_ready_for_send_fn)(struct ETCP_CONN*);) if (link->etcp->link_ready_for_send_fn) { link->etcp->link_ready_for_send_fn(link->etcp); @@ -165,6 +165,29 @@ void loadbalancer_link_ready(struct ETCP_LINK* link) { } } +// Get ETCP link status: 1 = at least one link is up, 0 = all links down or no links +int etcp_loadbalancer_get_link_status(struct ETCP_CONN* etcp) { + if (!etcp) { + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_loadbalancer_get_link_status: NULL etcp"); + return 0; + } + + struct ETCP_LINK* link = etcp->links; + int alive_count = 0; + + while (link) { + if (link->link_status == 1) { + alive_count++; + } + link = link->next; + } + + DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] link status check: %d alive links", + etcp->log_name ? etcp->log_name : "????→????", alive_count); + + return (alive_count > 0) ? 1 : 0; +} + // Shaper timer callback static void shaper_timer_cb(void* arg) { struct ETCP_LINK* link = (struct ETCP_LINK*)arg; diff --git a/src/etcp_loadbalancer.h b/src/etcp_loadbalancer.h index f83252fb..594484cd 100644 --- a/src/etcp_loadbalancer.h +++ b/src/etcp_loadbalancer.h @@ -1,4 +1,4 @@ -// etcp_loadbalancer.h - Load Balancer for ETCP Channels +// etcp_loadbalancer.h - Load Balancer for ETCP Channels #ifndef ETCP_LOADBALANCER_H #define ETCP_LOADBALANCER_H @@ -26,6 +26,8 @@ void etcp_loadbalancer_send(struct ETCP_DGRAM* dgram); // сообщаем в loadbalancer о готовности линка void loadbalancer_link_ready(struct ETCP_LINK* link); +// Получить состояние связи ETCP: 1 - есть живой линк, 0 - все недоступны +int etcp_loadbalancer_get_link_status(struct ETCP_CONN* etcp); #ifdef __cplusplus }