diff --git a/src/chat/db_sync.c b/src/chat/db_sync.c index be09cb7d..447936a0 100644 --- a/src/chat/db_sync.c +++ b/src/chat/db_sync.c @@ -691,7 +691,7 @@ static int db_sync_send_gid(struct DB_SYNC* db, uint64_t node_id, uint64_t group (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", + DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "sync: send → 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; } @@ -782,8 +782,8 @@ static void db_handle_error(struct DB_SYNC_INSTANCE* si, uint64_t src, const uin SI_SHRT(si), (unsigned long long)(src >> 16), p[0], what); struct SI_PEER* sp = si_peer_find(si, src); if (sp) { - 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); + DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← ERROR: state=%u→0, synced_pos=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), sp->sync_state, sp->synced_pos); sp->sync_state = 0; sp->sync_start_tb = 0; db_sync_resume_peer_check(si->db_sync); @@ -1304,13 +1304,26 @@ static int db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t pid) 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; + int ready = topo_group_peer_ready(group, pid); + if (si->sync_gated || !ready) { + DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "history waiting group=%016llx peer=%016llx members_ready=%d peer_ready=%d", + (unsigned long long)si->group_id, (unsigned long long)pid, !si->sync_gated, ready); + return 0; + } + if (!conn || !member_sync_validate_pubkey(db->inst, ch_id, conn->crypto_ctx.peer_public_key) || + db_get_ed25519_pubkey(db, pid, ek) != 0) { + DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "history READY peer lacks transport, membership or key group=%016llx peer=%016llx", + (unsigned long long)si->group_id, (unsigned long long)pid); + 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; + if (!p) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "history start without peer group=%016llx", si->group_id); return 0; } + if (db_sync_peer_throttled(db, p)) { + DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "history sleeping peer deferred group=%016llx peer=%016llx", si->group_id, pid); + db_sync_resume_peer_check(db); 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)); @@ -1500,6 +1513,7 @@ static void db_sync_resume_peer_check(struct DB_SYNC* db) if (!db || !db->enabled || db->peer_check_timer) return; db->peer_check_timer = uasync_set_timeout(db->inst->ua, DB_SYNC_PEER_CHECK_INTERVAL * 10000u, db, db_sync_peer_check_cb, "db_sync_peer"); + if (!db->peer_check_timer) DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "history retry timer allocation failed"); } // Периодическая проверка: запускает синхронизацию для пиров в состоянии sync_state==0, @@ -1674,10 +1688,12 @@ void db_sync_destroy(struct UTUN_INSTANCE* inst) // Создаёт локальное состояние таблицы, наблюдает группу. До 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; } + if (!inst || !inst->db_sync || !si_name_valid(table_name)) { + DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "history registration: invalid instance or table name"); 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; + if (!db->enabled || !db->db) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "history registration: module disabled"); 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; } @@ -1685,7 +1701,7 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const if (si) { DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "instance already exists group=%016llx", (unsigned long long)group_id); return si; } si = db_instance_alloc(db); - if (!si) return NULL; + if (!si) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "history instance allocation failed"); return NULL; } si->group_id = group_id; snprintf(si->table_name, sizeof(si->table_name), "%s", table_name); @@ -1707,7 +1723,10 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const " PRIMARY KEY (timestamp, author_signature))", si->table_name); int rc = sqlite3_exec(db->db, sql, NULL, NULL, NULL); - if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "CREATE TABLE %s: %s", si->table_name, sqlite3_errmsg(db->db)); si->enabled = 0; return si; } + if (rc != SQLITE_OK) { + DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "CREATE TABLE %s: %s", si->table_name, sqlite3_errmsg(db->db)); + db_sync_instance_remove(si); return NULL; + } snprintf(sql, sizeof(sql), "CREATE INDEX IF NOT EXISTS \"idx_%s_ttl\"" " ON \"%s\" (node_id, timestamp)",