|
|
|
@ -429,7 +429,18 @@ static int db_record_insert(struct DB_SYNC_INSTANCE* si, |
|
|
|
int rc; |
|
|
|
int rc; |
|
|
|
|
|
|
|
|
|
|
|
// Verify author signature (mandatory)
|
|
|
|
// Verify author signature (mandatory)
|
|
|
|
if (db_verify_author_sig(si, ts, author_node_id, json, jlen, author_sig) != 0) return -2; |
|
|
|
if (db_verify_author_sig(si, ts, author_node_id, json, jlen, author_sig) != 0) { |
|
|
|
|
|
|
|
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)ts); |
|
|
|
|
|
|
|
sqlite3_bind_blob(del, 2, author_sig, DB_SIG_SIZE, SQLITE_STATIC); |
|
|
|
|
|
|
|
sqlite3_step(del); |
|
|
|
|
|
|
|
int chg = sqlite3_changes(SI_DB(si)); |
|
|
|
|
|
|
|
sqlite3_finalize(del); |
|
|
|
|
|
|
|
if (chg > 0) db_cascade_from(si, si_find_pos(si, ts, author_sig)); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
return -2; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
rc = sqlite3_exec(db, "BEGIN IMMEDIATE", NULL, NULL, NULL); |
|
|
|
rc = sqlite3_exec(db, "BEGIN IMMEDIATE", NULL, NULL, NULL); |
|
|
|
if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "insert BEGIN: %s", sqlite3_errmsg(db)); return -1; } |
|
|
|
if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "insert BEGIN: %s", sqlite3_errmsg(db)); return -1; } |
|
|
|
@ -905,10 +916,12 @@ 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); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
if (received > 0) { |
|
|
|
uint16_t total_processed = received + duplicates + bad_sigs; |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
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 + received); |
|
|
|
if (sp) sp->synced_pos = from + received - 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) |
|
|
|
si->peers[j].synced_pos = rst > 0 ? rst - 1 : 0; |
|
|
|
si->peers[j].synced_pos = rst > 0 ? rst - 1 : 0; |
|
|
|
@ -1051,6 +1064,13 @@ static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t src, const |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
if (mc < pc) { |
|
|
|
if (mc < pc) { |
|
|
|
|
|
|
|
if (sp && sp->synced_pos + 1 >= pc) { |
|
|
|
|
|
|
|
if (sp) { sp->synced_pos = pc - 1; sp->sync_state = 2; } |
|
|
|
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, |
|
|
|
|
|
|
|
"sync [%s:%04llX] ← SYNC_DONE: count diff my=%u peer=%u but synced=%u covers tail — accept", |
|
|
|
|
|
|
|
SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, sp ? sp->synced_pos : 0); |
|
|
|
|
|
|
|
return; |
|
|
|
|
|
|
|
} |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, |
|
|
|
"sync [%s:%04llX] ← SYNC_DONE: count diff my=%u peer=%u — requesting tail from=%u", |
|
|
|
"sync [%s:%04llX] ← SYNC_DONE: count diff my=%u peer=%u — requesting tail from=%u", |
|
|
|
SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, mc); |
|
|
|
SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, mc); |
|
|
|
|