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.
 
 
 
 
 
 

2001 lines
65 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 "../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>
#define DEBUG_CATEGORY_DB_SYNC DEBUG_CATEGORY_DEBUG
// ---- 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_new_conn(struct ETCP_CONN* conn, 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);
// ============================================================
// Structures
// ============================================================
struct SI_PEER {
uint64_t node_id;
uint32_t synced_pos;
uint8_t sync_state;
};
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;
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)
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)
{
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);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: SQLite opened at %s", path);
return 0;
}
static void db_sqlite_close(struct DB_SYNC* db)
{
if (!db->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_datahash(const uint8_t* data, size_t len)
{
uint8_t hash[32];
db_sha256(data, len, hash);
uint64_t dh;
memcpy(&dh, hash, 8);
return dh;
}
static void db_chain_hash_compute(const uint8_t prev_chain_hash[32],
uint64_t id, uint64_t timestamp,
uint64_t datahash, uint8_t out[32])
{
uint8_t buf[56];
memcpy(buf, prev_chain_hash, 32);
memcpy(buf + 32, &id, 8);
memcpy(buf + 40, &timestamp, 8);
memcpy(buf + 48, &datahash, 8);
db_sha256(buf, 56, 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)
{
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_datahash(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++];
p->node_id = node_id;
p->synced_pos = 0;
p->sync_state = 0;
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, datahash"
" 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) {
sqlite3_finalize(stmt);
return -1;
}
const void* b = sqlite3_column_blob(stmt, 0);
if (!b || sqlite3_column_bytes(stmt, 0) < 32) {
sqlite3_finalize(stmt);
return -1;
}
memcpy(out, b, 32);
sqlite3_finalize(stmt);
return 0;
}
static int db_datahash_at(struct DB_SYNC_INSTANCE* si,
uint32_t pos, uint64_t* dh)
{
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT datahash FROM \"%s\""
" ORDER BY timestamp, datahash"
" LIMIT 1 OFFSET ?") != SQLITE_OK)
return -1;
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)pos);
if (sqlite3_step(stmt) != SQLITE_ROW) {
sqlite3_finalize(stmt);
return -1;
}
*dh = (uint64_t)sqlite3_column_int64(stmt, 0);
sqlite3_finalize(stmt);
return 0;
}
static uint32_t si_find_pos(struct DB_SYNC_INSTANCE* si,
uint64_t ts, uint64_t dh)
{
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT COUNT(*) FROM \"%s\""
" WHERE timestamp<?1"
" OR (timestamp=?1 AND datahash<?2)") != SQLITE_OK)
return 0;
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)dh);
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, uint64_t dh, uint8_t out[32])
{
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT chain_hash FROM \"%s\""
" WHERE timestamp<?1"
" OR (timestamp=?1 AND datahash<?2)"
" ORDER BY timestamp DESC, datahash DESC"
" LIMIT 1") != SQLITE_OK)
return -1;
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)dh);
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)
{
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, datahash"
" 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,datahash FROM \"%s\""
" ORDER BY timestamp,datahash"
" LIMIT -1 OFFSET ?") != SQLITE_OK)
{
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 datahash=?") != SQLITE_OK)
{
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_dh = sqlite3_column_int64(sel, 2);
db_chain_hash_compute(prev_ch, (uint64_t)n_id,
(uint64_t)n_ts, (uint64_t)n_dh, ch);
sqlite3_reset(upd);
sqlite3_bind_blob(upd, 1, ch, 32, SQLITE_STATIC);
sqlite3_bind_int64(upd, 2, n_ts);
sqlite3_bind_int64(upd, 3, n_dh);
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));
}
static int db_record_insert(struct DB_SYNC_INSTANCE* si,
uint64_t id, uint64_t ts, uint64_t dh,
const char* json, size_t jlen,
const uint8_t* sig, size_t sig_len,
int do_cascade)
{
sqlite3* db = SI_DB(si);
sqlite3_stmt* stmt;
int rc;
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
if (si_prep(si, &stmt,
"SELECT 1 FROM \"%s\" WHERE timestamp=? AND datahash=?")
!= SQLITE_OK)
{
sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL);
return -1;
}
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)dh);
int exists = (sqlite3_step(stmt) == SQLITE_ROW);
sqlite3_finalize(stmt);
if (exists) {
sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL);
return 1;
}
// Compute chain hash
uint8_t prev_ch[32];
db_prev_chain_hash(si, ts, dh, prev_ch);
uint8_t ch[32];
db_chain_hash_compute(prev_ch, id, ts, dh, ch);
// Insert new record
if (si_prep(si, &stmt,
"INSERT INTO \"%s\""
" (timestamp,datahash,id,chain_hash,author,flags,data,"
" author_signature,delivered_peers,delivery_chain)"
" VALUES (?,?,?,?,?,0,?,?,0,'')") != SQLITE_OK)
{
sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL);
return -1;
}
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)dh);
sqlite3_bind_int64(stmt, 3, (sqlite3_int64)id);
sqlite3_bind_blob(stmt, 4, ch, 32, SQLITE_STATIC);
sqlite3_bind_int64(stmt, 5, (sqlite3_int64)si->db_sync->inst->node_id);
if (jlen > 0)
sqlite3_bind_blob(stmt, 6, json, (int)jlen, SQLITE_STATIC);
else
sqlite3_bind_null(stmt, 6);
if (sig && sig_len > 0)
sqlite3_bind_blob(stmt, 7, sig, (int)sig_len, SQLITE_STATIC);
else
sqlite3_bind_null(stmt, 7);
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,datahash,id FROM \"%s\""
" WHERE timestamp>?1"
" OR (timestamp=?1 AND datahash>?2)"
" ORDER BY timestamp,datahash") != SQLITE_OK)
{
sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL);
return -1;
}
sqlite3_bind_int64(sel, 1, (sqlite3_int64)ts);
sqlite3_bind_int64(sel, 2, (sqlite3_int64)dh);
sqlite3_stmt* upd;
if (si_prep(si, &upd,
"UPDATE \"%s\" SET chain_hash=?"
" WHERE timestamp=? AND datahash=?") != SQLITE_OK)
{
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_dh = sqlite3_column_int64(sel, 1);
sqlite3_int64 n_id = sqlite3_column_int64(sel, 2);
uint8_t nch[32];
db_chain_hash_compute(rch, (uint64_t)n_id,
(uint64_t)n_ts, (uint64_t)n_dh, nch);
sqlite3_reset(upd);
sqlite3_bind_blob(upd, 1, nch, 32, SQLITE_STATIC);
sqlite3_bind_int64(upd, 2, n_ts);
sqlite3_bind_int64(upd, 3, n_dh);
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) {
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,
"db_sync: no direct conn to %016llx, dropping",
(unsigned long long)node_id);
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 dh, uint64_t peer_id)
{
char hex[17];
snprintf(hex, sizeof(hex), "%016llx", (unsigned long long)peer_id);
// Check if peer already in delivery_chain
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT delivery_chain FROM \"%s\""
" WHERE timestamp=? AND datahash=?") == SQLITE_OK)
{
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)dh);
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);
}
// Append peer to delivery_chain and increment delivered_peers
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 datahash = ?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)dh);
sqlite3_step(stmt);
sqlite3_finalize(stmt);
}
}
// ============================================================
// Sync protocol handlers
// ============================================================
static void db_handle_error(struct DB_SYNC* db, uint64_t src,
const uint8_t* p, size_t len)
{
if (len < 1)
return;
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC,
"db_sync: ERROR from %016llx code=%u",
(unsigned long long)src, p[0]);
(void)db;
}
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);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"INIT_SYNC from %016llx peer_count=%u my_count=%u tbl=%s",
(unsigned long long)src, pc, mc, SI_TBL(si));
uint32_t tp = (pc < mc ? pc : mc);
if (tp > 0)
tp--;
uint8_t resp[4096];
uint32_t off = 0;
resp[off++] = DB_MSG_INIT_RESP;
memcpy(resp + off, &tp, 4); off += 4;
uint64_t my_dh = 0;
if (mc > 0)
db_datahash_at(si, tp, &my_dh);
memcpy(resp + off, &my_dh, 8); off += 8;
uint32_t sp = off;
off++;
int scnt = 0;
for (uint32_t k = 0; k < 16; k++) {
uint32_t step = (uint32_t)(1u << k);
if (tp < step)
break;
uint32_t pos = tp - step;
uint64_t sdh;
if (db_datahash_at(si, pos, &sdh) != 0)
break;
if (off + 12 > sizeof(resp))
break;
memcpy(resp + off, &pos, 4); off += 4;
memcpy(resp + off, &sdh, 8); off += 8;
scnt++;
}
resp[sp] = (uint8_t)scnt;
db_sync_send(si, src, resp, off);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"INIT_RESP to %016llx tp=%u sparse=%d",
(unsigned long long)src, tp, scnt);
}
static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si,
uint64_t src,
const uint8_t* p, size_t len)
{
if (len < 13) {
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_dh;
memcpy(&peer_dh, p + 4, 8);
uint8_t sc = p[12];
uint64_t my_dh = 0;
db_datahash_at(si, tp, &my_dh);
if (peer_dh == 0 && sc == 0) {
uint32_t mc = db_count(si);
uint32_t batches = (mc + DB_SEND_DATA_MAX - 1) / DB_SEND_DATA_MAX;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"peer %016llx empty, sending all %u records in %u batches",
(unsigned long long)src, mc, batches);
uint32_t sent = 0;
while (sent < mc) {
uint32_t b = mc - sent;
if (b > DB_SEND_DATA_MAX)
b = DB_SEND_DATA_MAX;
uint8_t sdbuf[8192];
uint32_t off = 0;
sdbuf[off++] = DB_MSG_SEND_DATA;
memcpy(sdbuf + off, &sent, 4); off += 4;
uint16_t rc = 0;
uint16_t* rcp = (uint16_t*)(sdbuf + off); off += 2;
sqlite3_stmt* stmt;
si_prep(si, &stmt,
"SELECT id,timestamp,datahash,data,author_signature"
" FROM \"%s\" ORDER BY timestamp,datahash"
" LIMIT ? OFFSET ?");
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)b);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)sent);
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 rdh = (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;
int rec_sz = 28 + rdl + 1 + rsl;
if (off + rec_sz > (int)sizeof(sdbuf))
break;
memcpy(sdbuf + off, &rid, 8); off += 8;
memcpy(sdbuf + off, &rts, 8); off += 8;
memcpy(sdbuf + off, &rdh, 8); off += 8;
memcpy(sdbuf + off, &rdl, 4); off += 4;
if (rdl > 0) {
memcpy(sdbuf + off, rd, rdl); off += rdl;
}
sdbuf[off++] = (uint8_t)rsl;
if (rsl > 0) {
memcpy(sdbuf + off, rsig, rsl); off += rsl;
}
rc++;
}
sqlite3_finalize(stmt);
*rcp = rc;
sent += rc;
db_sync_send(si, src, sdbuf, off);
if (rc == 0)
break;
}
return;
}
if (my_dh == peer_dh) {
uint32_t mc = db_count(si);
struct SI_PEER* sp = si_peer_find(si, src);
uint32_t tail = tp + 1;
if (mc > tail) {
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"dh match at tp=%u my_mc=%u — sending tail [%u,%u) to %016llx",
tp, mc, tail, mc, (unsigned long long)src);
uint32_t sent = tail;
while (sent < mc) {
uint32_t b = mc - sent;
if (b > DB_SEND_DATA_MAX) b = DB_SEND_DATA_MAX;
uint8_t sdbuf[8192];
uint32_t off = 0;
sdbuf[off++] = DB_MSG_SEND_DATA;
memcpy(sdbuf + off, &sent, 4); off += 4;
uint16_t rc = 0;
uint16_t* rcp = (uint16_t*)(sdbuf + off); off += 2;
sqlite3_stmt* stmt;
si_prep(si, &stmt,
"SELECT id,timestamp,datahash,data,author_signature"
" FROM \"%s\" ORDER BY timestamp,datahash"
" LIMIT ? OFFSET ?");
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)b);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)sent);
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 rdh = (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;
int rec_sz = 28 + rdl + 1 + rsl;
if (off + rec_sz > (int)sizeof(sdbuf)) break;
memcpy(sdbuf + off, &rid, 8); off += 8;
memcpy(sdbuf + off, &rts, 8); off += 8;
memcpy(sdbuf + off, &rdh, 8); off += 8;
memcpy(sdbuf + off, &rdl, 4); off += 4;
if (rdl > 0) { memcpy(sdbuf + off, rd, rdl); off += rdl; }
sdbuf[off++] = (uint8_t)rsl;
if (rsl > 0) { memcpy(sdbuf + off, rsig, rsl); off += rsl; }
rc++;
}
sqlite3_finalize(stmt);
*rcp = rc;
sent += rc;
db_sync_send(si, src, sdbuf, off);
if (rc == 0) break;
}
}
if (sp) {
sp->synced_pos = tp;
sp->sync_state = 2;
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"sync complete with %016llx tp=%u my_mc=%u",
(unsigned long long)src, tp, mc);
return;
}
uint32_t ds = 0, de = tp;
const uint8_t* spr = p + 13;
for (int i = 0; i < sc && spr + 12 <= p + len; i++) {
uint32_t pos = *(uint32_t*)spr;
uint64_t pdh;
memcpy(&pdh, spr + 4, 8);
uint64_t mdh;
if (db_datahash_at(si, pos, &mdh) == 0) {
if (mdh == pdh) {
if (pos + 1 > ds)
ds = pos + 1;
}
else {
if (pos < de)
de = pos;
}
}
spr += 12;
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"divergence with %016llx range [%u,%u]",
(unsigned long long)src, ds, de);
if (de - ds <= 1) {
uint8_t ref[512];
uint32_t roff = 0;
ref[roff++] = DB_MSG_REFINE;
memcpy(ref + roff, &ds, 4); roff += 4;
memcpy(ref + roff, &de, 4); roff += 4;
ref[roff++] = 0;
db_sync_send(si, src, ref, roff);
}
else {
uint8_t ref[512];
uint32_t roff = 0;
ref[roff++] = DB_MSG_REFINE;
memcpy(ref + roff, &ds, 4); roff += 4;
memcpy(ref + roff, &de, 4); roff += 4;
uint32_t rng = de - ds;
uint8_t cnt = rng < DB_REFINE_HASHES
? (uint8_t)rng : DB_REFINE_HASHES;
ref[roff++] = cnt;
for (uint8_t i = 0; i < cnt; i++) {
uint32_t pos = ds + (rng * i / cnt);
uint64_t ddh;
if (db_datahash_at(si, pos, &ddh) == 0) {
memcpy(ref + roff, &pos, 4); roff += 4;
memcpy(ref + roff, &ddh, 8); roff += 8;
}
}
db_sync_send(si, src, ref, roff);
}
}
// ---- SEND_DATA batch helper ----
static void si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst,
uint32_t from, uint32_t count,
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;
sqlite3_stmt* stmt;
si_prep(si, &stmt,
"SELECT id,timestamp,datahash,data,author_signature"
" FROM \"%s\" ORDER BY timestamp,datahash"
" LIMIT ? OFFSET ?");
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)count);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)from);
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 rdh = (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;
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, &rdh, 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);
if (allocated_buf)
u_free(buf);
}
static void db_handle_refine(struct DB_SYNC_INSTANCE* si, uint64_t src,
const uint8_t* p, size_t len)
{
if (len < 9) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC,
"REFINE too short %zu", len);
return;
}
uint32_t from = *(uint32_t*)p;
uint32_t to = *(uint32_t*)(p + 4);
uint8_t hc = p[8];
if (hc == 0) {
uint32_t mc = db_count(si);
uint32_t scnt = (to - from + 1) < DB_SEND_DATA_MAX
? (to - from + 1) : DB_SEND_DATA_MAX;
if (from >= mc)
return;
si_send_data_batch(si, src, from, scnt, 1);
return;
}
uint32_t fm = to + 1;
const uint8_t* hp = p + 9;
for (uint8_t i = 0; i < hc && hp + 12 <= p + len; i++) {
uint32_t pos = *(uint32_t*)hp;
uint64_t pdh = *(uint64_t*)(hp + 4);
uint64_t mdh;
if (db_datahash_at(si, pos, &mdh) == 0 && mdh == pdh) {
if (pos >= from && pos + 1 < fm)
fm = pos + 1;
}
else {
if (pos < fm)
fm = pos;
}
hp += 12;
}
uint32_t mc = db_count(si);
uint32_t scnt = 4;
if (fm > to || fm >= mc)
scnt = 0;
si_send_data_batch(si, src, fm, scnt, 0);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"SEND_DATA to %016llx from=%u count=%u",
(unsigned long long)src, fm, scnt);
}
// ---- Parse one record from SEND_DATA/PUSH wire format ----
static int si_parse_record(const uint8_t** pp, const uint8_t* end,
uint64_t* rid, uint64_t* rts, uint64_t* rdh,
uint32_t* rdlen, const uint8_t** rdata,
const uint8_t** rsig, int* rsiglen)
{
if (*pp + 28 > end)
return -1;
*rid = *(uint64_t*)*pp; *pp += 8;
*rts = *(uint64_t*)*pp; *pp += 8;
*rdh = *(uint64_t*)*pp; *pp += 8;
*rdlen = *(uint32_t*)*pp; *pp += 4;
if (*pp + *rdlen > end)
return -1;
*rdata = (*rdlen > 0) ? *pp : NULL;
*pp += *rdlen;
if (*pp + 1 > end)
return -1;
*rsiglen = (int)(*(*pp)++);
if (*rsiglen > 0) {
if (*pp + *rsiglen > end)
return -1;
*rsig = *pp;
*pp += *rsiglen;
}
else {
*rsig = NULL;
}
return 0;
}
static void db_handle_send_data(struct DB_SYNC_INSTANCE* si,
uint64_t src, const uint8_t* p, size_t len)
{
if (len < 6)
return;
uint32_t from = *(uint32_t*)p;
uint16_t count = *(uint16_t*)(p + 4);
const uint8_t* ptr = p + 6;
struct SI_PEER* sp = si_peer_find(si, src);
uint32_t fix_from = sp ? sp->synced_pos : 0;
uint16_t received = 0;
for (uint16_t i = 0; i < count; i++) {
uint64_t rid, rts, rdh;
uint32_t rdlen;
const uint8_t* rdata, *rsig;
int rsiglen;
if (si_parse_record(&ptr, p + len,
&rid, &rts, &rdh, &rdlen,
&rdata, &rsig, &rsiglen) != 0)
break;
int ret = db_record_insert(si, rid, rts, rdh,
(const char*)rdata, rdlen,
rsig, rsiglen, 0);
if (ret >= 0)
received++;
if (ret == 0 && si->on_insert)
si->on_insert(si, (const char*)rdata, rdlen,
src, si->on_insert_arg);
}
db_cascade_from(si, fix_from);
if (sp && received > 0)
sp->synced_pos = from + received - 1;
uint32_t mc = db_count(si);
uint32_t pk = from + received;
if (mc > pk && sp && sp->sync_state == 1) {
uint32_t scnt = mc - pk;
if (scnt > DB_SEND_DATA_MAX)
scnt = DB_SEND_DATA_MAX;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"send_data continue to %016llx from=%u cnt=%u mc=%u",
(unsigned long long)src, pk, scnt, mc);
si_send_data_batch(si, src, pk, scnt, 1);
}
else if (mc > pk) {
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"send_data stopping to %016llx — mc=%u pk=%u state=%d",
(unsigned long long)src, mc, pk, sp ? sp->sync_state : -1);
}
uint32_t nc = db_count(si);
uint64_t ldh = 0;
if (nc > 0)
db_datahash_at(si, nc - 1, &ldh);
uint8_t sd[13];
sd[0] = DB_MSG_SYNC_DONE;
memcpy(sd + 1, &nc, 4);
memcpy(sd + 5, &ldh, 8);
db_sync_send(si, src, sd, 13);
if (sp)
sp->sync_state = 2;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"received %u/%u records from %016llx range=[%u,%u] → SYNC_DONE nc=%u dh=%016llx",
received, count, (unsigned long long)src, from, from + received - 1,
nc, (unsigned long long)ldh);
}
static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si,
uint64_t src, const uint8_t* p, size_t len)
{
if (len < 12)
return;
uint32_t pc = *(uint32_t*)p;
uint64_t pdh;
memcpy(&pdh, p + 4, 8);
uint32_t mc = db_count(si);
uint64_t mdh = 0;
if (mc > 0)
db_datahash_at(si, mc - 1, &mdh);
if (mc != pc || mdh != pdh) {
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"SYNC_DONE mismatch with %016llx my_count=%u peer_count=%u my_dh=%016llx peer_dh=%016llx — re-initiating",
(unsigned long long)src, mc, pc,
(unsigned long long)mdh, (unsigned long long)pdh);
db_sync_initiate_sync(si, src);
return;
}
struct SI_PEER* sp = si_peer_find(si, src);
if (sp) {
sp->synced_pos = (mc < pc ? mc : pc) > 0 ? (mc < pc ? mc : pc) - 1 : 0;
sp->sync_state = 2;
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"sync confirmed with %016llx count=%u dh=%016llx",
(unsigned long long)src, mc, (unsigned long long)mdh);
}
static void db_handle_ack_push(struct DB_SYNC_INSTANCE* si,
uint64_t src, const uint8_t* p, size_t len)
{
if (len < 16)
return;
uint64_t dh = *(uint64_t*)p;
uint64_t ts = *(uint64_t*)(p + 8);
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"UPDATE \"%s\" SET flags = flags | 1"
" WHERE timestamp=? AND datahash=?"
" AND author=? AND (flags & 1) = 0") == SQLITE_OK)
{
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)dh);
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,
"ACK_PUSH mark sent dh=%016llx ts=%llu from %016llx",
(unsigned long long)dh,
(unsigned long long)ts,
(unsigned long long)src);
}
si_delivery_update(si, ts, dh, src);
}
static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src,
const uint8_t* p, size_t len)
{
if (len < 28)
return;
const uint8_t* ptr = p;
uint64_t rid, rts, rdh;
uint32_t rdlen;
const uint8_t* rdata, *rsig;
int rsiglen;
if (si_parse_record(&ptr, p + len,
&rid, &rts, &rdh, &rdlen,
&rdata, &rsig, &rsiglen) != 0)
return;
int ret = db_record_insert(si, rid, rts, rdh,
(const char*)rdata, rdlen,
rsig, rsiglen, 0);
if (ret == 0) {
uint32_t ins_pos = si_find_pos(si, rts, rdh);
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;
}
if (si->on_insert)
si->on_insert(si, (const char*)rdata, rdlen,
src, si->on_insert_arg);
uint8_t ack[17];
ack[0] = DB_MSG_ACK_PUSH;
memcpy(ack + 1, &rdh, 8);
memcpy(ack + 9, &rts, 8);
db_sync_send(si, src, ack, 17);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"PUSH inserted from %016llx id=%llu dh=%016llx",
(unsigned long long)src,
(unsigned long long)rid,
(unsigned long long)rdh);
}
}
// ============================================================
// 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;
struct DB_SYNC_INSTANCE* si = db_instance_find(db, hash);
if (!si) {
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"recv msg type=0x%02x from %016llx hash=%016llx — instance NOT FOUND (not registered yet?); sending DB_ERR_NOT_FOUND",
type, (unsigned long long)src, (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:
db_handle_init_sync(si, src, payload, plen);
break;
case DB_MSG_INIT_RESP:
db_handle_init_resp(si, src, payload, plen);
break;
case DB_MSG_REFINE:
db_handle_refine(si, src, payload, plen);
break;
case DB_MSG_SEND_DATA:
db_handle_send_data(si, src, payload, plen);
break;
case DB_MSG_PUSH:
db_handle_push(si, src, payload, plen);
break;
case DB_MSG_ACK_PUSH:
db_handle_ack_push(si, src, payload, plen);
break;
case DB_MSG_SYNC_DONE:
db_handle_sync_done(si, src, payload, plen);
break;
case DB_MSG_ERROR:
db_handle_error(db, src, payload, plen);
break;
default:
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC,
"unknown msg type 0x%02x from %016llx",
type, (unsigned long long)src);
break;
}
queue_dgram_free(entry);
queue_entry_free(entry);
}
// ============================================================
// Connection callbacks
// ============================================================
static void db_sync_on_new_conn(struct ETCP_CONN* conn, void* arg)
{
(void)arg;
if (!conn || !conn->instance || !conn->instance->db_sync)
return;
etcp_conn_add_up_cbk(conn, db_sync_on_conn_up, NULL);
etcp_conn_add_down_cbk(conn, db_sync_on_conn_down, NULL);
}
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;
db->last_connected_tb = get_time_tb();
int synced = 0, skipped_state = 0, skipped_not_ready = 0;
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) { skipped_state++; continue; }
if (!conn->initialized || !conn->links_up) { skipped_not_ready++; continue; }
p->sync_state = 1;
db_sync_initiate_sync(si, pid);
synced++;
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"conn_up peer=%016llx init=%d links=%d instances=%d synced=%d skipped_state=%d skipped_not_ready=%d",
(unsigned long long)pid, conn->initialized, conn->links_up,
db->instance_count, synced, skipped_state, skipped_not_ready);
}
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;
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,
"peer down %016llx", (unsigned long long)pid);
}
// ============================================================
// Initiate sync
// ============================================================
static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si,
uint64_t pid)
{
uint32_t mc = db_count(si);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"INIT_SYNC → %016llx my=%u tbl=%s",
(unsigned long long)pid, mc, SI_TBL(si));
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,
"INIT_SYNC -> %016llx my=%u",
(unsigned long long)pid, mc);
}
static void db_verify_chain(struct DB_SYNC_INSTANCE* si)
{
uint32_t mc = db_count(si);
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,datahash,chain_hash FROM \"%s\""
" ORDER BY timestamp,datahash") != 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 rdh = (uint64_t)sqlite3_column_int64(stmt, 2);
const void* b = sqlite3_column_blob(stmt, 3);
if (b && sqlite3_column_bytes(stmt, 3) >= 32)
memcpy(stored_ch, b, 32);
else
memset(stored_ch, 0, 32);
db_chain_hash_compute(prev_ch, rid, rts, rdh, 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);
}
}
// ============================================================
// Timers
// ============================================================
static void db_sync_peer_check_cb(void* arg)
{
struct DB_SYNC* db = (struct DB_SYNC*)arg;
if (!db || !db->enabled)
return;
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;
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++;
}
}
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) {
best->sync_state = 1;
db_sync_initiate_sync(si, best->node_id);
total_synced++;
} else {
total_skipped++;
}
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"peer_check: instances=%d synced=%d skipped=%d",
db->instance_count, total_synced, total_skipped);
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;
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 author=? 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;
}
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;
char sp[512];
if (dp[0])
snprintf(sp, sizeof(sp), "%s/sync", 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_set_new_conn_cbk(inst, db_sync_on_new_conn, NULL);
// Attach to existing connections
{
struct ll_entry* entry = inst->connections->head;
while (entry) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data;
if (ce && ce->conn) {
etcp_conn_add_up_cbk(ce->conn, db_sync_on_conn_up, NULL);
etcp_conn_add_down_cbk(ce->conn, db_sync_on_conn_down, NULL);
}
entry = entry->next;
}
}
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;
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);
}
// Detach from existing connections
{
struct ll_entry* entry = inst->connections->head;
while (entry) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data;
if (ce && ce->conn) {
etcp_conn_remove_up_cbk(ce->conn, db_sync_on_conn_up, NULL);
etcp_conn_remove_down_cbk(ce->conn, db_sync_on_conn_down, NULL);
}
entry = entry->next;
}
}
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* name,
uint64_t id)
{
if (!inst || !inst->db_sync || !name || !name[0]) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC,
"instance_add invalid args");
return NULL;
}
struct DB_SYNC* db = inst->db_sync;
if (!db->enabled || !db->db)
return NULL;
if (!si_name_valid(name)) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC,
"invalid instance name '%s'", name);
return NULL;
}
uint64_t hash = db_instance_hash_compute(name, id);
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),
"db_sync_%s_%llx", name, (unsigned long long)id);
// Create table
{
char sql[512];
snprintf(sql, sizeof(sql),
"CREATE TABLE IF NOT EXISTS \"%s\" ("
" timestamp INTEGER NOT NULL,"
" datahash INTEGER NOT NULL,"
" id INTEGER NOT NULL,"
" chain_hash BLOB NOT NULL,"
" author INTEGER NOT NULL,"
" flags INTEGER NOT NULL DEFAULT 0,"
" data BLOB,"
" author_signature BLOB,"
" delivered_peers INTEGER NOT NULL DEFAULT 0,"
" delivery_chain TEXT NOT NULL DEFAULT '',"
" PRIMARY KEY (timestamp, datahash))",
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\" (author, 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);
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) {
p->sync_state = 1;
db_sync_initiate_sync(si, pid);
peers_synced++;
}
}
entry = entry->next;
}
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"instance added name=%s id=%llx tbl=%s hash=%016llx next_id=%llu peers_found=%d peers_synced=%d mc=%u",
name, (unsigned long long)id, 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;
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_len(struct DB_SYNC_INSTANCE* si,
const char* json_data, size_t len);
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);
int db_sync_insert_len(struct DB_SYNC_INSTANCE* si,
const char* json_data, size_t len)
{
return db_sync_insert_signed(si, json_data, len, NULL, 0);
}
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)
{
if (!si || !si->enabled || !json_data || len == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "insert invalid args");
return -1;
}
struct timeval tv;
utun_gettimeofday(&tv, NULL);
uint64_t nu = (uint64_t)tv.tv_sec * 1000ULL
+ (uint64_t)tv.tv_usec / 1000ULL;
if (nu <= si->last_timestamp_ms)
nu = si->last_timestamp_ms + 1;
si->last_timestamp_ms = nu;
uint64_t dh = db_datahash((const uint8_t*)json_data, len);
uint64_t ts = nu;
uint64_t id = si->next_id;
int ret = db_record_insert(si, id, ts, dh,
json_data, len, sig, sig_len, 1);
if (ret != 0)
return ret;
if (si->on_insert)
si->on_insert(si, json_data, len,
si->db_sync->inst->node_id,
si->on_insert_arg);
// Build PUSH payload: [DB_MSG_PUSH][id:8][ts:8][dh:8][dlen:4][data][sig_len:1][sig]
uint32_t sl = (sig && sig_len > 0) ? (uint32_t)sig_len : 0;
uint8_t pbuf[2304];
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, &dh, 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 (sl > 0 && off + sl <= sizeof(pbuf)) {
memcpy(pbuf + off, sig, sl); off += sl;
}
// Push to all synced peers
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, dh, si->peers[i].node_id);
}
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"insert id=%llu dh=%016llx ts=%llu len=%zu sig=%u",
(unsigned long long)id, (unsigned long long)dh,
(unsigned long long)ts, len, sl);
return 0;
}
int db_sync_insert(struct DB_SYNC_INSTANCE* si, const char* d)
{
if (!d)
return -1;
return db_sync_insert_len(si, d, strlen(d));
}
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;
}
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;
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT id,timestamp,author,data,"
"author_signature,delivered_peers,delivery_chain"
" FROM \"%s\" ORDER BY timestamp,datahash"
" 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;
}