Browse Source

db_sync: verify sig before send + fix total_processed double-count

topo_upd
Evgeny 2 months ago
parent
commit
34404403ef
  1. 26
      src/db_sync.c

26
src/db_sync.c

@ -772,6 +772,7 @@ static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_
} }
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)count); sqlite3_bind_int64(stmt, 1, (sqlite3_int64)count);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)from); sqlite3_bind_int64(stmt, 2, (sqlite3_int64)from);
int bad_del = 0;
while (sqlite3_step(stmt) == SQLITE_ROW) { while (sqlite3_step(stmt) == SQLITE_ROW) {
uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0); uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0);
uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1); uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1);
@ -780,6 +781,19 @@ static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_
uint32_t rdl = (uint32_t)sqlite3_column_bytes(stmt, 3); if (!rd) rdl = 0; 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); 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;
if (rsl == DB_SIG_SIZE && db_verify_author_sig(si, rts, rauth, (const char*)rd, rdl, rsig) != 0) {
uint32_t del_pos = si_find_pos(si, rts, rsig);
sqlite3_stmt* del;
if (si_prep(si, &del, "DELETE FROM \"%s\" WHERE timestamp=? AND author_signature=?") == SQLITE_OK) {
sqlite3_bind_int64(del, 1, (sqlite3_int64)rts);
sqlite3_bind_blob(del, 2, rsig, DB_SIG_SIZE, SQLITE_STATIC);
sqlite3_step(del);
sqlite3_finalize(del);
db_cascade_from(si, del_pos);
bad_del++;
}
continue;
}
int rec_sz = 28 + rdl + 1 + rsl; int rec_sz = 28 + rdl + 1 + rsl;
if (off + rec_sz > 8000) break; if (off + rec_sz > 8000) break;
memcpy(buf + off, &rid, 8); off += 8; memcpy(buf + off, &rid, 8); off += 8;
@ -795,9 +809,9 @@ static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_
*rcp = rc; *rcp = rc;
db_sync_send(si, dst, buf, off); db_sync_send(si, dst, buf, off);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SEND DATA: %u records from pos=%u", DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SEND DATA: %u records from pos=%u%s",
SI_SHRT(si), (unsigned long long)(dst >> 16), rc, from); SI_SHRT(si), (unsigned long long)(dst >> 16), rc, from, bad_del > 0 ? (rc > 0 ? " (bad_sig deleted)" : " (all bad_sig deleted)") : "");
if (rc == 0 && count > 0) { if (rc == 0 && count > 0 && bad_del == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] SEND_DATA batch: loaded 0 records from pos=%u count=%u — SQL error or empty range", DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] SEND_DATA batch: loaded 0 records from pos=%u count=%u — SQL error or empty range",
SI_SHRT(si), (unsigned long long)(dst >> 16), from, count); SI_SHRT(si), (unsigned long long)(dst >> 16), from, count);
} }
@ -916,11 +930,11 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const
if (ret == 0 && si->on_insert) si->on_insert(si, rts, (const char*)rdata, rdlen, src, si->on_insert_arg); if (ret == 0 && si->on_insert) si->on_insert(si, rts, (const char*)rdata, rdlen, src, si->on_insert_arg);
} }
uint16_t total_processed = received + duplicates + bad_sigs; uint16_t total_processed = count - parse_fails;
if (total_processed > 0) { if (total_processed > 0) {
uint32_t rst = (from < old_synced_pos) ? from : old_synced_pos; uint32_t rst = (from < old_synced_pos) ? from : old_synced_pos;
db_cascade_range(si, rst, from + received); db_cascade_range(si, rst, from + total_processed);
if (sp) sp->synced_pos = from + total_processed - 1; if (sp) sp->synced_pos = from + total_processed - 1;
for (int j = 0; j < si->peer_count; j++) { for (int j = 0; j < si->peer_count; j++) {
if (si->peers[j].synced_pos >= rst && &si->peers[j] != sp) if (si->peers[j].synced_pos >= rst && &si->peers[j] != sp)
@ -931,7 +945,7 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const
if (bad_sigs) { size_t el = strlen(extra_buf); snprintf(extra_buf+el, sizeof(extra_buf)-el, "%s%u bad_sig", extra_buf[0]?", ":" (", bad_sigs); if (!extra_buf[0]) extra_buf[0]=' '; } if (bad_sigs) { size_t el = strlen(extra_buf); snprintf(extra_buf+el, sizeof(extra_buf)-el, "%s%u bad_sig", extra_buf[0]?", ":" (", bad_sigs); if (!extra_buf[0]) extra_buf[0]=' '; }
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← DATA: from=%u count=%u vp=%u → inserted %u [id=%llu..%llu]%s cascade=[%u..%u] synced=%u→%u", DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← DATA: from=%u count=%u vp=%u → inserted %u [id=%llu..%llu]%s cascade=[%u..%u] synced=%u→%u",
SI_SHRT(si), (unsigned long long)(src >> 16), from, count, vp, received, (unsigned long long)first_id, (unsigned long long)last_id, SI_SHRT(si), (unsigned long long)(src >> 16), from, count, vp, received, (unsigned long long)first_id, (unsigned long long)last_id,
extra_buf, rst, from+received-1, old_synced_pos, sp?sp->synced_pos:0); extra_buf, rst, from+total_processed-1, old_synced_pos, sp?sp->synced_pos:0);
int relay_count = 0; int relay_count = 0;
uint8_t push_buf[2560]; uint32_t push_off; uint8_t push_buf[2560]; uint32_t push_off;

Loading…
Cancel
Save