From d8921f0e3d4c1562ba20921d3fc33094b0f0f64b Mon Sep 17 00:00:00 2001 From: evgeny Date: Sun, 9 Aug 2026 19:25:36 +0300 Subject: [PATCH] merkle_sync: fix infinite WAL growth from bg_check MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - _recompute_bucket: skip INSERT OR REPLACE when hash unchanged, also skip DELETE when row doesn't exist (was always writing) - bg_check interval: 100ms → 5000ms (was writing 10x/sec × 4KB/page into WAL = ~1.4GB over hours) - bg_timer: round-robin across channels instead of only first - bg_timer: periodic WAL checkpoint (every ~60s) to prevent unbounded WAL growth from any source - test: sync ms_session struct copy with real definition - stcp_link: handle race where conn already CLOSED before close - stcp_client: remove double-free c->free_on_close = cli - etcp_connections: stcp_link_close on TCP link down - u_async: round sub-ms timeouts to 1ms to avoid busy-loop --- lib/u_async.c | 13 +++-- src/chat/merkle_sync.c | 69 ++++++++++++++++++-------- src/transport_layer/etcp_connections.c | 1 + src/transport_layer/stcp_client.c | 1 - src/transport_layer/stcp_link.c | 25 ++++++++-- tests/test_merkle_sync.c | 8 ++- 6 files changed, 82 insertions(+), 35 deletions(-) diff --git a/lib/u_async.c b/lib/u_async.c index b87b99d7..75f3dac0 100644 --- a/lib/u_async.c +++ b/lib/u_async.c @@ -432,9 +432,12 @@ static void timeval_add_tb(struct timeval* tv, int dt) { tv->tv_usec %= 1000000; } -// Convert timeval to milliseconds (uint64_t) +// Convert timeval to milliseconds. Sub-ms values (>0, <1ms) round up to 1ms +// to avoid 0ms timeout → busy-loop in epoll_wait/select. static uint64_t timeval_to_ms(const struct timeval* tv) { - return (uint64_t)tv->tv_sec * 1000ULL + (uint64_t)tv->tv_usec / 1000ULL; + uint64_t ms = (uint64_t)tv->tv_sec * 1000ULL + (uint64_t)tv->tv_usec / 1000ULL; + if (ms == 0 && tv->tv_usec > 0) ms = 1; + return ms; } @@ -1033,11 +1036,11 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) { int timeout_ms; if (timeout_tb < 0 && (next_timeout.tv_sec > 0 || next_timeout.tv_usec > 0)) { - timeout_ms = (poll_timeout.tv_sec * 1000) + (poll_timeout.tv_usec / 1000); + timeout_ms = (int)timeval_to_ms(&poll_timeout); } else if (timeout_tb < 0) { timeout_ms = -1; // Infinite } else { - timeout_ms = (poll_timeout.tv_sec * 1000) + (poll_timeout.tv_usec / 1000); + timeout_ms = (int)timeval_to_ms(&poll_timeout); } // Count active sockets @@ -1045,7 +1048,7 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) { if (socket_count == 0 && timeout_ms == -1) { // No sockets and infinite wait - but we have timers? Wait for timer if (ua->timeout_heap->size > 0) { - timeout_ms = (poll_timeout.tv_sec * 1000) + (poll_timeout.tv_usec / 1000); + timeout_ms = (int)timeval_to_ms(&poll_timeout); } else { // Nothing to do - return immediately ua->last_poll_exit_us = get_time_us(); diff --git a/src/chat/merkle_sync.c b/src/chat/merkle_sync.c index 0e01cb60..5abdf72f 100644 --- a/src/chat/merkle_sync.c +++ b/src/chat/merkle_sync.c @@ -12,7 +12,7 @@ #include #define MS_ID "merkle_sync" -#define MS_BG_INTERVAL_MS 100 +#define MS_BG_INTERVAL_MS 5000 #define MS_PENDING_MAX 128 struct bucket_entry { uint8_t level; uint8_t prefix_bytes; uint64_t prefix; }; @@ -106,22 +106,37 @@ static int _recompute_bucket(struct merkle_sync* ms, const char* ns, if (count < 0) { DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: → error", MS_ID); EVP_MD_CTX_free(ctx); return -1; } if (count == 0) { - DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: → empty, deleted", MS_ID); + DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: → empty", MS_ID); EVP_MD_CTX_free(ctx); - sqlite3_stmt* ds = NULL; + sqlite3_stmt* cs = NULL; sqlite3_prepare_v2(db, - "DELETE FROM merkle_tree_hash WHERE namespace=? AND level=? AND prefix64=?", - -1, &ds, NULL); - if (ds) { sqlite3_bind_text(ds, 1, ns, -1, SQLITE_STATIC); sqlite3_bind_int(ds, 2, level); - sqlite3_bind_int64(ds, 3, (sqlite3_int64)prefix64); sqlite3_step(ds); sqlite3_finalize(ds); } + "SELECT COUNT(*) FROM merkle_tree_hash WHERE namespace=? AND level=? AND prefix64=?", -1, &cs, NULL); + if (cs) { sqlite3_bind_text(cs, 1, ns, -1, SQLITE_STATIC); sqlite3_bind_int(cs, 2, level); + sqlite3_bind_int64(cs, 3, (sqlite3_int64)prefix64); + if (sqlite3_step(cs) == SQLITE_ROW && sqlite3_column_int(cs, 0) > 0) { + sqlite3_finalize(cs); + DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: → cleared (was non-empty)", MS_ID); + sqlite3_stmt* ds = NULL; + sqlite3_prepare_v2(db, + "DELETE FROM merkle_tree_hash WHERE namespace=? AND level=? AND prefix64=?", -1, &ds, NULL); + if (ds) { sqlite3_bind_text(ds, 1, ns, -1, SQLITE_STATIC); sqlite3_bind_int(ds, 2, level); + sqlite3_bind_int64(ds, 3, (sqlite3_int64)prefix64); sqlite3_step(ds); sqlite3_finalize(ds); } + } else { sqlite3_finalize(cs); } + } return 0; } - uint8_t hash[MT_HASH_SIZE]; - EVP_DigestFinal_ex(ctx, hash, NULL); + uint8_t new_hash[MT_HASH_SIZE]; + EVP_DigestFinal_ex(ctx, new_hash, NULL); EVP_MD_CTX_free(ctx); - DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: → saved, count=%d", MS_ID, count); + const uint8_t* old_hash = merkle_sync_get_hash(ms->inst, ns, level, prefix64); + if (memcmp(old_hash, new_hash, MT_HASH_SIZE) == 0) { + DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: → unchanged (count=%d)", MS_ID, count); + return 0; + } + + DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: → saved, count=%d", MS_ID, count); sqlite3_stmt* is = NULL; sqlite3_prepare_v2(db, "INSERT OR REPLACE INTO merkle_tree_hash(namespace, level, prefix64, hash, member_count)" @@ -130,7 +145,7 @@ static int _recompute_bucket(struct merkle_sync* ms, const char* ns, sqlite3_bind_text(is, 1, ns, -1, SQLITE_STATIC); sqlite3_bind_int(is, 2, level); sqlite3_bind_int64(is, 3, (sqlite3_int64)prefix64); - sqlite3_bind_blob(is, 4, hash, MT_HASH_SIZE, SQLITE_STATIC); + sqlite3_bind_blob(is, 4, new_hash, MT_HASH_SIZE, SQLITE_STATIC); sqlite3_bind_int(is, 5, count); sqlite3_step(is); sqlite3_finalize(is); } @@ -689,21 +704,33 @@ static void _bg_timer_cb(void* arg) { if (!ms || !ms->initialized || !ms->inst) return; sqlite3* db = _db(ms->inst); if (!db) return; + static int tick = 0; tick++; + sqlite3_stmt* cs = NULL; sqlite3_prepare_v2(db, "SELECT channel_id FROM channels", -1, &cs, NULL); - if (!cs) return; - while (sqlite3_step(cs) == SQLITE_ROW) { - const char* ch_id = (const char*)sqlite3_column_text(cs, 0); - if (!ch_id) continue; - int rc = merkle_sync_bg_check(ms->inst, ch_id); - if (rc < 0) continue; - DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: bg_check ns=%s → %s", MS_ID, ch_id, rc == 0 ? "consistent" : "RECALCULATED"); - break; + if (cs) { + int ch_idx = 0; + while (sqlite3_step(cs) == SQLITE_ROW) { + const char* ch_id = (const char*)sqlite3_column_text(cs, 0); + if (!ch_id) { ch_idx++; continue; } + if ((tick % (ch_idx + 1)) == 0) { + int rc = merkle_sync_bg_check(ms->inst, ch_id); + if (rc >= 0) DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: bg_check ns=%s → %s", MS_ID, ch_id, rc == 0 ? "consistent" : "RECALCULATED"); + } + ch_idx++; + } + sqlite3_finalize(cs); + + if (tick % 12 == 0) { /* every ~60s: try to checkpoint WAL, keep it under control */ + int pn = 0, ckpt = 0; + sqlite3_wal_checkpoint_v2(db, NULL, SQLITE_CHECKPOINT_PASSIVE, &pn, &ckpt); + if (pn > 0 || ckpt > 0) + DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: WAL checkpoint done: %d frames in log, %d checkpointed", MS_ID, pn, ckpt); + } } - sqlite3_finalize(cs); ms->bg_timer = uasync_set_timeout(ms->inst->ua, (uint32_t)(MS_BG_INTERVAL_MS * 10), - ms, _bg_timer_cb, "ms_bg"); + ms, _bg_timer_cb, "ms_bg"); } /* ── Public lifecycle ── */ diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index 9d02adc5..ca992f26 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -1144,6 +1144,7 @@ static void tcp_link_close_cb(struct stcp_link *sl, int err, void *arg) { if (!link || !link->etcp) return; if (link->etcp->state == 2) return; DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] TCP link %d down err=%d, scheduling reconnect", link->etcp->log_name, link->local_link_id, err); + stcp_link_close(sl); link->tcp_link = NULL; etcp_tcp_link_start_reconnect(link); } diff --git a/src/transport_layer/stcp_client.c b/src/transport_layer/stcp_client.c index aacc1ad3..afae2544 100644 --- a/src/transport_layer/stcp_client.c +++ b/src/transport_layer/stcp_client.c @@ -187,7 +187,6 @@ struct stcp_client *stcp_client_connect(struct UASYNC *ua, const char *addr, uin client_send_handshake(c, cli->peer_pubkey, cli->my_ed25519_pubkey); stcp_recv_set(c, SC_PUBKEY_ENC_SIZE + STCP_HS_ENC_SERVER, 0, client_hs_cb); } - c->free_on_close = cli; return cli; } diff --git a/src/transport_layer/stcp_link.c b/src/transport_layer/stcp_link.c index 5bfb796f..1dff31a5 100644 --- a/src/transport_layer/stcp_link.c +++ b/src/transport_layer/stcp_link.c @@ -40,6 +40,9 @@ struct stcp_link { void *close_arg; struct ll_queue *tx_queue; // owned by this link + + struct ll_queue *saved_rx_queue; // saved rx_queue for pre-closed conn cleanup + uint8_t conn_pre_closed; // 1 = conn already CLOSED before stcp_link_close }; // ====== rx dispatch ====== @@ -222,9 +225,15 @@ static void stcp_link_close_impl(void *arg) { if (link->cli) { struct stcp_conn *c = stcp_client_get_conn(link->cli); - if (c && c->rx_queue) queue_free(c->rx_queue); - if (link->tx_queue) queue_free(link->tx_queue); - stcp_client_destroy(link->cli); + if (!link->conn_pre_closed) { + if (c && c->rx_queue) queue_free(c->rx_queue); + if (link->tx_queue) queue_free(link->tx_queue); + stcp_client_destroy(link->cli); + } else { + if (link->saved_rx_queue) queue_free(link->saved_rx_queue); + if (link->tx_queue) queue_free(link->tx_queue); + u_free(link->cli); + } } else if (link->conn) { struct stcp_conn *conn = link->conn; if (conn->rx_queue) { queue_free(conn->rx_queue); conn->rx_queue = NULL; } @@ -241,6 +250,16 @@ void stcp_link_close(struct stcp_link *link) { link->closing = 1; link->etcp_conn = NULL; link->etcp_link = NULL; + + if (link->cli) { + struct stcp_conn *c = stcp_client_get_conn(link->cli); + if (c) { + if (c->hs_timer) { uasync_cancel_timeout(link->cfg.ua, c->hs_timer); c->hs_timer = NULL; } + if (c->state == STCP_STATE_CLOSED || c->state == STCP_STATE_ERROR) { + link->saved_rx_queue = c->rx_queue; c->rx_queue = NULL; link->conn_pre_closed = 1; + } + } + } uasync_call_soon(link->cfg.ua, link, stcp_link_close_impl); } diff --git a/tests/test_merkle_sync.c b/tests/test_merkle_sync.c index 31454174..8ba688b4 100644 --- a/tests/test_merkle_sync.c +++ b/tests/test_merkle_sync.c @@ -19,15 +19,13 @@ #include #include -/* merkle_sync private struct — needed for Stage 2 unit tests */ +/* merkle_sync private struct — needed for Stage 2 unit tests. + * Must match struct ms_session prefix in merkle_sync.c exactly. */ struct ms_session { struct ms_session* next; char ns[64]; uint64_t peer; - uint8_t active; - uint8_t synced; - void* done_cb; /* merkle_sync_done_cb */ - void* cb_arg; + uint8_t sess_state; }; /* merkle_sync private struct — needed for Stage 2 unit tests.