11namespace PowerSync . Common . Client ;
22
3- using System . Runtime . CompilerServices ;
4- using System . Text . RegularExpressions ;
5- using System . Threading . Channels ;
63using System . Threading . Tasks ;
74
85using Microsoft . Data . Sqlite ;
96using Microsoft . Extensions . Logging ;
107using Microsoft . Extensions . Logging . Abstractions ;
118
129using Newtonsoft . Json ;
10+
1311using Nito . AsyncEx ;
14- using ThrottleDebounce ;
1512
1613using PowerSync . Common . Client . Connection ;
1714using PowerSync . Common . Client . Sync . Bucket ;
@@ -133,9 +130,6 @@ public class PowerSyncDatabase : IPowerSyncDatabase
133130 public IDBAdapter Database { get ; protected set ; }
134131 private Schema schema ;
135132
136- private const int DEFAULT_WATCH_THROTTLE_MS = 30 ;
137- private static readonly Regex POWERSYNC_TABLE_MATCH = new Regex ( @"(^ps_data__|^ps_data_local__)" , RegexOptions . Compiled ) ;
138-
139133 public bool Closed { get ; protected set ; }
140134 public bool Ready { get ; protected set ; }
141135
@@ -146,6 +140,8 @@ public class PowerSyncDatabase : IPowerSyncDatabase
146140
147141 private readonly InternalSubscriptionManager subscriptions ;
148142
143+ private readonly WatchManager watchManager ;
144+
149145 private StreamingSyncImplementation ? syncStreamImplementation ;
150146 public string SdkVersion { get ; protected set ; }
151147
@@ -205,6 +201,8 @@ public PowerSyncDatabase(PowerSyncDatabaseOptions options)
205201
206202 remoteFactory = options . RemoteFactory ?? ( connector => new Remote ( connector ) ) ;
207203
204+ watchManager = new WatchManager ( this , masterCts . Token ) ;
205+
208206 // Start async init
209207 subscriptions = new InternalSubscriptionManager (
210208 firstStatusMatching : WaitForStatus ,
@@ -346,7 +344,9 @@ protected async Task Initialize(PowerSyncDatabaseOptions options)
346344 await LoadVersion ( ) ;
347345 await Database . WriteTransaction ( tx => tx . Execute ( "SELECT powersync_init()" ) ) ;
348346
349- await UpdateSchema ( options . Schema ) ;
347+ // Note: no SchemaChangedEvent here - watched queries only resolve their source tables
348+ // once initialization has completed, so the initial schema is never stale for them.
349+ await ReplaceSchema ( options . Schema ) ;
350350 await ResolveOfflineSyncStatus ( ) ;
351351 await Database . Execute ( "PRAGMA RECURSIVE_TRIGGERS=TRUE" ) ;
352352 Ready = true ;
@@ -421,6 +421,12 @@ public async Task UpdateSchema(Schema schema)
421421 throw new Exception ( "Cannot update schema while connected" ) ;
422422 }
423423
424+ await ReplaceSchema ( schema ) ;
425+ Events . Emit ( new PowerSyncDBEvents . SchemaChangedEvent ( schema ) ) ;
426+ }
427+
428+ private async Task ReplaceSchema ( Schema schema )
429+ {
424430 try
425431 {
426432 schema . Validate ( ) ;
@@ -434,7 +440,6 @@ public async Task UpdateSchema(Schema schema)
434440 await Database . WriteTransaction ( tx =>
435441 tx . Execute ( "SELECT powersync_replace_schema(?)" , [ JsonConvert . SerializeObject ( schema ) ] ) ) ;
436442 await Database . RefreshSchema ( ) ;
437- Events . Emit ( new PowerSyncDBEvents . SchemaChangedEvent ( schema ) ) ;
438443 }
439444
440445 /// <summary>
@@ -763,192 +768,24 @@ public async Task<T> WriteTransaction<T>(Func<ITransaction, Task<T>> fn, DBLockO
763768 return await Database . WriteTransaction ( fn , options ) ;
764769 }
765770
766- public IAsyncEnumerable < WatchOnChangeEvent > OnChange ( SQLWatchOptions ? options = null )
767- {
768- options ??= new SQLWatchOptions ( ) ;
769-
770- var tables = options ? . Tables ?? [ ] ;
771- var powersyncTables = new HashSet < string > (
772- tables . SelectMany ( table => new [ ] { $ "ps_data__{ table } ", $ "ps_data_local__{ table } " } )
773- ) ;
774-
775- var signal = options ? . Signal != null
776- ? CancellationTokenSource . CreateLinkedTokenSource ( masterCts . Token , options . Signal . Value )
777- : CancellationTokenSource . CreateLinkedTokenSource ( masterCts . Token ) ;
778-
779- var listener = Database . Events . OnTablesUpdated . ListenAsync ( signal . Token ) ;
780-
781- // Return the actual IAsyncEnumerable here, using OnChange as a synchronous wrapper that blocks until the
782- // connection is established
783- var throttleMs = options ? . ThrottleMs ?? DEFAULT_WATCH_THROTTLE_MS ;
784- return OnChangeCore ( powersyncTables , listener , signal , options ? . TriggerImmediately == true , throttleMs ) ;
785- }
786-
787- private async IAsyncEnumerable < WatchOnChangeEvent > OnChangeCore (
788- HashSet < string > watchedTables ,
789- IAsyncEnumerable < DBAdapterEvents . TablesUpdatedEvent > listener ,
790- CancellationTokenSource signal ,
791- bool triggerImmediately ,
792- int throttleMs = DEFAULT_WATCH_THROTTLE_MS
793- )
794- {
795- try
796- {
797- await foreach ( var update in OnRawTableChange ( watchedTables , listener , signal . Token , triggerImmediately , throttleMs ) )
798- {
799- // Convert from 'ps_data__<name>' to '<name>'
800- for ( int i = 0 ; i < update . ChangedTables . Length ; i ++ )
801- {
802- update . ChangedTables [ i ] = InternalToFriendlyTableName ( update . ChangedTables [ i ] ) ;
803- }
804- yield return update ;
805- }
806- }
807- finally
808- {
809- signal . Dispose ( ) ;
810- }
811- }
812-
813- private static string InternalToFriendlyTableName ( string internalName )
814- {
815- const string PS_DATA_PREFIX = "ps_data__" ;
816- const string PS_DATA_LOCAL_PREFIX = "ps_data_local__" ;
817-
818- if ( internalName . StartsWith ( PS_DATA_PREFIX ) )
819- return internalName . Substring ( PS_DATA_PREFIX . Length ) ;
820-
821- if ( internalName . StartsWith ( PS_DATA_LOCAL_PREFIX ) )
822- return internalName . Substring ( PS_DATA_LOCAL_PREFIX . Length ) ;
823-
824- return internalName ;
825- }
826-
827- public IAsyncEnumerable < T [ ] > Watch < T > (
828- string sql ,
829- object ? [ ] ? parameters = null ,
830- SQLWatchOptions ? options = null
831- )
771+ /// <summary>
772+ /// Executes a read query every time the source tables are modified.
773+ ///
774+ /// The query's source tables are resolved automatically (or taken from
775+ /// <see cref="SQLWatchOptions.Tables"/> when provided), and are re-resolved whenever the
776+ /// schema changes, so watches keep working across <see cref="UpdateSchema"/> calls.
777+ /// </summary>
778+ public IAsyncEnumerable < T [ ] > Watch < T > ( string sql , object ? [ ] ? parameters = null , SQLWatchOptions ? options = null )
832779 {
833- options ??= new SQLWatchOptions ( ) ;
834-
835- // Stop watching on master CTS cancellation, or on user CTS cancellation
836- var signal = options . Signal != null
837- ? CancellationTokenSource . CreateLinkedTokenSource ( masterCts . Token , options . Signal . Value )
838- : CancellationTokenSource . CreateLinkedTokenSource ( masterCts . Token ) ;
839-
840- // Establish the initial DB listener synchronously before returning the IAsyncEnumerable,
841- // so that table changes between Watch() being called and iteration starting are not missed.
842- // This mirrors the pattern used in OnChange().
843- var initialRestartCts = CancellationTokenSource . CreateLinkedTokenSource ( signal . Token ) ;
844- var initialListener = Database . Events . OnTablesUpdated . ListenAsync ( initialRestartCts . Token ) ;
845-
846- return WatchCore < T > ( sql , parameters , options , signal , initialRestartCts , initialListener ) ;
780+ return watchManager . Watch < T > ( sql , parameters , options ) ;
847781 }
848782
849- private async IAsyncEnumerable < T [ ] > WatchCore < T > (
850- string sql ,
851- object ? [ ] ? parameters ,
852- SQLWatchOptions options ,
853- CancellationTokenSource signal ,
854- CancellationTokenSource initialRestartCts ,
855- IAsyncEnumerable < DBAdapterEvents . TablesUpdatedEvent > initialListener
856- )
783+ /// <summary>
784+ /// Emits an event whenever any of the tables in <see cref="SQLWatchOptions.Tables"/> are modified.
785+ /// </summary>
786+ public IAsyncEnumerable < WatchOnChangeEvent > OnChange ( SQLWatchOptions ? options = null )
857787 {
858- var schemaChanged = new TaskCompletionSource < bool > ( ) ;
859-
860- // Listen for schema changes in the background
861- var schemaListenerTask = Task . Run ( async ( ) =>
862- {
863- await foreach ( var update in Events . OnSchemaChanged . ListenAsync ( signal . Token ) )
864- {
865- // Swap schemaChanged with an unresolved TCS
866- var oldTcs = Interlocked . Exchange ( ref schemaChanged , new ( ) ) ;
867- oldTcs . TrySetResult ( true ) ;
868- }
869- } , signal . Token ) ;
870-
871- // Re-register query on schema updates
872- bool isRestart = false ;
873- var currentRestartCts = initialRestartCts ;
874- var currentListener = initialListener ;
875- var throttleMs = options ? . ThrottleMs ?? DEFAULT_WATCH_THROTTLE_MS ;
876-
877- try
878- {
879- while ( ! signal . Token . IsCancellationRequested )
880- {
881- // Resolve tables
882- HashSet < string > powersyncTables ;
883- if ( options ? . Tables != null )
884- {
885- powersyncTables = [ .. options
886- . Tables
887- . SelectMany < string , string > ( table => [ $ "ps_data__{ table } ", $ "ps_data_local__{ table } "]
888- ) ] ;
889- }
890- else
891- {
892- powersyncTables = await GetSourceTables ( sql , parameters ) ;
893- }
894-
895- var enumerator = OnRawTableChange (
896- powersyncTables ,
897- currentListener ,
898- currentRestartCts . Token ,
899- isRestart || ( options ? . TriggerImmediately == true ) ,
900- throttleMs
901- ) . GetAsyncEnumerator ( ) ;
902-
903- // Continually wait for either OnChange or SchemaChanged to fire
904- while ( true )
905- {
906- var currentSchemaTask = schemaChanged . Task ;
907- var onChangeTask = enumerator . MoveNextAsync ( ) . AsTask ( ) ;
908- var completedTask = await Task . WhenAny ( onChangeTask , currentSchemaTask ) ;
909-
910- if ( completedTask == currentSchemaTask )
911- {
912- var oldRestartCts = currentRestartCts ;
913- oldRestartCts . Cancel ( ) ;
914- isRestart = true ;
915- // Let the current task complete/cancel gracefully
916- try { await onChangeTask ; }
917- catch ( OperationCanceledException ) { }
918-
919- // Establish a new listener BEFORE resolving source tables in the next iteration,
920- // so that changes during the async GetSourceTables call are not missed.
921- currentRestartCts = CancellationTokenSource . CreateLinkedTokenSource ( signal . Token ) ;
922- currentListener = Database . Events . OnTablesUpdated . ListenAsync ( currentRestartCts . Token ) ;
923- oldRestartCts . Dispose ( ) ;
924-
925- break ;
926- }
927-
928- // Await onChangeTask to propagate cancellation and detect end-of-enumeration
929- bool hasNext ;
930- try { hasNext = await onChangeTask ; }
931- catch ( OperationCanceledException ) { yield break ; }
932-
933- if ( ! hasNext ) break ;
934-
935- var update = enumerator . Current ;
936- if ( update . ChangedTables != null )
937- {
938- yield return await GetAll < T > ( sql , parameters ) ;
939- }
940- }
941- }
942- }
943- finally
944- {
945- signal . Cancel ( ) ;
946- try { await schemaListenerTask ; }
947- catch ( OperationCanceledException ) { }
948-
949- currentRestartCts . Dispose ( ) ;
950- signal . Dispose ( ) ;
951- }
788+ return watchManager . OnChange ( options ) ;
952789 }
953790
954791 private class ExplainedResult
@@ -980,112 +817,6 @@ internal async Task<HashSet<string>> GetSourceTables(string sql, object?[]? para
980817
981818 return [ .. tables . Select ( row => row . tbl_name ) ] ;
982819 }
983-
984- private async IAsyncEnumerable < WatchOnChangeEvent > OnRawTableChange (
985- HashSet < string > watchedTables ,
986- IAsyncEnumerable < DBAdapterEvents . TablesUpdatedEvent > listener ,
987- [ EnumeratorCancellation ] CancellationToken signal ,
988- bool triggerImmediately = false ,
989- int throttleMs = DEFAULT_WATCH_THROTTLE_MS
990- )
991- {
992- if ( triggerImmediately )
993- {
994- yield return new WatchOnChangeEvent { ChangedTables = [ ] } ;
995- }
996-
997- if ( throttleMs <= 0 )
998- {
999- // No throttling
1000- HashSet < string > changedTables = new ( ) ;
1001- await foreach ( var e in listener )
1002- {
1003- GetTablesFromNotification ( e . TablesUpdated , changedTables ) ;
1004- changedTables . IntersectWith ( watchedTables ) ;
1005- if ( changedTables . Count == 0 ) continue ;
1006- yield return new WatchOnChangeEvent { ChangedTables = [ .. changedTables ] } ;
1007- }
1008- yield break ;
1009- }
1010-
1011- // Throttled - publish via throttled call to an action that flushes accumulated changes into this channel
1012- var channel = Channel . CreateUnbounded < WatchOnChangeEvent > ( ) ;
1013- var accumulatedTables = new HashSet < string > ( ) ;
1014-
1015- _ = Task . Run ( async ( ) =>
1016- {
1017- using var throttledFlush = Throttler . Throttle ( ( ) =>
1018- {
1019- // Safe to lock directly on accumulatedTables because it's a local variable
1020- lock ( accumulatedTables )
1021- {
1022- if ( accumulatedTables . Count == 0 ) return ;
1023- channel . Writer . TryWrite ( new WatchOnChangeEvent { ChangedTables = [ .. accumulatedTables ] } ) ;
1024- accumulatedTables . Clear ( ) ;
1025- }
1026- } ,
1027- TimeSpan . FromMilliseconds ( throttleMs ) ,
1028- leading : false ,
1029- trailing : true
1030- ) ;
1031-
1032- try
1033- {
1034- var changedTables = new HashSet < string > ( ) ;
1035- await foreach ( var e in listener )
1036- {
1037- GetTablesFromNotification ( e . TablesUpdated , changedTables ) ;
1038- changedTables . IntersectWith ( watchedTables ) ;
1039- if ( changedTables . Count == 0 ) continue ;
1040-
1041- lock ( accumulatedTables ) { accumulatedTables . UnionWith ( changedTables ) ; }
1042- throttledFlush . Invoke ( ) ;
1043- }
1044- }
1045- catch ( OperationCanceledException ) { }
1046- finally
1047- {
1048- // Flush any remaining events and close the channel
1049- lock ( accumulatedTables )
1050- {
1051- if ( accumulatedTables . Count > 0 )
1052- {
1053- channel . Writer . TryWrite ( new WatchOnChangeEvent { ChangedTables = [ .. accumulatedTables ] } ) ;
1054- accumulatedTables . Clear ( ) ;
1055- }
1056- }
1057- channel . Writer . Complete ( ) ;
1058- }
1059- } ) ;
1060-
1061- // Continuously pull values from channel and publish to the consumer
1062- while ( await channel . Reader . WaitToReadAsync ( CancellationToken . None ) )
1063- {
1064- while ( channel . Reader . TryRead ( out var evt ) )
1065- {
1066- yield return evt ;
1067- }
1068- }
1069- }
1070-
1071- private static void GetTablesFromNotification ( INotification updateNotification , HashSet < string > changedTables )
1072- {
1073- changedTables . Clear ( ) ;
1074- string [ ] tables = [ ] ;
1075- if ( updateNotification is BatchedUpdateNotification batchedUpdate )
1076- {
1077- tables = batchedUpdate . Tables ;
1078- }
1079- else if ( updateNotification is UpdateNotification singleUpdate )
1080- {
1081- tables = [ singleUpdate . Table ] ;
1082- }
1083-
1084- foreach ( var table in tables )
1085- {
1086- changedTables . Add ( table ) ;
1087- }
1088- }
1089820}
1090821
1091822public class SQLWatchOptions
0 commit comments