@@ -89,7 +89,166 @@ public string GetMaxOpId()
8989 {
9090 return MAX_OP_ID ;
9191 }
92-
92+
93+ public void StartSession ( ) { }
94+
95+ public async Task < BucketState [ ] > GetBucketStates ( )
96+ {
97+ return
98+ await db . GetAll < BucketState > ( "SELECT name as bucket, cast(last_op as TEXT) as op_id FROM ps_buckets WHERE pending_delete = 0 AND name != '$local'" ) ;
99+ }
100+
101+ public async Task SaveSyncData ( SyncDataBatch batch )
102+ {
103+ await db . WriteTransaction ( async tx =>
104+ {
105+ int count = 0 ;
106+ foreach ( var b in batch . Buckets )
107+ {
108+ var result = await tx . Execute ( "INSERT INTO powersync_operations(op, data) VALUES(?, ?)" ,
109+ [ "save" , JsonConvert . SerializeObject ( new { buckets = new [ ] { JsonConvert . DeserializeObject ( b . ToJSON ( ) ) } } ) ] ) ;
110+ logger . LogDebug ( "saveSyncData {message}" , JsonConvert . SerializeObject ( result ) ) ;
111+ count += b . Data . Length ;
112+ }
113+ compactCounter += count ;
114+ } ) ;
115+ }
116+
117+ public async Task RemoveBuckets ( string [ ] buckets )
118+ {
119+ foreach ( var bucket in buckets )
120+ {
121+ await DeleteBucket ( bucket ) ;
122+ }
123+ }
124+
125+ private async Task DeleteBucket ( string bucket )
126+ {
127+ await db . WriteTransaction ( async tx =>
128+ {
129+ await tx . Execute ( "INSERT INTO powersync_operations(op, data) VALUES(?, ?)" ,
130+ [ "delete_bucket" , bucket ] ) ;
131+ } ) ;
132+
133+ logger . LogDebug ( "Done deleting bucket" ) ;
134+ pendingBucketDeletes = true ;
135+ }
136+
137+ private record LastSyncedResult ( string ? synced_at ) ;
138+ public async Task < bool > HasCompletedSync ( )
139+ {
140+ if ( hasCompletedSync ) return true ;
141+
142+ var result = await db . Get < LastSyncedResult > ( "SELECT powersync_last_synced_at() as synced_at" ) ;
143+
144+ hasCompletedSync = result . synced_at != null ;
145+ return hasCompletedSync ;
146+ }
147+
148+ public async Task < SyncLocalDatabaseResult > SyncLocalDatabase ( Checkpoint checkpoint )
149+ {
150+ var validation = await ValidateChecksums ( checkpoint ) ;
151+ if ( ! validation . CheckpointValid )
152+ {
153+ logger . LogError ( "Checksums failed for {failures}" , JsonConvert . SerializeObject ( validation . CheckpointFailures ) ) ;
154+ foreach ( var failedBucket in validation . CheckpointFailures ?? [ ] )
155+ {
156+ await DeleteBucket ( failedBucket ) ;
157+ }
158+ return new SyncLocalDatabaseResult
159+ {
160+ Ready = false ,
161+ CheckpointValid = false ,
162+ CheckpointFailures = validation . CheckpointFailures
163+ } ;
164+ }
165+
166+ var bucketNames = checkpoint . Buckets . Select ( b => b . Bucket ) . ToArray ( ) ;
167+ await db . WriteTransaction ( async tx =>
168+ {
169+ await tx . Execute (
170+ "UPDATE ps_buckets SET last_op = ? WHERE name IN (SELECT json_each.value FROM json_each(?))" ,
171+ [ checkpoint . LastOpId , JsonConvert . SerializeObject ( bucketNames ) ]
172+ ) ;
173+
174+ if ( checkpoint . WriteCheckpoint != null )
175+ {
176+ await tx . Execute (
177+ "UPDATE ps_buckets SET last_op = ? WHERE name = '$local'" ,
178+ [ checkpoint . WriteCheckpoint ]
179+ ) ;
180+ }
181+ } ) ;
182+
183+ var valid = await UpdateObjectsFromBuckets ( checkpoint ) ;
184+ if ( ! valid )
185+ {
186+ logger . LogDebug ( "Not at a consistent checkpoint - cannot update local db" ) ;
187+ return new SyncLocalDatabaseResult
188+ {
189+ Ready = false ,
190+ CheckpointValid = true
191+ } ;
192+ }
193+
194+ await ForceCompact ( ) ;
195+
196+ return new SyncLocalDatabaseResult
197+ {
198+ Ready = true ,
199+ CheckpointValid = true
200+ } ;
201+ }
202+
203+ private async Task < bool > UpdateObjectsFromBuckets ( Checkpoint checkpoint )
204+ {
205+ return await db . WriteTransaction ( async tx =>
206+ {
207+ var result = await tx . Execute ( "INSERT INTO powersync_operations(op, data) VALUES(?, ?)" ,
208+ [ "sync_local" , "" ] ) ;
209+
210+ return result . InsertId == 1 ;
211+ } ) ;
212+ }
213+
214+ private record ResultResult ( object result ) ;
215+
216+ public class ResultDetail
217+ {
218+ [ JsonProperty ( "valid" ) ]
219+ public bool Valid { get ; set ; }
220+
221+ [ JsonProperty ( "failed_buckets" ) ]
222+ public List < string > ? FailedBuckets { get ; set ; }
223+ }
224+
225+ public async Task < SyncLocalDatabaseResult > ValidateChecksums (
226+ Checkpoint checkpoint )
227+ {
228+ var result = await db . Get < ResultResult > ( "SELECT powersync_validate_checkpoint(?) as result" ,
229+ [ JsonConvert . SerializeObject ( checkpoint ) ] ) ;
230+
231+ logger . LogDebug ( "validateChecksums result item {message}" , JsonConvert . SerializeObject ( result ) ) ;
232+
233+ if ( result == null ) return new SyncLocalDatabaseResult { CheckpointValid = false , Ready = false } ;
234+
235+ var resultDetail = JsonConvert . DeserializeObject < ResultDetail > ( result . result . ToString ( ) ?? "{}" ) ;
236+
237+ if ( resultDetail ? . Valid == true )
238+ {
239+ return new SyncLocalDatabaseResult { Ready = true , CheckpointValid = true } ;
240+ }
241+ else
242+ {
243+ return new SyncLocalDatabaseResult
244+ {
245+ CheckpointValid = false ,
246+ Ready = false ,
247+ CheckpointFailures = resultDetail ? . FailedBuckets ? . ToArray ( ) ?? [ ]
248+ } ;
249+ }
250+ }
251+
93252 /// <summary>
94253 /// Force a compact operation, primarily for testing purposes.
95254 /// </summary>
@@ -267,16 +426,15 @@ public async Task<bool> HasCrud()
267426 {
268427 return await db . GetOptional < object > ( "SELECT 1 as ignore FROM ps_crud LIMIT 1" ) != null ;
269428 }
270-
271429
272430 record ControlResult ( string ? r ) ;
273431
274432 public async Task < string > Control ( string op , object ? payload = null )
275433 {
276434 return await db . WriteTransaction ( async tx =>
277435 {
278- var result = await tx . Get < ControlResult > ( "SELECT powersync_control(?, ?) AS r" , [ op , payload ?? "" ] ) ;
436+ var result = await tx . Get < ControlResult > ( "SELECT powersync_control(?, ?) AS r" , [ op , payload ] ) ;
279437 return result . r ! ;
280438 } ) ;
281439 }
282- }
440+ }
0 commit comments