|
|
|
@ -25,6 +25,7 @@ |
|
|
|
// ---- Forward declarations ----
|
|
|
|
// ---- Forward declarations ----
|
|
|
|
struct DB_SYNC; |
|
|
|
struct DB_SYNC; |
|
|
|
struct DB_SYNC_INSTANCE; |
|
|
|
struct DB_SYNC_INSTANCE; |
|
|
|
|
|
|
|
struct SI_PEER; |
|
|
|
|
|
|
|
|
|
|
|
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); |
|
|
|
static void db_sync_on_conn_status(struct ETCP_CONN* conn, int status, void* arg); |
|
|
|
static void db_sync_on_conn_status(struct ETCP_CONN* conn, int status, void* arg); |
|
|
|
@ -38,6 +39,7 @@ static int db_sync_send_gid(struct DB_SYNC* db, uint64_t node_id, uint64_t grou |
|
|
|
static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t peer_node_id); |
|
|
|
static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t peer_node_id); |
|
|
|
static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_t from, uint32_t count, uint32_t vp, uint32_t want_from); |
|
|
|
static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_t from, uint32_t count, uint32_t vp, uint32_t want_from); |
|
|
|
static void si_delivery_update(struct DB_SYNC_INSTANCE* si, uint64_t ts, uint64_t author, uint64_t peer_id); |
|
|
|
static void si_delivery_update(struct DB_SYNC_INSTANCE* si, uint64_t ts, uint64_t author, uint64_t peer_id); |
|
|
|
|
|
|
|
static int db_sync_peer_throttled(struct DB_SYNC* db, struct SI_PEER* p); |
|
|
|
|
|
|
|
|
|
|
|
// ============================================================
|
|
|
|
// ============================================================
|
|
|
|
// Structures
|
|
|
|
// Structures
|
|
|
|
@ -674,6 +676,7 @@ static int db_record_insert_cascade(struct DB_SYNC_INSTANCE* si, |
|
|
|
for (int i = 0; i < si->peer_count; i++) { |
|
|
|
for (int i = 0; i < si->peer_count; i++) { |
|
|
|
uint64_t pid = si->peers[i].node_id; |
|
|
|
uint64_t pid = si->peers[i].node_id; |
|
|
|
if (si->peers[i].sync_state >= 1 && pid != source_peer && pid != my_id) { |
|
|
|
if (si->peers[i].sync_state >= 1 && pid != source_peer && pid != my_id) { |
|
|
|
|
|
|
|
if (db_sync_peer_throttled(si->db_sync, &si->peers[i])) continue; /* спящий пир в «тихом» чате — не пушим */ |
|
|
|
if (db_sync_send(si, pid, pbuf, poff) >= 0) |
|
|
|
if (db_sync_send(si, pid, pbuf, poff) >= 0) |
|
|
|
{ si_delivery_update(si, rts, rauthor, pid); push_count++; } |
|
|
|
{ si_delivery_update(si, rts, rauthor, pid); push_count++; } |
|
|
|
} |
|
|
|
} |
|
|
|
|