Browse Source

db_sync: work-driven peer_check timer (stops when all peers synced)

topo_upd
evgeny 2 months ago
parent
commit
a4e73acceb
  1. 104
      src/chat/db_sync.c

104
src/chat/db_sync.c

@ -27,6 +27,7 @@ static void db_sync_on_conn_status(struct ETCP_CONN* conn, int status, void* arg
static void db_sync_on_conn_up(struct ETCP_CONN* conn, void* arg);
static void db_sync_on_conn_down(struct ETCP_CONN* conn, void* arg);
static void db_sync_peer_check_cb(void* arg);
static void db_sync_resume_peer_check(struct DB_SYNC* db);
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);
@ -1347,6 +1348,7 @@ static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t pid)
uint8_t msg[1] = { DB_MSG_REQUEST_SYNC };
db_sync_send(si, pid, msg, 1);
}
db_sync_resume_peer_check(si->db_sync);
}
// Проверяет целостность цепочки хешей всей таблицы, при несовпадении запускает пересчёт с позиции ошибки
@ -1457,61 +1459,81 @@ int db_sync_last_chain_hash8(struct DB_SYNC_INSTANCE* si, uint64_t* out_hash8)
// Timers
// ============================================================
// Есть ли у любого активного инстанса пир, требующий синхронизации (sync_state != 2)
static int db_sync_has_pending_work(struct DB_SYNC* db)
{
if (!db || !db->enabled) return 0;
for (int i = 0; i < db->instance_count; i++) {
struct DB_SYNC_INSTANCE* si = &db->instances[i];
if (!si->enabled) continue;
for (int j = 0; j < si->peer_count; j++)
if (si->peers[j].sync_state != 2) return 1;
}
return 0;
}
// Взвести work-driven таймер проверки пиров, если он ещё не взведён и модуль активен
static void db_sync_resume_peer_check(struct DB_SYNC* db)
{
if (!db || !db->enabled || db->peer_check_timer) return;
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");
}
// Периодическая проверка: запускает синхронизацию для пиров в состоянии sync_state==0,
// детектирует таймауты синхронизации (DB_SYNC_SYNC_TIMEOUT секунд), логирует зависшие
// детектирует таймауты синхронизации (DB_SYNC_SYNC_TIMEOUT секунд), логирует зависшие.
// Work-driven: перевзводится только пока есть пиры с sync_state != 2.
static void db_sync_peer_check_cb(void* arg)
{
struct DB_SYNC* db = (struct DB_SYNC*)arg;
if (!db || !db->enabled) return;
db->peer_check_timer = NULL;
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "instances=%d", db->instance_count);
struct TOPO_GROUP* g = topo_groups_get_default(db->inst->topo_groups);
if (!g) {
db->peer_check_timer = uasync_set_timeout(db->inst->ua, DB_SYNC_PEER_CHECK_INTERVAL * 10000u,
db, db_sync_peer_check_cb, "db_sync_peer");
return;
}
int total_synced = 0, total_skipped = 0, any_peers = 0;
char launched_list[256] = "";
for (int i = 0; i < db->instance_count; i++) {
struct DB_SYNC_INSTANCE* si = &db->instances[i];
if (!si->enabled) continue;
if (g) {
for (int i = 0; i < db->instance_count; i++) {
struct DB_SYNC_INSTANCE* si = &db->instances[i];
if (!si->enabled) continue;
int peers_found = 0;
struct ll_entry* e = g->senders_list->head;
while (e) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data;
if (item->conn && item->conn->peer_node_id != 0 && item->conn->links_up && item->conn->initialized) {
uint64_t pid = item->conn->peer_node_id;
if (pid != db->inst->node_id) { si_peer_add(si, pid); peers_found++; any_peers = 1; }
int peers_found = 0;
struct ll_entry* e = g->senders_list->head;
while (e) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data;
if (item->conn && item->conn->peer_node_id != 0 && item->conn->links_up && item->conn->initialized) {
uint64_t pid = item->conn->peer_node_id;
if (pid != db->inst->node_id) { si_peer_add(si, pid); peers_found++; any_peers = 1; }
}
e = e->next;
}
e = e->next;
}
struct SI_PEER* best = NULL;
uint32_t min_pos = UINT32_MAX;
for (int j = 0; j < si->peer_count; j++) {
if (si->peers[j].sync_state == 0 && si->peers[j].synced_pos < min_pos) {
min_pos = si->peers[j].synced_pos;
best = &si->peers[j];
struct SI_PEER* best = NULL;
uint32_t min_pos = UINT32_MAX;
for (int j = 0; j < si->peer_count; j++) {
if (si->peers[j].sync_state == 0 && si->peers[j].synced_pos < min_pos) {
min_pos = si->peers[j].synced_pos;
best = &si->peers[j];
}
}
}
if (best) {
uint8_t ek[32];
if (db_get_ed25519_pubkey(db, best->node_id, ek) == 0) {
best->sync_state = 1;
db_sync_initiate_sync(si, best->node_id); total_synced++;
{ size_t tl = strlen(launched_list);
snprintf(launched_list + tl, sizeof(launched_list) - tl,
"%s%s:%04llX[sp=%u]", tl ? "," : "",
SI_SHRT(si), (unsigned long long)(best->node_id >> 16), min_pos); }
} else {
DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "peer_check skip peer=%016llx — no Ed25519 pubkey yet", (unsigned long long)best->node_id);
best = NULL; total_skipped++;
if (best) {
uint8_t ek[32];
if (db_get_ed25519_pubkey(db, best->node_id, ek) == 0) {
best->sync_state = 1;
db_sync_initiate_sync(si, best->node_id); total_synced++;
{ size_t tl = strlen(launched_list);
snprintf(launched_list + tl, sizeof(launched_list) - tl,
"%s%s:%04llX[sp=%u]", tl ? "," : "",
SI_SHRT(si), (unsigned long long)(best->node_id >> 16), min_pos); }
} else {
DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "peer_check skip peer=%016llx — no Ed25519 pubkey yet", (unsigned long long)best->node_id);
best = NULL; total_skipped++;
}
}
else { total_skipped++; }
}
else { total_skipped++; }
}
{
@ -1542,8 +1564,10 @@ static void db_sync_peer_check_cb(void* arg)
db->instance_count, total_synced, launched_list);
}
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");
if (db_sync_has_pending_work(db))
db_sync_resume_peer_check(db);
else
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "all peers synced — peer_check stopped");
}
// Периодическая TTL-очистка: удаляет неподтверждённые записи (flags&1==0) старше db_sync_ttl от локального узла

Loading…
Cancel
Save