You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

137 lines
7.2 KiB

// db_sync.h — Distributed content-addressed table with SQLite + peer sync via etcp_api (direct P2P)
//
// Назначение: децентрализованная реплицируемая таблица 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. 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() — вызывается автоматически при завершении
//
// Синхронизация:
// - При поднятии ETCP-соединения с пиром для каждого инстанса запускается полная синхронизация
// - Каждое сообщение содержит instance_hash (первые 64 бита SHA256(name||id_be)),
// что позволяет маршрутизировать сообщения к нужному инстансу на приёмной стороне
// - Если пир не имеет инстанса с таким хешем — возвращает DB_MSG_ERROR(0x01)
// - Если пир добавляет инстанс позже — сам инициирует sync в нашу сторону
// - Протокол: сравнение 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)
// - Дубликаты определяются по (timestamp, author)
// - БД хранится в SQLite, путь: <db_path>/sync
// - Синхронизация идёт через etcp_api (service ID 0x20), только прямые P2P соединения
#ifndef DB_SYNC_H
#define DB_SYNC_H
#ifdef __cplusplus
extern "C" {
#endif
#include <stdint.h>
#include <stddef.h>
struct UTUN_INSTANCE;
struct DB_SYNC_INSTANCE;
// etcp_router service ID
#define ETCP_RT_ID_DB_SYNC 0x20
// Message types
#define DB_MSG_INIT_SYNC 0x01
#define DB_MSG_INIT_RESP 0x02
#define DB_MSG_SEND_DATA 0x04
#define DB_MSG_PUSH 0x05
#define DB_MSG_ACK_PUSH 0x06
#define DB_MSG_SYNC_DONE 0x07
#define DB_MSG_ERROR 0x08
// Error codes for DB_MSG_ERROR
#define DB_ERR_NOT_FOUND 0x01 // instance not found
#define DB_ERR_DISABLED 0x02 // instance exists but disabled
// Defaults
#define DB_SYNC_DEFAULT_TTL 86400
#define DB_SYNC_DEFAULT_MAPSIZE (100UL * 1024 * 1024)
#define DB_SYNC_PEER_CHECK_INTERVAL 5
#define DB_SYNC_SYNC_TIMEOUT 15
#define DB_SYNC_TTL_INTERVAL 3600
// Record flags
#define DB_REC_FLAG_WAS_SENT 0x01
// Sync protocol constants
#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);
// Instance management
struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t hash, int auto_sync);
void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si);
// Data operations (per-instance)
// 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, uint64_t ts,
const char* local_attrs);
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, 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,
const uint8_t* author_sig, size_t sig_len,
int delivered_peers, const char* delivery_chain);
int db_sync_select(struct DB_SYNC_INSTANCE* si, uint32_t offset, uint32_t limit,
db_sync_select_cb cb, void* arg);
// 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);
// Verify chain integrity: returns 0 if all chain_hashes are correct, 1 if any mismatch found
int db_sync_chain_verify(struct DB_SYNC_INSTANCE* si);
// Peer state control (for testing — disable PUSH to isolate sync protocol)
void db_sync_peer_set_state(struct DB_SYNC_INSTANCE* si, uint64_t node_id, uint8_t state);
// Force re-initiate sync to a specific peer (for testing)
void db_sync_reinitiate(struct DB_SYNC_INSTANCE* si, uint64_t node_id);
// Get last chain hash8 (for cross-peer consistency check in tests)
int db_sync_last_chain_hash8(struct DB_SYNC_INSTANCE* si, uint64_t* out_hash8);
#ifdef __cplusplus
}
#endif
#endif // DB_SYNC_H