diff --git a/crates/core/src/sync/storage_adapter.rs b/crates/core/src/sync/storage_adapter.rs index 6721a92..8d35fb4 100644 --- a/crates/core/src/sync/storage_adapter.rs +++ b/crates/core/src/sync/storage_adapter.rs @@ -59,7 +59,7 @@ impl StorageAdapter { // language=SQLite let progress = - db.prepare_v2("SELECT name, count_at_last, count_since_last FROM ps_buckets")?; + db.prepare_v2("SELECT id, name, count_at_last, count_since_last FROM ps_buckets")?; // language=SQLite let time = db.prepare_v2("SELECT CAST(unixepoch('subsec') * 1000000 as integer)")?; @@ -192,11 +192,13 @@ WHERE bucket = ?1", pub fn step_progress(&'_ self) -> Result>> { if self.progress_stmt.step()? { - let bucket = self.progress_stmt.column_text(0)?; - let count_at_last = self.progress_stmt.column_int64(1); - let count_since_last = self.progress_stmt.column_int64(2); + let bucket_id = self.progress_stmt.column_int64(0); + let bucket = self.progress_stmt.column_text(1)?; + let count_at_last = self.progress_stmt.column_int64(2); + let count_since_last = self.progress_stmt.column_int64(3); Ok(Some(PersistedBucketProgress { + bucket_id, bucket, count_at_last, count_since_last, @@ -208,12 +210,6 @@ WHERE bucket = ?1", } } - pub fn reset_progress(&self) -> Result<()> { - self.db - .exec_safe(c"UPDATE ps_buckets SET count_since_last = 0, count_at_last = 0;")?; - Ok(()) - } - pub fn lookup_bucket(&self, bucket: &str) -> Result { // We do an ON CONFLICT UPDATE simply so that the RETURNING bit works for existing rows. // We can consider splitting this into separate SELECT and INSERT statements. @@ -681,6 +677,7 @@ pub enum SyncLocalResult { /// operations have been inserted in the meantime. pub struct PersistedBucketProgress<'a> { pub bucket: &'a str, + pub bucket_id: i64, pub count_at_last: i64, pub count_since_last: i64, } diff --git a/crates/core/src/sync/streaming_sync.rs b/crates/core/src/sync/streaming_sync.rs index 6994bb4..ab69f98 100644 --- a/crates/core/src/sync/streaming_sync.rs +++ b/crates/core/src/sync/streaming_sync.rs @@ -36,7 +36,7 @@ use super::{ line::{Checkpoint, CheckpointDiff, SyncLine}, operations::insert_bucket_operations, storage_adapter::{StorageAdapter, SyncLocalResult}, - sync_status::{SyncDownloadProgress, SyncProgressFromCheckpoint, SyncStatusContainer}, + sync_status::{SyncDownloadProgress, SyncStatusContainer}, }; /// The sync client implementation, responsible for parsing lines received by the sync service and @@ -490,16 +490,7 @@ impl StreamingSyncIteration { } fn load_progress(&self, checkpoint: &OwnedCheckpoint) -> Result { - let SyncProgressFromCheckpoint { - progress, - needs_counter_reset, - } = SyncDownloadProgress::for_checkpoint(checkpoint, &self.adapter)?; - - if needs_counter_reset { - self.adapter.reset_progress()?; - } - - Ok(progress) + SyncDownloadProgress::for_checkpoint(checkpoint, &self.adapter) } fn try_applying_write_after_completed_upload<'a>( diff --git a/crates/core/src/sync/sync_status.rs b/crates/core/src/sync/sync_status.rs index a100aa5..0d3d26f 100644 --- a/crates/core/src/sync/sync_status.rs +++ b/crates/core/src/sync/sync_status.rs @@ -296,6 +296,8 @@ pub struct BucketProgress { pub at_last: i64, pub since_last: i64, pub target_count: i64, + #[serde(skip_serializing)] + pub reset_counter: bool, } #[derive(Hash)] @@ -336,6 +338,7 @@ impl Serialize for SyncDownloadProgress { at_last: 0, since_last: progress.downloaded, target_count: progress.total, + reset_counter: false, // ignored }, )?; } @@ -349,18 +352,12 @@ impl Serialize for SyncDownloadProgress { } } -pub struct SyncProgressFromCheckpoint { - pub progress: SyncDownloadProgress, - pub needs_counter_reset: bool, -} - impl SyncDownloadProgress { pub fn for_checkpoint<'a>( checkpoint: &OwnedCheckpoint, adapter: &StorageAdapter, - ) -> Result { + ) -> Result { let mut buckets = BTreeMap::::new(); - let mut needs_reset = false; for bucket in checkpoint.buckets.values() { buckets.insert( bucket.bucket.clone(), @@ -370,6 +367,7 @@ impl SyncDownloadProgress { // Will be filled out later by iterating local_progress at_last: 0, since_last: 0, + reset_counter: false, }, ); } @@ -377,6 +375,9 @@ impl SyncDownloadProgress { // Ignore errors here - SQLite seems to report errors from an earlier statement iteration // sometimes. let _ = adapter.progress_stmt.reset(); + let reset_progress = adapter.db.prepare_v2( + "UPDATE ps_buckets SET count_since_last = 0, count_at_last = 0 WHERE id = ?;", + )?; // Go through local bucket states to detect pending progress from previous sync iterations // that may have been interrupted. @@ -389,24 +390,21 @@ impl SyncDownloadProgress { progress.since_last = row.count_since_last; if progress.target_count < row.count_at_last + row.count_since_last { - needs_reset = true; - // Either due to a defrag / sync rule deploy or a compactioon operation, the size + // Either due to a defrag / sync rule deploy or a compaction operation, the size // of the bucket shrank so much that the local ops exceed the ops in the updated // bucket. We can't possibly report progress in this case (it would overshoot 100%). - for (_, progress) in &mut buckets { - progress.at_last = 0; - progress.since_last = 0; - } - break; + progress.reset_counter = true; + progress.at_last = 0; + progress.since_last = 0; + + reset_progress.bind_int64(1, row.bucket_id)?; + reset_progress.exec()?; } } adapter.progress_stmt.reset()?; - Ok(SyncProgressFromCheckpoint { - progress: Self { buckets }, - needs_counter_reset: needs_reset, - }) + Ok(Self { buckets }) } pub fn increment_download_count(&mut self, line: &DataLine) { diff --git a/dart/test/sync_test.dart b/dart/test/sync_test.dart index b7c969b..202ffe9 100644 --- a/dart/test/sync_test.dart +++ b/dart/test/sync_test.dart @@ -1226,27 +1226,37 @@ void _syncTests({ test('interrupt and defrag', () { applyInstructions(invokeControl('start', null)); applyInstructions(pushCheckpoint( - buckets: [bucketDescription('a', count: 10)], lastOpId: 10)); - expect(totalProgress(), (0, 10)); + buckets: [ + bucketDescription('a', count: 10), + bucketDescription('b', count: 5), + ], + lastOpId: 10, + )); + expect(totalProgress(), (0, 15)); pushSyncData('a', 5); - expect(totalProgress(), (5, 10)); + pushSyncData('b', 4); + expect(totalProgress(), (9, 15)); // Emulate stream closing applyInstructions(invokeControl('stop', null)); expect(progress, isNull); applyInstructions(invokeControl('start', null)); - // A defrag in the meantime shrank the bucket. - applyInstructions(pushCheckpoint( - buckets: [bucketDescription('a', count: 4)], lastOpId: 14)); - // So we shouldn't report 5/4. - expect(totalProgress(), (0, 4)); + // A defrag in the meantime shrank bucket a. + applyInstructions(pushCheckpoint(buckets: [ + bucketDescription('a', count: 4), + bucketDescription('b', count: 5), + ], lastOpId: 14)); + // The progress in a should no longer count, e.g. we shouldn't report 9/9. + expect(totalProgress(), (4, 9)); // This should also reset the persisted progress counters. - final [bucket] = db.select('SELECT * FROM ps_buckets'); - expect(bucket, containsPair('count_since_last', 0)); - expect(bucket, containsPair('count_at_last', 0)); + final [a, b] = db.select('SELECT * FROM ps_buckets ORDER BY name'); + expect(a, containsPair('count_since_last', 0)); + expect(a, containsPair('count_at_last', 0)); + expect(b, containsPair('count_since_last', 4)); + expect(b, containsPair('count_at_last', 0)); }); test('different priorities', () {