|
|
|
|
@ -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)", |
|
|
|
|
|