From 93871c592b5d4f5e0ccf9e50550b04a109bdf6ae Mon Sep 17 00:00:00 2001 From: Evgeny Date: Sat, 18 Jul 2026 10:52:54 +0300 Subject: [PATCH] db_sync: timeout for stuck sync_state=1, handle ERROR properly - Add sync_start_tb to SI_PEER, set in db_sync_initiate_sync - peer_check_cb: after 15s timeout reset sync_state=0 with WARN - db_handle_error: accept si, reset sync_state=0, DEBUG_ERROR with details - DB_SYNC_SYNC_TIMEOUT=15 added to db_sync.h --- src/db_sync.c | 39 +++++++++++++++++++++++++++++------ src/db_sync.h | 1 + tests/test_icmp_proxy.c | 2 +- tests/test_socks_http_proxy.c | 2 +- tests/test_udp_proxy.c | 2 +- 5 files changed, 37 insertions(+), 9 deletions(-) diff --git a/src/db_sync.c b/src/db_sync.c index d8e90417..b3354fbc 100644 --- a/src/db_sync.c +++ b/src/db_sync.c @@ -38,6 +38,7 @@ struct SI_PEER { uint64_t node_id; uint32_t synced_pos; uint8_t sync_state; + uint64_t sync_start_tb; }; struct DB_SYNC_INSTANCE { @@ -224,6 +225,7 @@ static struct SI_PEER* si_peer_add(struct DB_SYNC_INSTANCE* si, uint64_t node_id p->node_id = node_id; p->synced_pos = 0; p->sync_state = 0; + p->sync_start_tb = 0; return p; } @@ -581,12 +583,16 @@ static void si_delivery_update(struct DB_SYNC_INSTANCE* si, uint64_t ts, uint64_ // Sync protocol handlers // ============================================================ -static void db_handle_error(struct DB_SYNC* db, uint64_t src, const uint8_t* p, size_t len) +static void db_handle_error(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { - if (len < 1) return; - DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "from=%016llx code=%u", (unsigned long long)src, p[0]); - DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "db_sync: ERROR from %016llx code=%u", (unsigned long long)src, p[0]); - (void)db; + if (len < 1 || !si) return; + const char* what = (p[0] == DB_ERR_NOT_FOUND) ? "NOT_FOUND" : + (p[0] == DB_ERR_DISABLED) ? "DISABLED" : "UNKNOWN"; + DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, + "db_sync: ERROR from %016llx code=%u (%s) tbl=%s hash=%016llx", + (unsigned long long)src, p[0], what, SI_TBL(si), (unsigned long long)si->hash); + struct SI_PEER* sp = si_peer_find(si, src); + if (sp) sp->sync_state = 0; } static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) @@ -1075,7 +1081,7 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) case DB_MSG_PUSH: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle PUSH"); db_handle_push(si, src, payload, plen); break; case DB_MSG_ACK_PUSH: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle ACK_PUSH"); db_handle_ack_push(si, src, payload, plen); break; case DB_MSG_SYNC_DONE: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle SYNC_DONE"); db_handle_sync_done(si, src, payload, plen); break; - case DB_MSG_ERROR: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle ERROR"); db_handle_error(db, src, payload, plen); break; + case DB_MSG_ERROR: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle ERROR"); db_handle_error(si, src, payload, plen); break; default: DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "unknown msg type 0x%02x from %016llx", type, (unsigned long long)src); break; @@ -1160,6 +1166,8 @@ static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t pid) { uint32_t mc = db_count(si); DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "pid=%016llx my=%u tbl=%s", (unsigned long long)pid, mc, SI_TBL(si)); + struct SI_PEER* p = si_peer_find(si, pid); + if (p) p->sync_start_tb = get_time_tb(); uint8_t msg[5]; msg[0] = DB_MSG_INIT_SYNC; @@ -1252,6 +1260,25 @@ static void db_sync_peer_check_cb(void* arg) else { total_skipped++; } } + { + uint64_t now = get_time_tb(); + uint64_t to_tb = DB_SYNC_SYNC_TIMEOUT * 10000u; + for (int i = 0; i < db->instance_count; i++) { + struct DB_SYNC_INSTANCE* si = &db->instances[i]; + if (!si->enabled) continue; + for (int j = 0; j < si->peer_count; j++) { + struct SI_PEER* p = &si->peers[j]; + if (p->sync_state == 1 && p->sync_start_tb > 0 && now - p->sync_start_tb > to_tb) { + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, + "db_sync: no response for peer=%016llx tbl=%s elapsed=%llu ms, resetting sync_state", + (unsigned long long)p->node_id, SI_TBL(si), + (unsigned long long)((now - p->sync_start_tb) / 10)); + p->sync_state = 0; + } + } + } + } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "peer_check: instances=%d synced=%d skipped=%d", db->instance_count, total_synced, total_skipped); diff --git a/src/db_sync.h b/src/db_sync.h index 521cea11..c620db54 100644 --- a/src/db_sync.h +++ b/src/db_sync.h @@ -72,6 +72,7 @@ struct DB_SYNC_INSTANCE; #define DB_SYNC_DEFAULT_TTL 86400 #define DB_SYNC_DEFAULT_MAPSIZE (100UL * 1024 * 1024) #define DB_SYNC_PEER_CHECK_INTERVAL 5 +#define DB_SYNC_SYNC_TIMEOUT 15 #define DB_SYNC_TTL_INTERVAL 3600 // Record flags diff --git a/tests/test_icmp_proxy.c b/tests/test_icmp_proxy.c index 692a7dfa..43137fb4 100644 --- a/tests/test_icmp_proxy.c +++ b/tests/test_icmp_proxy.c @@ -51,7 +51,7 @@ static const char* cfg_node_client(void) { "enabled=yes\n" "tun_name=tun_tcp\n" "tun_ip=10.99.0.1\n" - "via_node=0x2b2d74c71e38c7fe\n"); + "via_node=0x5a77374f79310824\n"); return buf; } diff --git a/tests/test_socks_http_proxy.c b/tests/test_socks_http_proxy.c index bbfa0d6d..832ca54e 100644 --- a/tests/test_socks_http_proxy.c +++ b/tests/test_socks_http_proxy.c @@ -207,7 +207,7 @@ static char* make_cfg_client(void) { "socks_addr=127.0.0.1:%d\n" "http_proxy_enabled=yes\n" "http_proxy_addr=127.0.0.1:%d\n" - "via_node=0x2b2d74c71e38c7fe\n", + "via_node=0x5a77374f79310824\n", g_cli_port, g_srv_port, g_socks_port, g_http_proxy_port); return buf; } diff --git a/tests/test_udp_proxy.c b/tests/test_udp_proxy.c index 7ade1bf0..8e6654f0 100644 --- a/tests/test_udp_proxy.c +++ b/tests/test_udp_proxy.c @@ -50,7 +50,7 @@ static const char* cfg_node_client(void) { "enabled=yes\n" "tun_name=tun_tcp\n" "tun_ip=10.99.0.1\n" - "via_node=0x2b2d74c71e38c7fe\n"); + "via_node=0x5a77374f79310824\n"); return buf; }