Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 7 additions & 10 deletions crates/core/src/sync/storage_adapter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)")?;
Expand Down Expand Up @@ -192,11 +192,13 @@ WHERE bucket = ?1",

pub fn step_progress(&'_ self) -> Result<Option<PersistedBucketProgress<'_>>> {
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,
Expand All @@ -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<BucketInfo> {
// 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.
Expand Down Expand Up @@ -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,
}
13 changes: 2 additions & 11 deletions crates/core/src/sync/streaming_sync.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -490,16 +490,7 @@ impl StreamingSyncIteration {
}

fn load_progress(&self, checkpoint: &OwnedCheckpoint) -> Result<SyncDownloadProgress> {
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>(
Expand Down
34 changes: 16 additions & 18 deletions crates/core/src/sync/sync_status.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Expand Down Expand Up @@ -336,6 +338,7 @@ impl Serialize for SyncDownloadProgress {
at_last: 0,
since_last: progress.downloaded,
target_count: progress.total,
reset_counter: false, // ignored
},
)?;
}
Expand All @@ -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<SyncProgressFromCheckpoint, PowerSyncError> {
) -> Result<Self, PowerSyncError> {
let mut buckets = BTreeMap::<String, BucketProgress>::new();
let mut needs_reset = false;
for bucket in checkpoint.buckets.values() {
buckets.insert(
bucket.bucket.clone(),
Expand All @@ -370,13 +367,17 @@ impl SyncDownloadProgress {
// Will be filled out later by iterating local_progress
at_last: 0,
since_last: 0,
reset_counter: false,
},
);
}

// 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.
Expand All @@ -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) {
Expand Down
32 changes: 21 additions & 11 deletions dart/test/sync_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -1226,27 +1226,37 @@ void _syncTests<T>({
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', () {
Expand Down
Loading