From e86944d661547ba124d253a59623428364a24626 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Wed, 15 Jul 2026 16:32:36 +0300 Subject: [PATCH] =?UTF-8?q?style:=20reformat=20db=5Fsync.c=20=E2=80=94=20b?= =?UTF-8?q?reak=20500-char=20lines=20into=20readable=20~130-char=20max?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/db_sync.c | 1803 +++++++++++++++++++++++++++++++++++++++++++------ 1 file changed, 1601 insertions(+), 202 deletions(-) diff --git a/src/db_sync.c b/src/db_sync.c index 5b4927c0..1d30e835 100644 --- a/src/db_sync.c +++ b/src/db_sync.c @@ -25,9 +25,13 @@ 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 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 @@ -48,7 +52,8 @@ struct DB_SYNC_INSTANCE { void* ttl_timer; struct SI_PEER* peers; int peer_count, peer_capacity; - db_sync_insert_cb on_insert; void* on_insert_arg; + db_sync_insert_cb on_insert; + void* on_insert_arg; }; struct DB_SYNC { @@ -65,97 +70,395 @@ struct DB_SYNC { // SQL helpers // ============================================================ -#define SI_DB(si) ((si)->db_sync->db) -#define SI_TBL(si) ((si)->table_name) +#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); } +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) { +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; } + 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); } + if (rc != SQLITE_OK) { + DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: WAL pragma: %s", err); + sqlite3_free(err); + } 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"); } +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, ×tamp, 8); memcpy(buf + 48, &datahash, 8); db_sha256(buf, 56, out); +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, ×tamp, 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; iinstance_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 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); } -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; ipeer_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=0; return p; +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 = 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 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=32) memcpy(out,b,32); else memset(out,0,32); } else memset(out,0,32); sqlite3_finalize(stmt); return 0; } +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_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) { - 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; } - 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; } +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; +} - 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); - 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; } +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; +} - { 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),n_dh=sqlite3_column_int64(sel,1),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); } +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= 32) + memcpy(out, b, 32); + else + memset(out, 0, 32); + } + else { + memset(out, 0, 32); + } + sqlite3_finalize(stmt); + return 0; +} - 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; } +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) +{ + 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 + 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; } @@ -163,221 +466,1317 @@ static int db_record_insert(struct DB_SYNC_INSTANCE* si, uint64_t id, uint64_t t // 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 = db->inst->connections; while (conn) { if (conn->peer_node_id==node_id && conn->links_up>0) break; conn=conn->next; } - 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; } +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 = db->inst->connections; + while (conn) { + if (conn->peer_node_id == node_id && conn->links_up > 0) + break; + conn = conn->next; + } + 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); } + +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 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 +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,"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); } + 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<36) { 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; const uint8_t* plc = p+4; uint32_t mc = db_count(si); uint32_t tp = (pc0) tp--; - uint8_t resp[4096]; uint32_t off=0; resp[off++]=DB_MSG_INIT_RESP; memcpy(resp+off,&tp,4); off+=4; - uint8_t my_ch[32]; if (mc==0) memset(my_ch,0,32); else if (db_chain_hash_at(si,tp,my_ch)!=0) memset(my_ch,0,32); memcpy(resp+off,my_ch,32); off+=32; - if (memcmp(my_ch,plc,32)==0) { resp[off++]=0; db_sync_send(si,src,resp,off); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,"INIT_RESP complete to %016llx tp=%u",(unsigned long long)src,tp); return; } - uint32_t sp=off; off++; int scnt=0; - for (uint32_t k=0;k<16;k++) { uint32_t step=(uint32_t)(1u<sizeof(resp)) break; memcpy(resp+off,&pos,4); off+=4; memcpy(resp+off,my_ch,32); off+=32; 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<37) { 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; const uint8_t* pc=p+4; uint8_t sc=p[36]; uint8_t my_ch[32]; if (db_chain_hash_at(si,tp,my_ch)!=0) memset(my_ch,0,32); - if (sc==0 && memcmp(my_ch,pc,32)==0) { struct SI_PEER* sp=si_peer_find(si,src); if (sp) sp->synced=2; DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,"sync complete with %016llx",(unsigned long long)src); return; } - { uint8_t z[32]; memset(z,0,32); - if (memcmp(pc,z,32)==0 && sc==0) { uint32_t mc=db_count(si); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,"peer %016llx empty, sending all %u",(unsigned long long)src,mc); uint32_t sent=0; - while (sentDB_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),rts=(uint64_t)sqlite3_column_int64(stmt,1),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>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 (memcmp(my_ch,pc,32)==0) { struct SI_PEER* sp=si_peer_find(si,src); if (sp) sp->synced=2; DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,"chain_hash match with %016llx at %u",(unsigned long long)src,tp); return; } - uint32_t ds=0,de=tp; const uint8_t* spr=p+37; - for (int i=0; ids) ds=pos+1; } else { if (pos 0) + tp--; + + uint8_t resp[4096]; + uint32_t off = 0; + resp[off++] = DB_MSG_INIT_RESP; + memcpy(resp + off, &tp, 4); off += 4; + + uint8_t my_ch[32]; + if (mc == 0) + memset(my_ch, 0, 32); + else if (db_chain_hash_at(si, tp, my_ch) != 0) + memset(my_ch, 0, 32); + memcpy(resp + off, my_ch, 32); off += 32; + + // Exact match — sync complete + if (memcmp(my_ch, plc, 32) == 0) { + resp[off++] = 0; + db_sync_send(si, src, resp, off); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + "INIT_RESP complete to %016llx tp=%u", + (unsigned long long)src, tp); + return; + } + + // Sparse hashes — binary search + 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; + if (db_chain_hash_at(si, pos, my_ch) != 0) + break; + if (off + 36 > sizeof(resp)) + break; + memcpy(resp + off, &pos, 4); off += 4; + memcpy(resp + off, my_ch, 32); off += 32; + 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 < 37) { + 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; + const uint8_t* pc = p + 4; + uint8_t sc = p[36]; + + uint8_t my_ch[32]; + if (db_chain_hash_at(si, tp, my_ch) != 0) + memset(my_ch, 0, 32); + + // Exact match at tp — sync complete + if (sc == 0 && memcmp(my_ch, pc, 32) == 0) { + struct SI_PEER* sp = si_peer_find(si, src); + if (sp) + sp->synced = 2; + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + "sync complete with %016llx", + (unsigned long long)src); + return; + } + + // Peer has empty DB — send all records + { + uint8_t z[32]; + memset(z, 0, 32); + if (memcmp(pc, z, 32) == 0 && sc == 0) { + uint32_t mc = db_count(si); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + "peer %016llx empty, sending all %u", + (unsigned long long)src, mc); + 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; + } + } + + // Prefix match — sync confirmed + if (memcmp(my_ch, pc, 32) == 0) { + struct SI_PEER* sp = si_peer_find(si, src); + if (sp) + sp->synced = 2; + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + "chain_hash match with %016llx at %u", + (unsigned long long)src, tp); + return; + } + + // Binary search for divergence point + uint32_t ds = 0, de = tp; + const uint8_t* spr = p + 37; + for (int i = 0; i < sc && spr + 36 <= p + len; i++) { + uint32_t pos = *(uint32_t*)spr; + const uint8_t* pch = spr + 4; + uint8_t mch[32]; + if (db_chain_hash_at(si, pos, mch) == 0) { + if (memcmp(mch, pch, 32) == 0) { + if (pos + 1 > ds) + ds = pos + 1; + } + else { + if (pos < de) + de = pos; + } + } + spr += 36; + } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + "divergence with %016llx range [%u,%u]", + (unsigned long long)src, ds, de); + + if (de - ds <= 1) { + // Small divergence — REFINE with 0 hashes + 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 { + // Sparse REFINE + 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* buf = allocated_buf ? u_malloc(8192) : NULL; if (!buf) { uint8_t sbuf[8192]; buf = sbuf; } - uint32_t off=0; buf[off++]=DB_MSG_SEND_DATA; memcpy(buf+off,&from,4); off+=4; uint16_t rc=0; uint16_t* rcp=(uint16_t*)(buf+off); off+=2; - 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),rts=(uint64_t)sqlite3_column_int64(stmt,1),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 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, to=*(uint32_t*)(p+4); uint8_t hc=p[8]; - if (hc==0) { uint32_t mc=db_count(si), scnt=(to-from+1)=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=from&&pos+1to||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); +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; +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; - for (uint16_t i=0;ion_insert) si->on_insert(si,(const char*)rdata,rdlen,src,si->on_insert_arg); } - uint32_t mc=db_count(si); struct SI_PEER* sp=si_peer_find(si,src); uint32_t pk=from+count; - if (mc>pk && sp && sp->synced==1) { uint32_t scnt=mc-pk; if (scnt>DB_SEND_DATA_MAX) scnt=DB_SEND_DATA_MAX; si_send_data_batch(si,src,pk,scnt,1); } - uint32_t nc=db_count(si); uint8_t lch[32]; if (nc>0&&db_chain_hash_at(si,nc-1,lch)==0) { uint8_t sd[37]; sd[0]=DB_MSG_SYNC_DONE; memcpy(sd+1,&nc,4); memcpy(sd+5,lch,32); db_sync_send(si,src,sd,37); } - if (sp) sp->synced=2; DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,"received %u records from %016llx",count,(unsigned long long)src); +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; + + 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); + if (ret == 0 && si->on_insert) + si->on_insert(si, (const char*)rdata, rdlen, + src, si->on_insert_arg); + } + + uint32_t mc = db_count(si); + struct SI_PEER* sp = si_peer_find(si, src); + uint32_t pk = from + count; + + // Request next batch if more records expected + if (mc > pk && sp && sp->synced == 1) { + uint32_t scnt = mc - pk; + if (scnt > DB_SEND_DATA_MAX) + scnt = DB_SEND_DATA_MAX; + si_send_data_batch(si, src, pk, scnt, 1); + } + + // Send SYNC_DONE when we've received all + uint32_t nc = db_count(si); + uint8_t lch[32]; + if (nc > 0 && db_chain_hash_at(si, nc - 1, lch) == 0) { + uint8_t sd[37]; + sd[0] = DB_MSG_SYNC_DONE; + memcpy(sd + 1, &nc, 4); + memcpy(sd + 5, lch, 32); + db_sync_send(si, src, sd, 37); + } + if (sp) + sp->synced = 2; + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + "received %u records from %016llx", + count, (unsigned long long)src); } -static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { - if (len<36) return; uint32_t pc=*(uint32_t*)p; const uint8_t* pch=p+4; uint32_t mc=db_count(si); uint8_t mch[32]; if (mc>0) db_chain_hash_at(si,mc-1,mch); else memset(mch,0,32); - if (mc!=pc||memcmp(mch,pch,32)!=0) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC,"SYNC_DONE mismatch %016llx my=%u peer=%u",(unsigned long long)src,mc,pc); db_sync_initiate_sync(si,src); return; } - struct SI_PEER* sp=si_peer_find(si,src); if (sp) sp->synced=2; DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,"sync confirmed with %016llx count=%u",(unsigned long long)src,mc); +static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, + uint64_t src, const uint8_t* p, size_t len) +{ + if (len < 36) + return; + uint32_t pc = *(uint32_t*)p; + const uint8_t* pch = p + 4; + uint32_t mc = db_count(si); + uint8_t mch[32]; + if (mc > 0) + db_chain_hash_at(si, mc - 1, mch); + else + memset(mch, 0, 32); + + if (mc != pc || memcmp(mch, pch, 32) != 0) { + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, + "SYNC_DONE mismatch %016llx my=%u peer=%u", + (unsigned long long)src, mc, pc); + db_sync_initiate_sync(si, src); + return; + } + struct SI_PEER* sp = si_peer_find(si, src); + if (sp) + sp->synced = 2; + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + "sync confirmed with %016llx count=%u", + (unsigned long long)src, mc); } -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, 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); } } +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); - if (ret==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); } +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); + if (ret == 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, 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; +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) { 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); + if (!si) { + 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(); for (int i=0;iinstance_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&&p->synced==0) { p->synced=1; db_sync_initiate_sync(si,pid); } } DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,"peer up %016llx",(unsigned long long)pid); } -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;iinstance_count;i++) { struct SI_PEER* p=si_peer_find(&db->instances[i],pid); if (p) p->synced=0; } int any=0; for (int i=0;iinstance_count;i++) for (int j=0;jinstances[i].peer_count;j++) if (db->instances[i].peers[j].synced>=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); } +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(); + 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 && p->synced == 0) { + p->synced = 1; + db_sync_initiate_sync(si, pid); + } + } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + "peer up %016llx", (unsigned long long)pid); +} + +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->synced = 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].synced >= 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,"initiate_sync to %016llx my=%u",(unsigned long long)pid,mc); uint8_t msg[37]; msg[0]=DB_MSG_INIT_SYNC; memcpy(msg+1,&mc,4); if (mc>0) { if (db_chain_hash_at(si,mc-1,msg+5)!=0) memset(msg+5,0,32); } else memset(msg+5,0,32); db_sync_send(si,pid,msg,37); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,"INIT_SYNC -> %016llx my=%u",(unsigned long long)pid,mc); } +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, + "initiate_sync to %016llx my=%u", + (unsigned long long)pid, mc); + + uint8_t msg[37]; + msg[0] = DB_MSG_INIT_SYNC; + memcpy(msg + 1, &mc, 4); + if (mc > 0) { + if (db_chain_hash_at(si, mc - 1, msg + 5) != 0) + memset(msg + 5, 0, 32); + } + else { + memset(msg + 5, 0, 32); + } + db_sync_send(si, pid, msg, 37); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + "INIT_SYNC -> %016llx my=%u", + (unsigned long long)pid, mc); +} // ============================================================ // 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; } for (int i=0;iinstance_count;i++) { struct DB_SYNC_INSTANCE* si=&db->instances[i]; if (!si->enabled) continue; 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) { uint64_t pid=item->conn->peer_node_id; if (pid!=db->inst->node_id) { struct SI_PEER* p=si_peer_add(si,pid); if (p&&p->synced==0) { p->synced=1; db_sync_initiate_sync(si,pid); } } } e=e->next; } } 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; if (!si||!si->enabled||!si->db_sync||!si->db_sync->db) { if (si) si->ttl_timer=uasync_set_timeout(si->db_sync->inst->ua,DB_SYNC_TTL_INTERVAL*10000u,si,db_sync_instance_ttl_cb,"db_sync_ttl"); return; } uint64_t my_id=si->db_sync->inst->node_id,nu=get_time_us(),ttl=(uint64_t)(si->db_sync->inst->config->global.db_sync_ttl)*1000000uLL,cut=nu-ttl; sqlite3_stmt* stmt; if (si_prep(si,&stmt,"DELETE FROM \"%s\" WHERE author=? AND (flags & 1)=0 AND timestampttl_timer=uasync_set_timeout(si->db_sync->inst->ua,DB_SYNC_TTL_INTERVAL*10000u,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,DB_SYNC_TTL_INTERVAL*10000u,si,db_sync_instance_ttl_cb,"db_sync_ttl"); } +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; + } + + for (int i = 0; i < db->instance_count; i++) { + struct DB_SYNC_INSTANCE* si = &db->instances[i]; + if (!si->enabled) + continue; + + 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) + { + uint64_t pid = item->conn->peer_node_id; + if (pid != db->inst->node_id) { + struct SI_PEER* p = si_peer_add(si, pid); + if (p && p->synced == 0) { + p->synced = 1; + db_sync_initiate_sync(si, pid); + } + } + } + e = e->next; + } + } + + 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 timestampttl_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); { struct ETCP_CONN* conn=inst->connections; while (conn) { etcp_conn_add_up_cbk(conn,db_sync_on_conn_up,NULL); etcp_conn_add_down_cbk(conn,db_sync_on_conn_down,NULL); conn=conn->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;iinstance_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); } { struct ETCP_CONN* conn=inst->connections; while (conn) { etcp_conn_remove_up_cbk(conn,db_sync_on_conn_up,NULL); etcp_conn_remove_down_cbk(conn,db_sync_on_conn_down,NULL); conn=conn->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); - { 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)); } - { 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); } - si->ttl_timer=uasync_set_timeout(inst->ua,DB_SYNC_TTL_INTERVAL*10000u,si,db_sync_instance_ttl_cb,"db_sync_ttl"); - { struct ETCP_CONN* conn=inst->connections; while (conn) { uint64_t pid=conn->peer_node_id; if (pid!=0&&pid!=inst->node_id&&conn->links_up>0) { struct SI_PEER* p=si_peer_add(si,pid); if (p&&p->synced==0) { p->synced=1; db_sync_initiate_sync(si,pid); } } conn=conn->next; } } - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,"instance added name=%s id=%llx tbl=%s hash=%016llx next_id=%llu",name,(unsigned long long)id,si->table_name,(unsigned long long)hash,(unsigned long long)si->next_id); +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 ETCP_CONN* conn = inst->connections; + while (conn) { + etcp_conn_add_up_cbk(conn, db_sync_on_conn_up, NULL); + etcp_conn_add_down_cbk(conn, db_sync_on_conn_down, NULL); + conn = conn->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 ETCP_CONN* conn = inst->connections; + while (conn) { + etcp_conn_remove_up_cbk(conn, db_sync_on_conn_up, NULL); + etcp_conn_remove_down_cbk(conn, db_sync_on_conn_down, NULL); + conn = conn->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); + } + + si->ttl_timer = + uasync_set_timeout(inst->ua, + DB_SYNC_TTL_INTERVAL * 10000u, + si, db_sync_instance_ttl_cb, + "db_sync_ttl"); + + // Initiate sync with already connected peers + { + struct ETCP_CONN* conn = inst->connections; + while (conn) { + uint64_t pid = conn->peer_node_id; + if (pid != 0 && pid != inst->node_id + && conn->links_up > 0) + { + struct SI_PEER* p = si_peer_add(si, pid); + if (p && p->synced == 0) { + p->synced = 1; + db_sync_initiate_sync(si, pid); + } + } + conn = conn->next; + } + } + + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + "instance added name=%s id=%llx tbl=%s" + " hash=%016llx next_id=%llu", + name, (unsigned long long)id, si->table_name, + (unsigned long long)hash, + (unsigned long long)si->next_id); 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&&idxinstance_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), ts=nu, id=si->next_id; - int ret=db_record_insert(si,id,ts,dh,json_data,len,sig,sig_len); 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: [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; } - for (int i=0;ipeer_count;i++) if (si->peers[i].synced>=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); +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); + 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].synced >= 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)); } +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); } +uint32_t db_sync_count(struct DB_SYNC_INSTANCE* si) +{ + if (!si || !si->enabled) + return 0; + return db_count(si); +} -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; } +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; +} -uint64_t db_sync_get_last_timestamp(struct DB_SYNC_INSTANCE* si) { return si ? si->last_timestamp_ms : 0; } +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; -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),rts=(uint64_t)sqlite3_column_int64(stmt,1),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; + 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; }