Browse Source

Gate history on member synchronization and group readiness

master
evgeny 3 days ago
parent
commit
eaac611465
  1. 378
      src/chat/db_sync.c

378
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_<ch_id>, у которой 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; }

Loading…
Cancel
Save