Browse Source

merkle_sync: fix infinite WAL growth from bg_check

- _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
topo_upd
evgeny 2 months ago
parent
commit
d8921f0e3d
  1. 13
      lib/u_async.c
  2. 69
      src/chat/merkle_sync.c
  3. 1
      src/transport_layer/etcp_connections.c
  4. 1
      src/transport_layer/stcp_client.c
  5. 25
      src/transport_layer/stcp_link.c
  6. 8
      tests/test_merkle_sync.c

13
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();

69
src/chat/merkle_sync.c

@ -12,7 +12,7 @@
#include <sqlite3.h>
#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 ── */

1
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);
}

1
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;
}

25
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);
}

8
tests/test_merkle_sync.c

@ -19,15 +19,13 @@
#include <sqlite3.h>
#include <openssl/evp.h>
/* 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.

Loading…
Cancel
Save