@@ -63,10 +63,13 @@ public interface IPowerSyncDatabase : IEventStream<PowerSyncDBEvent>
6363 Task < NonQueryResult > ExecuteBatch ( string query , object ? [ ] [ ] ? parameters = null ) ;
6464
6565 Task < T [ ] > GetAll < T > ( string sql , object ? [ ] ? parameters = null ) ;
66+ Task < dynamic [ ] > GetAll ( string sql , object ? [ ] ? parameters = null ) ;
6667
6768 Task < T ? > GetOptional < T > ( string sql , object ? [ ] ? parameters = null ) ;
69+ Task < dynamic ? > GetOptional ( string sql , object ? [ ] ? parameters = null ) ;
6870
6971 Task < T > Get < T > ( string sql , object ? [ ] ? parameters = null ) ;
72+ Task < dynamic > Get ( string sql , object ? [ ] ? parameters = null ) ;
7073
7174 Task < T > ReadLock < T > ( Func < ILockContext , Task < T > > fn , DBLockOptions ? options = null ) ;
7275
@@ -422,7 +425,8 @@ await Database.WriteTransaction(async tx =>
422425 Closed = true ;
423426 }
424427
425- private record UploadQueueStatsResult ( int size , int count ) ;
428+ private record UploadQueueStatsSizeCountResult ( long size , long count ) ;
429+ private record UploadQueueStatsCountResult ( long count ) ;
426430 /// <summary>
427431 /// Get upload queue size estimate and count.
428432 /// </summary>
@@ -432,23 +436,22 @@ public async Task<UploadQueueStats> GetUploadQueueStats(bool includeSize = false
432436 {
433437 if ( includeSize )
434438 {
435- var result = await tx . Get < UploadQueueStatsResult > (
439+ var result = await tx . Get < UploadQueueStatsSizeCountResult > (
436440 $ "SELECT SUM(cast(data as blob) + 20) as size, count(*) as count FROM { PSInternalTable . CRUD } "
437441 ) ;
438442
439443 return new UploadQueueStats ( result . count , result . size ) ;
440444 }
441445 else
442446 {
443- var result = await tx . Get < UploadQueueStatsResult > (
447+ var result = await tx . Get < UploadQueueStatsCountResult > (
444448 $ "SELECT count(*) as count FROM { PSInternalTable . CRUD } "
445449 ) ;
446450 return new UploadQueueStats ( result . count ) ;
447451 }
448452 } ) ;
449453 }
450454
451-
452455 /// <summary>
453456 /// Get a batch of crud data to upload.
454457 /// <para />
@@ -468,7 +471,7 @@ public async Task<UploadQueueStats> GetUploadQueueStats(bool includeSize = false
468471 /// </summary>
469472 public async Task < CrudBatch ? > GetCrudBatch ( int limit = 100 )
470473 {
471- var crudResult = await GetAll < CrudEntryJSON > ( $ "SELECT id, tx_id, data FROM { PSInternalTable . CRUD } ORDER BY id ASC LIMIT ?", [ limit + 1 ] ) ;
474+ var crudResult = await GetAll < CrudEntryJSON > ( $ "SELECT id, tx_id as transactionId , data FROM { PSInternalTable . CRUD } ORDER BY id ASC LIMIT ?", [ limit + 1 ] ) ;
472475
473476 var all = crudResult . Select ( CrudEntry . FromRow ) . ToList ( ) ;
474477
@@ -510,7 +513,7 @@ public async Task<UploadQueueStats> GetUploadQueueStats(bool includeSize = false
510513 return await ReadTransaction ( async tx =>
511514 {
512515 var first = await tx . GetOptional < CrudEntryJSON > (
513- $ "SELECT id, tx_id, data FROM { PSInternalTable . CRUD } ORDER BY id ASC LIMIT 1") ;
516+ $ "SELECT id, tx_id AS transactionId , data FROM { PSInternalTable . CRUD } ORDER BY id ASC LIMIT 1") ;
514517
515518 if ( first == null )
516519 {
@@ -528,7 +531,7 @@ public async Task<UploadQueueStats> GetUploadQueueStats(bool includeSize = false
528531 else
529532 {
530533 var result = await tx . GetAll < CrudEntryJSON > (
531- $ "SELECT id, tx_id, data FROM { PSInternalTable . CRUD } WHERE tx_id = ? ORDER BY id ASC",
534+ $ "SELECT id, tx_id as transactionId , data FROM { PSInternalTable . CRUD } WHERE tx_id = ? ORDER BY id ASC",
532535 [ txId ] ) ;
533536
534537 all = result . Select ( CrudEntry . FromRow ) . ToList ( ) ;
@@ -596,17 +599,36 @@ public async Task<T[]> GetAll<T>(string query, object?[]? parameters = null)
596599 return await Database . GetAll < T > ( query , parameters ) ;
597600 }
598601
602+ public async Task < dynamic [ ] > GetAll ( string query , object ? [ ] ? parameters = null )
603+ {
604+ await WaitForReady ( ) ;
605+ return await Database . GetAll ( query , parameters ) ;
606+ }
607+
599608 public async Task < T ? > GetOptional < T > ( string query , object ? [ ] ? parameters = null )
600609 {
601610 await WaitForReady ( ) ;
602611 return await Database . GetOptional < T > ( query , parameters ) ;
603612 }
613+
614+ public async Task < dynamic ? > GetOptional ( string query , object ? [ ] ? parameters = null )
615+ {
616+ await WaitForReady ( ) ;
617+ return await Database . GetOptional ( query , parameters ) ;
618+ }
619+
604620 public async Task < T > Get < T > ( string query , object ? [ ] ? parameters = null )
605621 {
606622 await WaitForReady ( ) ;
607623 return await Database . Get < T > ( query , parameters ) ;
608624 }
609625
626+ public async Task < dynamic > Get ( string query , object ? [ ] ? parameters = null )
627+ {
628+ await WaitForReady ( ) ;
629+ return await Database . Get ( query , parameters ) ;
630+ }
631+
610632 public async Task < T > ReadLock < T > ( Func < ILockContext , Task < T > > fn , DBLockOptions ? options = null )
611633 {
612634 await WaitForReady ( ) ;
@@ -650,14 +672,32 @@ public async Task<T> WriteTransaction<T>(Func<ITransaction, Task<T>> fn, DBLockO
650672 /// Source tables are automatically detected using <c>EXPLAIN QUERY PLAN</c>.
651673 /// </summary>
652674 public Task Watch < T > ( string query , object ? [ ] ? parameters , WatchHandler < T > handler , SQLWatchOptions ? options = null )
675+ => WatchInternal ( query , parameters , handler , options , GetAll < T > ) ;
676+
677+ /// <summary>
678+ /// Executes a read query every time the source tables are modified.
679+ /// <para />
680+ /// Use <see cref="SQLWatchOptions.ThrottleMs"/> to specify the minimum interval between queries.
681+ /// Source tables are automatically detected using <c>EXPLAIN QUERY PLAN</c>.
682+ /// </summary>
683+ public Task Watch ( string query , object ? [ ] ? parameters , WatchHandler < dynamic > handler , SQLWatchOptions ? options = null )
684+ => WatchInternal ( query , parameters , handler , options , GetAll ) ;
685+
686+ private Task WatchInternal < T > (
687+ string query ,
688+ object ? [ ] ? parameters ,
689+ WatchHandler < T > handler ,
690+ SQLWatchOptions ? options ,
691+ Func < string , object ? [ ] ? , Task < T [ ] > > getter
692+ )
653693 {
654694 var tcs = new TaskCompletionSource < bool > ( ) ;
655695 Task . Run ( async ( ) =>
656696 {
657697 try
658698 {
659699 var resolvedTables = await ResolveTables ( query , parameters , options ) ;
660- var result = await GetAll < T > ( query , parameters ) ;
700+ var result = await getter ( query , parameters ) ;
661701 handler . OnResult ( result ) ;
662702
663703 OnChange ( new WatchOnChangeHandler
@@ -666,7 +706,7 @@ public Task Watch<T>(string query, object?[]? parameters, WatchHandler<T> handle
666706 {
667707 try
668708 {
669- var result = await GetAll < T > ( query , parameters ) ;
709+ var result = await getter ( query , parameters ) ;
670710 handler . OnResult ( result ) ;
671711 }
672712 catch ( Exception ex )
@@ -691,7 +731,16 @@ public Task Watch<T>(string query, object?[]? parameters, WatchHandler<T> handle
691731 return tcs . Task ;
692732 }
693733
694- private record ExplainedResult ( string opcode , int p2 , int p3 ) ;
734+ private class ExplainedResult
735+ {
736+ public int addr = 0 ;
737+ public string opcode = "" ;
738+ public int p1 = 0 ;
739+ public int p2 = 0 ;
740+ public int p3 = 0 ;
741+ public string p4 = "" ;
742+ public int p5 = 0 ;
743+ }
695744 private record TableSelectResult ( string tbl_name ) ;
696745 public async Task < string [ ] > ResolveTables ( string sql , object ? [ ] ? parameters = null , SQLWatchOptions ? options = null )
697746 {
@@ -718,7 +767,6 @@ public async Task<string[]> ResolveTables(string sql, object?[]? parameters = nu
718767 resolvedTables . Add ( POWERSYNC_TABLE_MATCH . Replace ( table . tbl_name , "" ) ) ;
719768 }
720769 }
721-
722770 return [ .. resolvedTables ] ;
723771 }
724772
0 commit comments