You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

1949 lines
92 KiB

// db_sync.c — Distributed content-addressed table with SQLite + peer sync via etcp_api (direct P2P)
#include "db_sync.h"
#include "../transport_layer/etcp_api.h"
#include "../transport_layer/etcp.h"
#include "../utun_instance.h"
#include "../routing_layer/topo_group.h"
#include "../routing_layer/topo_node_sqlite.h"
#include "../transport_layer/etcp_dump.h"
#include "../transport_layer/secure_channel.h"
#include "../../lib/debug_config.h"
#include "../../lib/mem.h"
#include "../../lib/u_async.h"
#ifdef UTUN_HAVE_STANDBY
#include "standby.h"
#endif
#include <openssl/sha.h>
#include "../../lib/sqlite3.h"
#include "../../lib/platform_compat.h"
#include <string.h>
#include <openssl/evp.h>
#include <unistd.h>
// ---- Forward declarations ----
struct DB_SYNC;
struct DB_SYNC_INSTANCE;
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_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 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);
// ============================================================
// Structures
// ============================================================
struct SI_PEER {
uint64_t node_id;
uint32_t synced_pos;
uint32_t verified_pos;
uint32_t last_peer_count;
uint8_t sync_state;
uint64_t sync_start_tb;
uint32_t retry_count;
};
struct DB_SYNC_INSTANCE {
struct DB_SYNC* db_sync;
uint64_t group_id;
char table_name[64];
uint64_t next_id;
uint64_t last_timestamp_ms;
uint8_t enabled;
void* ttl_timer;
struct {
uint32_t pos;
void* timer;
} recalc;
struct SI_PEER* peers;
int peer_count, peer_capacity;
db_sync_insert_cb on_insert;
void* on_insert_arg;
uint8_t sync_gated; /* 1 = не запускать авто-синхронизацию (ждём member_sync/pubkeys) */
};
struct db_sync_done_cbk_entry {
db_sync_done_fn fn;
void* arg;
struct db_sync_done_cbk_entry* next;
};
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;
void* peer_check_wait; /* standby_wait handle (Android) */
uint8_t enabled;
struct db_sync_done_cbk_entry* done_cbks;
};
// ============================================================
// Done callback chain (global, like etcp_status_cbk_entry)
// ============================================================
static void db_sync_done_cbk_add_chain(struct db_sync_done_cbk_entry** head, db_sync_done_fn fn, void* arg) {
if (!head || !fn) return;
struct db_sync_done_cbk_entry* e = u_malloc(sizeof(struct db_sync_done_cbk_entry));
if (!e) return;
e->fn = fn; e->arg = arg; e->next = *head;
*head = e;
}
static void db_sync_done_cbk_remove_chain(struct db_sync_done_cbk_entry** head, db_sync_done_fn fn, void* arg) {
if (!head || !fn) return;
struct db_sync_done_cbk_entry** p = head;
while (*p) {
if ((*p)->fn == fn && (*p)->arg == arg) { struct db_sync_done_cbk_entry* rm = *p; *p = rm->next; u_free(rm); return; }
p = &(*p)->next;
}
}
static void si_fire_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t peer_node_id) {
struct DB_SYNC* db = si->db_sync;
struct db_sync_done_cbk_entry* cb = db->done_cbks;
while (cb) { struct db_sync_done_cbk_entry* n = cb->next; cb->fn(si, peer_node_id, cb->arg); cb = n; }
}
void db_sync_add_done_cbk(struct UTUN_INSTANCE* inst, db_sync_done_fn fn, void* arg) {
if (inst && inst->db_sync) db_sync_done_cbk_add_chain(&inst->db_sync->done_cbks, fn, arg);
}
void db_sync_remove_done_cbk(struct UTUN_INSTANCE* inst, db_sync_done_fn fn, void* arg) {
if (inst && inst->db_sync) db_sync_done_cbk_remove_chain(&inst->db_sync->done_cbks, fn, arg);
}
// ============================================================
// SQL helpers
// ============================================================
#define SI_DB(si) ((si)->db_sync->db)
#define SI_TBL(si) ((si)->table_name)
#define SI_SHRT(si) ({ \
static char _tbuf[28]; \
snprintf(_tbuf, sizeof(_tbuf), "%s.%04llX", SI_TBL(si), (unsigned long long)((si)->group_id & 0xFFFF)); \
_tbuf; \
})
// Форматирует SQL-запрос, подставляя имя таблицы из инстанса (через %s)
static void si_sql(char* buf, size_t sz, struct DB_SYNC_INSTANCE* si, const char* fmt)
{
snprintf(buf, sz, fmt, SI_TBL(si));
}
// Готовит SQLite statement с подстановкой имени таблицы — обёртка над si_sql + sqlite3_prepare_v2
static int si_prep(struct DB_SYNC_INSTANCE* si, sqlite3_stmt** stmt, const char* fmt)
{
char sql[512];
si_sql(sql, sizeof(sql), si, fmt);
return sqlite3_prepare_v2(SI_DB(si), sql, -1, stmt, NULL);
}
// ============================================================
// SQLite open/close
// ============================================================
// Открывает SQLite БД, включает WAL, mmap и прочие оптимизации для высокой производительности
static int db_sqlite_open(struct DB_SYNC* db, const char* path)
{
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "path=%s", path);
int rc = sqlite3_open(path, &db->db);
if (rc != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "db_sync: sqlite3_open(%s): %s", path, sqlite3_errmsg(db->db));
sqlite3_close(db->db);
db->db = NULL;
return -1;
}
char* err = NULL;
rc = sqlite3_exec(db->db, "PRAGMA journal_mode=WAL", NULL, NULL, &err);
if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "db_sync: WAL pragma: %s", err); sqlite3_free(err); }
sqlite3_exec(db->db, "PRAGMA synchronous=NORMAL", NULL, NULL, NULL);
sqlite3_exec(db->db, "PRAGMA wal_autocheckpoint=10000", NULL, NULL, NULL);
sqlite3_exec(db->db, "PRAGMA cache_size=-32768", NULL, NULL, NULL);
sqlite3_exec(db->db, "PRAGMA mmap_size=134217728", NULL, NULL, NULL);
sqlite3_exec(db->db, "PRAGMA temp_store=MEMORY", NULL, NULL, NULL);
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "db_sync: SQLite opened at %s", path);
return 0;
}
// Закрывает SQLite БД (только свою, shared не трогает)
static void db_sqlite_close(struct DB_SYNC* db)
{
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "");
if (!db->db || db->shared_db) return;
sqlite3_close(db->db);
db->db = NULL;
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "db_sync: SQLite closed");
}
// ============================================================
// SHA256 helpers
// ============================================================
// Вычисляет SHA256 хеш данных (обёртка над OpenSSL SHA256)
static void db_sha256(const uint8_t* data, size_t len, uint8_t hash[32])
{
SHA256(data, len, hash);
}
// Возвращает первые 8 байт SHA256 как uint64 — короткий идентификатор для сравнения
static uint64_t db_hash64(const uint8_t* data, size_t len)
{
uint8_t hash[32]; db_sha256(data, len, hash);
uint64_t h; memcpy(&h, hash, 8);
return h;
}
// Вычисляет цепной хеш: SHA256(prev_chain_hash[32] || id[8] || timestamp[8] || author[8] || sig[64])
static void db_chain_hash_compute(const uint8_t prev_chain_hash[32],
uint64_t id, uint64_t timestamp,
uint64_t author, const uint8_t author_signature[DB_SIG_SIZE],
uint8_t out[32])
{
uint8_t buf[32 + 8 + 8 + 8 + DB_SIG_SIZE]; // 32 + 8 + 8 + 8 + 64 = 120
memcpy(buf, prev_chain_hash, 32);
memcpy(buf + 32, &id, 8);
memcpy(buf + 40, &timestamp, 8);
memcpy(buf + 48, &author, 8);
memcpy(buf + 56, author_signature, DB_SIG_SIZE);
db_sha256(buf, sizeof(buf), out);
}
// ============================================================
// Instance management
// ============================================================
// Ищет активный инстанс синхронизации по group_id канала
static struct DB_SYNC_INSTANCE* db_instance_find(struct DB_SYNC* db, uint64_t group_id)
{
for (int i = 0; i < db->instance_count; i++) {
if (db->instances[i]->group_id == group_id && db->instances[i]->enabled) return db->instances[i];
}
return NULL;
}
struct DB_SYNC_INSTANCE* db_sync_instance_find(struct UTUN_INSTANCE* inst, uint64_t group_id)
{
if (!inst || !inst->db_sync) return NULL;
return db_instance_find(inst->db_sync, group_id);
}
// Выделяет новый инстанс в куче (стабильный указатель) и регистрирует его в массиве
static struct DB_SYNC_INSTANCE* db_instance_alloc(struct DB_SYNC* db)
{
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "count=%d cap=%d", db->instance_count, db->instance_capacity);
if (db->instance_count >= db->instance_capacity) {
int nc = db->instance_capacity ? db->instance_capacity * 2 : 4;
struct DB_SYNC_INSTANCE** np = u_realloc(db->instances, nc * sizeof(*db->instances));
if (!np) return NULL;
db->instances = np;
db->instance_capacity = nc;
}
struct DB_SYNC_INSTANCE* si = u_calloc(1, sizeof(*si));
if (!si) return NULL;
si->db_sync = db;
si->enabled = 1;
db->instances[db->instance_count++] = si;
return si;
}
// Проверяет допустимость имени таблицы: только A-Za-z0-9_, до 48 символов
static int si_name_valid(const char* name)
{
if (!name || !name[0] || strlen(name) > 48) return 0;
for (const char* p = name; *p; p++) {
if (!((*p >= 'a' && *p <= 'z') || (*p >= 'A' && *p <= 'Z')
|| (*p >= '0' && *p <= '9') || *p == '_')) return 0;
}
return 1;
}
// Вычисляет 64-битный хеш инстанса: db_hash64(name[256] || id_be[8])
static uint64_t db_instance_hash_compute(const char* name, uint64_t id)
{
size_t nl = strlen(name);
uint8_t buf[264]; if (nl > 256) nl = 256;
memcpy(buf, name, nl);
uint64_t id_be = htobe64(id);
memcpy(buf + nl, &id_be, 8);
return db_hash64(buf, nl + 8);
}
// ============================================================
// Peer management
// ============================================================
// Ищет пира в инстансе по node_id
static struct SI_PEER* si_peer_find(struct DB_SYNC_INSTANCE* si, uint64_t node_id)
{
for (int i = 0; i < si->peer_count; i++) {
if (si->peers[i].node_id == node_id) return &si->peers[i];
}
return NULL;
}
// Добавляет пира в инстанс (если уже существует — возвращает существующего)
static struct SI_PEER* si_peer_add(struct DB_SYNC_INSTANCE* si, uint64_t node_id)
{
struct SI_PEER* p = si_peer_find(si, node_id);
if (p) return p;
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;
si->peers = np;
si->peer_capacity = nc;
}
p = &si->peers[si->peer_count++];
memset(p, 0, sizeof(*p));
p->node_id = node_id;
return p;
}
// ============================================================
// Data access
// ============================================================
// Возвращает количество записей в таблице инстанса
static uint32_t db_count(struct DB_SYNC_INSTANCE* si)
{
sqlite3_stmt* stmt;
if (si_prep(si, &stmt, "SELECT COUNT(*) FROM \"%s\"") != SQLITE_OK) return 0;
uint32_t c = (sqlite3_step(stmt) == SQLITE_ROW) ? (uint32_t)sqlite3_column_int64(stmt, 0) : 0;
sqlite3_finalize(stmt);
return c;
}
// Читает цепной хеш (32 байта) записи на позиции pos в отсортированной таблице
static int db_chain_hash_at(struct DB_SYNC_INSTANCE* si, uint32_t pos, uint8_t out[32])
{
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT chain_hash FROM \"%s\""
" ORDER BY timestamp, author_signature"
" LIMIT 1 OFFSET ?") != SQLITE_OK)
{
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "chain_hash_at prepare: %s", sqlite3_errmsg(SI_DB(si)));
return -1;
}
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)pos);
if (sqlite3_step(stmt) != SQLITE_ROW) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "chain_hash_at no row pos=%u tbl=%s", pos, SI_TBL(si)); sqlite3_finalize(stmt); return -1; }
const void* b = sqlite3_column_blob(stmt, 0);
if (!b || sqlite3_column_bytes(stmt, 0) < 32) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "chain_hash_at blob too small pos=%u tbl=%s", pos, SI_TBL(si)); sqlite3_finalize(stmt); return -1; }
memcpy(out, b, 32);
sqlite3_finalize(stmt);
return 0;
}
// Находит позицию записи (0-based) в отсортированной таблице по (timestamp, author_signature)
static uint32_t si_find_pos(struct DB_SYNC_INSTANCE* si, uint64_t ts, const uint8_t* author_sig)
{
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT COUNT(*) FROM \"%s\""
" WHERE timestamp<?1 OR (timestamp=?1 AND author_signature<?2)") != SQLITE_OK)
return 0;
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_blob(stmt, 2, author_sig, DB_SIG_SIZE, SQLITE_STATIC);
uint32_t pos = (sqlite3_step(stmt) == SQLITE_ROW) ? (uint32_t)sqlite3_column_int64(stmt, 0) : 0;
sqlite3_finalize(stmt);
return pos;
}
// Возвращает цепной хеш записи, предшествующей заданной по (ts, author_sig)
static int db_prev_chain_hash(struct DB_SYNC_INSTANCE* si, uint64_t ts, const uint8_t* author_sig, uint8_t out[32])
{
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT chain_hash FROM \"%s\""
" WHERE timestamp<?1 OR (timestamp=?1 AND author_signature<?2)"
" ORDER BY timestamp DESC, author_signature DESC"
" LIMIT 1") != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "db_prev_chain_hash prep: %s", sqlite3_errmsg(SI_DB(si)));
return -1; }
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_blob(stmt, 2, author_sig, DB_SIG_SIZE, SQLITE_STATIC);
if (sqlite3_step(stmt) == SQLITE_ROW) {
const void* b = sqlite3_column_blob(stmt, 0);
if (b && sqlite3_column_bytes(stmt, 0) >= 32) memcpy(out, b, 32);
else memset(out, 0, 32);
} else {
memset(out, 0, 32);
}
sqlite3_finalize(stmt);
return 0;
}
// ── Chain hash fragment (first 8 bytes) for sync protocol comparisons ──
// Возвращает первые 8 байт цепного хеша на позиции pos — для быстрого сравнения в протоколе sync
static int db_chain_hash8_at(struct DB_SYNC_INSTANCE* si, uint32_t pos, uint64_t* h8)
{
uint8_t ch[32];
if (db_chain_hash_at(si, pos, ch) != 0) return -1;
memcpy(h8, ch, 8);
return 0;
}
// ── Recalc scheduler (batch of 50, async via timer) ──
static void db_recalc_tick(void* arg);
// Планирует асинхронный пересчёт цепных хешей начиная с позиции from_pos (батчами по 50)
static void db_sync_schedule_recalc(struct DB_SYNC_INSTANCE* si, uint32_t from_pos)
{
if (!si->recalc.timer) {
si->recalc.pos = from_pos;
si->recalc.timer = uasync_set_timeout(si->db_sync->inst->ua, 0, si, db_recalc_tick, "db_recalc");
return;
}
if (from_pos < si->recalc.pos)
si->recalc.pos = from_pos;
}
// Обрабатывает один батч (до 50 записей) пересчёта цепных хешей, вызывается по таймеру
static void db_recalc_tick(void* arg)
{
struct DB_SYNC_INSTANCE* si = (struct DB_SYNC_INSTANCE*)arg;
sqlite3* db = SI_DB(si);
uint32_t total = db_count(si);
if (si->recalc.pos >= total) { si->recalc.timer = NULL; return; }
uint8_t prev_ch[32];
if (si->recalc.pos > 0) {
if (db_chain_hash_at(si, si->recalc.pos - 1, prev_ch) != 0)
memset(prev_ch, 0, 32);
} else {
memset(prev_ch, 0, 32);
}
uint32_t batch = (total - si->recalc.pos > 50) ? 50 : total - si->recalc.pos;
int rc = sqlite3_exec(db, "BEGIN IMMEDIATE", NULL, NULL, NULL);
if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "recalc BEGIN: %s", sqlite3_errmsg(db)); return; }
sqlite3_stmt* sel, *upd;
if (si_prep(si, &sel,
"SELECT id,timestamp,node_id,author_signature FROM \"%s\""
" ORDER BY timestamp, author_signature"
" LIMIT -1 OFFSET ?") != SQLITE_OK)
{ DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "recalc SELECT prep: %s", sqlite3_errmsg(db)); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return; }
sqlite3_bind_int64(sel, 1, (sqlite3_int64)si->recalc.pos);
if (si_prep(si, &upd,
"UPDATE \"%s\" SET chain_hash=?"
" WHERE timestamp=? AND author_signature=?") != SQLITE_OK)
{ DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "recalc UPDATE prep: %s", sqlite3_errmsg(db)); sqlite3_finalize(sel); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return; }
uint8_t ch[32]; uint32_t processed = 0;
while (sqlite3_step(sel) == SQLITE_ROW && processed < batch) {
sqlite3_int64 n_id = sqlite3_column_int64(sel, 0);
sqlite3_int64 n_ts = sqlite3_column_int64(sel, 1);
sqlite3_int64 n_auth = sqlite3_column_int64(sel, 2);
const void* sig_blob = sqlite3_column_blob(sel, 3);
uint8_t sig[DB_SIG_SIZE];
if (sig_blob) memcpy(sig, sig_blob, DB_SIG_SIZE); else memset(sig, 0, DB_SIG_SIZE);
db_chain_hash_compute(prev_ch, (uint64_t)n_id, (uint64_t)n_ts, (uint64_t)n_auth, sig, ch);
sqlite3_reset(upd);
sqlite3_bind_blob(upd, 1, ch, 32, SQLITE_STATIC);
sqlite3_bind_int64(upd, 2, n_ts);
sqlite3_bind_blob(upd, 3, sig, DB_SIG_SIZE, SQLITE_STATIC);
sqlite3_step(upd);
memcpy(prev_ch, ch, 32);
processed++;
}
sqlite3_finalize(upd);
sqlite3_finalize(sel);
rc = sqlite3_exec(db, "COMMIT", NULL, NULL, NULL);
if (rc != SQLITE_OK) DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "recalc COMMIT: %s", sqlite3_errmsg(db));
si->recalc.pos += processed;
DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "recalc tick tbl=%s pos=%u batch=%u/%u total=%u",
SI_TBL(si), si->recalc.pos - processed, processed, batch, total);
if (si->recalc.pos >= total) {
si->recalc.timer = NULL;
DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "recalc complete tbl=%s pos=%u(total)", SI_TBL(si), si->recalc.pos);
} else {
si->recalc.timer = uasync_set_timeout(si->db_sync->inst->ua, 100, si, db_recalc_tick, "db_recalc");
}
}
// Синхронно завершает все отложенные пересчёты хешей — вызывается перед операциями, требующими актуальных хешей
static void db_sync_flush_recalc(struct DB_SYNC_INSTANCE* si)
{
struct UASYNC* ua = si->db_sync->inst->ua;
while (si->recalc.timer) {
uasync_cancel_timeout(ua, si->recalc.timer);
si->recalc.timer = NULL;
db_recalc_tick(si);
}
}
// ── Ed25519 verification ──
// Получает Ed25519 публичный ключ узла: сначала из локального, затем из topo_node БД, затем из ETCP-соединения
static int db_get_ed25519_pubkey(struct DB_SYNC* db, uint64_t node_id, uint8_t out[32])
{
if (node_id == db->inst->node_id) { memcpy(out, db->inst->my_ed25519_pubkey, 32); return 0; }
/* из in-memory registry (BGP/topo, self-sig верифицирован) — первичный источник */
if (db->inst->topo_groups) {
struct TOPO_NODE* ni = topo_node_registry_find(db->inst->topo_groups, node_id);
if (ni) {
uint64_t chk; memcpy(&chk, ni->ed25519_public_key, 8);
if (chk != 0) { memcpy(out, ni->ed25519_public_key, 32); return 0; }
}
}
if (db->inst->topo_sqlite_db) {
if (topo_node_sqlite_get_ed25519_pubkey(db->inst->topo_sqlite_db, node_id, out) == 0) return 0;
}
struct ETCP_CONN* conn = instance_find_conn(db->inst, node_id);
if (conn) {
memcpy(out, conn->peer_ed25519_pubkey, 32);
uint64_t chk; memcpy(&chk, out, 8);
if (chk != 0) return 0;
}
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,
const char* json, size_t jlen,
const uint8_t* sig)
{
uint8_t pubkey[32];
if (db_get_ed25519_pubkey(si->db_sync, author_node_id, pubkey) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC,
"cannot get Ed25519 pubkey for author=%016llx — discarding",
(unsigned long long)author_node_id);
return -1;
}
uint8_t msg[8192]; size_t off = 0;
memcpy(msg + off, &ts, 8); off += 8;
if (jlen > 0) { memcpy(msg + off, json, jlen); off += jlen; }
if (sc_ed25519_verify(pubkey, msg, off, sig) != SC_OK) {
DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC,
"author_sig VERIFY FAIL author=%016llx ed_pubkey=%016llx... msg_len=%zu sig=%016llx... — discarding as forgery",
(unsigned long long)author_node_id, *(uint64_t*)pubkey, off, *(uint64_t*)sig);
return -1;
}
return 0;
}
// Вставляет запись в таблицу: проверяет подпись, дубликаты, вычисляет chain_hash, возвращает 0/1/-1/-2
static int db_record_insert(struct DB_SYNC_INSTANCE* si,
uint64_t id, uint64_t ts, uint64_t author_node_id,
const char* json, size_t jlen,
const uint8_t* author_sig,
const char* local_attrs)
{
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "id=%llu ts=%llu author=%N len=%zu", (unsigned long long)id, (unsigned long long)ts, (unsigned long long)author_node_id, jlen);
sqlite3* db = SI_DB(si);
sqlite3_stmt* stmt;
int rc;
// Verify author signature (mandatory)
if (db_verify_author_sig(si, ts, author_node_id, json, jlen, author_sig) != 0) {
sqlite3_stmt* del;
if (si_prep(si, &del, "DELETE FROM \"%s\" WHERE timestamp=? AND author_signature=?") == SQLITE_OK) {
sqlite3_bind_int64(del, 1, (sqlite3_int64)ts);
sqlite3_bind_blob(del, 2, author_sig, DB_SIG_SIZE, SQLITE_STATIC);
sqlite3_step(del);
int chg = sqlite3_changes(SI_DB(si));
sqlite3_finalize(del);
if (chg > 0) db_sync_schedule_recalc(si, si_find_pos(si, ts, author_sig));
}
return -2;
}
rc = sqlite3_exec(db, "BEGIN IMMEDIATE", NULL, NULL, NULL);
if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "insert BEGIN: %s", sqlite3_errmsg(db)); return -1; }
// Check if record already exists by (timestamp, author_signature)
if (si_prep(si, &stmt, "SELECT 1 FROM \"%s\" WHERE timestamp=? AND author_signature=?") != SQLITE_OK)
{ DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "db_record_insert SELECT prep: %s", sqlite3_errmsg(db)); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return -1; }
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_blob(stmt, 2, author_sig, DB_SIG_SIZE, SQLITE_STATIC);
int exists = (sqlite3_step(stmt) == SQLITE_ROW);
sqlite3_finalize(stmt);
if (exists) { DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ already exists, skip"); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return 1; }
// Compute chain hash
uint8_t prev_ch[32];
db_prev_chain_hash(si, ts, author_sig, prev_ch);
uint8_t ch[32];
db_chain_hash_compute(prev_ch, id, ts, author_node_id, author_sig, ch);
// Insert new record
if (si_prep(si, &stmt,
"INSERT INTO \"%s\""
" (timestamp,node_id,id,chain_hash,flags,data,"
" author_signature,local_attrs,delivered_peers,delivery_chain)"
" VALUES (?,?,?,?,0,?,?,?,0,'')") != SQLITE_OK)
{ DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "db_record_insert INSERT prep: %s", sqlite3_errmsg(db)); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return -1; }
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)author_node_id);
sqlite3_bind_int64(stmt, 3, (sqlite3_int64)id);
sqlite3_bind_blob(stmt, 4, ch, 32, SQLITE_STATIC);
if (jlen > 0) sqlite3_bind_blob(stmt, 5, json, (int)jlen, SQLITE_STATIC);
else sqlite3_bind_null(stmt, 5);
sqlite3_bind_blob(stmt, 6, author_sig, DB_SIG_SIZE, SQLITE_STATIC);
if (local_attrs && local_attrs[0]) sqlite3_bind_text(stmt, 7, local_attrs, -1, SQLITE_STATIC);
else sqlite3_bind_text(stmt, 7, "", -1, SQLITE_STATIC);
rc = sqlite3_step(stmt);
sqlite3_finalize(stmt);
if (rc != SQLITE_DONE) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "INSERT: %s", sqlite3_errmsg(db)); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return -1; }
si->next_id = id + 1;
rc = sqlite3_exec(db, "COMMIT", NULL, NULL, NULL);
if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "insert COMMIT: %s", sqlite3_errmsg(db)); return -1; }
uint32_t ins_pos = si_find_pos(si, ts, author_sig);
db_sync_schedule_recalc(si, ins_pos);
return 0;
}
// Вставляет запись + каскад: on_insert + PUSH всем connected пирам (кроме источника)
static int db_record_insert_cascade(struct DB_SYNC_INSTANCE* si,
uint64_t rid, uint64_t rts, uint64_t rauthor,
const char* rdata, uint32_t rdlen,
const uint8_t* rsig, uint64_t source_peer,
const char* local_attrs)
{
int ret = db_record_insert(si, rid, rts, rauthor, rdata, rdlen, rsig, local_attrs);
if (ret != 0) return ret;
if (si->on_insert) si->on_insert(si, rts, rdata, rdlen, rauthor, si->on_insert_arg);
uint8_t pbuf[4096]; uint32_t poff = 1; pbuf[0] = DB_MSG_PUSH;
memcpy(pbuf + poff, &rid, 8); poff += 8;
memcpy(pbuf + poff, &rts, 8); poff += 8;
memcpy(pbuf + poff, &rauthor, 8); poff += 8;
memcpy(pbuf + poff, &rdlen, 4); poff += 4;
if (rdlen > 0) { memcpy(pbuf + poff, rdata, rdlen); poff += rdlen; }
uint32_t sl = DB_SIG_SIZE; pbuf[poff++] = (uint8_t)sl;
memcpy(pbuf + poff, rsig, sl); poff += sl;
uint64_t my_id = si->db_sync->inst->node_id;
int push_count = 0;
for (int i = 0; i < si->peer_count; i++) {
uint64_t pid = si->peers[i].node_id;
if (si->peers[i].sync_state >= 1 && pid != source_peer && pid != my_id) {
if (db_sync_send(si, pid, pbuf, poff) >= 0)
{ si_delivery_update(si, rts, rauthor, pid); push_count++; }
}
}
if (push_count > 0)
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:ME--] → PUSH: id=%llu => %d peers",
SI_SHRT(si), (unsigned long long)rid, push_count);
return 0;
}
// ============================================================
// Send
// ============================================================
// Отправляет сообщение узлу через ETCP с заголовком: [ETCP_RT_ID_DB_SYNC:1][group_id_be:8][payload]. Возвращает код etcp_send
static int db_sync_send_gid(struct DB_SYNC* db, uint64_t node_id, uint64_t group_id, const uint8_t* payload, size_t plen)
{
struct ll_entry* entry = queue_entry_new(0);
if (!entry) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "queue_entry_new"); return -1; }
uint8_t* buf = u_malloc(plen + 9);
if (!buf) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "db_sync_send_gid: u_malloc(%zu) failed for dst=%016llx", plen + 9, (unsigned long long)node_id); queue_entry_free(entry); return -1; }
buf[0] = ETCP_RT_ID_DB_SYNC;
uint64_t gb = htobe64(group_id);
memcpy(buf + 1, &gb, 8);
memcpy(buf + 9, payload, plen);
entry->dgram = buf;
entry->len = plen + 9;
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);
queue_entry_free(entry);
return -1;
}
int ret = etcp_send(conn, 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;
}
// Отправляет сообщение через 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);
}
// ============================================================
// Delivery chain helpers (local fields)
// ============================================================
// Обновляет delivery_chain (добавляет peer_id hex) и счётчик delivered_peers для записи после отправки пиру
static void si_delivery_update(struct DB_SYNC_INSTANCE* si, uint64_t ts, uint64_t author, uint64_t peer_id)
{
char hex[17]; snprintf(hex, sizeof(hex), "%016llx", (unsigned long long)peer_id);
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT delivery_chain FROM \"%s\""
" WHERE timestamp=? AND node_id=?") == SQLITE_OK)
{
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)author);
if (sqlite3_step(stmt) == SQLITE_ROW) {
const char* dc = (const char*)sqlite3_column_text(stmt, 0);
if (dc && strstr(dc, hex)) { sqlite3_finalize(stmt); return; }
}
sqlite3_finalize(stmt);
}
if (si_prep(si, &stmt,
"UPDATE \"%s\""
" SET delivered_peers = delivered_peers + 1,"
" delivery_chain = CASE"
" WHEN delivery_chain = '' THEN ?1"
" ELSE delivery_chain || ',' || ?1"
" END"
" WHERE timestamp = ?2 AND author = ?3") == SQLITE_OK)
{
sqlite3_bind_text(stmt, 1, hex, -1, SQLITE_STATIC);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)ts);
sqlite3_bind_int64(stmt, 3, (sqlite3_int64)author);
sqlite3_step(stmt);
sqlite3_finalize(stmt);
}
}
// ============================================================
// Sync protocol handlers
// ============================================================
// Определяет направление инициации синхронизации: узел с бОльшим node_id — мастер
static int is_master(uint64_t a, uint64_t b) { return a > b; }
// Обрабатывает REQUEST_SYNC: запускает синхронизацию как слейв (мастер сам пришлёт INIT_SYNC)
static void db_handle_request_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len)
{
(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();
db_sync_initiate_sync(si, src);
}
}
// Обрабатывает ERROR от пира: логирует причину, сбрасывает sync_state пира
static void db_handle_error(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len)
{
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,
"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);
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);
sp->sync_state = 0;
}
}
// Обрабатывает INIT_SYNC: сбрасывает pending recalc, определяет точку расхождения (tp), формирует INIT_RESP
// с разреженными хешами (интервалы: 1,1,1,2,2,4,4,4,4,4,4,4,4,4,4,4) и флагом has_tail
static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len)
{
if (len < 4) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "INIT_SYNC too short %zu from %016llx", len, (unsigned long long)src); return; }
db_sync_flush_recalc(si);
uint32_t pc = *(uint32_t*)p;
uint32_t mc = db_count(si);
uint32_t tp = (pc < mc ? pc : mc);
if (tp > 0) tp--;
uint64_t my_ch8_at_tp = 0;
if (mc > 0) db_chain_hash8_at(si, tp, &my_ch8_at_tp);
uint8_t has_tail = (mc > tp + 1) ? 1 : 0;
static const uint16_t iv[16] = {1, 1, 1, 2, 2, 4, 4, 4, 4, 4, 4, 4, 4, 4, 4, 4};
uint32_t accum = 0; int sparse_cnt = 0;
for (int i = 0; i < 16; i++) { if (tp < accum + iv[i]) break; accum += iv[i]; sparse_cnt++; }
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← INIT_SYNC: peer=%u my=%u → tp=%u ch8=%016llX has_tail=%u sparse=%d",
SI_SHRT(si), (unsigned long long)(src >> 16), pc, mc, tp, my_ch8_at_tp, has_tail, sparse_cnt);
uint8_t resp[4096];
uint32_t off = 0;
resp[off++] = DB_MSG_INIT_RESP;
memcpy(resp + off, &tp, 4); off += 4;
uint64_t my_ch8 = 0;
if (mc > 0) db_chain_hash8_at(si, tp, &my_ch8);
memcpy(resp + off, &my_ch8, 8); off += 8;
resp[off++] = has_tail;
uint32_t sp = off;
off++;
int scnt = 0;
accum = 0;
for (int i = 0; i < 16; i++) {
if (tp < accum + iv[i]) break;
accum += iv[i];
uint32_t pos = tp - accum;
uint64_t sch8;
if (db_chain_hash8_at(si, pos, &sch8) != 0) break;
if (off + 12 > sizeof(resp)) break;
memcpy(resp + off, &pos, 4); off += 4;
memcpy(resp + off, &sch8, 8); off += 8;
scnt++;
}
resp[sp] = (uint8_t)scnt;
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(); }
}
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,
resp[0], resp[1]);
}
// Обрабатывает INIT_RESP: сравнивает хеши на точке tp, находит первую расходящуюся позицию,
// отправляет наши данные начиная с неё и запрашивает данные пира с той же позиции
static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len)
{
if (len < 14) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "INIT_RESP too short %zu from %016llx", len, (unsigned long long)src); return; }
db_sync_flush_recalc(si);
uint32_t tp = *(uint32_t*)p;
uint64_t peer_ch8; memcpy(&peer_ch8, p + 4, 8);
uint8_t has_tail = p[12];
uint8_t sc = p[13];
const uint8_t* spr = p + 14;
uint64_t my_ch8 = 0;
uint32_t mc = db_count(si);
if (mc > 0 && tp < mc) db_chain_hash8_at(si, tp, &my_ch8);
DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← INIT_RESP: tp=%u my=%u my_ch8=%016llX peer_ch8=%016llX has_tail=%u sc=%d",
SI_SHRT(si), (unsigned long long)(src >> 16), tp, mc, my_ch8, peer_ch8, has_tail, sc);
struct SI_PEER* sp = si_peer_find(si, src);
// Peer empty → send ALL our data, then SYNC_DONE will be sent from DATA handler when slave responds.
// Если мы тоже пусты (mc==0) — обе стороны пусты: не отправляем ничего, а завершаем sync через
// SYNC_DONE в ветке "exact match" ниже (иначе обе стороны вечно ждут ответа → TIMEOUT).
if (peer_ch8 == 0 && sc == 0 && mc > 0) {
uint32_t vp = (uint32_t)-1; if (sp) sp->verified_pos = vp;
uint32_t sent = 0;
while (sent < mc) {
uint32_t b = mc - sent; if (b > DB_SEND_DATA_MAX) b = DB_SEND_DATA_MAX;
uint32_t want = (sent + b >= mc) ? mc : DB_WANT_FROM_NONE;
if (si_send_data_batch(si, src, sent, b, vp, want) <= 0) break;
sent += b;
}
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← INIT_RESP: peer empty → sent %u/%u records, awaiting response",
SI_SHRT(si), (unsigned long long)(src >> 16), sent, mc);
return;
}
// Exact match: same hash at tp
if (my_ch8 == peer_ch8) {
if (!has_tail && mc <= tp + 1) {
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);
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);
return;
}
// Has tail (slave or us): send our records after tp, request peer's after tp
uint32_t vp = tp; if (sp) sp->verified_pos = vp;
uint32_t rfrom = tp + 1;
uint32_t scnt = mc - rfrom; if (scnt > DB_SEND_DATA_MAX) scnt = DB_SEND_DATA_MAX;
if (scnt > 0) si_send_data_batch(si, src, rfrom, scnt, vp, rfrom);
else si_send_data_batch(si, src, rfrom, 0, vp, rfrom);
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← INIT_RESP: MATCH+tail → SEND_DATA from=%u want=%u",
SI_SHRT(si), (unsigned long long)(src >> 16), rfrom, rfrom);
return;
}
// Mismatch: find first divergent position
uint32_t fm = tp;
for (int i = 0; i < sc && spr + 12 <= p + len; i++) {
uint32_t pos = *(uint32_t*)spr; uint64_t pch8; memcpy(&pch8, spr + 4, 8); spr += 12;
uint64_t mch8;
if (db_chain_hash8_at(si, pos, &mch8) == 0) {
if (mch8 != pch8 && pos < fm) fm = pos;
} else if (pos < fm) fm = pos;
}
uint32_t vp = fm > 0 ? fm - 1 : (uint32_t)-1; if (sp) sp->verified_pos = vp;
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← INIT_RESP: MISMATCH fm=%u vp=%u mc=%u",
SI_SHRT(si), (unsigned long long)(src >> 16), fm, vp, mc);
// Send our records from fm, request peer's from fm
uint32_t scnt = mc - fm; if (scnt > DB_SEND_DATA_MAX) scnt = DB_SEND_DATA_MAX;
if (scnt > 0) si_send_data_batch(si, src, fm, scnt, vp, fm);
else si_send_data_batch(si, src, fm, 0, vp, fm);
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] → SEND_DATA from=%u vp=%u cnt=%u want=%u",
SI_SHRT(si), (unsigned long long)(src >> 16), fm, vp, scnt, fm);
}
// ---- SEND_DATA batch helper ----
// Wire: [from:4][count:2][vp:4][want_from:4][records...]
// Record: [id:8][ts:8][node_id:8][dlen:4][data][sig_len:1=64][sig:64]
// Returns: number of records sent, or -1 on SQL error
// Формирует и отправляет батч записей (до 32) в формате SEND_DATA:
// [from:4][count:2][vp:4][want_from:4][records...]. Возвращает количество отправленных записей
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)
{
uint8_t sbuf[8192]; uint8_t* buf = sbuf;
uint32_t off = 0;
buf[off++] = DB_MSG_SEND_DATA;
memcpy(buf + off, &from, 4); off += 4;
uint16_t rc = 0;
uint16_t* rcp = (uint16_t*)(buf + off); off += 2;
memcpy(buf + off, &vp, 4); off += 4;
memcpy(buf + off, &want_from, 4); off += 4; // NEW: request peer's data from this position, DB_WANT_FROM_NONE=none
if (count > 0) {
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT id,timestamp,node_id,data,author_signature"
" FROM \"%s\" ORDER BY timestamp, author_signature"
" LIMIT ? OFFSET ?") != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] SEND_DATA: si_prep failed tbl=%s", SI_SHRT(si), (unsigned long long)(dst >> 16), SI_TBL(si)); return -1; }
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)count);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)from);
int bad_del = 0;
while (sqlite3_step(stmt) == SQLITE_ROW) {
uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0);
uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1);
uint64_t rauth = (uint64_t)sqlite3_column_int64(stmt, 2);
const uint8_t* rd = (const uint8_t*)sqlite3_column_blob(stmt, 3);
uint32_t rdl = (uint32_t)sqlite3_column_bytes(stmt, 3); if (!rd) rdl = 0;
const uint8_t* rsig = (const uint8_t*)sqlite3_column_blob(stmt, 4);
int rsl = sqlite3_column_bytes(stmt, 4); if (!rsig) rsl = 0;
if (rsl == DB_SIG_SIZE && db_verify_author_sig(si, rts, rauth, (const char*)rd, rdl, rsig) != 0) {
uint32_t del_pos = si_find_pos(si, rts, rsig);
sqlite3_stmt* del;
if (si_prep(si, &del, "DELETE FROM \"%s\" WHERE timestamp=? AND author_signature=?") == SQLITE_OK) {
sqlite3_bind_int64(del, 1, (sqlite3_int64)rts);
sqlite3_bind_blob(del, 2, rsig, DB_SIG_SIZE, SQLITE_STATIC);
sqlite3_step(del); sqlite3_finalize(del);
db_sync_schedule_recalc(si, del_pos); bad_del++;
}
continue;
}
int rec_sz = 28 + rdl + 1 + rsl;
if (off + rec_sz > 8000) break;
memcpy(buf + off, &rid, 8); off += 8;
memcpy(buf + off, &rts, 8); off += 8;
memcpy(buf + off, &rauth, 8); off += 8;
memcpy(buf + off, &rdl, 4); off += 4;
if (rdl > 0) { memcpy(buf + off, rd, rdl); off += rdl; }
buf[off++] = (uint8_t)rsl;
if (rsl > 0) { memcpy(buf + off, rsig, rsl); off += rsl; }
rc++;
}
sqlite3_finalize(stmt);
if (rc == 0 && bad_del == 0)
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] SEND_DATA: 0 records from=%u count=%u",
SI_SHRT(si), (unsigned long long)(dst >> 16), from, count);
}
*rcp = rc;
db_sync_send(si, dst, buf, off);
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;
}
// ---- Parse one record from SEND_DATA/PUSH wire format ----
// Wire: [id:8][ts:8][author:8][dlen:4][data][sig_len:1][sig]
// Разбирает одну запись из бинарного формата: [id:8][ts:8][author:8][dlen:4][data][sig_len:1][sig]
// Продвигает указатель *pp, возвращает 0 при успехе, -1 при ошибке парсинга
static int si_parse_record(const uint8_t** pp, const uint8_t* end,
uint64_t* rid, uint64_t* rts, uint64_t* rauthor,
uint32_t* rdlen, const uint8_t** rdata,
const uint8_t** rsig, int* rsiglen)
{
if (*pp + 28 > end) { DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "si_parse_record: truncated header, need 28 have %td", end - *pp); return -1; }
*rid = *(uint64_t*)*pp; *pp += 8;
*rts = *(uint64_t*)*pp; *pp += 8;
*rauthor = *(uint64_t*)*pp; *pp += 8;
*rdlen = *(uint32_t*)*pp; *pp += 4;
if (*pp + *rdlen > end) { DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "si_parse_record: data overrun dlen=%u have=%td", *rdlen, end - *pp); return -1; }
*rdata = (*rdlen > 0) ? *pp : NULL;
*pp += *rdlen;
if (*pp + 1 > end) { DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "si_parse_record: sig_len overrun"); return -1; }
*rsiglen = (int)(*(*pp)++);
if (*rsiglen > 0) {
if (*pp + *rsiglen > end) { DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "si_parse_record: sig data overrun slen=%d have=%td", *rsiglen, end - *pp); return -1; }
*rsig = *pp;
*pp += *rsiglen;
} else {
*rsig = NULL;
}
return 0;
}
// ---- Atomic SEND_DATA exchange ----
// Обрабатывает SEND_DATA: вставляет полученные записи (до 32), проверяет подписи,
// корректирует synced_pos, при want_from шлёт встречные данные
static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len)
{
if (len < 14) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← SEND_DATA: truncated len=%zu (need >=14)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; }
db_sync_flush_recalc(si);
uint32_t from = *(uint32_t*)p;
uint16_t count = *(uint16_t*)(p + 4);
uint32_t vp = *(uint32_t*)(p + 6);
uint32_t want_from = *(uint32_t*)(p + 10);
struct SI_PEER* sp = si_peer_find(si, src);
const uint8_t* ptr = p + 14;
uint16_t received = 0, duplicates = 0, bad_sigs = 0;
for (uint16_t i = 0; i < count && i < 32; i++) {
uint64_t rid, rts, rauthor; uint32_t rdlen; const uint8_t* rdata, *rsig; int rsiglen;
if (si_parse_record(&ptr, p + len, &rid, &rts, &rauthor, &rdlen, &rdata, &rsig, &rsiglen) != 0) break;
int ret = db_record_insert_cascade(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, src, NULL);
if (ret >= 0) received++;
if (ret == 1) duplicates++;
if (ret == -2) bad_sigs++;
}
uint32_t old_spos = sp ? sp->synced_pos : 0;
if (received > 0) {
uint32_t rst = (from < old_spos) ? from : old_spos;
if (sp) sp->synced_pos = from + (count - bad_sigs) - 1;
for (int j = 0; j < si->peer_count; j++)
if (si->peers[j].synced_pos >= rst && &si->peers[j] != sp)
si->peers[j].synced_pos = rst > 0 ? rst - 1 : 0;
}
char extra[32] = ""; if (duplicates) snprintf(extra, sizeof(extra), " (%u dup)", duplicates);
if (bad_sigs) { size_t el = strlen(extra); snprintf(extra+el, sizeof(extra)-el, "%s%u bad", extra[0]?", ":" (", bad_sigs); }
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← DATA: from=%u cnt=%u vp=%u want=%u → ins=%u%s mc=%u sp=%u→%u",
SI_SHRT(si), (unsigned long long)(src >> 16), from, count, vp, want_from, received, extra, db_count(si),
old_spos, sp ? sp->synced_pos : 0);
// Respond with our data if peer requested (want_from != DB_WANT_FROM_NONE)
if (want_from != DB_WANT_FROM_NONE) {
uint32_t mc = db_count(si);
if (want_from < mc) {
uint32_t scnt = mc - want_from; if (scnt > DB_SEND_DATA_MAX) scnt = DB_SEND_DATA_MAX;
si_send_data_batch(si, src, want_from, scnt, vp, DB_WANT_FROM_NONE);
} else {
si_send_data_batch(si, src, want_from, 0, vp, DB_WANT_FROM_NONE);
}
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← DATA: responded with %u recs from=%u",
SI_SHRT(si), (unsigned long long)(src >> 16), want_from < mc ? (mc - want_from > DB_SEND_DATA_MAX ? DB_SEND_DATA_MAX : mc - want_from) : 0, want_from);
}
// If slave responded with want=DB_WANT_FROM_NONE and we are in active sync — complete the exchange with SYNC_DONE
if (want_from == DB_WANT_FROM_NONE && sp && sp->sync_state == 1) {
uint32_t mc = db_count(si);
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);
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",
SI_SHRT(si), (unsigned long long)(src >> 16), mc);
}
}
// Обрабатывает SYNC_DONE: финальная сверка хешей. mc==pc → sync завершён, mc<pc → запрашиваем хвост, mc>pc → отправляем хвост
static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len)
{
if (len < 12) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "SYNC_DONE too short %zu from %016llx", len, (unsigned long long)src); return; }
db_sync_flush_recalc(si);
uint32_t pc = *(uint32_t*)p;
uint64_t pch8; memcpy(&pch8, p + 4, 8);
uint32_t mc = db_count(si);
uint64_t mch8 = 0; if (mc > 0) db_chain_hash8_at(si, mc - 1, &mch8);
struct SI_PEER* sp = si_peer_find(si, src);
if (mc == pc && mch8 == pch8) {
if (sp) { sp->synced_pos = mc > 0 ? mc - 1 : 0; sp->sync_state = 2; sp->sync_start_tb = 0; }
si_fire_sync_done(si, src);
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← SYNC_DONE: matched ✓ mc=%u ss→2",
SI_SHRT(si), (unsigned long long)(src >> 16), mc);
} else if (mc < pc) {
// Request tail
uint32_t vp_send = mc > 0 ? mc - 1 : (uint32_t)-1;
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;
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) {
// mc > pc: send tail
if (sp) sp->last_peer_count = pc;
uint32_t scnt = mc - pc; if (scnt > DB_SEND_DATA_MAX) scnt = DB_SEND_DATA_MAX;
si_send_data_batch(si, src, pc, scnt, pc > 0 ? pc - 1 : (uint32_t)-1, DB_WANT_FROM_NONE);
if (sp) sp->sync_state = 1;
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← SYNC_DONE: mc=%u > pc=%u → send tail from=%u",
SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, pc);
} else {
// mc == pc but chain hash differs — accept as converged
if (sp) { sp->synced_pos = mc > 0 ? mc - 1 : 0; sp->sync_state = 2; sp->sync_start_tb = 0; }
si_fire_sync_done(si, src);
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← SYNC_DONE: mc=%u == pc=%u (hash !match) → accept as converged",
SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc);
}
}
// Обрабатывает ACK_PUSH: помечает запись флагом DB_REC_FLAG_WAS_SENT и обновляет delivery_chain
static void db_handle_ack_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len)
{
if (len < 16) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← PUSH ACK: truncated len=%zu (need >=16)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; }
uint64_t ts = *(uint64_t*)p;
uint64_t author = *(uint64_t*)(p + 8);
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "from=%016llx ts=%llu author=%016llx", (unsigned long long)src, (unsigned long long)ts, (unsigned long long)author);
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"UPDATE \"%s\" SET flags = flags | 1"
" WHERE timestamp=? AND node_id=?"
" AND node_id=? AND (flags & 1) = 0") == SQLITE_OK)
{
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)author);
sqlite3_bind_int64(stmt, 3, (sqlite3_int64)si->db_sync->inst->node_id);
sqlite3_step(stmt);
int changed = sqlite3_changes(SI_DB(si));
sqlite3_finalize(stmt);
if (changed > 0)
DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← PUSH ACK: ts=%llu — delivery confirmed by peer",
SI_SHRT(si), (unsigned long long)(src >> 16), (unsigned long long)ts);
}
si_delivery_update(si, ts, author, src);
}
// Обрабатывает PUSH: вставляет новую запись от пира, шлёт ACK_PUSH, корректирует synced_pos остальных пиров
// если запись вставлена не в конец (в середину)
static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len)
{
if (len < 30) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← PUSH: truncated len=%zu (need >=30)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; }
const uint8_t* ptr = p;
uint64_t rid, rts, rauthor;
uint32_t rdlen;
const uint8_t* rdata, *rsig;
int rsiglen;
if (si_parse_record(&ptr, p + len, &rid, &rts, &rauthor, &rdlen, &rdata, &rsig, &rsiglen) != 0) return;
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "from=%016llx id=%llu author=%016llx ts=%llu len=%u", (unsigned long long)src, (unsigned long long)rid, (unsigned long long)rauthor, (unsigned long long)rts, rdlen);
int ret = db_record_insert_cascade(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, src, NULL);
if (ret == 0) {
uint32_t ins_pos = si_find_pos(si, rts, rsig);
uint32_t total = db_count(si);
int adj_count = 0;
for (int j = 0; j < si->peer_count; j++) {
if (si->peers[j].synced_pos >= ins_pos) { si->peers[j].synced_pos = ins_pos > 0 ? ins_pos - 1 : 0; adj_count++; }
}
uint8_t ack[17];
ack[0] = DB_MSG_ACK_PUSH;
memcpy(ack + 1, &rts, 8);
memcpy(ack + 9, &rauthor, 8);
db_sync_send(si, src, ack, 17);
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:%04llX] ← PUSH: msg id=%llu author=%016llX ts=%llu → inserted at pos=%u (of %u total)",
SI_SHRT(si), (unsigned long long)(src >> 16),
(unsigned long long)rid, (unsigned long long)rauthor, (unsigned long long)rts, ins_pos, total);
if (adj_count > 0 && ins_pos < total - 1) {
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s:----] PUSH: inserted at pos=%u NOT at tail → reset synced_pos of %d peers from ≥%u back to %u",
SI_SHRT(si), ins_pos, adj_count, ins_pos, ins_pos > 0 ? ins_pos - 1 : 0);
}
}
}
// ============================================================
// Receive callback
// ============================================================
// Главный колбэк приёма ETCP: извлекает хеш инстанса и тип сообщения, находит/создаёт инстанс,
// маршрутизирует в соответствующий обработчик (INIT_SYNC, INIT_RESP, SEND_DATA, PUSH, ACK_PUSH, SYNC_DONE, ERROR, REQUEST_SYNC)
static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry)
{
if (!entry || entry->len < 10) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return; }
struct DB_SYNC* db = conn ? conn->instance->db_sync : NULL;
if (!db || !db->enabled) { queue_dgram_free(entry); queue_entry_free(entry); return; }
uint64_t src = conn ? conn->peer_node_id : 0;
uint64_t group_id = be64toh(*(uint64_t*)(entry->dgram + 1));
uint8_t type = entry->dgram[9];
const uint8_t* payload = entry->dgram + 10;
size_t plen = entry->len - 10;
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "recv type=%02x from=%016llx group=%016llx len=%zu", type, (unsigned long long)src, (unsigned long long)group_id, plen);
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);
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;
}
}
if (!si->enabled) {
uint8_t err[2]; err[0] = DB_MSG_ERROR; err[1] = DB_ERR_DISABLED;
db_sync_send_gid(db, src, group_id, err, 2);
queue_dgram_free(entry); queue_entry_free(entry); return;
}
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;
case DB_MSG_SEND_DATA: DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ handle SEND_DATA"); db_handle_send_data(si, src, payload, plen); break;
case DB_MSG_PUSH: DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ handle PUSH"); db_handle_push(si, src, payload, plen); break;
case DB_MSG_ACK_PUSH: DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ handle ACK_PUSH"); db_handle_ack_push(si, src, payload, plen); break;
case DB_MSG_REQUEST_SYNC: DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ handle REQUEST_SYNC"); db_handle_request_sync(si, src, payload, plen); break;
case DB_MSG_SYNC_DONE: DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ handle SYNC_DONE"); db_handle_sync_done(si, src, payload, plen); break;
case DB_MSG_ERROR: DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ handle ERROR"); db_handle_error(si, src, payload, plen); break;
default:
DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "sync: unknown msg type 0x%02x from %04llX (group=%016llx) — dropped",
type, (unsigned long long)(src >> 16), group_id);
break;
}
queue_dgram_free(entry);
queue_entry_free(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);
}
// ============================================================
// Initiate sync
// ============================================================
// Запускает синхронизацию: мастер шлёт INIT_SYNC со своим количеством записей, слейв шлёт REQUEST_SYNC
static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t pid)
{
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);
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);
} else {
if (p) { p->sync_state = 1; }
uint8_t msg[1] = { DB_MSG_REQUEST_SYNC };
db_sync_send(si, pid, msg, 1);
}
db_sync_resume_peer_check(si->db_sync);
}
// Проверяет целостность цепочки хешей всей таблицы, при несовпадении запускает пересчёт с позиции ошибки
static void db_verify_chain(struct DB_SYNC_INSTANCE* si)
{
uint32_t mc = db_count(si);
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "tbl=%s mc=%u", SI_TBL(si), mc);
if (mc == 0) return;
uint8_t prev_ch[32]; memset(prev_ch, 0, 32);
uint8_t exp_ch[32], stored_ch[32];
int bad_pos = -1;
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT id,timestamp,node_id,author_signature,chain_hash FROM \"%s\""
" ORDER BY timestamp, author_signature") != SQLITE_OK)
return;
for (uint32_t pos = 0; sqlite3_step(stmt) == SQLITE_ROW; pos++) {
uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0);
uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1);
uint64_t rauth = (uint64_t)sqlite3_column_int64(stmt, 2);
const void* sig_blob = sqlite3_column_blob(stmt, 3);
uint8_t sig[DB_SIG_SIZE];
if (sig_blob) memcpy(sig, sig_blob, DB_SIG_SIZE); else memset(sig, 0, DB_SIG_SIZE);
const void* b = sqlite3_column_blob(stmt, 4);
if (b && sqlite3_column_bytes(stmt, 4) >= 32) memcpy(stored_ch, b, 32);
else memset(stored_ch, 0, 32);
db_chain_hash_compute(prev_ch, rid, rts, rauth, sig, exp_ch);
if (memcmp(exp_ch, stored_ch, 32) != 0) { bad_pos = (int)pos; break; }
memcpy(prev_ch, exp_ch, 32);
}
sqlite3_finalize(stmt);
if (bad_pos >= 0) {
DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "chain_hash mismatch at pos %d/%u in %s, recalculating",
bad_pos, mc, SI_TBL(si));
db_sync_schedule_recalc(si, (uint32_t)bad_pos);
}
}
// Публичная проверка целостности цепных хешей: возвращает 0 если всё верно, 1 при несовпадении, -1 при ошибке
int db_sync_chain_verify(struct DB_SYNC_INSTANCE* si)
{
if (!si || !si->enabled) return 0;
db_sync_flush_recalc(si);
uint32_t mc = db_count(si);
if (mc == 0) return 0;
uint8_t prev_ch[32]; memset(prev_ch, 0, 32);
uint8_t exp_ch[32], stored_ch[32];
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT id,timestamp,node_id,author_signature,chain_hash FROM \"%s\""
" ORDER BY timestamp, author_signature") != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "verify_chain_hash: si_prep failed tbl=%s", SI_TBL(si));
return -1;
}
for (uint32_t pos = 0; sqlite3_step(stmt) == SQLITE_ROW; pos++) {
uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0);
uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1);
uint64_t rauth = (uint64_t)sqlite3_column_int64(stmt, 2);
const void* sig_blob = sqlite3_column_blob(stmt, 3);
uint8_t sig[DB_SIG_SIZE];
if (sig_blob) memcpy(sig, sig_blob, DB_SIG_SIZE); else memset(sig, 0, DB_SIG_SIZE);
const void* b = sqlite3_column_blob(stmt, 4);
if (b && sqlite3_column_bytes(stmt, 4) >= 32) memcpy(stored_ch, b, 32);
else memset(stored_ch, 0, 32);
db_chain_hash_compute(prev_ch, rid, rts, rauth, sig, exp_ch);
if (memcmp(exp_ch, stored_ch, 32) != 0) { sqlite3_finalize(stmt); return 1; }
memcpy(prev_ch, exp_ch, 32);
}
sqlite3_finalize(stmt);
return 0;
}
// Устанавливает состояние синхронизации пира (для тестов)
void db_sync_peer_set_state(struct DB_SYNC_INSTANCE* si, uint64_t node_id, uint8_t state)
{
struct SI_PEER* p = si_peer_find(si, node_id);
if (p) p->sync_state = state;
}
// Гейт авто-синхронизации: пока gated — инстанс не запускает sync автоматически
// (ждёт member_sync/pubkeys). При снятии гейта запускает sync к пирам в sync_state==0.
void db_sync_instance_set_gated(struct DB_SYNC_INSTANCE* si, int gated)
{
if (!si || !si->enabled) return;
gated = gated ? 1 : 0;
if (si->sync_gated == (uint8_t)gated) return;
si->sync_gated = (uint8_t)gated;
if (!gated) {
struct DB_SYNC* db = si->db_sync;
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++;
}
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 {
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s] → gated: auto-sync suspended (mc=%u)",
SI_SHRT(si), db_count(si));
}
}
// Принудительно перезапускает синхронизацию с указанным пиром (для тестов)
void db_sync_reinitiate(struct DB_SYNC_INSTANCE* si, uint64_t node_id)
{
struct SI_PEER* p = si_peer_find(si, node_id);
if (p) {
p->sync_state = 1; p->sync_start_tb = get_time_tb();
db_sync_initiate_sync(si, node_id);
}
}
// Возвращает последний chain_hash8 (для сверки синхронизации в тестах)
int db_sync_last_chain_hash8(struct DB_SYNC_INSTANCE* si, uint64_t* out_hash8)
{
if (!si || !out_hash8) return -1;
db_sync_flush_recalc(si);
uint32_t mc = db_count(si);
if (mc == 0) { *out_hash8 = 0; return 0; }
return db_chain_hash8_at(si, mc - 1, out_hash8);
}
// ============================================================
// Timers
// ============================================================
// Есть ли у любого активного инстанса пир, требующий синхронизации (sync_state != 2)
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;
for (int j = 0; j < si->peer_count; j++)
if (si->peers[j].sync_state != 2) return 1;
}
return 0;
}
// Взвести work-driven таймер проверки пиров, если он ещё не взведён и модуль активен
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");
}
// Периодическая проверка: запускает синхронизацию для пиров в состоянии sync_state==0,
// детектирует таймауты синхронизации (DB_SYNC_SYNC_TIMEOUT секунд), логирует зависшие.
// Work-driven: перевзводится только пока есть пиры с sync_state != 2.
static void db_sync_peer_check_cb(void* arg)
{
struct DB_SYNC* db = (struct DB_SYNC*)arg;
if (!db || !db->enabled) return;
db->peer_check_timer = NULL;
#ifdef UTUN_HAVE_STANDBY
if (standby_get_sleep_tb() > 0) {
db->peer_check_wait = standby_wait(db, db_sync_peer_check_cb);
return;
}
db->peer_check_wait = NULL;
#endif
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "instances=%d", db->instance_count);
struct TOPO_GROUP* g = topo_groups_get_default(db->inst->topo_groups);
int total_synced = 0, total_skipped = 0, any_peers = 0;
char launched_list[256] = "";
if (g) {
for (int i = 0; i < db->instance_count; i++) {
struct DB_SYNC_INSTANCE* si = db->instances[i];
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) {
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);
}
}
}
}
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");
}
// Периодическая TTL-очистка: удаляет неподтверждённые записи (flags&1==0) старше db_sync_ttl от локального узла
static void db_sync_instance_ttl_cb(void* arg)
{
struct DB_SYNC_INSTANCE* si = (struct DB_SYNC_INSTANCE*)arg;
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "tbl=%s enabled=%d", SI_TBL(si), si->enabled);
uint64_t interval = DB_SYNC_TTL_INTERVAL * 10000u;
if (!si || !si->enabled || !si->db_sync || !si->db_sync->db) {
if (si) si->ttl_timer = uasync_set_timeout(si->db_sync->inst->ua, interval, si, db_sync_instance_ttl_cb, "db_sync_ttl");
return;
}
uint64_t my_id = si->db_sync->inst->node_id;
uint64_t nu = get_time_us();
uint64_t ttl = (uint64_t)(si->db_sync->inst->config->global.db_sync_ttl) * 1000000uLL;
uint64_t cut = nu - ttl;
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"DELETE FROM \"%s\""
" WHERE node_id=? AND (flags & 1) = 0"
" AND timestamp<?") != SQLITE_OK)
{
si->ttl_timer = uasync_set_timeout(si->db_sync->inst->ua, interval, si, db_sync_instance_ttl_cb, "db_sync_ttl");
return;
}
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)my_id);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)cut);
sqlite3_step(stmt);
int d = sqlite3_changes(SI_DB(si));
sqlite3_finalize(stmt);
if (d > 0) DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "TTL cleanup deleted %d unsent records (tbl=%s)", d, SI_TBL(si));
si->ttl_timer = uasync_set_timeout(si->db_sync->inst->ua, interval, si, db_sync_instance_ttl_cb, "db_sync_ttl");
}
// ============================================================
// Public API
// ============================================================
// Инициализирует модуль db_sync: открывает БД, регистрирует ETCP-биндинг, запускает таймер проверки пиров
int db_sync_init(struct UTUN_INSTANCE* inst)
{
if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "NULL instance"); return -1; }
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "inst=%p enabled=%d", (void*)inst, inst->config->global.db_sync_enabled);
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;
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "db_sync: disabled by config");
return 0;
}
db->enabled = 1;
inst->db_sync = db;
const char* dp = inst->config->global.db_path;
if (inst->topo_sqlite_db) {
db->db = inst->topo_sqlite_db; db->shared_db = 1;
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "db_sync: using shared SQLite db=%p", (void*)db->db);
} else {
char sp[512];
if (dp[0]) snprintf(sp, sizeof(sp), "%s/chats.db", dp);
else snprintf(sp, sizeof(sp), "/tmp/utun_db_sync");
if (db_sqlite_open(db, sp) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "SQLite open failed, sync disabled");
db->enabled = 0;
}
}
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;
}
// Уничтожает модуль db_sync: отменяет таймеры, закрывает БД, освобождает память
void db_sync_destroy(struct UTUN_INSTANCE* inst)
{
if (!inst || !inst->db_sync) return;
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "inst=%p instances=%d", (void*)inst, inst->db_sync->instance_count);
struct DB_SYNC* db = inst->db_sync;
inst->db_sync = NULL;
etcp_unbind(inst, ETCP_RT_ID_DB_SYNC);
if (db->peer_check_timer) { uasync_cancel_timeout(inst->ua, db->peer_check_timer); db->peer_check_timer = NULL; }
#ifdef UTUN_HAVE_STANDBY
if (db->peer_check_wait) { standby_wait_cancel(db->peer_check_wait); db->peer_check_wait = NULL; }
#endif
for (int i = 0; i < db->instance_count; i++) {
struct DB_SYNC_INSTANCE* si = db->instances[i];
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; } }
db_sqlite_close(db);
if (db->instances) u_free(db->instances);
u_free(db);
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)
{
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 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; }
si = db_instance_alloc(db);
if (!si) return NULL;
si->group_id = group_id;
snprintf(si->table_name, sizeof(si->table_name), "%s", table_name);
// Create table — PK is (timestamp, node_id), author_signature NOT NULL
{
char sql[512];
snprintf(sql, sizeof(sql),
"CREATE TABLE IF NOT EXISTS \"%s\" ("
" timestamp INTEGER NOT NULL,"
" node_id INTEGER NOT NULL,"
" id INTEGER NOT NULL,"
" chain_hash BLOB NOT NULL,"
" flags INTEGER NOT NULL DEFAULT 0,"
" data BLOB,"
" author_signature BLOB NOT NULL,"
" local_attrs TEXT DEFAULT '',"
" delivered_peers INTEGER NOT NULL DEFAULT 0,"
" delivery_chain TEXT NOT NULL DEFAULT '',"
" 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; }
snprintf(sql, sizeof(sql),
"CREATE INDEX IF NOT EXISTS \"idx_%s_ttl\""
" ON \"%s\" (node_id, timestamp)",
si->table_name, si->table_name);
rc = sqlite3_exec(db->db, sql, NULL, NULL, NULL);
if (rc != SQLITE_OK) DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "CREATE INDEX %s: %s", si->table_name, sqlite3_errmsg(db->db));
}
// Get next id
{
sqlite3_stmt* stmt;
if (si_prep(si, &stmt, "SELECT COALESCE(MAX(id),0)+1 FROM \"%s\"") == SQLITE_OK
&& sqlite3_step(stmt) == SQLITE_ROW)
si->next_id = (uint64_t)sqlite3_column_int64(stmt, 0);
else
si->next_id = 1;
sqlite3_finalize(stmt);
}
db_verify_chain(si);
DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "instance_add: tbl=%s mc=%u next_id=%llu group=%016llx", SI_TBL(si), db_count(si), (unsigned long long)si->next_id, (unsigned long long)si->group_id);
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;
}
}
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",
table_name, si->table_name, (unsigned long long)group_id,
(unsigned long long)si->next_id, peers_found, peers_synced, db_count(si));
return si;
}
// Удаляет инстанс: отменяет таймеры, освобождает пиров, сдвигает массив (таблица БД не удаляется)
void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si)
{
if (!si) return;
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "tbl=%s group=%016llx", si->table_name, (unsigned long long)si->group_id);
struct DB_SYNC* db = si->db_sync;
struct UTUN_INSTANCE* inst = db->inst;
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; }
if (si->peers) { u_free(si->peers); si->peers = NULL; si->peer_count = si->peer_capacity = 0; }
int idx = -1;
for (int i = 0; i < db->instance_count; i++) {
if (db->instances[i] == si) { idx = i; break; }
}
if (idx >= 0) {
memmove(&db->instances[idx], &db->instances[idx + 1],
(db->instance_count - idx - 1) * sizeof(*db->instances));
db->instance_count--;
}
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "instance removed tbl=%s group=%016llx", si->table_name, (unsigned long long)si->group_id);
u_free(si);
}
// Вставляет подписанную запись: проверяет подпись, вставляет в БД, рассылает PUSH всем синхронизированным пирам
int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si, const char* json_data, size_t len,
const uint8_t* sig, size_t sig_len, uint64_t ts,
const char* local_attrs)
{
if (!si || !si->enabled || !json_data || len == 0) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "insert invalid args"); return -1; }
if (!sig || sig_len != DB_SIG_SIZE) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "insert_signed: signature required (64 bytes Ed25519)"); return -1; }
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "tbl=%s id=%llu ts=%llu len=%zu", SI_TBL(si), (unsigned long long)si->next_id, (unsigned long long)ts, len);
uint64_t id = si->next_id;
uint64_t author_node_id = si->db_sync->inst->node_id;
return db_record_insert_cascade(si, id, ts, author_node_id, json_data, (uint32_t)len, sig, author_node_id, local_attrs);
}
// Возвращает количество записей в инстансе
uint32_t db_sync_count(struct DB_SYNC_INSTANCE* si)
{
if (!si || !si->enabled) return 0;
return db_count(si);
}
// Возвращает последний использованный timestamp инстанса
uint64_t db_sync_get_last_timestamp(struct DB_SYNC_INSTANCE* si)
{
return si ? si->last_timestamp_ms : 0;
}
// Возвращает следующий монотонно возрастающий timestamp (мкс → мс), использует NTP если синхронизировано
uint64_t db_sync_next_timestamp(struct DB_SYNC_INSTANCE* si)
{
if (!si) return 0;
int64_t now_us;
if (si->db_sync && si->db_sync->inst && si->db_sync->inst->ntp.synced)
now_us = ntp_time_get_us(si->db_sync->inst);
else {
struct timeval tv;
utun_gettimeofday(&tv, NULL);
now_us = (int64_t)tv.tv_sec * 1000000LL + tv.tv_usec;
}
uint64_t nu = (uint64_t)(now_us / 1000LL);
if (nu <= si->last_timestamp_ms) nu = si->last_timestamp_ms + 1;
si->last_timestamp_ms = nu;
return nu;
}
// Устанавливает колбэк, вызываемый при каждой успешной вставке записи (локальной или от пира)
void db_sync_set_insert_cb(struct DB_SYNC_INSTANCE* si, db_sync_insert_cb cb, void* arg)
{
if (!si) return;
si->on_insert = cb;
si->on_insert_arg = arg;
}
// Итерирует записи с offset/limit, передавая каждую в колбэк. Возвращает количество переданных записей
int db_sync_select(struct DB_SYNC_INSTANCE* si, uint32_t offset, uint32_t limit,
db_sync_select_cb cb, void* arg)
{
if (!si || !si->enabled || !cb) return 0;
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "tbl=%s off=%u lim=%u", SI_TBL(si), offset, limit);
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT id,timestamp,node_id,data,"
"author_signature,delivered_peers,delivery_chain"
" FROM \"%s\" ORDER BY timestamp, author_signature"
" LIMIT ? OFFSET ?") != SQLITE_OK)
return 0;
sqlite3_bind_int64(stmt, 1, limit > 0 ? (sqlite3_int64)limit : -1);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)offset);
int cnt = 0;
while (sqlite3_step(stmt) == SQLITE_ROW) {
uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0);
uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1);
uint64_t rauth = (uint64_t)sqlite3_column_int64(stmt, 2);
const char* rd = (const char*)sqlite3_column_blob(stmt, 3);
int rdl = sqlite3_column_bytes(stmt, 3);
const uint8_t* rsig = (const uint8_t*)sqlite3_column_blob(stmt, 4);
int rsl = sqlite3_column_bytes(stmt, 4); if (!rsig) rsl = 0;
int rdp = sqlite3_column_int(stmt, 5);
const char* rdc = (const char*)sqlite3_column_text(stmt, 6); if (!rdc) rdc = "";
cb(arg, rid, rts, rd, (size_t)(rdl > 0 ? rdl : 0), rauth, rsig, (size_t)rsl, rdp, rdc);
cnt++;
}
sqlite3_finalize(stmt);
return cnt;
}