Browse Source

db_sync: datahash removed, Ed25519 signature mandatory, PK=(timestamp,author_signature), chain_hash includes signature, signed message=ts||json

topo_upd
Evgeny 3 months ago
parent
commit
f0bb03c2cf
  1. 1442
      src/db_sync.c
  2. 28
      src/db_sync.h
  3. 41
      src/secure_channel.c
  4. 3
      src/secure_channel.h
  5. 17
      src/topo_node_sqlite.c
  6. 1
      src/topo_node_sqlite.h
  7. 17
      tests/test_chat_sync_stress.c
  8. 17
      tests/test_db_sync.c
  9. 97
      tools/chatgui/transport/chat_core.c

1442
src/db_sync.c

File diff suppressed because it is too large Load Diff

28
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, путь: <db_path>/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);

41
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) {

3
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);

17
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;
}

1
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

17
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;
}

17
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; }

97
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":<uint64>,"ch":"<str>","ct":"<str>","d":"<str>"} */
@ -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);

Loading…
Cancel
Save