diff --git a/src/db_sync.c b/src/db_sync.c index aacce73b..2f3c0768 100644 --- a/src/db_sync.c +++ b/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, 2, (sqlite3_int64)from); + int bad_del = 0; while (sqlite3_step(stmt) == SQLITE_ROW) { uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0); 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; const uint8_t* rsig = (const uint8_t*)sqlite3_column_blob(stmt, 4); 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; if (off + rec_sz > 8000) break; 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; db_sync_send(si, dst, buf, off); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SEND DATA: %u records from pos=%u", - SI_SHRT(si), (unsigned long long)(dst >> 16), rc, from); - if (rc == 0 && count > 0) { + 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, bad_del > 0 ? (rc > 0 ? " (bad_sig deleted)" : " (all bad_sig deleted)") : ""); + 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", 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); } - uint16_t total_processed = received + duplicates + bad_sigs; + uint16_t total_processed = count - parse_fails; if (total_processed > 0) { 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; for (int j = 0; j < si->peer_count; j++) { 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]=' '; } 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, - 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; uint8_t push_buf[2560]; uint32_t push_off;