Browse Source

ETCP: реализация восстановления связи при падении всех линков

- Добавлена функция etcp_loadbalancer_get_link_status() для проверки состояния связи
- Клиент: при падении всех линков запускается отправка ETCP_INIT_REQUEST_NOINIT (0x04)
- Сервер: при приеме NOINIT отвечает ETCP_INIT_RESPONSE_NOINIT (0x05) без сброса если линк initialized
- При приеме любого пакета отменяется таймер восстановления
- ETCP_INIT_RESPONSE (0x03) вызывает сброс всего ETCP_CONN
- Переименованы константы: ETCP_CHANNEL_INIT -> ETCP_INIT_REQUEST_NOINIT,
  ETCP_CHANNEL_RESPONSE -> ETCP_INIT_RESPONSE_NOINIT

Все 23 теста проходят
nodeinfo-routing-update
Evgeny 8 months ago
parent
commit
00e09c767e
  1. 183
      src/etcp_connections.c
  2. 4
      src/etcp_connections.h
  3. 27
      src/etcp_loadbalancer.c
  4. 4
      src/etcp_loadbalancer.h

183
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++) {

4
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)

27
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;

4
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
}

Loading…
Cancel
Save