From f0bb03c2cf1cdc70c8904e0dde58201e7b56f760 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Fri, 17 Jul 2026 15:44:19 +0300 Subject: [PATCH] db_sync: datahash removed, Ed25519 signature mandatory, PK=(timestamp,author_signature), chain_hash includes signature, signed message=ts||json --- src/db_sync.c | 1442 ++++++++++----------------- src/db_sync.h | 28 +- src/secure_channel.c | 41 + src/secure_channel.h | 3 + src/topo_node_sqlite.c | 17 + src/topo_node_sqlite.h | 1 + tests/test_chat_sync_stress.c | 17 +- tests/test_db_sync.c | 17 +- tools/chatgui/transport/chat_core.c | 97 +- 9 files changed, 643 insertions(+), 1020 deletions(-) diff --git a/src/db_sync.c b/src/db_sync.c index e87127ad..cd28da47 100644 --- a/src/db_sync.c +++ b/src/db_sync.c @@ -5,6 +5,8 @@ #include "etcp.h" #include "utun_instance.h" #include "topo_group.h" +#include "topo_node_sqlite.h" +#include "secure_channel.h" #include "../lib/debug_config.h" #include "../lib/mem.h" #include "../lib/u_async.h" @@ -12,6 +14,7 @@ #include "../lib/sqlite3.h" #include "../lib/platform_compat.h" #include +#include #define DEBUG_CATEGORY_DB_SYNC DEBUG_CATEGORY_DEBUG @@ -25,13 +28,9 @@ 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 @@ -74,14 +73,12 @@ struct DB_SYNC { #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) +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) +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); @@ -96,19 +93,14 @@ 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)); + 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); } 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); @@ -117,8 +109,7 @@ static int db_sqlite_open(struct DB_SYNC* db, const char* path) static void db_sqlite_close(struct DB_SYNC* db) { - if (!db->db) - return; + if (!db->db) return; sqlite3_close(db->db); db->db = NULL; DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: SQLite closed"); @@ -136,37 +127,35 @@ static void db_sha256(const uint8_t* data, size_t len, uint8_t hash[32]) sc_sha256_final(&ctx, hash); } -static uint64_t db_datahash(const uint8_t* data, size_t len) +static uint64_t db_hash64(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; + uint8_t hash[32]; db_sha256(data, len, hash); + uint64_t h; memcpy(&h, hash, 8); + return h; } static void db_chain_hash_compute(const uint8_t prev_chain_hash[32], uint64_t id, uint64_t timestamp, - uint64_t datahash, uint8_t out[32]) + uint64_t author, const uint8_t author_signature[DB_SIG_SIZE], + uint8_t out[32]) { - uint8_t buf[56]; + uint8_t buf[32 + 8 + 8 + 8 + DB_SIG_SIZE]; // 32 + 8 + 8 + 8 + 64 = 120 memcpy(buf, prev_chain_hash, 32); memcpy(buf + 32, &id, 8); memcpy(buf + 40, ×tamp, 8); - memcpy(buf + 48, &datahash, 8); - db_sha256(buf, 56, out); + memcpy(buf + 48, &author, 8); + memcpy(buf + 56, author_signature, DB_SIG_SIZE); + db_sha256(buf, sizeof(buf), out); } // ============================================================ // Instance management // ============================================================ -static struct DB_SYNC_INSTANCE* db_instance_find(struct DB_SYNC* db, - uint64_t hash) +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]; + if (db->instances[i].hash == hash && db->instances[i].enabled) return &db->instances[i]; } return NULL; } @@ -175,10 +164,8 @@ 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; + struct DB_SYNC_INSTANCE* np = u_realloc(db->instances, nc * sizeof(*db->instances)); + if (!np) return NULL; db->instances = np; db->instance_capacity = nc; } @@ -191,14 +178,10 @@ static struct DB_SYNC_INSTANCE* db_instance_alloc(struct DB_SYNC* db) static int si_name_valid(const char* name) { - if (!name || !name[0] || strlen(name) > 48) - return 0; + 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; + if (!((*p >= 'a' && *p <= 'z') || (*p >= 'A' && *p <= 'Z') + || (*p >= '0' && *p <= '9') || *p == '_')) return 0; } return 1; } @@ -206,41 +189,33 @@ static int si_name_valid(const char* name) 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; + 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); + return db_hash64(buf, nl + 8); } // ============================================================ // Peer management // ============================================================ -static struct SI_PEER* si_peer_find(struct DB_SYNC_INSTANCE* si, - uint64_t node_id) +static struct SI_PEER* si_peer_find(struct DB_SYNC_INSTANCE* si, uint64_t node_id) { for (int i = 0; i < si->peer_count; i++) { - if (si->peers[i].node_id == node_id) - return &si->peers[i]; + 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) +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 (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; + struct SI_PEER* np = u_realloc(si->peers, nc * sizeof(*si->peers)); + if (!np) return NULL; si->peers = np; si->peer_capacity = nc; } @@ -258,101 +233,62 @@ static struct SI_PEER* si_peer_add(struct DB_SYNC_INSTANCE* si, 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; + 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]) +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" + " ORDER BY timestamp, author_signature" " LIMIT 1 OFFSET ?") != SQLITE_OK) { - DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, - "chain_hash_at prepare: %s", - sqlite3_errmsg(SI_DB(si))); + 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; - } + 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; - } + 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) +static uint32_t si_find_pos(struct DB_SYNC_INSTANCE* si, uint64_t ts, const uint8_t* author_sig) { sqlite3_stmt* stmt; if (si_prep(si, &stmt, "SELECT COUNT(*) FROM \"%s\"" - " WHERE timestamp= 32) - memcpy(out, b, 32); - else - memset(out, 0, 32); - } - else { + 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); @@ -363,26 +299,20 @@ 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; - } + 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); + 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" + " ORDER BY timestamp, author_signature" " LIMIT 1 OFFSET ?") == SQLITE_OK) { sqlite3_bind_int64(stmt, 1, (sqlite3_int64)(from_pos - 1)); if (sqlite3_step(stmt) == SQLITE_ROW) { const void* b = sqlite3_column_blob(stmt, 0); - if (b && sqlite3_column_bytes(stmt, 0) >= 32) - memcpy(prev_ch, b, 32); + if (b && sqlite3_column_bytes(stmt, 0) >= 32) memcpy(prev_ch, b, 32); } sqlite3_finalize(stmt); } @@ -390,35 +320,30 @@ static void db_cascade_from(struct DB_SYNC_INSTANCE* si, uint32_t from_pos) sqlite3_stmt* sel, *upd; if (si_prep(si, &sel, - "SELECT id,timestamp,datahash FROM \"%s\"" - " ORDER BY timestamp,datahash" + "SELECT id,timestamp,author,author_signature FROM \"%s\"" + " ORDER BY timestamp, author_signature" " LIMIT -1 OFFSET ?") != SQLITE_OK) - { - sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); - return; - } + { 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; - } + " WHERE timestamp=? AND author_signature=?") != 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_int64 n_id = sqlite3_column_int64(sel, 0); + sqlite3_int64 n_ts = sqlite3_column_int64(sel, 1); + sqlite3_int64 n_auth = sqlite3_column_int64(sel, 2); + const void* sig_blob = sqlite3_column_blob(sel, 3); + uint8_t sig[DB_SIG_SIZE]; + if (sig_blob) memcpy(sig, sig_blob, DB_SIG_SIZE); else memset(sig, 0, DB_SIG_SIZE); + db_chain_hash_compute(prev_ch, (uint64_t)n_id, (uint64_t)n_ts, (uint64_t)n_auth, sig, ch); sqlite3_reset(upd); sqlite3_bind_blob(upd, 1, ch, 32, SQLITE_STATIC); sqlite3_bind_int64(upd, 2, n_ts); - sqlite3_bind_int64(upd, 3, n_dh); + sqlite3_bind_blob(upd, 3, sig, DB_SIG_SIZE, SQLITE_STATIC); sqlite3_step(upd); memcpy(prev_ch, ch, 32); } @@ -426,135 +351,152 @@ static void db_cascade_from(struct DB_SYNC_INSTANCE* si, uint32_t from_pos) 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)); + if (rc != SQLITE_OK) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "cascade_from COMMIT: %s", sqlite3_errmsg(db)); +} + +// ── Chain hash fragment (first 8 bytes) for sync protocol comparisons ── + +static int db_chain_hash8_at(struct DB_SYNC_INSTANCE* si, uint32_t pos, uint64_t* h8) +{ + uint8_t ch[32]; + if (db_chain_hash_at(si, pos, ch) != 0) return -1; + memcpy(h8, ch, 8); + return 0; +} + +// ── Ed25519 verification ── + +static int db_get_ed25519_pubkey(struct DB_SYNC* db, uint64_t node_id, uint8_t out[32]) +{ + if (node_id == db->inst->node_id) { memcpy(out, db->inst->my_ed25519_pubkey, 32); return 0; } + if (db->inst->topo_groups && db->inst->topo_groups->topo_sqlite_db) { + if (topo_node_sqlite_get_ed25519_pubkey(db->inst->topo_groups->topo_sqlite_db, node_id, out) == 0) return 0; + } + struct ll_entry* e = queue_find_data_by_index(db->inst->connections, (const uint8_t*)&node_id); + if (e) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + if (ce->conn) { + memcpy(out, ce->conn->peer_ed25519_pubkey, 32); + uint64_t chk; memcpy(&chk, out, 8); + if (chk != 0) return 0; + } + } + return -1; +} + +static int db_verify_author_sig(struct DB_SYNC_INSTANCE* si, + uint64_t ts, uint64_t author_node_id, + const char* json, size_t jlen, + const uint8_t* sig) +{ + uint8_t pubkey[32]; + if (db_get_ed25519_pubkey(si->db_sync, author_node_id, pubkey) != 0) { + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, + "cannot get Ed25519 pubkey for author=%016llx", (unsigned long long)author_node_id); + return -1; + } + uint8_t msg[8192]; size_t off = 0; + memcpy(msg + off, &ts, 8); off += 8; + if (jlen > 0) { memcpy(msg + off, json, jlen); off += jlen; } + if (sc_ed25519_verify(pubkey, msg, off, sig) != SC_OK) { + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, + "author_sig VERIFY FAIL node=%016llx — discarding as forgery", + (unsigned long long)author_node_id); + return -1; + } + return 0; } static int db_record_insert(struct DB_SYNC_INSTANCE* si, - uint64_t id, uint64_t ts, uint64_t dh, + uint64_t id, uint64_t ts, uint64_t author_node_id, const char* json, size_t jlen, - const uint8_t* sig, size_t sig_len, - int do_cascade) + const uint8_t* author_sig, int do_cascade) { sqlite3* db = SI_DB(si); sqlite3_stmt* stmt; int rc; + // Verify author signature (mandatory) + if (db_verify_author_sig(si, ts, author_node_id, json, jlen, author_sig) != 0) return -2; + rc = sqlite3_exec(db, "BEGIN IMMEDIATE", NULL, NULL, NULL); - if (rc != SQLITE_OK) { - DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, - "insert BEGIN: %s", sqlite3_errmsg(db)); - return -1; - } + 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; - } + // Check if record already exists by (timestamp, author_signature) + if (si_prep(si, &stmt, "SELECT 1 FROM \"%s\" WHERE timestamp=? AND author_signature=?") != SQLITE_OK) + { 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_blob(stmt, 2, author_sig, DB_SIG_SIZE, SQLITE_STATIC); int exists = (sqlite3_step(stmt) == SQLITE_ROW); sqlite3_finalize(stmt); - if (exists) { - sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); - return 1; - } + 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); + db_prev_chain_hash(si, ts, author_sig, prev_ch); uint8_t ch[32]; - db_chain_hash_compute(prev_ch, id, ts, dh, ch); + db_chain_hash_compute(prev_ch, id, ts, author_node_id, author_sig, ch); // Insert new record if (si_prep(si, &stmt, "INSERT INTO \"%s\"" - " (timestamp,datahash,id,chain_hash,author,flags,data," + " (timestamp,author,id,chain_hash,flags,data," " author_signature,delivered_peers,delivery_chain)" - " VALUES (?,?,?,?,?,0,?,?,0,'')") != SQLITE_OK) - { - sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); - return -1; - } + " 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, 2, (sqlite3_int64)author_node_id); 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); + if (jlen > 0) sqlite3_bind_blob(stmt, 5, json, (int)jlen, SQLITE_STATIC); + else sqlite3_bind_null(stmt, 5); + sqlite3_bind_blob(stmt, 6, author_sig, DB_SIG_SIZE, SQLITE_STATIC); rc = sqlite3_step(stmt); sqlite3_finalize(stmt); - if (rc != SQLITE_DONE) { - DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, - "INSERT: %s", sqlite3_errmsg(db)); - sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); - return -1; - } + 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_stmt* sel; + if (si_prep(si, &sel, + "SELECT timestamp,id,author,author_signature FROM \"%s\"" + " WHERE timestamp>?1" + " OR (timestamp=?1 AND author_signature>?2)" + " ORDER BY timestamp, author_signature") != SQLITE_OK) + { sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return -1; } + sqlite3_bind_int64(sel, 1, (sqlite3_int64)ts); + sqlite3_bind_blob(sel, 2, author_sig, DB_SIG_SIZE, SQLITE_STATIC); + + sqlite3_stmt* upd; + if (si_prep(si, &upd, + "UPDATE \"%s\" SET chain_hash=?" + " WHERE timestamp=? AND author_signature=?") != SQLITE_OK) + { sqlite3_finalize(sel); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return -1; } + + uint8_t rch[32]; memcpy(rch, ch, 32); + while (sqlite3_step(sel) == SQLITE_ROW) { + sqlite3_int64 n_ts = sqlite3_column_int64(sel, 0); + sqlite3_int64 n_id = sqlite3_column_int64(sel, 1); + sqlite3_int64 n_auth = sqlite3_column_int64(sel, 2); + const void* sig_blob = sqlite3_column_blob(sel, 3); + uint8_t sig[DB_SIG_SIZE]; + if (sig_blob) memcpy(sig, sig_blob, DB_SIG_SIZE); else memset(sig, 0, DB_SIG_SIZE); + uint8_t nch[32]; + db_chain_hash_compute(rch, (uint64_t)n_id, (uint64_t)n_ts, (uint64_t)n_auth, sig, nch); + sqlite3_reset(upd); + sqlite3_bind_blob(upd, 1, nch, 32, SQLITE_STATIC); + sqlite3_bind_int64(upd, 2, n_ts); + sqlite3_bind_blob(upd, 3, sig, DB_SIG_SIZE, SQLITE_STATIC); + sqlite3_step(upd); + memcpy(rch, nch, 32); + } + sqlite3_finalize(upd); sqlite3_finalize(sel); - 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; - } + if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "insert COMMIT: %s", sqlite3_errmsg(db)); return -1; } return 0; } @@ -562,20 +504,12 @@ static int db_record_insert(struct DB_SYNC_INSTANCE* si, // 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) +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; - } + 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; - } + 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); @@ -589,52 +523,40 @@ static int db_sync_send_hash(struct DB_SYNC* db, uint64_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); + 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) +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); + 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) +static void si_delivery_update(struct DB_SYNC_INSTANCE* si, uint64_t ts, uint64_t author, uint64_t peer_id) { - char hex[17]; - snprintf(hex, sizeof(hex), "%016llx", (unsigned long long)peer_id); + 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) + " WHERE timestamp=? AND author=?") == SQLITE_OK) { sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts); - sqlite3_bind_int64(stmt, 2, (sqlite3_int64)dh); + sqlite3_bind_int64(stmt, 2, (sqlite3_int64)author); if (sqlite3_step(stmt) == SQLITE_ROW) { const char* dc = (const char*)sqlite3_column_text(stmt, 0); - if (dc && strstr(dc, hex)) { - sqlite3_finalize(stmt); - return; - } + 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," @@ -642,11 +564,11 @@ static void si_delivery_update(struct DB_SYNC_INSTANCE* si, " WHEN delivery_chain = '' THEN ?1" " ELSE delivery_chain || ',' || ?1" " END" - " WHERE timestamp = ?2 AND datahash = ?3") == SQLITE_OK) + " WHERE timestamp = ?2 AND author = ?3") == SQLITE_OK) { sqlite3_bind_text(stmt, 1, hex, -1, SQLITE_STATIC); sqlite3_bind_int64(stmt, 2, (sqlite3_int64)ts); - sqlite3_bind_int64(stmt, 3, (sqlite3_int64)dh); + sqlite3_bind_int64(stmt, 3, (sqlite3_int64)author); sqlite3_step(stmt); sqlite3_finalize(stmt); } @@ -656,100 +578,69 @@ static void si_delivery_update(struct DB_SYNC_INSTANCE* si, // Sync protocol handlers // ============================================================ -static void db_handle_error(struct DB_SYNC* db, uint64_t src, - const uint8_t* p, size_t len) +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]); + 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) +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; - } + 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", + 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--; + 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; + uint64_t my_ch8 = 0; + if (mc > 0) db_chain_hash8_at(si, tp, &my_ch8); + memcpy(resp + off, &my_ch8, 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; + 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; + uint64_t sch8; + if (db_chain_hash8_at(si, pos, &sch8) != 0) break; + if (off + 12 > sizeof(resp)) break; memcpy(resp + off, &pos, 4); off += 4; - memcpy(resp + off, &sdh, 8); off += 8; + memcpy(resp + off, &sch8, 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); + 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) +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; - } + 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); + uint64_t peer_ch8; memcpy(&peer_ch8, p + 4, 8); uint8_t sc = p[12]; - uint64_t my_dh = 0; - db_datahash_at(si, tp, &my_dh); + uint64_t my_ch8 = 0; + db_chain_hash8_at(si, tp, &my_ch8); - if (peer_dh == 0 && sc == 0) { + if (peer_ch8 == 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", + 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; - + 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; @@ -759,64 +650,50 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, sqlite3_stmt* stmt; si_prep(si, &stmt, - "SELECT id,timestamp,datahash,data,author_signature" - " FROM \"%s\" ORDER BY timestamp,datahash" + "SELECT id,timestamp,author,data,author_signature" + " FROM \"%s\" ORDER BY timestamp, author_signature" " 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; - + uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0); + uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1); + uint64_t rauth = (uint64_t)sqlite3_column_int64(stmt, 2); + const uint8_t* rd = (const uint8_t*)sqlite3_column_blob(stmt, 3); + uint32_t rdl = (uint32_t)sqlite3_column_bytes(stmt, 3); if (!rd) rdl = 0; + const uint8_t* rsig = (const uint8_t*)sqlite3_column_blob(stmt, 4); + int rsl = sqlite3_column_bytes(stmt, 4); if (!rsig) rsl = 0; + // Wire: [id:8][ts:8][author:8][dlen:4][data][sig_len:1][sig] int rec_sz = 28 + rdl + 1 + rsl; - if (off + rec_sz > (int)sizeof(sdbuf)) - break; - + 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, &rauth, 8); off += 8; memcpy(sdbuf + off, &rdl, 4); off += 4; - if (rdl > 0) { - memcpy(sdbuf + off, rd, rdl); off += rdl; - } + 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; - } + 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 (rc == 0) break; } return; } - if (my_dh == peer_dh) { + if (my_ch8 == peer_ch8) { 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", + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "ch8 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; + 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; @@ -825,26 +702,24 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, 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" + "SELECT id,timestamp,author,data,author_signature" + " FROM \"%s\" ORDER BY timestamp, author_signature" " 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; + uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0); + uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1); + uint64_t rauth = (uint64_t)sqlite3_column_int64(stmt, 2); + const uint8_t* rd = (const uint8_t*)sqlite3_column_blob(stmt, 3); + uint32_t rdl = (uint32_t)sqlite3_column_bytes(stmt, 3); if (!rd) rdl = 0; const uint8_t* rsig = (const uint8_t*)sqlite3_column_blob(stmt, 4); - int rsl = sqlite3_column_bytes(stmt, 4); - if (!rsig) rsl = 0; + 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, &rauth, 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; @@ -858,12 +733,8 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, 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", + 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; } @@ -872,24 +743,15 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, 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; - } + uint64_t pch8; memcpy(&pch8, spr + 4, 8); + uint64_t mch8; + if (db_chain_hash8_at(si, pos, &mch8) == 0) { + if (mch8 == pch8) { if (pos + 1 > ds) ds = pos + 1; } + else { if (pos < de) de = pos; } } spr += 12; } - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, - "divergence with %016llx range [%u,%u]", - (unsigned long long)src, ds, de); + 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]; @@ -899,23 +761,21 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, memcpy(ref + roff, &de, 4); roff += 4; ref[roff++] = 0; db_sync_send(si, src, ref, roff); - } - else { + } 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; + 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; + uint64_t ch8; + if (db_chain_hash8_at(si, pos, &ch8) == 0) { + memcpy(ref + roff, &pos, 4); roff += 4; + memcpy(ref + roff, &ch8, 8); roff += 8; } } db_sync_send(si, src, ref, roff); @@ -923,14 +783,12 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, } // ---- 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) +// Wire format: [id:8][ts:8][author:8][dlen:4][data][sig_len:1=64][sig:64] +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; + if (!buf) buf = sbuf_stack; uint32_t off = 0; buf[off++] = DB_MSG_SEND_DATA; @@ -940,67 +798,48 @@ static void si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, sqlite3_stmt* stmt; si_prep(si, &stmt, - "SELECT id,timestamp,datahash,data,author_signature" - " FROM \"%s\" ORDER BY timestamp,datahash" + "SELECT id,timestamp,author,data,author_signature" + " FROM \"%s\" ORDER BY timestamp, author_signature" " 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; - + uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0); + uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1); + uint64_t rauth = (uint64_t)sqlite3_column_int64(stmt, 2); + const uint8_t* rd = (const uint8_t*)sqlite3_column_blob(stmt, 3); + uint32_t rdl = (uint32_t)sqlite3_column_bytes(stmt, 3); if (!rd) rdl = 0; + const uint8_t* rsig = (const uint8_t*)sqlite3_column_blob(stmt, 4); + int rsl = sqlite3_column_bytes(stmt, 4); if (!rsig) rsl = 0; 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; - } + if (off + rec_sz > 8000) break; + memcpy(buf + off, &rid, 8); off += 8; + memcpy(buf + off, &rts, 8); off += 8; + memcpy(buf + off, &rauth, 8); off += 8; + memcpy(buf + off, &rdl, 4); off += 4; + if (rdl > 0) { memcpy(buf + off, rd, rdl); off += rdl; } buf[off++] = (uint8_t)rsl; - if (rsl > 0) { - memcpy(buf + off, rsig, rsl); off += rsl; - } + 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); + 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) +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; - } + 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; + 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; } @@ -1009,66 +848,54 @@ static void db_handle_refine(struct DB_SYNC_INSTANCE* si, uint64_t src, 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; + uint64_t pch8 = *(uint64_t*)(hp + 4); + uint64_t mch8; + if (db_chain_hash8_at(si, pos, &mch8) == 0 && mch8 == pch8) { + 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; + 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); + 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 ---- +// Wire: [id:8][ts:8][author:8][dlen:4][data][sig_len:1][sig] static int si_parse_record(const uint8_t** pp, const uint8_t* end, - uint64_t* rid, uint64_t* rts, uint64_t* rdh, + uint64_t* rid, uint64_t* rts, uint64_t* rauthor, uint32_t* rdlen, const uint8_t** rdata, const uint8_t** rsig, int* rsiglen) { - if (*pp + 28 > end) - 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 + 28 > end) return -1; // 8+8+8+4 = 28 + *rid = *(uint64_t*)*pp; *pp += 8; + *rts = *(uint64_t*)*pp; *pp += 8; + *rauthor = *(uint64_t*)*pp; *pp += 8; + *rdlen = *(uint32_t*)*pp; *pp += 4; - if (*pp + *rdlen > end) - return -1; + if (*pp + *rdlen > end) return -1; *rdata = (*rdlen > 0) ? *pp : NULL; *pp += *rdlen; - if (*pp + 1 > end) - return -1; + if (*pp + 1 > end) return -1; *rsiglen = (int)(*(*pp)++); if (*rsiglen > 0) { - if (*pp + *rsiglen > end) - return -1; + if (*pp + *rsiglen > end) return -1; *rsig = *pp; *pp += *rsiglen; - } - else { + } 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) +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; + if (len < 6) return; uint32_t from = *(uint32_t*)p; uint16_t count = *(uint16_t*)(p + 4); const uint8_t* ptr = p + 6; @@ -1078,163 +905,119 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint16_t received = 0; for (uint16_t i = 0; i < count; i++) { - uint64_t rid, rts, rdh; + uint64_t rid, rts, rauthor; 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 (si_parse_record(&ptr, p + len, &rid, &rts, &rauthor, &rdlen, &rdata, &rsig, &rsiglen) != 0) break; + int ret = db_record_insert(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, 0); + if (ret >= 0) received++; if (ret == 0 && si->on_insert) - si->on_insert(si, (const char*)rdata, rdlen, - src, si->on_insert_arg); + si->on_insert(si, rts, (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; + 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", + 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", + } 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); + uint64_t lch8 = 0; + if (nc > 0) db_chain_hash8_at(si, nc - 1, &lch8); uint8_t sd[13]; sd[0] = DB_MSG_SYNC_DONE; memcpy(sd + 1, &nc, 4); - memcpy(sd + 5, &ldh, 8); + memcpy(sd + 5, &lch8, 8); db_sync_send(si, src, sd, 13); - if (sp) - sp->sync_state = 2; + 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); + "received %u/%u records from %016llx range=[%u,%u] → SYNC_DONE nc=%u ch8=%016llx", + received, count, (unsigned long long)src, from, from + received - 1, nc, (unsigned long long)lch8); } -static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, - uint64_t src, const uint8_t* p, size_t len) +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; + if (len < 12) return; uint32_t pc = *(uint32_t*)p; - uint64_t pdh; - memcpy(&pdh, p + 4, 8); + uint64_t pch8; memcpy(&pch8, p + 4, 8); uint32_t mc = db_count(si); - uint64_t mdh = 0; - if (mc > 0) - db_datahash_at(si, mc - 1, &mdh); + uint64_t mch8 = 0; + if (mc > 0) db_chain_hash8_at(si, mc - 1, &mch8); - if (mc != pc || mdh != pdh) { + if (mc != pc || mch8 != pch8) { 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); + "SYNC_DONE mismatch with %016llx my_count=%u peer_count=%u my_ch8=%016llx peer_ch8=%016llx — re-initiating", + (unsigned long long)src, mc, pc, (unsigned long long)mch8, (unsigned long long)pch8); 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); + if (sp) { sp->synced_pos = (mc > 0) ? mc - 1 : 0; sp->sync_state = 2; } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync confirmed with %016llx count=%u ch8=%016llx", + (unsigned long long)src, mc, (unsigned long long)mch8); } -static void db_handle_ack_push(struct DB_SYNC_INSTANCE* si, - uint64_t src, const uint8_t* p, size_t len) +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); + if (len < 16) return; + uint64_t ts = *(uint64_t*)p; + uint64_t author = *(uint64_t*)(p + 8); sqlite3_stmt* stmt; if (si_prep(si, &stmt, "UPDATE \"%s\" SET flags = flags | 1" - " WHERE timestamp=? AND datahash=?" + " WHERE timestamp=? AND author=?" " 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_bind_int64(stmt, 2, (sqlite3_int64)author); + sqlite3_bind_int64(stmt, 3, (sqlite3_int64)si->db_sync->inst->node_id); sqlite3_step(stmt); int changed = sqlite3_changes(SI_DB(si)); sqlite3_finalize(stmt); if (changed > 0) - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, - "ACK_PUSH mark sent dh=%016llx ts=%llu from %016llx", - (unsigned long long)dh, - (unsigned long long)ts, - (unsigned long long)src); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "ACK_PUSH mark sent author=%016llx ts=%llu from %016llx", + (unsigned long long)author, (unsigned long long)ts, (unsigned long long)src); } - si_delivery_update(si, ts, dh, src); + si_delivery_update(si, ts, author, src); } -static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src, - const uint8_t* p, size_t len) +static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { - if (len < 28) - return; + if (len < 30) return; const uint8_t* ptr = p; - uint64_t rid, rts, rdh; + uint64_t rid, rts, rauthor; 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; + if (si_parse_record(&ptr, p + len, &rid, &rts, &rauthor, &rdlen, &rdata, &rsig, &rsiglen) != 0) return; - int ret = db_record_insert(si, rid, rts, rdh, - (const char*)rdata, rdlen, - rsig, rsiglen, 0); + int ret = db_record_insert(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, 0); if (ret == 0) { - uint32_t ins_pos = si_find_pos(si, rts, rdh); + uint32_t ins_pos = si_find_pos(si, rts, rsig); 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->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); + if (si->on_insert) si->on_insert(si, rts, (const char*)rdata, rdlen, src, si->on_insert_arg); uint8_t ack[17]; ack[0] = DB_MSG_ACK_PUSH; - memcpy(ack + 1, &rdh, 8); - memcpy(ack + 9, &rts, 8); + memcpy(ack + 1, &rts, 8); + memcpy(ack + 9, &rauthor, 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); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "PUSH inserted from %016llx id=%llu author=%016llx ts=%llu", + (unsigned long long)src, (unsigned long long)rid, (unsigned long long)rauthor, (unsigned long long)rts); } } @@ -1242,22 +1025,11 @@ static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src, // Receive callback // ============================================================ -static void db_sync_recv_cb(struct ETCP_CONN* conn, - struct ll_entry* entry) +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; - } + 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; - } + 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)); @@ -1268,55 +1040,29 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, 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", + "recv msg type=0x%02x from %016llx hash=%016llx — instance NOT FOUND; sending DB_ERR_NOT_FOUND", type, (unsigned long long)src, (unsigned long long)hash); - uint8_t err[2]; - err[0] = DB_MSG_ERROR; - err[1] = DB_ERR_NOT_FOUND; + 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; + 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; + 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; + 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; + 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); + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "unknown msg type 0x%02x from %016llx", type, (unsigned long long)src); break; } queue_dgram_free(entry); @@ -1330,8 +1076,7 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, static void db_sync_on_new_conn(struct ETCP_CONN* conn, void* arg) { (void)arg; - if (!conn || !conn->instance || !conn->instance->db_sync) - return; + 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); } @@ -1339,14 +1084,11 @@ 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) { (void)arg; - if (!conn || !conn->instance || !conn->instance->db_sync) - return; + if (!conn || !conn->instance || !conn->instance->db_sync) return; struct DB_SYNC* db = conn->instance->db_sync; - if (!db->enabled) - return; + if (!db->enabled) return; uint64_t pid = conn->peer_node_id; - if (pid == 0 || pid == db->inst->node_id) - return; + 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; @@ -1370,63 +1112,48 @@ 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) { (void)arg; - if (!conn || !conn->instance || !conn->instance->db_sync) - return; + if (!conn || !conn->instance || !conn->instance->db_sync) return; struct DB_SYNC* db = conn->instance->db_sync; - if (!db->enabled) - return; + if (!db->enabled) return; uint64_t pid = conn->peer_node_id; - if (pid == 0) - return; + 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; + 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; - } + 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); + 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) +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)); + 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); + 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; + if (mc == 0) return; uint8_t prev_ch[32]; memset(prev_ch, 0, 32); uint8_t exp_ch[32], stored_ch[32]; @@ -1434,31 +1161,28 @@ static void db_verify_chain(struct DB_SYNC_INSTANCE* si) sqlite3_stmt* stmt; if (si_prep(si, &stmt, - "SELECT id,timestamp,datahash,chain_hash FROM \"%s\"" - " ORDER BY timestamp,datahash") != SQLITE_OK) + "SELECT id,timestamp,author,author_signature,chain_hash FROM \"%s\"" + " ORDER BY timestamp, author_signature") != SQLITE_OK) return; for (uint32_t pos = 0; sqlite3_step(stmt) == SQLITE_ROW; pos++) { - uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0); - uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1); - uint64_t 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; - } + uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0); + uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1); + uint64_t rauth = (uint64_t)sqlite3_column_int64(stmt, 2); + const void* sig_blob = sqlite3_column_blob(stmt, 3); + uint8_t sig[DB_SIG_SIZE]; + if (sig_blob) memcpy(sig, sig_blob, DB_SIG_SIZE); else memset(sig, 0, DB_SIG_SIZE); + const void* b = sqlite3_column_blob(stmt, 4); + if (b && sqlite3_column_bytes(stmt, 4) >= 32) memcpy(stored_ch, b, 32); + else memset(stored_ch, 0, 32); + + db_chain_hash_compute(prev_ch, rid, rts, rauth, sig, exp_ch); + if (memcmp(exp_ch, stored_ch, 32) != 0) { bad_pos = (int)pos; break; } memcpy(prev_ch, exp_ch, 32); } sqlite3_finalize(stmt); if (bad_pos >= 0) { - DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, - "chain_hash mismatch at pos %d/%u in %s, recalculating", + 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); } @@ -1471,39 +1195,27 @@ static void db_verify_chain(struct DB_SYNC_INSTANCE* si) static void db_sync_peer_check_cb(void* arg) { struct DB_SYNC* db = (struct DB_SYNC*)arg; - if (!db || !db->enabled) - return; + if (!db || !db->enabled) return; - struct TOPO_GROUP* g = - topo_groups_get_default(db->inst->topo_groups); + 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"); + 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; + 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) - { + 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++; - } + if (pid != db->inst->node_id) { si_peer_add(si, pid); peers_found++; } } e = e->next; } @@ -1516,24 +1228,15 @@ static void db_sync_peer_check_cb(void* arg) best = &si->peers[j]; } } - if (best) { - best->sync_state = 1; - db_sync_initiate_sync(si, best->node_id); - total_synced++; - } else { - total_skipped++; - } + 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", + 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"); + 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) @@ -1542,20 +1245,13 @@ static void db_sync_instance_ttl_cb(void* 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"); + 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 ttl = (uint64_t)(si->db_sync->inst->config->global.db_sync_ttl) * 1000000uLL; uint64_t cut = nu - ttl; sqlite3_stmt* stmt; @@ -1564,11 +1260,7 @@ static void db_sync_instance_ttl_cb(void* arg) " 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"); + 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); @@ -1577,16 +1269,9 @@ static void db_sync_instance_ttl_cb(void* arg) 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"); + 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"); } // ============================================================ @@ -1595,23 +1280,15 @@ static void db_sync_instance_ttl_cb(void* arg) int db_sync_init(struct UTUN_INSTANCE* inst) { - if (!inst) { - DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "NULL instance"); - return -1; - } + 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; - } + 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"); + db->enabled = 0; inst->db_sync = db; + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: disabled by config"); return 0; } @@ -1620,21 +1297,17 @@ int db_sync_init(struct UTUN_INSTANCE* inst) 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 (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"); + 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) { @@ -1648,42 +1321,29 @@ int db_sync_init(struct UTUN_INSTANCE* inst) } 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"); + 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); + 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; + 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; - } + 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); + 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) { @@ -1697,88 +1357,58 @@ void db_sync_destroy(struct UTUN_INSTANCE* inst) } db_sqlite_close(db); - if (db->instances) - u_free(db->instances); + 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) +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; - } + 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 (!db->enabled || !db->db) return NULL; - if (!si_name_valid(name)) { - DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, - "invalid instance name '%s'", name); - 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; - } + 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; + if (!si) return NULL; si->hash = hash; - snprintf(si->table_name, sizeof(si->table_name), - "db_sync_%s_%llx", name, (unsigned long long)id); + snprintf(si->table_name, sizeof(si->table_name), "db_sync_%s_%llx", name, (unsigned long long)id); - // Create table + // Create table — PK is (timestamp, author), author_signature NOT NULL { char sql[512]; snprintf(sql, sizeof(sql), "CREATE TABLE IF NOT EXISTS \"%s\" (" " timestamp INTEGER NOT NULL," - " datahash INTEGER NOT NULL," + " author 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," + " author_signature BLOB NOT NULL," " delivered_peers INTEGER NOT NULL DEFAULT 0," " delivery_chain TEXT NOT NULL DEFAULT ''," - " PRIMARY KEY (timestamp, datahash))", + " PRIMARY KEY (timestamp, author_signature))", si->table_name); int rc = sqlite3_exec(db->db, sql, NULL, NULL, NULL); - if (rc != SQLITE_OK) { - DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, - "CREATE TABLE %s: %s", - si->table_name, sqlite3_errmsg(db->db)); - si->enabled = 0; - return si; - } + 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)); + 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 + 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 @@ -1788,11 +1418,7 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, 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"); + 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; { @@ -1800,16 +1426,10 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, 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) - { + 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++; - } + if (p && p->sync_state == 0) { p->sync_state = 1; db_sync_initiate_sync(si, pid); peers_synced++; } } entry = entry->next; } @@ -1817,130 +1437,72 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, 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)); + 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; + 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; - } + 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], + 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); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "instance removed tbl=%s hash=%016llx", si->table_name, (unsigned long long)si->hash); } -int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si, - const char* json_data, size_t len, - const uint8_t* sig, size_t sig_len) +int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si, const char* json_data, size_t len, + const uint8_t* sig, size_t sig_len, uint64_t ts) { - if (!si || !si->enabled || !json_data || len == 0) { - DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "insert invalid args"); - return -1; - } - - 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; + if (!si || !si->enabled || !json_data || len == 0) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "insert invalid args"); return -1; } + if (!sig || sig_len != DB_SIG_SIZE) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "insert_signed: signature required (64 bytes Ed25519)"); return -1; } - uint64_t dh = db_datahash((const uint8_t*)json_data, len); - uint64_t ts = nu; uint64_t id = si->next_id; + uint64_t author_node_id = si->db_sync->inst->node_id; + + int ret = db_record_insert(si, id, ts, author_node_id, json_data, len, sig, 1); + if (ret != 0) return ret; + if (si->on_insert) si->on_insert(si, ts, json_data, len, author_node_id, si->on_insert_arg); - 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]; + // Build PUSH: [DB_MSG_PUSH][id:8][ts:8][author:8][dlen:4][data][sig_len:1=64][sig:64] + uint32_t sl = DB_SIG_SIZE; + uint8_t pbuf[2560]; uint32_t off = 0; pbuf[off++] = DB_MSG_PUSH; - memcpy(pbuf + off, &id, 8); off += 8; - memcpy(pbuf + off, &ts, 8); off += 8; - memcpy(pbuf + off, &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; - } + memcpy(pbuf + off, &id, 8); off += 8; + memcpy(pbuf + off, &ts, 8); off += 8; + memcpy(pbuf + off, &author_node_id, 8); off += 8; + memcpy(pbuf + off, &len, 4); off += 4; + if (len > 0 && off + len <= sizeof(pbuf)) { memcpy(pbuf + off, json_data, len); off += len; } pbuf[off++] = (uint8_t)sl; - if (sl > 0 && off + sl <= sizeof(pbuf)) { - memcpy(pbuf + off, sig, sl); off += sl; - } + if (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 (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); + si_delivery_update(si, ts, author_node_id, 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); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "insert id=%llu author=%016llx ts=%llu len=%zu", + (unsigned long long)id, (unsigned long long)author_node_id, (unsigned long long)ts, len); 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; + if (!si || !si->enabled) return 0; return db_count(si); } @@ -1949,32 +1511,38 @@ 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) +uint64_t db_sync_next_timestamp(struct DB_SYNC_INSTANCE* si) { - if (!si) - return; + if (!si) return 0; + 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; + return nu; +} + +void db_sync_set_insert_cb(struct DB_SYNC_INSTANCE* si, db_sync_insert_cb cb, void* arg) +{ + if (!si) return; si->on_insert = cb; si->on_insert_arg = arg; } -int db_sync_select(struct DB_SYNC_INSTANCE* si, - uint32_t offset, uint32_t limit, - db_sync_select_cb cb, void* arg) +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; + 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" + " FROM \"%s\" ORDER BY timestamp, author_signature" " LIMIT ? OFFSET ?") != SQLITE_OK) return 0; - sqlite3_bind_int64(stmt, 1, - limit > 0 ? (sqlite3_int64)limit : -1); + sqlite3_bind_int64(stmt, 1, limit > 0 ? (sqlite3_int64)limit : -1); sqlite3_bind_int64(stmt, 2, (sqlite3_int64)offset); int cnt = 0; @@ -1984,16 +1552,12 @@ int db_sync_select(struct DB_SYNC_INSTANCE* si, 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; + 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 = ""; + 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); + cb(arg, rid, rts, rd, (size_t)(rdl > 0 ? rdl : 0), rauth, rsig, (size_t)rsl, rdp, rdc); cnt++; } sqlite3_finalize(stmt); diff --git a/src/db_sync.h b/src/db_sync.h index 43a12e0b..521cea11 100644 --- a/src/db_sync.h +++ b/src/db_sync.h @@ -3,13 +3,21 @@ // Назначение: децентрализованная реплицируемая таблица JSON-записей между всеми узлами сети. // Модуль поддерживает несколько независимых инстансов (таблиц), каждый идентифицируется парой (name, id). // +// Каждая запись ОБЯЗАТЕЛЬНО содержит Ed25519-подпись автора. +// chain_hash[pos] = SHA256(chain_hash[pos-1] || id || timestamp || author || author_signature[64]) +// Первичный ключ: (timestamp, author) — тот же автор в ту же ms = дубликат. +// Упорядочение: ORDER BY timestamp, author. +// // Использование: // 1. Включить в конфиге: db_sync_enabled = 1 // Опционально: db_sync_ttl = 86400 (по умолчанию) // 2. db_sync_init() — вызывается автоматически при старте utun_instance // 3. struct DB_SYNC_INSTANCE* si = db_sync_instance_add(inst, "chats", 1); // Создаёт/регистрирует инстанс. При поднятии соединений автоматически запускает sync. -// 4. db_sync_insert(si, json_data) — вставить запись. Автоматически PUSH всем synced-пирам. +// 4. uint64_t ts = db_sync_next_timestamp(si); +// sig = Ed25519(ts || my_node_id || json_data) +// db_sync_insert_signed(si, json_data, len, sig, 64, ts) — вставить запись. +// sig=NULL → ошибка. // 5. db_sync_count(si) — количество записей в локальной БД для этого инстанса // 6. db_sync_instance_remove(si) — удалить инстанс (таблица БД не удаляется) // 7. db_sync_destroy() — вызывается автоматически при завершении @@ -23,9 +31,11 @@ // - Протокол: сравнение chain_hash для поиска расхождений (7 типов сообщений) // - Новые записи немедленно рассылаются подключённым пирам через PUSH // +// Wire-формат записи (SEND_DATA/PUSH): [id:8][ts:8][author:8][dlen:4][data][sig_len:1=64][sig:64] +// // Нюансы: // - Записи не редактируются и не удаляются явно — только TTL-очистка (per-instance) -// - Дубликаты определяются по (instance_hash, timestamp, datahash) +// - Дубликаты определяются по (timestamp, author) // - БД хранится в SQLite, путь: /sync // - Синхронизация идёт через etcp_api (service ID 0x20), только прямые P2P соединения #ifndef DB_SYNC_H @@ -71,6 +81,9 @@ struct DB_SYNC_INSTANCE; #define DB_REFINE_HASHES 16 #define DB_SEND_DATA_MAX 32 +// Ed25519 signature size +#define DB_SIG_SIZE 64 + // Global lifecycle int db_sync_init(struct UTUN_INSTANCE* inst); void db_sync_destroy(struct UTUN_INSTANCE* inst); @@ -80,14 +93,16 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si); // Data operations (per-instance) -int db_sync_insert_len(struct DB_SYNC_INSTANCE* si, const char* json_data, size_t len); -int db_sync_insert(struct DB_SYNC_INSTANCE* si, const char* json_data); +// db_sync_insert_signed: sig MUST be non-NULL, 64 bytes. +// sig = Ed25519(ts[8] || author[8] || json). +// ts — call db_sync_next_timestamp(si) before signing to reserve monotonically increasing timestamp. 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); + const uint8_t* sig, size_t sig_len, uint64_t ts); uint32_t db_sync_count(struct DB_SYNC_INSTANCE* si); uint64_t db_sync_get_last_timestamp(struct DB_SYNC_INSTANCE* si); +uint64_t db_sync_next_timestamp(struct DB_SYNC_INSTANCE* si); -// Select: iterate records ordered by (timestamp, datahash), starting at offset, max limit (0=unlimited). +// Select: iterate records ordered by (timestamp, author), starting at offset, max limit (0=unlimited). // Returns number of records passed to callback. typedef void (*db_sync_select_cb)(void* arg, uint64_t id, uint64_t timestamp, const char* data, size_t data_len, uint64_t author, @@ -99,6 +114,7 @@ int db_sync_select(struct DB_SYNC_INSTANCE* si, uint32_t offset, uint32_t l // Insert callback: fired after local or peer insert succeeds. // author_node_id = self for local inserts, peer node_id for remote. typedef void (*db_sync_insert_cb)(struct DB_SYNC_INSTANCE* si, + uint64_t record_timestamp, const char* json_data, size_t len, uint64_t author_node_id, void* arg); void db_sync_set_insert_cb(struct DB_SYNC_INSTANCE* si, db_sync_insert_cb cb, void* arg); diff --git a/src/secure_channel.c b/src/secure_channel.c index e10e65f6..ce78eb35 100644 --- a/src/secure_channel.c +++ b/src/secure_channel.c @@ -526,6 +526,47 @@ void sc_stream_sign_cleanup(struct sc_stream_sign_state *state) { state->initialized = 0; } +// ── One-shot Ed25519 helpers (raw keys, any context) ── + +sc_status_t sc_ed25519_sign(const uint8_t privkey[32], const uint8_t* msg, size_t msg_len, uint8_t sig_out[64]) +{ + if (!privkey || !msg || !sig_out) return SC_ERR_INVALID_ARG; + EVP_PKEY* pkey = EVP_PKEY_new_raw_private_key(EVP_PKEY_ED25519, NULL, privkey, 32); + if (!pkey) { DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "EVP_PKEY_new_raw_private_key failed"); return SC_ERR_CRYPTO; } + EVP_MD_CTX* ctx = EVP_MD_CTX_new(); + int rc = SC_ERR_CRYPTO; + if (ctx) { + size_t slen = 64; + if (EVP_DigestSignInit(ctx, NULL, NULL, NULL, pkey) == 1 + && EVP_DigestSign(ctx, sig_out, &slen, msg, msg_len) == 1) + rc = SC_OK; + else + DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "EVP_DigestSign failed"); + EVP_MD_CTX_free(ctx); + } + EVP_PKEY_free(pkey); + return rc; +} + +sc_status_t sc_ed25519_verify(const uint8_t pubkey[32], const uint8_t* msg, size_t msg_len, const uint8_t sig[64]) +{ + if (!pubkey || !msg || !sig) return SC_ERR_INVALID_ARG; + EVP_PKEY* pkey = EVP_PKEY_new_raw_public_key(EVP_PKEY_ED25519, NULL, pubkey, 32); + if (!pkey) { DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "EVP_PKEY_new_raw_public_key failed"); return SC_ERR_CRYPTO; } + EVP_MD_CTX* ctx = EVP_MD_CTX_new(); + int rc = SC_ERR_CRYPTO; + if (ctx) { + if (EVP_DigestVerifyInit(ctx, NULL, NULL, NULL, pkey) == 1 + && EVP_DigestVerify(ctx, sig, 64, msg, msg_len) == 1) + rc = SC_OK; + else + rc = SC_ERR_AUTH_FAILED; + EVP_MD_CTX_free(ctx); + } + EVP_PKEY_free(pkey); + return rc; +} + // --- Common crypto utilities --- sc_status_t sc_sha_transcode(const uint8_t *key, size_t key_len, uint8_t *data, size_t data_len) { diff --git a/src/secure_channel.h b/src/secure_channel.h index e94bee0f..686b4505 100644 --- a/src/secure_channel.h +++ b/src/secure_channel.h @@ -109,6 +109,9 @@ sc_status_t sc_stream_sign_final(struct sc_stream_sign_state *state, uint8_t *si sc_status_t sc_stream_sign_verify(struct sc_stream_sign_state *state, const uint8_t *sig, size_t sig_len); void sc_stream_sign_cleanup(struct sc_stream_sign_state *state); +sc_status_t sc_ed25519_sign(const uint8_t privkey[32], const uint8_t* msg, size_t msg_len, uint8_t sig_out[64]); +sc_status_t sc_ed25519_verify(const uint8_t pubkey[32], const uint8_t* msg, size_t msg_len, const uint8_t sig[64]); + sc_status_t sc_derive_ed25519_pubkey(const uint8_t *x25519_privkey, uint8_t *ed25519_pubkey_out); uint64_t sc_derive_node_id(const uint8_t *private_key); diff --git a/src/topo_node_sqlite.c b/src/topo_node_sqlite.c index d0b89f79..98af5747 100644 --- a/src/topo_node_sqlite.c +++ b/src/topo_node_sqlite.c @@ -394,3 +394,20 @@ int topo_node_sqlite_node_get_online(sqlite3* db, uint64_t node_id) { sqlite3_finalize(stmt); return online; } + +int topo_node_sqlite_get_ed25519_pubkey(sqlite3* db, uint64_t node_id, uint8_t pubkey_out[32]) +{ + if (!db || !pubkey_out) return -1; + sqlite3_stmt* stmt = NULL; + if (sqlite3_prepare_v2(db, "SELECT ed25519_pubkey FROM nodes WHERE node_id=?", + -1, &stmt, NULL) != SQLITE_OK) return -1; + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)node_id); + int rc = -1; + if (sqlite3_step(stmt) == SQLITE_ROW) { + const void* b = sqlite3_column_blob(stmt, 0); + int bytes = sqlite3_column_bytes(stmt, 0); + if (b && bytes >= 32) { memcpy(pubkey_out, b, 32); rc = 0; } + } + sqlite3_finalize(stmt); + return rc; +} diff --git a/src/topo_node_sqlite.h b/src/topo_node_sqlite.h index 13c691e7..0ce3c1ed 100644 --- a/src/topo_node_sqlite.h +++ b/src/topo_node_sqlite.h @@ -47,5 +47,6 @@ int topo_node_sqlite_channel_peers_all(sqlite3* db, const char* channel_id, int topo_node_sqlite_node_set_online(sqlite3* db, uint64_t node_id, int online); int topo_node_sqlite_node_get_online(sqlite3* db, uint64_t node_id); +int topo_node_sqlite_get_ed25519_pubkey(sqlite3* db, uint64_t node_id, uint8_t pubkey_out[32]); #endif diff --git a/tests/test_chat_sync_stress.c b/tests/test_chat_sync_stress.c index 908733d9..a88a451e 100644 --- a/tests/test_chat_sync_stress.c +++ b/tests/test_chat_sync_stress.c @@ -17,6 +17,7 @@ #include "../src/config_parser.h" #include "../src/config_updater.h" #include "../src/db_sync.h" +#include "../src/secure_channel.h" #include "../lib/u_async.h" #include "../lib/debug_config.h" @@ -118,6 +119,16 @@ static int wait_for(const char* desc, int timeout_tb, uint64_t* elapsed_out) { return ok; } +static int insert_record(struct DB_SYNC_INSTANCE* si, struct UTUN_INSTANCE* inst, const char* data, size_t len) { + uint64_t ts = db_sync_next_timestamp(si); + uint8_t sig_msg[4096]; size_t off = 0; + memcpy(sig_msg + off, &ts, 8); off += 8; + memcpy(sig_msg + off, data, len); off += len; + uint8_t sig[64]; + if (sc_ed25519_sign(inst->my_ed25519_privkey, sig_msg, off, sig) != SC_OK) return -1; + return db_sync_insert_signed(si, data, len, sig, 64, ts); +} + int main(void) { g_seed = (unsigned int)time(NULL); srand(g_seed); @@ -146,8 +157,8 @@ int main(void) { uint64_t t0 = get_time_tb(); // Phase 1: local inserts on A - if (db_sync_insert(si_a, "{\"test\":1}") != 0 - || db_sync_insert(si_a, "{\"test\":2}") != 0 + if (insert_record(si_a, inst_a, "{\"test\":1}", 11) != 0 + || insert_record(si_a, inst_a, "{\"test\":2}", 11) != 0 || db_sync_count(si_a) != 2) { DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "PHASE1_FAIL A=%u", db_sync_count(si_a)); test_phase = 2; goto cleanup; @@ -180,7 +191,7 @@ int main(void) { "{\"seq\":%d,\"n\":%llu,\"ch\":\"test\",\"ct\":\"text\",\"d\":\"msg_%d\"," "\"pad\":\"xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx\"}", global_seq + i, (unsigned long long)node, global_seq + i); - if (db_sync_insert(si, buf) < 0) { + if (insert_record(si, side ? inst_b : inst_a, buf, strlen(buf)) < 0) { DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "INSERT_FAIL seq=%d", global_seq + i); test_phase = 2; break; } diff --git a/tests/test_db_sync.c b/tests/test_db_sync.c index 8ccdc81f..41bba303 100644 --- a/tests/test_db_sync.c +++ b/tests/test_db_sync.c @@ -22,6 +22,7 @@ #include "../src/routing.h" #include "../src/tun_if.h" #include "../src/secure_channel.h" +#include "../src/secure_channel.h" #include "../src/db_sync.h" #include "../lib/u_async.h" #include "../lib/debug_config.h" @@ -123,12 +124,18 @@ static uint32_t ca_target, cb_target; static int _cond_ca(void) { return si_a && db_sync_count(si_a) == ca_target; } static int _cond_cb(void) { return si_b && db_sync_count(si_b) == cb_target; } -static int insert_many(struct DB_SYNC_INSTANCE* si, int start, int count) { +static int insert_many(struct DB_SYNC_INSTANCE* si, struct UTUN_INSTANCE* inst, int start, int count) { char buf[128]; for (int i = start; i < start + count && test_phase == 0; i++) { snprintf(buf, sizeof(buf), "{\"idx\":%d,\"val\":\"data_%d\",\"pad\":\"%.50s\"}", i, i, "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"); - if (db_sync_insert(si, buf) < 0) { fprintf(stderr,"insert_many fail %d\n", i); return -1; } + uint64_t ts = db_sync_next_timestamp(si); + uint8_t sig_msg[256]; size_t off = 0; + memcpy(sig_msg + off, &ts, 8); off += 8; + size_t jl = strlen(buf); memcpy(sig_msg + off, buf, jl); off += jl; + uint8_t sig[64]; + if (sc_ed25519_sign(inst->my_ed25519_privkey, sig_msg, off, sig) != SC_OK) { fprintf(stderr,"sign fail %d\n", i); return -1; } + if (db_sync_insert_signed(si, buf, jl, sig, 64, ts) < 0) { fprintf(stderr,"insert_many fail %d\n", i); return -1; } } return 0; } @@ -164,7 +171,7 @@ int main(void) { if (!wait_for("links up", cond_links_init, PHASE_TIMEOUT_TB)) { test_phase=2; goto done; } si_a = db_sync_instance_add(inst_a, "test", 1); - if (!si_a || insert_many(si_a, 0, 50) != 0) { test_phase=2; goto done; } + if (!si_a || insert_many(si_a, inst_a, 0, 50) != 0) { test_phase=2; goto done; } if (db_sync_count(si_a) != 50) { fprintf(stderr, "FAIL: A!=50\n"); test_phase=2; goto done; } printf(" A has 50 records, creating B now (B empty, A conn_up already fired)\n"); @@ -181,7 +188,7 @@ int main(void) { // =================================================================== printf("Phase 2: recreate A — dh_match at tp=49 (sc>0 sparse) + tail-send 30 records\n"); - if (insert_many(si_a, 50, 30) != 0) { test_phase=2; goto done; } + if (insert_many(si_a, inst_a, 50, 30) != 0) { test_phase=2; goto done; } if (db_sync_count(si_a) != 80) { fprintf(stderr, "FAIL: A!=80\n"); test_phase=2; goto done; } remove_si(&si_a); @@ -207,7 +214,7 @@ int main(void) { si_b = db_sync_instance_add(inst_b, "test2", 2); if (!si_a || !si_b) { test_phase=2; goto done; } - if (insert_many(si_a, 0, 30) != 0 || db_sync_count(si_a) != 30) { test_phase=2; goto done; } + if (insert_many(si_a, inst_a, 0, 30) != 0 || db_sync_count(si_a) != 30) { test_phase=2; goto done; } cb_target = 30; if (!wait_for("B count=30 (peer empty)", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase=2; goto done; } diff --git a/tools/chatgui/transport/chat_core.c b/tools/chatgui/transport/chat_core.c index ad74df9b..f2f5067f 100644 --- a/tools/chatgui/transport/chat_core.c +++ b/tools/chatgui/transport/chat_core.c @@ -51,7 +51,7 @@ static struct chat_core_ctx { static struct DB_SYNC_INSTANCE* si_find(const char* ch_id); static void si_register(struct DB_SYNC_INSTANCE* si, const char* ch_id); -static void on_db_sync_insert(struct DB_SYNC_INSTANCE* si, const char* data, size_t len, uint64_t author, void* arg); +static void on_db_sync_insert(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char* data, size_t len, uint64_t author, void* arg); /* ─── утилиты ─── */ @@ -77,40 +77,6 @@ static void peers_table_name(const char* ch_id, char* buf, size_t sz) { snprintf(buf, sz, "peers_%s", san); } -static void compute_datahash(const uint8_t* data, size_t len, uint64_t* dh) { - uint8_t hash[32]; SHA256(data, len, hash); - memcpy(dh, hash, 8); -} - -static void compute_chain_hash(const uint8_t* prev_chain, int64_t ts, - uint64_t dh, uint8_t* out) { - uint8_t buf[32 + 8 + 8]; - memcpy(buf, prev_chain, 32); - memcpy(buf + 32, &ts, 8); - memcpy(buf + 40, &dh, 8); - SHA256(buf, 48, out); -} - -static int get_last_chain_hash(const char* ch_id, uint8_t* out) { - char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); - char sql[256]; snprintf(sql, sizeof(sql), - "SELECT chain_hash FROM \"%s\" ORDER BY timestamp, datahash DESC LIMIT 1", tbl); - sqlite3_stmt* stmt = NULL; - if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) { - memset(out, 0, 32); return -1; - } - if (sqlite3_step(stmt) == SQLITE_ROW) { - const void* blob = sqlite3_column_blob(stmt, 0); - int len = sqlite3_column_bytes(stmt, 0); - if (blob && len == 32) memcpy(out, blob, 32); - else memset(out, 0, 32); - } else { - memset(out, 0, 32); - } - sqlite3_finalize(stmt); - return 0; -} - static int db_exec(const char* sql) { char* err = NULL; int rc = sqlite3_exec(g_cc.db, sql, NULL, NULL, &err); @@ -359,25 +325,31 @@ void chat_core_submit_message(struct chat_msg_submit* req) { "{\"n\":%llu,\"ch\":\"%s\",\"ct\":\"%s\",\"d\":\"%.*s\"}", (unsigned long long)g_cc.my_node_id, req->channel_id, req->content_type, (int)req->data_len, (const char*)req->data); - /* FIXME: JSON escaping for d (quotes/backslashes in data) */ - /* For now, naive format; will break on special chars. TODO: use base64. */ - int ret = db_sync_insert_signed(si, json, strlen(json), NULL, 0); + /* Generate Ed25519 signature: ts || json */ + uint64_t ts = db_sync_next_timestamp(si); + uint8_t sig_msg[8192]; size_t soff = 0; + memcpy(sig_msg + soff, &ts, 8); soff += 8; + size_t jl = strlen(json); memcpy(sig_msg + soff, json, jl); soff += jl; + uint8_t sig[64]; + if (sc_ed25519_sign(g_cc.inst->my_ed25519_privkey, sig_msg, soff, sig) != SC_OK) { + DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: Ed25519 sign failed", CC_ID); + return; + } + + int ret = db_sync_insert_signed(si, json, strlen(json), sig, 64, ts); if (ret != 0) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: db_sync_insert_signed failed ret=%d", CC_ID, ret); /* Insert directly to msg_ table as fallback (dual-write not triggered by callback) */ - uint64_t dh; compute_datahash((const uint8_t*)req->data, req->data_len, &dh); char tbl[80]; msg_table_name(req->channel_id, tbl, sizeof(tbl)); - char sql[512]; snprintf(sql, sizeof(sql), "INSERT OR IGNORE INTO \"%s\" (node_id,content_type,data,timestamp,datahash,chain_hash,signature,is_outgoing,is_read) VALUES(?,?,?,?,?,?,?,1,1)", tbl); + char sql[512]; snprintf(sql, sizeof(sql), "INSERT OR IGNORE INTO \"%s\" (node_id,content_type,data,timestamp,signature,is_outgoing,is_read) VALUES(?,?,?,?,?,1,1)", tbl); sqlite3_stmt* st=NULL; if (sqlite3_prepare_v2(g_cc.db,sql,-1,&st,NULL)==SQLITE_OK) { sqlite3_bind_int64(st,1,(sqlite3_int64)g_cc.my_node_id); sqlite3_bind_text(st,2,req->content_type,-1,SQLITE_STATIC); sqlite3_bind_blob(st,3,req->data,(int)req->data_len,SQLITE_STATIC); sqlite3_bind_int64(st,4,(sqlite3_int64)req->timestamp); - sqlite3_bind_int64(st,5,(sqlite3_int64)dh); - static const uint8_t ch32[32]={0}; sqlite3_bind_blob(st,6,ch32,32,SQLITE_STATIC); - static const uint8_t sig64[64]={0}; sqlite3_bind_blob(st,7,sig64,64,SQLITE_STATIC); + sqlite3_bind_blob(st,5,sig,64,SQLITE_STATIC); sqlite3_step(st); sqlite3_finalize(st); } /* notify GUI anyway */ @@ -408,7 +380,7 @@ int chat_core_chain_hash_at(const char* ch_id, uint32_t pos, uint8_t* hash_out) if (!g_cc.initialized || !hash_out) return -1; char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); char sql[200]; snprintf(sql, sizeof(sql), - "SELECT chain_hash FROM \"%s\" ORDER BY timestamp, datahash ASC LIMIT 1 OFFSET ?", + "SELECT chain_hash FROM \"%s\" ORDER BY timestamp, node_id ASC LIMIT 1 OFFSET ?", tbl); sqlite3_stmt* stmt = NULL; if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) { @@ -1005,16 +977,13 @@ void chat_core_ensure_channel_ready(const char* ch_id) { " content_type TEXT NOT NULL," " data BLOB NOT NULL," " timestamp INTEGER NOT NULL," - " datahash INTEGER NOT NULL," - " chain_hash BLOB NOT NULL," " signature BLOB," " is_outgoing INTEGER DEFAULT 0," " is_read INTEGER DEFAULT 0," - " sync_flags INTEGER DEFAULT 0," - " UNIQUE(timestamp, datahash))", tbl_msg); + " UNIQUE(timestamp, node_id))", tbl_msg); db_exec(sql); snprintf(sql, sizeof(sql), - "CREATE INDEX IF NOT EXISTS \"idx_%s_ts_dh\" ON \"%s\"(timestamp, datahash)", + "CREATE INDEX IF NOT EXISTS \"idx_%s_ts_node\" ON \"%s\"(timestamp, node_id)", tbl_msg, tbl_msg); db_exec(sql); @@ -1134,7 +1103,7 @@ static void si_register(struct DB_SYNC_INSTANCE* si, const char* ch_id) { g_cc.si[g_cc.si_count]=si; g_cc.si_ch_id[g_cc.si_count]=u_strdup(ch_id); g_cc.si_count++; } -static void on_db_sync_insert(struct DB_SYNC_INSTANCE* si, const char* data, size_t len, uint64_t author, void* arg) { +static void on_db_sync_insert(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char* data, size_t len, uint64_t author, void* arg) { const char* ch_id = (const char*)arg; static int insert_count = 0; /* parse JSON: {"n":,"ch":"","ct":"","d":""} */ @@ -1153,37 +1122,31 @@ static void on_db_sync_insert(struct DB_SYNC_INSTANCE* si, const char* data, siz } if (!jd_start) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert JSON parse fail: no '\"d\"' field ch=%s", CC_ID, ch_id); return; } - uint64_t dh; compute_datahash((const uint8_t*)jd_start, jd_len, &dh); - uint64_t ts = db_sync_get_last_timestamp(si); char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); - uint8_t prev_ch[32]; get_last_chain_hash(ch_id, prev_ch); - uint8_t chain_h[32]; compute_chain_hash(prev_ch, (int64_t)ts, dh, chain_h); char sql[512]; snprintf(sql, sizeof(sql), - "INSERT OR IGNORE INTO \"%s\" (node_id,content_type,data,timestamp,datahash,chain_hash,signature,is_outgoing,is_read,sync_flags)" - " VALUES(?,?,?,?,?,?,?,?,?,?)", tbl); + "INSERT OR IGNORE INTO \"%s\" (node_id,content_type,data,timestamp,signature,is_outgoing,is_read)" + " VALUES(?,?,?,?,?,?,?)", tbl); sqlite3_stmt* st=NULL; if (sqlite3_prepare_v2(g_cc.db,sql,-1,&st,NULL)!=SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert prepare failed ch=%s", CC_ID, ch_id); return; } sqlite3_bind_int64(st,1,(sqlite3_int64)jn); sqlite3_bind_text(st,2,jct,(int)jct_len,SQLITE_STATIC); sqlite3_bind_blob(st,3,jd_start,(int)jd_len,SQLITE_STATIC); - sqlite3_bind_int64(st,4,(sqlite3_int64)ts); - sqlite3_bind_int64(st,5,(sqlite3_int64)dh); - sqlite3_bind_blob(st,6,chain_h,32,SQLITE_STATIC); - { static const uint8_t z64[64]={0}; sqlite3_bind_blob(st,7,z64,64,SQLITE_STATIC); } - sqlite3_bind_int(st,8,(uint64_t)jn==g_cc.my_node_id?1:0); - sqlite3_bind_int(st,9,1); sqlite3_bind_int(st,10,0); + sqlite3_bind_int64(st,4,(sqlite3_int64)record_ts); + { static const uint8_t z64[64]={0}; sqlite3_bind_blob(st,5,z64,64,SQLITE_STATIC); } + sqlite3_bind_int(st,6,(uint64_t)jn==g_cc.my_node_id?1:0); + sqlite3_bind_int(st,7,1); int rc=sqlite3_step(st); sqlite3_finalize(st); if (rc==SQLITE_DONE) { insert_count++; if (insert_count <= 3 || insert_count % 10 == 0) - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert #%d ch=%s n=%llu ts=%llu dh=%016llx ct=%.*s", - CC_ID, insert_count, ch_id, (unsigned long long)jn, (unsigned long long)ts, - (unsigned long long)dh, (int)jct_len, jct); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert #%d ch=%s n=%llu ts=%llu ct=%.*s", + CC_ID, insert_count, ch_id, (unsigned long long)jn, (unsigned long long)record_ts, + (int)jct_len, jct); uint8_t evt[65]; uint8_t cl=(uint8_t)strlen(ch_id); evt[0]=cl; memcpy(evt+1,ch_id,cl); gui_bridge_post(GUI_EVT_MSG_RECEIVED, evt, 1+cl); } else if (rc == SQLITE_CONSTRAINT) { - DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert duplicate ch=%s ts=%lld dh=0x%016llx", - CC_ID, ch_id, (long long)ts, (unsigned long long)dh); + DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert duplicate ch=%s ts=%lld", + CC_ID, ch_id, (long long)record_ts); } else { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert step failed rc=%d ch=%s", CC_ID, rc, ch_id);