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.
1855 lines
84 KiB
1855 lines
84 KiB
// db_sync.c — Distributed content-addressed table with SQLite + peer sync via etcp_api (direct P2P) |
|
|
|
#include "db_sync.h" |
|
#include "etcp_api.h" |
|
#include "etcp.h" |
|
#include "utun_instance.h" |
|
#include "topo_group.h" |
|
#include "topo_node_sqlite.h" |
|
#include "secure_channel.h" |
|
#include "../lib/debug_config.h" |
|
#include "../lib/mem.h" |
|
#include "../lib/u_async.h" |
|
#include "../lib/sha256.h" |
|
#include "../lib/sqlite3.h" |
|
#include "../lib/platform_compat.h" |
|
#include <string.h> |
|
#include <openssl/evp.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_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_hash(struct DB_SYNC* db, uint64_t node_id, uint64_t hash, 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, int allocated_buf); |
|
|
|
// ============================================================ |
|
// Structures |
|
// ============================================================ |
|
|
|
struct SI_PEER { |
|
uint64_t node_id; |
|
uint32_t synced_pos; |
|
uint32_t verified_pos; |
|
uint8_t sync_state; |
|
uint64_t sync_start_tb; |
|
}; |
|
|
|
struct DB_SYNC_INSTANCE { |
|
struct DB_SYNC* db_sync; |
|
uint64_t hash; |
|
char table_name[64]; |
|
uint64_t next_id; |
|
uint64_t last_timestamp_ms; |
|
uint8_t enabled; |
|
void* ttl_timer; |
|
struct SI_PEER* peers; |
|
int peer_count, peer_capacity; |
|
db_sync_insert_cb on_insert; |
|
void* on_insert_arg; |
|
}; |
|
|
|
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; |
|
uint8_t enabled; |
|
}; |
|
|
|
// ============================================================ |
|
// SQL helpers |
|
// ============================================================ |
|
|
|
#define SI_DB(si) ((si)->db_sync->db) |
|
#define SI_TBL(si) ((si)->table_name) |
|
#define SI_SHRT(si) ({ const char* _t = SI_TBL(si); const char* _u = strrchr(_t, '_'); _u ? _u + 1 : _t; }) |
|
|
|
static void si_sql(char* buf, size_t sz, struct DB_SYNC_INSTANCE* si, const char* fmt) |
|
{ |
|
snprintf(buf, sz, fmt, SI_TBL(si)); |
|
} |
|
|
|
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 |
|
// ============================================================ |
|
|
|
static int db_sqlite_open(struct DB_SYNC* db, const char* path) |
|
{ |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "path=%s", path); |
|
int rc = sqlite3_open(path, &db->db); |
|
if (rc != SQLITE_OK) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_DB_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_DB_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_DB_SYNC, "db_sync: SQLite opened at %s", path); |
|
return 0; |
|
} |
|
|
|
static void db_sqlite_close(struct DB_SYNC* db) |
|
{ |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, ""); |
|
if (!db->db || db->shared_db) return; |
|
sqlite3_close(db->db); |
|
db->db = NULL; |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: SQLite closed"); |
|
} |
|
|
|
// ============================================================ |
|
// SHA256 helpers |
|
// ============================================================ |
|
|
|
static void db_sha256(const uint8_t* data, size_t len, uint8_t hash[32]) |
|
{ |
|
SC_SHA256_CTX ctx; |
|
sc_sha256_init(&ctx); |
|
sc_sha256_update(&ctx, data, len); |
|
sc_sha256_final(&ctx, hash); |
|
} |
|
|
|
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; |
|
} |
|
|
|
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, ×tamp, 8); |
|
memcpy(buf + 48, &author, 8); |
|
memcpy(buf + 56, author_signature, DB_SIG_SIZE); |
|
db_sha256(buf, sizeof(buf), out); |
|
} |
|
|
|
// ============================================================ |
|
// Instance management |
|
// ============================================================ |
|
|
|
static struct DB_SYNC_INSTANCE* db_instance_find(struct DB_SYNC* db, uint64_t hash) |
|
{ |
|
for (int i = 0; i < db->instance_count; i++) { |
|
if (db->instances[i].hash == hash && db->instances[i].enabled) return &db->instances[i]; |
|
} |
|
return NULL; |
|
} |
|
|
|
static struct DB_SYNC_INSTANCE* db_instance_alloc(struct DB_SYNC* db) |
|
{ |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_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 = &db->instances[db->instance_count++]; |
|
memset(si, 0, sizeof(*si)); |
|
si->db_sync = db; |
|
si->enabled = 1; |
|
return si; |
|
} |
|
|
|
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; |
|
} |
|
|
|
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 |
|
// ============================================================ |
|
|
|
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; |
|
} |
|
|
|
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_DB_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_DB_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_DB_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; |
|
} |
|
|
|
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; |
|
} |
|
|
|
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_DB_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; |
|
} |
|
|
|
static void db_cascade_from(struct DB_SYNC_INSTANCE* si, uint32_t from_pos) |
|
{ |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "tbl=%s from=%u", SI_TBL(si), from_pos); |
|
sqlite3* db = SI_DB(si); |
|
int rc = sqlite3_exec(db, "BEGIN IMMEDIATE", NULL, NULL, NULL); |
|
if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "cascade_from BEGIN: %s", sqlite3_errmsg(db)); return; } |
|
|
|
uint8_t prev_ch[32]; memset(prev_ch, 0, 32); |
|
if (from_pos > 0) { |
|
sqlite3_stmt* stmt; |
|
if (si_prep(si, &stmt, |
|
"SELECT chain_hash FROM \"%s\"" |
|
" ORDER BY timestamp, author_signature" |
|
" LIMIT 1 OFFSET ?") == SQLITE_OK) |
|
{ |
|
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)(from_pos - 1)); |
|
if (sqlite3_step(stmt) == SQLITE_ROW) { |
|
const void* b = sqlite3_column_blob(stmt, 0); |
|
if (b && sqlite3_column_bytes(stmt, 0) >= 32) memcpy(prev_ch, b, 32); |
|
} |
|
sqlite3_finalize(stmt); |
|
} |
|
} |
|
|
|
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_DB_SYNC, "db_cascade_from SELECT prep: %s", sqlite3_errmsg(db)); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return; } |
|
sqlite3_bind_int64(sel, 1, (sqlite3_int64)from_pos); |
|
|
|
if (si_prep(si, &upd, |
|
"UPDATE \"%s\" SET chain_hash=?" |
|
" WHERE timestamp=? AND author_signature=?") != SQLITE_OK) |
|
{ DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_cascade_from UPDATE prep: %s", sqlite3_errmsg(db)); sqlite3_finalize(sel); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return; } |
|
|
|
uint8_t ch[32]; |
|
while (sqlite3_step(sel) == SQLITE_ROW) { |
|
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); |
|
} |
|
sqlite3_finalize(upd); |
|
sqlite3_finalize(sel); |
|
|
|
rc = sqlite3_exec(db, "COMMIT", NULL, NULL, NULL); |
|
if (rc != SQLITE_OK) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "cascade_from COMMIT: %s", sqlite3_errmsg(db)); |
|
} |
|
|
|
// ── Chain hash fragment (first 8 bytes) for sync protocol comparisons ── |
|
|
|
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; |
|
} |
|
|
|
// ── Ed25519 verification ── |
|
|
|
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; } |
|
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 ll_entry* e = queue_find_data_by_index(db->inst->connections, (const uint8_t*)&node_id); |
|
if (e) { |
|
struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; |
|
if (ce->conn) { |
|
memcpy(out, ce->conn->peer_ed25519_pubkey, 32); |
|
uint64_t chk; memcpy(&chk, out, 8); |
|
if (chk != 0) return 0; |
|
} |
|
} |
|
return -1; |
|
} |
|
|
|
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_DB_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_DB_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; |
|
} |
|
|
|
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, int do_cascade) |
|
{ |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "id=%llu ts=%llu author=%N len=%zu cascade=%d", (unsigned long long)id, (unsigned long long)ts, (unsigned long long)author_node_id, jlen, do_cascade); |
|
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_cascade_from(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_DB_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_DB_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_DB_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,delivered_peers,delivery_chain)" |
|
" VALUES (?,?,?,?,0,?,?,0,'')") != SQLITE_OK) |
|
{ DEBUG_ERROR(DEBUG_CATEGORY_DB_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); |
|
rc = sqlite3_step(stmt); |
|
sqlite3_finalize(stmt); |
|
if (rc != SQLITE_DONE) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "INSERT: %s", sqlite3_errmsg(db)); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return -1; } |
|
|
|
// Recompute chain hashes for all subsequent records |
|
if (do_cascade) { |
|
sqlite3_stmt* sel; |
|
if (si_prep(si, &sel, |
|
"SELECT timestamp,id,node_id,author_signature FROM \"%s\"" |
|
" WHERE timestamp>?1" |
|
" OR (timestamp=?1 AND author_signature>?2)" |
|
" ORDER BY timestamp, author_signature") != SQLITE_OK) |
|
{ DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_record_insert cascade SELECT prep: %s", sqlite3_errmsg(db)); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return -1; } |
|
sqlite3_bind_int64(sel, 1, (sqlite3_int64)ts); |
|
sqlite3_bind_blob(sel, 2, author_sig, DB_SIG_SIZE, SQLITE_STATIC); |
|
|
|
sqlite3_stmt* upd; |
|
if (si_prep(si, &upd, |
|
"UPDATE \"%s\" SET chain_hash=?" |
|
" WHERE timestamp=? AND author_signature=?") != SQLITE_OK) |
|
{ DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_record_insert cascade UPDATE prep: %s", sqlite3_errmsg(db)); sqlite3_finalize(sel); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return -1; } |
|
|
|
uint8_t rch[32]; memcpy(rch, ch, 32); |
|
while (sqlite3_step(sel) == SQLITE_ROW) { |
|
sqlite3_int64 n_ts = sqlite3_column_int64(sel, 0); |
|
sqlite3_int64 n_id = 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); |
|
uint8_t nch[32]; |
|
db_chain_hash_compute(rch, (uint64_t)n_id, (uint64_t)n_ts, (uint64_t)n_auth, sig, nch); |
|
sqlite3_reset(upd); |
|
sqlite3_bind_blob(upd, 1, nch, 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(rch, nch, 32); |
|
} |
|
sqlite3_finalize(upd); |
|
sqlite3_finalize(sel); |
|
} |
|
|
|
si->next_id = id + 1; |
|
rc = sqlite3_exec(db, "COMMIT", NULL, NULL, NULL); |
|
if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "insert COMMIT: %s", sqlite3_errmsg(db)); return -1; } |
|
return 0; |
|
} |
|
|
|
// ============================================================ |
|
// Send |
|
// ============================================================ |
|
|
|
static int db_sync_send_hash(struct DB_SYNC* db, uint64_t node_id, uint64_t hash, const uint8_t* payload, size_t plen) |
|
{ |
|
struct ll_entry* entry = queue_entry_new(0); |
|
if (!entry) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "queue_entry_new"); return -1; } |
|
uint8_t* buf = u_malloc(plen + 9); |
|
if (!buf) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync_send_hash: 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 hb = htobe64(hash); |
|
memcpy(buf + 1, &hb, 8); |
|
memcpy(buf + 9, payload, plen); |
|
entry->dgram = buf; |
|
entry->len = plen + 9; |
|
|
|
struct ETCP_CONN* conn = NULL; |
|
{ |
|
struct ll_entry* e = queue_find_data_by_index(db->inst->connections, (const uint8_t*)&node_id); |
|
if (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; if (ce->conn->links_up) conn = ce->conn; } |
|
} |
|
if (!conn) { |
|
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "sync: send DROP — no direct conn to %04llX (hash=%016llx, type=%02x)", |
|
(unsigned long long)(node_id >> 16), hash, plen > 0 ? payload[0] : 0); |
|
queue_entry_free(entry); |
|
return -1; |
|
} |
|
return etcp_send(conn, entry); |
|
} |
|
|
|
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_hash(si->db_sync, dst_node_id, si->hash, payload, len); |
|
} |
|
|
|
// ============================================================ |
|
// Delivery chain helpers (local fields) |
|
// ============================================================ |
|
|
|
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 |
|
// ============================================================ |
|
|
|
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_DB_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_DB_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; |
|
} |
|
} |
|
|
|
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_DB_SYNC, "INIT_SYNC too short %zu from %016llx", len, (unsigned long long)src); return; } |
|
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_DB_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; |
|
db_sync_send(si, src, resp, off); |
|
} |
|
|
|
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_DB_SYNC, "INIT_RESP too short %zu from %016llx", len, (unsigned long long)src); return; } |
|
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_INFO(DEBUG_CATEGORY_DB_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); |
|
|
|
if (peer_ch8 == 0 && sc == 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; |
|
int pushed = si_send_data_batch(si, src, sent, b, vp, 1); |
|
if (pushed <= 0) break; |
|
sent += (uint32_t)pushed; |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: peer empty → sent %u/%u records vp=%u", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), sent, mc, vp); |
|
return; |
|
} |
|
|
|
if (my_ch8 == peer_ch8) { |
|
uint32_t old_sp = sp ? sp->synced_pos : 0; |
|
if (!has_tail) { |
|
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); |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: MATCH+no_tail ✓ synced_pos %u→%u → SYNC_DONE", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), old_sp, tp); |
|
return; |
|
} |
|
uint32_t vp = tp; |
|
if (sp) sp->verified_pos = vp; |
|
uint8_t req[11]; |
|
req[0] = DB_MSG_SEND_DATA; uint32_t rfrom = tp + 1; memcpy(req + 1, &rfrom, 4); |
|
uint16_t z = 0; memcpy(req + 5, &z, 2); memcpy(req + 7, &vp, 4); |
|
db_sync_send(si, src, req, 11); |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: MATCH but has_tail=1 → requesting SEND_DATA from=%u vp=%u", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), rfrom, vp); |
|
return; |
|
} |
|
|
|
uint32_t ds = 0, de = tp, 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) { if (pos + 1 > ds) ds = pos + 1; } |
|
else { if (pos < de) de = pos; if (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_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: MISMATCH tp=%u ds=%u de=%u fm=%u vp=%u", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), tp, ds, de, fm, vp); |
|
|
|
{ |
|
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, 0); |
|
else { |
|
uint8_t req[11]; |
|
req[0] = DB_MSG_SEND_DATA; memcpy(req + 1, &fm, 4); uint16_t z = 0; memcpy(req + 5, &z, 2); memcpy(req + 7, &vp, 4); |
|
db_sync_send(si, src, req, 11); |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SEND_DATA from=%u vp=%u count=%u (ds=%u de=%u fm=%u)", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), fm, vp, scnt, ds, de, fm); |
|
} |
|
} |
|
|
|
// ---- SEND_DATA batch helper ---- |
|
// Wire format: [from:4][count:2][vp: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 |
|
static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_t from, uint32_t count, uint32_t vp, int allocated_buf) |
|
{ |
|
uint8_t sbuf_stack[8192]; |
|
uint8_t* buf = allocated_buf ? u_malloc(8192) : NULL; |
|
if (!buf) buf = sbuf_stack; |
|
|
|
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; |
|
|
|
sqlite3_stmt* stmt; |
|
int prep_rc = si_prep(si, &stmt, |
|
"SELECT id,timestamp,node_id,data,author_signature" |
|
" FROM \"%s\" ORDER BY timestamp, author_signature" |
|
" LIMIT ? OFFSET ?"); |
|
if (prep_rc != SQLITE_OK || !stmt) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] SEND_DATA batch: si_prep failed rc=%d", |
|
SI_SHRT(si), (unsigned long long)(dst >> 16), prep_rc); |
|
if (allocated_buf) u_free(buf); |
|
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_cascade_from(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); |
|
|
|
*rcp = rc; |
|
db_sync_send(si, dst, buf, off); |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SEND DATA: %u records from pos=%u%s", |
|
SI_SHRT(si), (unsigned long long)(dst >> 16), rc, from, bad_del > 0 ? (rc > 0 ? " (bad_sig deleted)" : " (all bad_sig deleted)") : ""); |
|
if (rc == 0 && count > 0 && bad_del == 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] SEND_DATA batch: loaded 0 records from pos=%u count=%u — SQL error or empty range", |
|
SI_SHRT(si), (unsigned long long)(dst >> 16), from, count); |
|
} |
|
if (allocated_buf) u_free(buf); |
|
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] |
|
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_DB_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_DB_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_DB_SYNC, "si_parse_record: sig_len overrun"); return -1; } |
|
*rsiglen = (int)(*(*pp)++); |
|
if (*rsiglen > 0) { |
|
if (*pp + *rsiglen > end) { DEBUG_WARN(DEBUG_CATEGORY_DB_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; |
|
} |
|
|
|
// ---- Cascade chain_hash for a specific range, then full cascade from end of range ---- |
|
static void db_cascade_range(struct DB_SYNC_INSTANCE* si, uint32_t range_start, uint32_t range_end) |
|
{ |
|
uint32_t mc = db_count(si); |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "tbl=%s start=%u end=%u mc=%u", SI_TBL(si), range_start, range_end, mc); |
|
if (range_start >= range_end || range_start >= mc) return; |
|
sqlite3* db = SI_DB(si); |
|
int rc = sqlite3_exec(db, "BEGIN IMMEDIATE", NULL, NULL, NULL); |
|
if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "cascade_range BEGIN: %s", sqlite3_errmsg(db)); return; } |
|
|
|
uint8_t prev_ch[32]; memset(prev_ch, 0, 32); |
|
if (range_start > 0) { sqlite3_stmt* s; if (si_prep(si, &s, "SELECT chain_hash FROM \"%s\" ORDER BY timestamp, author_signature LIMIT 1 OFFSET ?") == SQLITE_OK) { sqlite3_bind_int64(s, 1, (sqlite3_int64)(range_start - 1)); if (sqlite3_step(s) == SQLITE_ROW) { const void* b = sqlite3_column_blob(s, 0); if (b && sqlite3_column_bytes(s, 0) >= 32) memcpy(prev_ch, b, 32); } sqlite3_finalize(s); } } |
|
|
|
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) { sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return; } |
|
sqlite3_bind_int64(sel, 1, (sqlite3_int64)range_start); |
|
if (si_prep(si, &upd, "UPDATE \"%s\" SET chain_hash=? WHERE timestamp=? AND author_signature=?") != SQLITE_OK) { sqlite3_finalize(sel); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return; } |
|
|
|
uint8_t ch[32]; uint32_t count = 0; |
|
while (sqlite3_step(sel) == SQLITE_ROW) { |
|
sqlite3_int64 n_id = sqlite3_column_int64(sel, 0), n_ts = sqlite3_column_int64(sel, 1), 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); count++; |
|
} |
|
sqlite3_finalize(upd); sqlite3_finalize(sel); |
|
rc = sqlite3_exec(db, "COMMIT", NULL, NULL, NULL); |
|
if (rc != SQLITE_OK) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "cascade_range COMMIT: %s", sqlite3_errmsg(db)); |
|
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "cascade_range [%s] from=%u → %u records recalculated", SI_TBL(si), range_start, count); |
|
} |
|
|
|
static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) |
|
{ |
|
if (len < 10) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← SEND_DATA: truncated len=%zu (need >=10)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; } |
|
uint32_t from = *(uint32_t*)p; |
|
uint16_t count = *(uint16_t*)(p + 4); |
|
uint32_t vp = *(uint32_t*)(p + 6); |
|
|
|
struct SI_PEER* sp = si_peer_find(si, src); |
|
|
|
// count=0: peer is requesting OUR data from position "from" |
|
if (count == 0) { |
|
uint32_t mc = db_count(si); |
|
uint32_t eff_from = from; |
|
if (vp != (uint32_t)-1 && vp + 1 > from) eff_from = vp + 1; |
|
if (eff_from >= mc) return; |
|
uint32_t scnt = mc - eff_from; if (scnt > DB_SEND_DATA_MAX) scnt = DB_SEND_DATA_MAX; |
|
int sent = si_send_data_batch(si, src, eff_from, scnt, vp, 0); |
|
uint64_t nc = db_count(si); uint64_t lch8 = 0; if (nc > 0) db_chain_hash8_at(si, nc - 1, &lch8); |
|
uint8_t sd[13]; |
|
sd[0] = DB_MSG_SYNC_DONE; memcpy(sd + 1, &nc, 4); memcpy(sd + 5, &lch8, 8); |
|
db_sync_send(si, src, sd, 13); |
|
if (sp) { sp->synced_pos = nc > 0 ? nc - 1 : 0; sp->sync_state = 2; sp->sync_start_tb = 0; } |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← SEND_DATA request from=%u (eff=%u) → sent %d records mc=%u + SYNC_DONE", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), from, eff_from, sent, mc); |
|
return; |
|
} |
|
|
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "from=%016llx from=%u count=%u vp=%u", (unsigned long long)src, from, count, vp); |
|
const uint8_t* ptr = p + 10; |
|
|
|
uint32_t old_synced_pos = sp ? sp->synced_pos : 0; |
|
uint16_t received = 0, parse_fails = 0, duplicates = 0, bad_sigs = 0; |
|
uint64_t first_id = 0, last_id = 0; |
|
uint64_t recv_ts[32], recv_sig8[32]; uint16_t recv_cnt = 0; |
|
|
|
for (uint16_t i = 0; i < count && i < 32; i++) { |
|
const uint8_t* save = ptr; |
|
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) { parse_fails++; break; } |
|
int ret = db_record_insert(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, 0); |
|
if (ret >= 0) { recv_ts[recv_cnt] = rts; if (rsig && rsiglen >= 8) memcpy(&recv_sig8[recv_cnt], rsig, 8); else recv_sig8[recv_cnt] = 0; recv_cnt++; } |
|
if (ret >= 0) { if (received == 0) first_id = rid; last_id = rid; received++; } |
|
if (ret == 1) duplicates++; |
|
if (ret == -2) bad_sigs++; |
|
if (ret == 0 && si->on_insert) si->on_insert(si, rts, (const char*)rdata, rdlen, src, si->on_insert_arg); |
|
} |
|
|
|
uint16_t total_processed = count - parse_fails; |
|
|
|
if (total_processed > 0) { |
|
uint32_t rst = (from < old_synced_pos) ? from : old_synced_pos; |
|
db_cascade_range(si, rst, from + total_processed); |
|
if (sp) sp->synced_pos = from + total_processed - 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; |
|
} |
|
uint32_t nc = db_count(si); |
|
char extra_buf[64] = ""; if (duplicates) snprintf(extra_buf, sizeof(extra_buf), " (%u dup)", duplicates); |
|
if (bad_sigs) { size_t el = strlen(extra_buf); snprintf(extra_buf+el, sizeof(extra_buf)-el, "%s%u bad_sig", extra_buf[0]?", ":" (", bad_sigs); if (!extra_buf[0]) extra_buf[0]=' '; } |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← DATA: from=%u count=%u vp=%u → inserted %u [id=%llu..%llu]%s cascade=[%u..%u] synced=%u→%u", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), from, count, vp, received, (unsigned long long)first_id, (unsigned long long)last_id, |
|
extra_buf, rst, from+total_processed-1, old_synced_pos, sp?sp->synced_pos:0); |
|
|
|
int relay_count = 0; |
|
uint8_t push_buf[2560]; uint32_t push_off; |
|
sqlite3_stmt* pstmt; |
|
if (si_prep(si, &pstmt, |
|
"SELECT id,timestamp,node_id,data,author_signature" |
|
" FROM \"%s\" ORDER BY timestamp, author_signature" |
|
" LIMIT ? OFFSET ?") == SQLITE_OK) |
|
{ |
|
sqlite3_bind_int64(pstmt, 1, (sqlite3_int64)received); |
|
sqlite3_bind_int64(pstmt, 2, (sqlite3_int64)from); |
|
while (sqlite3_step(pstmt) == SQLITE_ROW) { |
|
push_off = 1; |
|
uint64_t prid = (uint64_t)sqlite3_column_int64(pstmt, 0); |
|
uint64_t prts = (uint64_t)sqlite3_column_int64(pstmt, 1); |
|
uint64_t prauth = (uint64_t)sqlite3_column_int64(pstmt, 2); |
|
const uint8_t* prd = (const uint8_t*)sqlite3_column_blob(pstmt, 3); |
|
uint32_t prdl = (uint32_t)sqlite3_column_bytes(pstmt, 3); if (!prd) prdl = 0; |
|
const uint8_t* prsig = (const uint8_t*)sqlite3_column_blob(pstmt, 4); |
|
int prsl = sqlite3_column_bytes(pstmt, 4); if (!prsig) prsl = 0; |
|
memcpy(push_buf + push_off, &prid, 8); push_off += 8; |
|
memcpy(push_buf + push_off, &prts, 8); push_off += 8; |
|
memcpy(push_buf + push_off, &prauth, 8); push_off += 8; |
|
memcpy(push_buf + push_off, &prdl, 4); push_off += 4; |
|
if (prdl > 0) { memcpy(push_buf + push_off, prd, prdl); push_off += prdl; } |
|
push_buf[push_off++] = (uint8_t)prsl; |
|
if (prsl > 0) { memcpy(push_buf + push_off, prsig, prsl); push_off += prsl; } |
|
push_buf[0] = DB_MSG_PUSH; |
|
for (int j = 0; j < si->peer_count; j++) { |
|
if (si->peers[j].sync_state >= 1 && si->peers[j].node_id != src |
|
&& si->peers[j].node_id != si->db_sync->inst->node_id) { |
|
if (db_sync_send(si, si->peers[j].node_id, push_buf, push_off) >= 0) { |
|
si_delivery_update(si, prts, prauth, si->peers[j].node_id); relay_count++; |
|
} |
|
} |
|
} |
|
} |
|
sqlite3_finalize(pstmt); |
|
} |
|
if (relay_count > 0) |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← DATA: relayed %u records to %d peers", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), received, relay_count / (int)received); |
|
} else { |
|
uint32_t nc = db_count(si); |
|
if (old_synced_pos < nc) db_cascade_range(si, old_synced_pos, nc); |
|
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← DATA: from=%u count=%u vp=%u → 0 inserted (%u dup %u bad_sig %u parse_err) cascade_from=%u", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), from, count, vp, duplicates, bad_sigs, parse_fails, old_synced_pos); |
|
} |
|
|
|
// Send back own records after vp EXCEPT those just received |
|
uint32_t mc = db_count(si); |
|
uint32_t resp_from = vp; |
|
uint32_t resp_start = (vp == (uint32_t)-1) ? 0 : (vp + 1 > from + received ? vp + 1 : from + received); |
|
int resp_count = 0; |
|
if (resp_start < mc) { |
|
uint8_t rbuf[8192]; |
|
uint32_t roff = 0; |
|
rbuf[roff++] = DB_MSG_SEND_DATA; |
|
memcpy(rbuf + roff, &resp_start, 4); roff += 4; |
|
uint16_t* rrcp = (uint16_t*)(rbuf + roff); roff += 2; |
|
memcpy(rbuf + roff, &vp, 4); roff += 4; |
|
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) |
|
{ |
|
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)(mc - resp_start)); |
|
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)resp_start); |
|
while (sqlite3_step(stmt) == SQLITE_ROW && resp_count < DB_SEND_DATA_MAX) { |
|
uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1); |
|
const uint8_t* rsig = (const uint8_t*)sqlite3_column_blob(stmt, 4); |
|
uint64_t chk_sig8 = 0; if (rsig) memcpy(&chk_sig8, rsig, 8); |
|
int is_dup = 0; |
|
for (int k = 0; k < recv_cnt; k++) { if (recv_ts[k] == rts && recv_sig8[k] == chk_sig8) { is_dup = 1; break; } } |
|
if (is_dup) continue; |
|
uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0); |
|
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; |
|
int rsl = sqlite3_column_bytes(stmt, 4); if (!rsig) rsl = 0; |
|
int rec_sz = 28 + rdl + 1 + rsl; |
|
if (roff + rec_sz > 8000) break; |
|
memcpy(rbuf + roff, &rid, 8); roff += 8; |
|
memcpy(rbuf + roff, &rts, 8); roff += 8; |
|
memcpy(rbuf + roff, &rauth, 8); roff += 8; |
|
memcpy(rbuf + roff, &rdl, 4); roff += 4; |
|
if (rdl > 0) { memcpy(rbuf + roff, rd, rdl); roff += rdl; } |
|
rbuf[roff++] = (uint8_t)rsl; |
|
if (rsl > 0) { memcpy(rbuf + roff, rsig, rsl); roff += rsl; } |
|
resp_count++; |
|
} |
|
sqlite3_finalize(stmt); |
|
} |
|
*rrcp = resp_count; |
|
if (resp_count > 0) { |
|
db_sync_send(si, src, rbuf, roff); |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → RESP DATA: %u own records after vp=%u (from=%u)", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), resp_count, vp, resp_start); |
|
} |
|
} |
|
|
|
uint64_t lch8 = 0; if (mc > 0) db_chain_hash8_at(si, mc - 1, &lch8); |
|
if (received == 0 && resp_count == 0) { |
|
if (sp) { uint32_t nsp = mc > 0 ? mc - 1 : 0; if (sp->synced_pos < nsp) sp->synced_pos = nsp; sp->sync_state = 2; sp->sync_start_tb = 0; } |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SEND_DATA: 0 new records, skip SYNC_DONE (synced=%u)", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), sp ? sp->synced_pos : 0); |
|
return; |
|
} |
|
uint8_t sd[13]; |
|
sd[0] = DB_MSG_SYNC_DONE; memcpy(sd + 1, &mc, 4); memcpy(sd + 5, &lch8, 8); |
|
db_sync_send(si, src, sd, 13); |
|
if (sp) { uint32_t nsp = mc > 0 ? mc - 1 : 0; if (sp->synced_pos < nsp) sp->synced_pos = nsp; sp->sync_state = 2; sp->sync_start_tb = 0; } |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SYNC_DONE: fc=%u ch8=%016llX synced=%u⇥2", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), mc, lch8, sp ? sp->synced_pos : 0); |
|
} |
|
|
|
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_DB_SYNC, "sync [%s:%04llX] ← SYNC_DONE: truncated len=%zu (need >=12)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; } |
|
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); |
|
uint32_t old_spos = sp ? sp->synced_pos : 0; |
|
|
|
if (mc == pc && mch8 == pch8) { |
|
if (sp) { sp->synced_pos = mc > 0 ? mc - 1 : 0; sp->sync_state = 2; } |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, |
|
"sync [%s:%04llX] ← SYNC_DONE: matched (%u=%u, %016llX=%016llX) ✓ synced_pos %u→%u", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, mch8, pch8, old_spos, mc > 0 ? mc - 1 : 0); |
|
return; |
|
} |
|
|
|
if (mc < pc) { |
|
if (sp && sp->synced_pos + 1 >= pc) { |
|
if (sp) { sp->synced_pos = pc - 1; sp->sync_state = 2; } |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, |
|
"sync [%s:%04llX] ← SYNC_DONE: count diff my=%u peer=%u but synced=%u covers tail — accept", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, sp ? sp->synced_pos : 0); |
|
return; |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, |
|
"sync [%s:%04llX] ← SYNC_DONE: count diff my=%u peer=%u — requesting tail from=%u", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, mc); |
|
uint32_t vp_send = mc > 0 ? mc - 1 : (uint32_t)-1; |
|
uint8_t req[11]; |
|
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); |
|
db_sync_send(si, src, req, 11); |
|
if (sp) sp->sync_state = 1; |
|
return; |
|
} |
|
|
|
if (mc > pc) { |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, |
|
"sync [%s:%04llX] ← SYNC_DONE: count diff my=%u peer=%u — sending tail from=%u", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, pc); |
|
uint32_t vp_send = pc > 0 ? pc - 1 : (uint32_t)-1; |
|
uint32_t scnt = mc - pc; if (scnt > DB_SEND_DATA_MAX) scnt = DB_SEND_DATA_MAX; |
|
si_send_data_batch(si, src, pc, scnt, vp_send, 0); |
|
if (sp) sp->sync_state = 1; |
|
return; |
|
} |
|
|
|
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, |
|
"sync [%s:%04llX] ← SYNC_DONE: HASH MISMATCH my=%u/%016llX vs peer=%u/%016llX (same count) — re-initiating sync", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), mc, mch8, pc, pch8); |
|
if (sp) { sp->sync_state = 1; sp->sync_start_tb = get_time_tb(); } |
|
db_sync_initiate_sync(si, src); |
|
} |
|
|
|
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_DB_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_DB_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_INFO(DEBUG_CATEGORY_DB_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); |
|
} |
|
|
|
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_DB_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_DB_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(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, 0); |
|
if (ret == 0) { |
|
uint32_t ins_pos = si_find_pos(si, rts, rsig); |
|
uint32_t total = db_count(si); |
|
db_cascade_from(si, ins_pos); |
|
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++; } |
|
} |
|
if (si->on_insert) si->on_insert(si, rts, (const char*)rdata, rdlen, src, si->on_insert_arg); |
|
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); |
|
uint8_t rbuf[2560]; |
|
rbuf[0] = DB_MSG_PUSH; |
|
memcpy(rbuf + 1, p, len); |
|
int relay_count = 0; |
|
for (int i = 0; i < si->peer_count; i++) { |
|
if (si->peers[i].sync_state >= 1 && si->peers[i].node_id != src && si->peers[i].node_id != si->db_sync->inst->node_id) { |
|
if (db_sync_send(si, si->peers[i].node_id, rbuf, len + 1) >= 0) |
|
{ si_delivery_update(si, rts, rauthor, si->peers[i].node_id); relay_count++; } |
|
} |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← PUSH: msg id=%llu author=%016llX ts=%llu → inserted at pos=%u (of %u total), cascade from pos=%u, ACK sent → relayed to %d peers", |
|
SI_SHRT(si), (unsigned long long)(src >> 16), |
|
(unsigned long long)rid, (unsigned long long)rauthor, (unsigned long long)rts, ins_pos, total, ins_pos, relay_count); |
|
if (adj_count > 0 && ins_pos < total - 1) { |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_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 |
|
// ============================================================ |
|
|
|
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 hash = 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_DB_SYNC, "recv type=%02x from=%016llx hash=%016llx len=%zu", type, (unsigned long long)src, (unsigned long long)hash, plen); |
|
|
|
struct DB_SYNC_INSTANCE* si = db_instance_find(db, hash); |
|
if (!si) { |
|
/* lazy-register: find table msg_<ch_id> with matching SHA256 */ |
|
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 h; { uint8_t sh[32]; SC_SHA256_CTX ctx; sc_sha256_init(&ctx); sc_sha256_update(&ctx, (const uint8_t*)ch_id, strlen(ch_id)); sc_sha256_final(&ctx, sh); memcpy(&h, sh, 8); } |
|
if (h == hash) { |
|
si = db_sync_instance_add(db->inst, tbl, hash, 1); |
|
if (si) { si_peer_add(si, src); } |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "lazy-register: tbl=%s ch=%s hash=%016llx", tbl, ch_id, (unsigned long long)hash); |
|
break; |
|
} |
|
} |
|
sqlite3_finalize(st); |
|
} |
|
if (!si) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, |
|
"recv msg type=0x%02x from %016llx hash=%016llx — instance NOT FOUND; sending DB_ERR_NOT_FOUND", |
|
type, (unsigned long long)src, (unsigned long long)hash); |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ instance not found for hash=%016llx", (unsigned long long)hash); |
|
uint8_t err[2]; err[0] = DB_MSG_ERROR; err[1] = DB_ERR_NOT_FOUND; |
|
db_sync_send_hash(db, src, hash, 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_hash(db, src, hash, err, 2); |
|
queue_dgram_free(entry); queue_entry_free(entry); return; |
|
} |
|
|
|
switch (type) { |
|
case DB_MSG_INIT_SYNC: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle INIT_SYNC"); db_handle_init_sync(si, src, payload, plen); break; |
|
case DB_MSG_INIT_RESP: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle INIT_RESP"); db_handle_init_resp(si, src, payload, plen); break; |
|
case DB_MSG_SEND_DATA: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle SEND_DATA"); db_handle_send_data(si, src, payload, plen); break; |
|
case DB_MSG_PUSH: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle PUSH"); db_handle_push(si, src, payload, plen); break; |
|
case DB_MSG_ACK_PUSH: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle ACK_PUSH"); db_handle_ack_push(si, src, payload, plen); break; |
|
case DB_MSG_SYNC_DONE: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle SYNC_DONE"); db_handle_sync_done(si, src, payload, plen); break; |
|
case DB_MSG_ERROR: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle ERROR"); db_handle_error(si, src, payload, plen); break; |
|
default: |
|
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "sync: unknown msg type 0x%02x from %04llX (hash=%016llx) — dropped", |
|
type, (unsigned long long)(src >> 16), hash); |
|
break; |
|
} |
|
queue_dgram_free(entry); |
|
queue_entry_free(entry); |
|
} |
|
|
|
// ============================================================ |
|
// Connection callbacks |
|
// ============================================================ |
|
|
|
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: 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_DB_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 (!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_DB_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); |
|
} |
|
|
|
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_DB_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_DB_SYNC, "sync: CONN DOWN peer=%04llX → reset sync state for %d tables", |
|
(unsigned long long)(pid >> 16), db->instance_count); |
|
} |
|
|
|
// ============================================================ |
|
// Initiate sync |
|
// ============================================================ |
|
|
|
static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t pid) |
|
{ |
|
uint32_t mc = db_count(si); |
|
struct SI_PEER* p = si_peer_find(si, pid); |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "pid=%016llx my=%u tbl=%s", (unsigned long long)pid, mc, SI_TBL(si)); |
|
if (p) { p->sync_start_tb = get_time_tb(); } |
|
|
|
uint8_t msg[5]; |
|
msg[0] = DB_MSG_INIT_SYNC; |
|
memcpy(msg + 1, &mc, 4); |
|
db_sync_send(si, pid, msg, 5); |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → sent INIT_SYNC: my_count=%u — \"let's compare chains\"", |
|
SI_SHRT(si), (unsigned long long)(pid >> 16), mc); |
|
} |
|
|
|
static void db_verify_chain(struct DB_SYNC_INSTANCE* si) |
|
{ |
|
uint32_t mc = db_count(si); |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_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_DB_SYNC, "chain_hash mismatch at pos %d/%u in %s, recalculating", |
|
bad_pos, mc, SI_TBL(si)); |
|
db_cascade_from(si, (uint32_t)bad_pos); |
|
} |
|
} |
|
|
|
int db_sync_chain_verify(struct DB_SYNC_INSTANCE* si) |
|
{ |
|
if (!si || !si->enabled) return 0; |
|
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) |
|
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; |
|
} |
|
|
|
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); |
|
} |
|
} |
|
|
|
int db_sync_last_chain_hash8(struct DB_SYNC_INSTANCE* si, uint64_t* out_hash8) |
|
{ |
|
if (!si || !out_hash8) return -1; |
|
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 |
|
// ============================================================ |
|
|
|
static void db_sync_peer_check_cb(void* arg) |
|
{ |
|
struct DB_SYNC* db = (struct DB_SYNC*)arg; |
|
if (!db || !db->enabled) return; |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "instances=%d", db->instance_count); |
|
|
|
struct TOPO_GROUP* g = topo_groups_get_default(db->inst->topo_groups); |
|
if (!g) { |
|
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"); |
|
return; |
|
} |
|
|
|
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]; |
|
if (!si->enabled) 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_DB_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_WARN(DEBUG_CATEGORY_DB_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); |
|
p->sync_state = 0; |
|
} |
|
} |
|
} |
|
} |
|
|
|
if (total_synced > 0) { |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync: PEER CHECK → %d tables, launched sync for %d: %s", |
|
db->instance_count, total_synced, launched_list); |
|
} |
|
|
|
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"); |
|
} |
|
|
|
static void db_sync_instance_ttl_cb(void* arg) |
|
{ |
|
struct DB_SYNC_INSTANCE* si = (struct DB_SYNC_INSTANCE*)arg; |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_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_DB_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 |
|
// ============================================================ |
|
|
|
int db_sync_init(struct UTUN_INSTANCE* inst) |
|
{ |
|
if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "NULL instance"); return -1; } |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_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_DB_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_DB_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_DB_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_DB_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_DB_SYNC, "db_sync: initialized (enabled=%d)", db->enabled); |
|
return 0; |
|
} |
|
|
|
void db_sync_destroy(struct UTUN_INSTANCE* inst) |
|
{ |
|
if (!inst || !inst->db_sync) return; |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_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; } |
|
|
|
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->peers) u_free(si->peers); |
|
} |
|
|
|
etcp_remove_conn_status_cbk(inst, db_sync_on_conn_status, NULL); |
|
|
|
db_sqlite_close(db); |
|
if (db->instances) u_free(db->instances); |
|
u_free(db); |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: destroyed"); |
|
} |
|
|
|
struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t hash, int auto_sync) |
|
{ |
|
if (!inst || !inst->db_sync || !table_name || !table_name[0]) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "instance_add invalid args"); return NULL; } |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "table=%s hash=%016llx", table_name, (unsigned long long)hash); |
|
struct DB_SYNC* db = inst->db_sync; |
|
if (!db->enabled || !db->db) return NULL; |
|
|
|
struct DB_SYNC_INSTANCE* si = db_instance_find(db, hash); |
|
if (si) { DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "instance already exists hash=%016llx", (unsigned long long)hash); return si; } |
|
|
|
si = db_instance_alloc(db); |
|
if (!si) return NULL; |
|
si->hash = hash; |
|
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," |
|
" 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_DB_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_DB_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_DEBUG, "instance_add: tbl=%s mc=%u next_id=%llu hash=%016llx", SI_TBL(si), db_count(si), (unsigned long long)si->next_id, (unsigned long long)si->hash); |
|
|
|
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 && ce->conn->initialized) { |
|
peers_found++; |
|
struct SI_PEER* p = si_peer_add(si, pid); |
|
if (p && p->sync_state == 0 && auto_sync) { p->sync_state = 1; db_sync_initiate_sync(si, pid); peers_synced++; } |
|
} |
|
entry = entry->next; |
|
} |
|
} |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, |
|
"instance added table=%s tbl=%s hash=%016llx next_id=%llu peers_found=%d peers_synced=%d mc=%u", |
|
table_name, si->table_name, (unsigned long long)hash, |
|
(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_DB_SYNC, "tbl=%s hash=%016llx", si->table_name, (unsigned long long)si->hash); |
|
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->peers) { u_free(si->peers); si->peers = NULL; si->peer_count = si->peer_capacity = 0; } |
|
|
|
int idx = (int)(si - db->instances); |
|
if (idx >= 0 && idx < db->instance_count) { |
|
memmove(&db->instances[idx], &db->instances[idx + 1], |
|
(db->instance_count - idx - 1) * sizeof(*db->instances)); |
|
db->instance_count--; |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "instance removed tbl=%s hash=%016llx", si->table_name, (unsigned long long)si->hash); |
|
} |
|
|
|
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) |
|
{ |
|
if (!si || !si->enabled || !json_data || len == 0) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "insert invalid args"); return -1; } |
|
if (!sig || sig_len != DB_SIG_SIZE) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "insert_signed: signature required (64 bytes Ed25519)"); return -1; } |
|
DEBUG_TRACE(DEBUG_CATEGORY_DB_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; |
|
|
|
int ret = db_record_insert(si, id, ts, author_node_id, json_data, len, sig, 1); |
|
if (ret != 0) return ret; |
|
if (si->on_insert) si->on_insert(si, ts, json_data, len, author_node_id, si->on_insert_arg); |
|
|
|
// Build PUSH: [DB_MSG_PUSH][id:8][ts:8][author:8][dlen:4][data][sig_len:1=64][sig:64] |
|
uint32_t sl = DB_SIG_SIZE; |
|
uint8_t pbuf[2560]; |
|
uint32_t off = 0; |
|
pbuf[off++] = DB_MSG_PUSH; |
|
memcpy(pbuf + off, &id, 8); off += 8; |
|
memcpy(pbuf + off, &ts, 8); off += 8; |
|
memcpy(pbuf + off, &author_node_id, 8); off += 8; |
|
memcpy(pbuf + off, &len, 4); off += 4; |
|
if (len > 0 && off + len <= sizeof(pbuf)) { memcpy(pbuf + off, json_data, len); off += len; } |
|
pbuf[off++] = (uint8_t)sl; |
|
if (off + sl <= sizeof(pbuf)) { memcpy(pbuf + off, sig, sl); off += sl; } |
|
|
|
// Push to all synced peers |
|
int push_count = 0; |
|
for (int i = 0; i < si->peer_count; i++) { |
|
if (si->peers[i].sync_state >= 1 && si->peers[i].node_id != si->db_sync->inst->node_id) { |
|
if (db_sync_send(si, si->peers[i].node_id, pbuf, off) >= 0) |
|
{ si_delivery_update(si, ts, author_node_id, si->peers[i].node_id); push_count++; } |
|
} |
|
} |
|
if (push_count > 0) |
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:ME--] → PUSH out: new msg id=%llu ts=%llu → forwarded to %d synced peers", |
|
SI_SHRT(si), (unsigned long long)id, (unsigned long long)ts, push_count); |
|
return 0; |
|
} |
|
|
|
uint32_t db_sync_count(struct DB_SYNC_INSTANCE* si) |
|
{ |
|
if (!si || !si->enabled) return 0; |
|
return db_count(si); |
|
} |
|
|
|
uint64_t db_sync_get_last_timestamp(struct DB_SYNC_INSTANCE* si) |
|
{ |
|
return si ? si->last_timestamp_ms : 0; |
|
} |
|
|
|
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; |
|
} |
|
|
|
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_DB_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; |
|
}
|
|
|