From eaac6114655636089e0992ed0928652df0f9a941 Mon Sep 17 00:00:00 2001 From: evgeny Date: Sun, 4 Oct 2026 16:45:47 +0200 Subject: [PATCH] Gate history on member synchronization and group readiness --- src/chat/db_sync.c | 378 ++++++++++++++------------------------------- 1 file changed, 119 insertions(+), 259 deletions(-) diff --git a/src/chat/db_sync.c b/src/chat/db_sync.c index 9ee34083..34d09322 100644 --- a/src/chat/db_sync.c +++ b/src/chat/db_sync.c @@ -28,15 +28,13 @@ struct DB_SYNC_INSTANCE; struct SI_PEER; static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry); -static void db_sync_on_conn_status(struct ETCP_CONN* conn, int status, void* arg); -static void db_sync_on_conn_up(struct ETCP_CONN* conn, void* arg); -static void db_sync_on_conn_down(struct ETCP_CONN* conn, void* arg); +static void db_sync_on_peer_ready(struct TOPO_GROUP* group, uint64_t pid, int ready, void* arg); static void db_sync_peer_check_cb(void* arg); static void db_sync_resume_peer_check(struct DB_SYNC* db); static void db_sync_instance_ttl_cb(void* arg); static int db_sync_send(struct DB_SYNC_INSTANCE* si, uint64_t dst_node_id, const uint8_t* payload, size_t len); static int db_sync_send_gid(struct DB_SYNC* db, uint64_t node_id, uint64_t group_id, const uint8_t* payload, size_t len); -static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t peer_node_id); +static int db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t peer_node_id); static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_t from, uint32_t count, uint32_t vp, uint32_t want_from); static void si_delivery_update(struct DB_SYNC_INSTANCE* si, uint64_t ts, uint64_t author, uint64_t peer_id); static int db_sync_peer_throttled(struct DB_SYNC* db, struct SI_PEER* p); @@ -85,7 +83,6 @@ struct DB_SYNC { struct UTUN_INSTANCE* inst; sqlite3* db; uint8_t shared_db; /* db == inst->topo_sqlite_db, do not close */ - uint64_t last_connected_tb; struct DB_SYNC_INSTANCE** instances; int instance_count, instance_capacity; void* peer_check_timer; @@ -260,6 +257,7 @@ static struct DB_SYNC_INSTANCE* db_instance_alloc(struct DB_SYNC* db) if (!si) return NULL; si->db_sync = db; si->enabled = 1; + si->sync_gated = 1; db->instances[db->instance_count++] = si; return si; } @@ -307,7 +305,7 @@ static struct SI_PEER* si_peer_add(struct DB_SYNC_INSTANCE* si, uint64_t node_id if (si->peer_count >= si->peer_capacity) { int nc = si->peer_capacity ? si->peer_capacity * 2 : 8; struct SI_PEER* np = u_realloc(si->peers, nc * sizeof(*si->peers)); - if (!np) return NULL; + if (!np) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "history peer allocation failed group=%016llx", si->group_id); return NULL; } si->peers = np; si->peer_capacity = nc; } @@ -526,32 +524,6 @@ static int db_get_ed25519_pubkey(struct DB_SYNC* db, uint64_t node_id, uint8_t o return -1; } -/* dump table: compact one-line per record with pos, id, ts, author, ch8 */ -// Отладочный дамп содержимого таблицы: одна строка на запись (pos, id, ts, author, ch8), макс. 200 записей -static void si_dump_table(struct DB_SYNC_INSTANCE* si, const char* tag) -{ - uint32_t mc = db_count(si); - if (mc > 200) return; /* skip very large tables */ - sqlite3_stmt* stmt; - if (si_prep(si, &stmt, - "SELECT id,timestamp,node_id,chain_hash FROM \"%s\" ORDER BY timestamp,author_signature") != SQLITE_OK) return; - char buf[256]; int total = 0; - while (sqlite3_step(stmt) == SQLITE_ROW) { - uint64_t id = (uint64_t)sqlite3_column_int64(stmt, 0); - uint64_t ts = (uint64_t)sqlite3_column_int64(stmt, 1); - uint64_t auth = (uint64_t)sqlite3_column_int64(stmt, 2); - const uint8_t* ch = (const uint8_t*)sqlite3_column_blob(stmt, 3); - uint64_t ch8 = 0; if (ch) memcpy(&ch8, ch, 8); - snprintf(buf, sizeof(buf), "[%2u] id=%-2llu ts=%-10llu auth=%04llX ch8=%016llX", - total, (unsigned long long)id, (unsigned long long)ts, - (unsigned long long)(auth >> 16), ch8); - DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, " %s %s", tag, buf); - total++; - } - sqlite3_finalize(stmt); - DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, " %s: table=%s mc=%u", tag, SI_TBL(si), mc); -} - // Проверяет Ed25519 подпись автора: verify(pubkey, ts[8] || json, sig). Возвращает 0 если подпись верна static int db_verify_author_sig(struct DB_SYNC_INSTANCE* si, uint64_t ts, uint64_t author_node_id, @@ -708,14 +680,17 @@ static int db_sync_send_gid(struct DB_SYNC* db, uint64_t node_id, uint64_t group struct ETCP_CONN* conn = instance_find_conn(db->inst, node_id); if (!conn || !conn->links_up) { - DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync: send DROP — no conn to %04llX (group=%016llx, type=%02x)", - (unsigned long long)(node_id >> 16), group_id, plen > 0 ? payload[0] : 0, - (void*)conn, conn ? conn->links_up : -1); + DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "history send: transport unavailable peer=%016llx group=%016llx type=%02x", + (unsigned long long)node_id, (unsigned long long)group_id, plen > 0 ? payload[0] : 0); queue_dgram_free(entry); queue_entry_free(entry); return -1; } int ret = etcp_send(conn, entry); - if (ret != 0) { queue_dgram_free(entry); queue_entry_free(entry); } + if (ret != 0) { + DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "history send rejected peer=%016llx group=%016llx type=%02x ret=%d", + (unsigned long long)node_id, (unsigned long long)group_id, plen > 0 ? payload[0] : 0, ret); + queue_dgram_free(entry); queue_entry_free(entry); + } DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "sync: send OK → etcp_send ret=%d conn=%s l_up=%d init=%d q=%p", ret, conn->log_name, conn->links_up, conn->initialized, (void*)conn->send_input_q); return ret; @@ -724,7 +699,18 @@ static int db_sync_send_gid(struct DB_SYNC* db, uint64_t node_id, uint64_t group // Отправляет сообщение через ETCP от имени конкретного инстанса (подставляет group_id инстанса) static int db_sync_send(struct DB_SYNC_INSTANCE* si, uint64_t dst_node_id, const uint8_t* payload, size_t len) { - return db_sync_send_gid(si->db_sync, dst_node_id, si->group_id, payload, len); + struct TOPO_GROUP* group = topo_groups_find(si->db_sync->inst->topo_groups, si->group_id); + int ret = -1; + if (!si->sync_gated && topo_group_peer_ready(group, dst_node_id)) + ret = db_sync_send_gid(si->db_sync, dst_node_id, si->group_id, payload, len); + else DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "history send blocked group=%016llx peer=%016llx members_ready=%d peer_ready=%d", + si->group_id, dst_node_id, !si->sync_gated, topo_group_peer_ready(group, dst_node_id)); + if (ret != 0) { + struct SI_PEER* p = si_peer_find(si, dst_node_id); + if (p) { p->sync_state = 0; p->sync_start_tb = 0; } + db_sync_resume_peer_check(si->db_sync); + } + return ret; } // ============================================================ @@ -779,9 +765,7 @@ static void db_handle_request_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, co { (void)p; (void)len; struct SI_PEER* sp = si_peer_add(si, src); - if (sp && sp->sync_state != 1) { - sp->sync_state = 1; - sp->sync_start_tb = get_time_tb(); + if (sp && is_master(si->db_sync->inst->node_id, src)) { db_sync_initiate_sync(si, src); } } @@ -791,8 +775,9 @@ static void db_handle_error(struct DB_SYNC_INSTANCE* si, uint64_t src, const uin { 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_CHAT_SYNC, + (p[0] == DB_ERR_DISABLED) ? "DISABLED" : + (p[0] == DB_ERR_NOT_READY) ? "NOT_READY" : "UNKNOWN"; + DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← ERROR: code=%u (%s) — peer rejected our sync", SI_SHRT(si), (unsigned long long)(src >> 16), p[0], what); struct SI_PEER* sp = si_peer_find(si, src); @@ -800,6 +785,8 @@ static void db_handle_error(struct DB_SYNC_INSTANCE* si, uint64_t src, const uin DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← ERROR: reset sync_state 2→0, synced_pos stays at %u", SI_SHRT(si), (unsigned long long)(src >> 16), sp->synced_pos); sp->sync_state = 0; + sp->sync_start_tb = 0; + db_sync_resume_peer_check(si->db_sync); } } @@ -849,7 +836,7 @@ static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, const int sret = db_sync_send(si, src, resp, off); { struct SI_PEER* sp = si_peer_add(si, src); - if (sp && sp->sync_state == 0) { sp->sync_state = 1; sp->sync_start_tb = get_time_tb(); } + if (sp && sret == 0) { sp->sync_state = 1; sp->sync_start_tb = get_time_tb(); } } DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] → INIT_RESP: sent %u bytes to peer=%016llx ret=%d flow=%u,%u", SI_SHRT(si), (unsigned long long)(src >> 16), off, (unsigned long long)src, sret, @@ -899,7 +886,7 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const if (sp) { sp->verified_pos = tp; sp->synced_pos = tp; sp->sync_state = 2; sp->sync_start_tb = 0; } uint8_t sd[13]; sd[0] = DB_MSG_SYNC_DONE; memcpy(sd + 1, &mc, 4); memcpy(sd + 5, &my_ch8, 8); - db_sync_send(si, src, sd, 13); + if (db_sync_send(si, src, sd, 13) != 0) return; si_fire_sync_done(si, src); DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← INIT_RESP: MATCH+no_tail ss→2 → SYNC_DONE tp=%u", SI_SHRT(si), (unsigned long long)(src >> 16), tp); @@ -1001,7 +988,7 @@ static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_ } memcpy(rcp, &rc, 2); - db_sync_send(si, dst, buf, off); + if (db_sync_send(si, dst, buf, off) != 0) return -1; DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] → SEND_DATA: %u recs from=%u want=%u vp=%u", SI_SHRT(si), (unsigned long long)(dst >> 16), rc, from, want_from, vp); return (int)rc; @@ -1097,7 +1084,7 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint64_t mch8 = 0; if (mc > 0) db_chain_hash8_at(si, mc - 1, &mch8); uint8_t sd[13]; sd[0] = DB_MSG_SYNC_DONE; memcpy(sd + 1, &mc, 4); memcpy(sd + 5, &mch8, 8); - db_sync_send(si, src, sd, 13); + if (db_sync_send(si, src, sd, 13) != 0) return; sp->sync_state = 2; sp->sync_start_tb = 0; si_fire_sync_done(si, src); DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← DATA: exchange complete, → SYNC_DONE mc=%u", @@ -1128,8 +1115,8 @@ static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t req[15]; req[0] = DB_MSG_SEND_DATA; memcpy(req + 1, &mc, 4); uint16_t z = 0; memcpy(req + 5, &z, 2); memcpy(req + 7, &vp_send, 4); memcpy(req + 11, &mc, 4); - db_sync_send(si, src, req, 15); - if (sp) sp->sync_state = 1; + if (db_sync_send(si, src, req, 15) != 0) return; + if (sp) { sp->sync_state = 1; sp->sync_start_tb = get_time_tb(); } DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← SYNC_DONE: mc=%u < pc=%u → request tail from=%u", SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, mc); } else if (mc > pc) { @@ -1233,34 +1220,13 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) struct DB_SYNC_INSTANCE* si = db_instance_find(db, group_id); if (!si) { - /* lazy-register: найти таблицу msg_, у которой strtoull(ch_id) == group_id */ - sqlite3_stmt* st = NULL; - if (sqlite3_prepare_v2(db->db, - "SELECT name, substr(name,5) FROM sqlite_master WHERE type='table' AND name LIKE 'msg_%'", - -1, &st, NULL) == SQLITE_OK) { - while (sqlite3_step(st) == SQLITE_ROW) { - const char* tbl = (const char*)sqlite3_column_text(st, 0); - const char* ch_id = (const char*)sqlite3_column_text(st, 1); - if (!ch_id || !ch_id[0]) continue; - uint64_t gid = strtoull(ch_id, NULL, 10); - if (gid == group_id) { - si = db_sync_instance_add(db->inst, tbl, group_id, 1); - if (si) { si_peer_add(si, src); } - DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "lazy-register: tbl=%s ch=%s group=%016llx", tbl, ch_id, (unsigned long long)group_id); - break; - } - } - sqlite3_finalize(st); - } - if (!si) { - DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, - "recv msg type=0x%02x from %016llx group=%016llx — instance NOT FOUND; sending DB_ERR_NOT_FOUND", - type, (unsigned long long)src, (unsigned long long)group_id); - DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ instance not found for group=%016llx", (unsigned long long)group_id); + DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "history table not registered group=%016llx peer=%016llx type=%02x", + (unsigned long long)group_id, (unsigned long long)src, type); + if (type != DB_MSG_ERROR) { uint8_t err[2]; err[0] = DB_MSG_ERROR; err[1] = DB_ERR_NOT_FOUND; db_sync_send_gid(db, src, group_id, err, 2); - queue_dgram_free(entry); queue_entry_free(entry); return; } + queue_dgram_free(entry); queue_entry_free(entry); return; } if (!si->enabled) { uint8_t err[2]; err[0] = DB_MSG_ERROR; err[1] = DB_ERR_DISABLED; @@ -1278,6 +1244,19 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) } } + struct TOPO_GROUP* group = topo_groups_find(db->inst->topo_groups, group_id); + if (type != DB_MSG_ERROR && (si->sync_gated || !topo_group_peer_ready(group, src))) { + DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "history receive deferred group=%016llx peer=%016llx type=%02x members_ready=%d peer_ready=%d", + group_id, src, type, !si->sync_gated, topo_group_peer_ready(group, src)); + uint8_t err[2] = { DB_MSG_ERROR, DB_ERR_NOT_READY }; + db_sync_send_gid(db, src, group_id, err, sizeof(err)); + queue_dgram_free(entry); queue_entry_free(entry); return; + } + struct SI_PEER* progress = si_peer_find(si, src); + if (progress && progress->sync_state == 1 && + ((type == DB_MSG_INIT_RESP && plen >= 14) || (type == DB_MSG_SEND_DATA && plen >= 14))) + progress->sync_start_tb = get_time_tb(); + switch (type) { case DB_MSG_INIT_SYNC: DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ handle INIT_SYNC"); db_handle_init_sync(si, src, payload, plen); break; case DB_MSG_INIT_RESP: DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ handle INIT_RESP"); db_handle_init_resp(si, src, payload, plen); break; @@ -1300,75 +1279,17 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) // Connection callbacks // ============================================================ -// Маршрутизирует статус соединения (UP/DOWN) в соответствующие обработчики -static void db_sync_on_conn_status(struct ETCP_CONN* conn, int status, void* arg) { - switch (status) { - case ETCP_CONN_STATUS_UP: db_sync_on_conn_up(conn, arg); break; - case ETCP_CONN_STATUS_DOWN: - case ETCP_CONN_STATUS_DELETE: db_sync_on_conn_down(conn, arg); break; - default: break; - } -} - -// При поднятии соединения с пиром: для всех активных инстансов добавляет пира и запускает синхронизацию -static void db_sync_on_conn_up(struct ETCP_CONN* conn, void* arg) -{ - (void)arg; - if (!conn || !conn->instance || !conn->instance->db_sync) return; - struct DB_SYNC* db = conn->instance->db_sync; - if (!db->enabled) return; - uint64_t pid = conn->peer_node_id; - if (pid == 0 || pid == db->inst->node_id) return; - DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "peer=%016llx init=%d links=%d", (unsigned long long)pid, conn->initialized, conn->links_up); - - db->last_connected_tb = get_time_tb(); - int synced = 0; - char tbl_list[256] = ""; - for (int i = 0; i < db->instance_count; i++) { - struct DB_SYNC_INSTANCE* si = db->instances[i]; - if (!si->enabled) continue; - struct SI_PEER* p = si_peer_add(si, pid); - if (!p) continue; - if (p->sync_state != 0) continue; - if (si->sync_gated) continue; - if (!conn->initialized || !conn->links_up) continue; - { uint8_t ek[32]; if (db_get_ed25519_pubkey(db, pid, ek) != 0) continue; } - p->sync_state = 1; - db_sync_initiate_sync(si, pid); - synced++; - if (tbl_list[0]) { size_t tl = strlen(tbl_list); snprintf(tbl_list + tl, sizeof(tbl_list) - tl, ",%s", SI_SHRT(si)); } - else snprintf(tbl_list, sizeof(tbl_list), "%s", SI_SHRT(si)); - } - DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync: CONN UP peer=%04llX → %d tables: [%s] — initiating sync for %d", - (unsigned long long)(pid >> 16), db->instance_count, synced > 0 ? tbl_list : "none", synced); -} - -// При разрыве соединения с пиром: сбрасывает sync_state для всех инстансов, если нет других активных пиров — обнуляет last_connected_tb -static void db_sync_on_conn_down(struct ETCP_CONN* conn, void* arg) -{ - (void)arg; - if (!conn || !conn->instance || !conn->instance->db_sync) return; - struct DB_SYNC* db = conn->instance->db_sync; - if (!db->enabled) return; - uint64_t pid = conn->peer_node_id; - if (pid == 0) return; - DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "peer=%016llx", (unsigned long long)pid); - - for (int i = 0; i < db->instance_count; i++) { - struct SI_PEER* p = si_peer_find(db->instances[i], pid); - if (p) p->sync_state = 0; - } - - int any = 0; - for (int i = 0; i < db->instance_count; i++) { - for (int j = 0; j < db->instances[i]->peer_count; j++) { - if (db->instances[i]->peers[j].sync_state >= 1) { any = 1; goto cd_done; } - } - } -cd_done: - if (!any) db->last_connected_tb = 0; - DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync: CONN DOWN peer=%04llX → reset sync state for %d tables", - (unsigned long long)(pid >> 16), db->instance_count); +// История наблюдает сессию своей группы и не владеет дополнительным транспортом. +static void db_sync_on_peer_ready(struct TOPO_GROUP* group, uint64_t pid, int ready, void* arg) { + struct DB_SYNC_INSTANCE* si = arg; + if (!si->enabled) return; + struct SI_PEER* p = ready ? si_peer_add(si, pid) : si_peer_find(si, pid); + if (!p) return; + if (!ready) { + DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "history peer unavailable group=%016llx peer=%016llx state=%u", + group->group_id, pid, p->sync_state); + p->sync_state = 0; p->sync_start_tb = 0; + } else if (p->sync_state == 0) db_sync_initiate_sync(si, pid); } // ============================================================ @@ -1376,25 +1297,37 @@ cd_done: // ============================================================ // Запускает синхронизацию: мастер шлёт INIT_SYNC со своим количеством записей, слейв шлёт REQUEST_SYNC -static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t pid) +static int db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t pid) { + struct DB_SYNC* db = si->db_sync; + struct TOPO_GROUP* group = topo_groups_find(db->inst->topo_groups, si->group_id); + struct ETCP_CONN* conn = instance_find_conn(db->inst, pid); + uint8_t ek[32]; + char ch_id[32]; snprintf(ch_id, sizeof(ch_id), "%llu", (unsigned long long)si->group_id); + if (si->sync_gated || !topo_group_peer_ready(group, pid) || !conn || + !member_sync_validate_pubkey(db->inst, ch_id, conn->crypto_ctx.peer_public_key) || + db_get_ed25519_pubkey(db, pid, ek) != 0) return 0; uint32_t mc = db_count(si); uint64_t my_id = si->db_sync->inst->node_id; struct SI_PEER* p = si_peer_find(si, pid); + if (!p || db_sync_peer_throttled(db, p)) return 0; + p->sync_start_tb = get_time_tb(); p->sync_state = 1; + int ret; DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "pid=%016llx my=%u tbl=%s", (unsigned long long)pid, mc, SI_TBL(si)); if (is_master(my_id, pid)) { - if (p) { p->sync_start_tb = get_time_tb(); p->sync_state = 1; } uint8_t msg[5]; msg[0] = DB_MSG_INIT_SYNC; memcpy(msg + 1, &mc, 4); - db_sync_send(si, pid, msg, 5); + ret = db_sync_send(si, pid, msg, 5); } else { - if (p) { p->sync_state = 1; } uint8_t msg[1] = { DB_MSG_REQUEST_SYNC }; - db_sync_send(si, pid, msg, 1); + ret = db_sync_send(si, pid, msg, 1); } + DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "history start group=%016llx peer=%016llx role=%s count=%u timeout_s=%d ret=%d", + si->group_id, pid, is_master(my_id, pid) ? "master" : "requester", mc, DB_SYNC_SYNC_TIMEOUT, ret); db_sync_resume_peer_check(si->db_sync); + return ret == 0; } // Проверяет целостность цепочки хешей всей таблицы, при несовпадении запускает пересчёт с позиции ошибки @@ -1497,8 +1430,7 @@ void db_sync_set_peer_sound(struct UTUN_INSTANCE* inst, uint64_t group_id, uint6 } } -// Гейт авто-синхронизации: пока gated — инстанс не запускает sync автоматически -// (ждёт member_sync/pubkeys). При снятии гейта запускает sync к пирам в sync_state==0. +// После member_sync запускаем обмен только через READY-сессии этой группы. void db_sync_instance_set_gated(struct DB_SYNC_INSTANCE* si, int gated) { if (!si || !si->enabled) return; @@ -1507,20 +1439,19 @@ void db_sync_instance_set_gated(struct DB_SYNC_INSTANCE* si, int gated) si->sync_gated = (uint8_t)gated; if (!gated) { struct DB_SYNC* db = si->db_sync; + struct TOPO_GROUP* group = topo_groups_find(db->inst->topo_groups, si->group_id); int launched = 0; - for (int i = 0; i < si->peer_count; i++) { - struct SI_PEER* p = &si->peers[i]; - if (p->sync_state != 0) continue; - uint8_t ek[32]; - if (db_get_ed25519_pubkey(db, p->node_id, ek) != 0) continue; - p->sync_state = 1; p->sync_start_tb = get_time_tb(); - db_sync_initiate_sync(si, p->node_id); - launched++; + for (struct ll_entry* e = group && group->senders_list ? group->senders_list->head : NULL; e; e = e->next) { + struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; + if (!topo_group_peer_ready(group, item->node_id)) continue; + struct SI_PEER* p = si_peer_add(si, item->node_id); + if (p && p->sync_state == 0) launched += db_sync_initiate_sync(si, p->node_id); } if (launched == 0) db_sync_resume_peer_check(db); DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s] → un-gated: launched sync for %d peers (mc=%u)", SI_SHRT(si), launched, db_count(si)); } else { + for (int i = 0; i < si->peer_count; i++) { si->peers[i].sync_state = 0; si->peers[i].sync_start_tb = 0; } DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s] → gated: auto-sync suspended (mc=%u)", SI_SHRT(si), db_count(si)); } @@ -1529,9 +1460,8 @@ void db_sync_instance_set_gated(struct DB_SYNC_INSTANCE* si, int gated) // Принудительно перезапускает синхронизацию с указанным пиром (для тестов) void db_sync_reinitiate(struct DB_SYNC_INSTANCE* si, uint64_t node_id) { - struct SI_PEER* p = si_peer_find(si, node_id); + struct SI_PEER* p = si_peer_add(si, node_id); if (p) { - p->sync_state = 1; p->sync_start_tb = get_time_tb(); db_sync_initiate_sync(si, node_id); } } @@ -1556,9 +1486,10 @@ static int db_sync_has_pending_work(struct DB_SYNC* db) if (!db || !db->enabled) return 0; for (int i = 0; i < db->instance_count; i++) { struct DB_SYNC_INSTANCE* si = db->instances[i]; - if (!si->enabled) continue; + if (!si->enabled || si->sync_gated) continue; + struct TOPO_GROUP* group = topo_groups_find(db->inst->topo_groups, si->group_id); for (int j = 0; j < si->peer_count; j++) - if (si->peers[j].sync_state != 2) return 1; + if (si->peers[j].sync_state != 2 && topo_group_peer_ready(group, si->peers[j].node_id)) return 1; } return 0; } @@ -1597,86 +1528,28 @@ static void db_sync_peer_check_cb(void* arg) #endif DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "instances=%d", db->instance_count); - int total_synced = 0, total_skipped = 0, any_peers = 0; - char launched_list[256] = ""; - { - for (int i = 0; i < db->instance_count; i++) { - struct DB_SYNC_INSTANCE* si = db->instances[i]; - struct TOPO_GROUP* g = topo_groups_find(db->inst->topo_groups, si->group_id); - if (!g) continue; - if (!si->enabled) continue; - if (si->sync_gated) continue; - - int peers_found = 0; - struct ll_entry* e = g->senders_list->head; - while (e) { - struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; - if (item->conn && item->conn->peer_node_id != 0 && item->conn->links_up && item->conn->initialized) { - uint64_t pid = item->conn->peer_node_id; - if (pid != db->inst->node_id) { si_peer_add(si, pid); peers_found++; any_peers = 1; } - } - e = e->next; - } - - struct SI_PEER* best = NULL; - uint32_t min_pos = UINT32_MAX; - for (int j = 0; j < si->peer_count; j++) { - if (si->peers[j].sync_state == 0 && si->peers[j].synced_pos < min_pos) { - if (db_sync_peer_throttled(db, &si->peers[j])) continue; /* спящий пир в «тихом» чате — пропускаем */ - min_pos = si->peers[j].synced_pos; - best = &si->peers[j]; - } - } - if (best) { - uint8_t ek[32]; - if (db_get_ed25519_pubkey(db, best->node_id, ek) == 0) { - best->sync_state = 1; - db_sync_initiate_sync(si, best->node_id); total_synced++; - { size_t tl = strlen(launched_list); - snprintf(launched_list + tl, sizeof(launched_list) - tl, - "%s%s:%04llX[sp=%u]", tl ? "," : "", - SI_SHRT(si), (unsigned long long)(best->node_id >> 16), min_pos); } - } else { - DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "peer_check skip peer=%016llx — no Ed25519 pubkey yet", (unsigned long long)best->node_id); - best = NULL; total_skipped++; - } - } - 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_INFO(DEBUG_CATEGORY_CHAT_SYNC, - "sync [%s:%04llX] TIMEOUT: no response for %llu ms, sync_state 1→0 (synced_pos was %u)", - SI_SHRT(si), (unsigned long long)(p->node_id >> 16), - (unsigned long long)((now - p->sync_start_tb) / 10), p->synced_pos); - si_dump_table(si, "STUCK_SYNC_TIMEOUT"); - p->sync_state = 0; - static int dump_ctr = 0; - if (dump_ctr++ == 0 || (dump_ctr % 10) == 0) - etcp_dump_all(db->inst); - } + uint64_t now = get_time_tb(); + for (int i = 0; i < db->instance_count; i++) { + struct DB_SYNC_INSTANCE* si = db->instances[i]; + if (!si->enabled || si->sync_gated) continue; + struct TOPO_GROUP* group = topo_groups_find(db->inst->topo_groups, si->group_id); + for (int j = 0; j < si->peer_count; j++) { + struct SI_PEER* p = &si->peers[j]; + if (!topo_group_peer_ready(group, p->node_id)) continue; + if (p->sync_state == 1 && now - p->sync_start_tb >= DB_SYNC_SYNC_TIMEOUT * 10000ULL) { + DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "history timeout group=%016llx peer=%016llx idle_ms=%llu pos=%u retry=%u", + si->group_id, p->node_id, (unsigned long long)((now - p->sync_start_tb) / 10), + p->synced_pos, ++p->retry_count); + p->sync_state = 0; p->sync_start_tb = 0; } + if (p->sync_state == 0) db_sync_initiate_sync(si, p->node_id); } } - if (total_synced > 0) { - DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync: PEER CHECK → %d tables, launched sync for %d: %s", - db->instance_count, total_synced, launched_list); - } - if (db_sync_has_pending_work(db)) db_sync_resume_peer_check(db); else - DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "all peers synced — peer_check stopped"); + DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "history peer_check stopped: no pending READY peers"); } // Периодическая TTL-очистка: удаляет неподтверждённые записи (flags&1==0) старше db_sync_ttl от локального узла @@ -1728,7 +1601,6 @@ int db_sync_init(struct UTUN_INSTANCE* inst) struct DB_SYNC* db = u_calloc(1, sizeof(struct DB_SYNC)); if (!db) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "u_calloc failed"); return -1; } db->inst = inst; - db->last_connected_tb = 0; if (!inst->config->global.db_sync_enabled) { db->enabled = 0; inst->db_sync = db; @@ -1760,11 +1632,6 @@ int db_sync_enable(struct UTUN_INSTANCE* inst) { db->enabled = 1; etcp_bind(inst, ETCP_RT_ID_DB_SYNC, db_sync_recv_cb); - etcp_add_conn_status_cbk(inst, db_sync_on_conn_status, NULL); - - if (db->enabled) - db->peer_check_timer = uasync_set_timeout(inst->ua, DB_SYNC_PEER_CHECK_INTERVAL * 10000u, - db, db_sync_peer_check_cb, "db_sync_peer"); DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "db_sync: initialized (enabled=%d)", db->enabled); return 0; @@ -1787,13 +1654,14 @@ void db_sync_destroy(struct UTUN_INSTANCE* inst) for (int i = 0; i < db->instance_count; i++) { struct DB_SYNC_INSTANCE* si = db->instances[i]; + struct TOPO_GROUP* group = topo_groups_find(inst->topo_groups, si->group_id); + topo_group_remove_peer_ready_cbk(group, db_sync_on_peer_ready, si); if (si->ttl_timer) { uasync_cancel_timeout(inst->ua, si->ttl_timer); si->ttl_timer = NULL; } if (si->recalc.timer) { uasync_cancel_timeout(inst->ua, si->recalc.timer); si->recalc.timer = NULL; } if (si->peers) u_free(si->peers); u_free(si); } - etcp_remove_conn_status_cbk(inst, db_sync_on_conn_status, NULL); { struct db_sync_done_cbk_entry* cb = db->done_cbks; while (cb) { struct db_sync_done_cbk_entry* n = cb->next; u_free(cb); cb = n; } } @@ -1803,13 +1671,15 @@ void db_sync_destroy(struct UTUN_INSTANCE* inst) DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "db_sync: destroyed"); } -// Создаёт SQLite-таблицу для синхронизации, регистрирует пиров, запускает первичную синхронизацию и TTL-таймер -struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t group_id, int auto_sync) +// Создаёт локальное состояние таблицы, наблюдает группу. До member_sync сетевого обмена нет. +struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t group_id) { if (!inst || !inst->db_sync || !table_name || !table_name[0]) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "instance_add invalid args"); return NULL; } DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "table=%s group=%016llx", table_name, (unsigned long long)group_id); struct DB_SYNC* db = inst->db_sync; if (!db->enabled || !db->db) return NULL; + struct TOPO_GROUP* group = topo_groups_find(inst->topo_groups, group_id); + if (!group) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "history registration requires group=%016llx", group_id); return NULL; } struct DB_SYNC_INSTANCE* si = db_instance_find(db, group_id); if (si) { DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "instance already exists group=%016llx", (unsigned long long)group_id); return si; } @@ -1863,26 +1733,14 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const si->ttl_timer = uasync_set_timeout(inst->ua, DB_SYNC_TTL_INTERVAL * 10000u, si, db_sync_instance_ttl_cb, "db_sync_ttl"); - int peers_found = 0, peers_synced = 0; - { - struct ll_entry* entry = inst->connections->head; - while (entry) { - struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; - uint64_t pid = ce->conn->peer_node_id; - if (pid != 0 && pid != inst->node_id && ce->conn->links_up > 0) { - peers_found++; - struct SI_PEER* p = si_peer_add(si, pid); - if (!p) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "si_peer_add failed for %N in %s", (unsigned long long)pid, SI_TBL(si)); continue; } - if (p->sync_state == 0 && auto_sync && !si->sync_gated) { p->sync_state = 1; db_sync_initiate_sync(si, pid); peers_synced++; } - } - entry = entry->next; - } + if (topo_group_add_peer_ready_cbk(group, db_sync_on_peer_ready, si) != 0) { + db_sync_instance_remove(si); return NULL; } DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, - "instance added table=%s tbl=%s group=%016llx next_id=%llu peers_found=%d peers_synced=%d mc=%u", + "history registered table=%s tbl=%s group=%016llx next_id=%llu gated=1 mc=%u", table_name, si->table_name, (unsigned long long)group_id, - (unsigned long long)si->next_id, peers_found, peers_synced, db_count(si)); + (unsigned long long)si->next_id, db_count(si)); return si; } @@ -1894,6 +1752,8 @@ void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si) struct DB_SYNC* db = si->db_sync; struct UTUN_INSTANCE* inst = db->inst; + struct TOPO_GROUP* group = topo_groups_find(inst->topo_groups, si->group_id); + topo_group_remove_peer_ready_cbk(group, db_sync_on_peer_ready, si); si->enabled = 0; if (si->ttl_timer) { uasync_cancel_timeout(inst->ua, si->ttl_timer); si->ttl_timer = NULL; } if (si->recalc.timer) { uasync_cancel_timeout(inst->ua, si->recalc.timer); si->recalc.timer = NULL; }