@@ -569,106 +569,193 @@ await tx.Execute(
569569 await watched . Task ;
570570 }
571571
572+ private class WrappedTCS
573+ {
574+ public TaskCompletionSource < bool > TCS { get ; protected set ; } = new ( ) ;
575+ public Task < bool > Task => TCS . Task ;
576+
577+ public void Complete ( ) => TCS . TrySetResult ( true ) ;
578+ public void Reset ( ) => TCS = new ( ) ;
579+ }
580+
572581 [ Fact ( Timeout = 2000 ) ]
573582 public async void WatchDisposableSubscriptionTest ( )
574583 {
575584 int callCount = 0 ;
585+ var tcs = new WrappedTCS ( ) ;
576586
577- var subscription = await db . Watch ( "select id, description, make from assets" , null , new ( )
587+ var query = await db . Watch ( "select id from assets" , null , new ( )
578588 {
579- OnResult = ( results ) => callCount ++ ,
580- OnError = ( ex ) => Assert . Fail ( "An exception occurred: " + ex . ToString ( ) )
589+ OnResult = ( results ) =>
590+ {
591+ callCount ++ ;
592+ tcs . Complete ( ) ;
593+ } ,
594+ OnError = ( ex ) => Assert . Fail ( ex . ToString ( ) )
581595 } ) ;
582- Thread . Sleep ( 200 ) ;
596+ await tcs . Task ;
597+ tcs . Reset ( ) ;
583598 Assert . Equal ( 1 , callCount ) ;
584599
585- // Bump callCount to 2
586600 await db . Execute (
587601 "insert into assets(id, description, make) values (?, ?, ?)" ,
588602 [ Guid . NewGuid ( ) . ToString ( ) , "some desc" , "some make" ]
589603 ) ;
590- Thread . Sleep ( 200 ) ;
604+ await tcs . Task ;
605+ tcs . Reset ( ) ;
591606 Assert . Equal ( 2 , callCount ) ;
592607
593- subscription . Dispose ( ) ;
608+ query . Dispose ( ) ;
609+
594610 await db . Execute (
595611 "insert into assets(id, description, make) values (?, ?, ?)" ,
596612 [ Guid . NewGuid ( ) . ToString ( ) , "some desc" , "some make" ]
597613 ) ;
598- Thread . Sleep ( 200 ) ;
599- Assert . Equal ( 2 , callCount ) ;
614+ await Task . Delay ( 100 ) ;
615+ Assert . Equal ( 2 , callCount ) ; // Same value
600616 }
601617
602- [ Fact ( Timeout = 2500 ) ]
618+ [ Fact ( Timeout = 2000 ) ]
603619 public async void WatchDisposableCustomTokenTest ( )
604620 {
605621 var customTokenSource = new CancellationTokenSource ( ) ;
606622 int callCount = 0 ;
623+ var tcs = new WrappedTCS ( ) ;
607624
608- using var subscription = await db . Watch ( "select id, description, make from assets" , null , new ( )
625+ using var query = await db . Watch ( "select id, description, make from assets" , null , new ( )
609626 {
610- OnResult = ( results ) => callCount ++ ,
611- OnError = ( ex ) => Assert . Fail ( "An exception occurred: " + ex . ToString ( ) )
627+ OnResult = ( results ) =>
628+ {
629+ callCount ++ ;
630+ tcs . Complete ( ) ;
631+ } ,
632+ OnError = ( ex ) => Assert . Fail ( ex . ToString ( ) )
612633 } , new ( )
613634 {
614635 Signal = customTokenSource . Token
615636 } ) ;
616- Thread . Sleep ( 200 ) ;
637+ await tcs . Task ;
638+ tcs . Reset ( ) ;
617639 Assert . Equal ( 1 , callCount ) ;
618640
619641 await db . Execute (
620642 "insert into assets(id, description, make) values (?, ?, ?)" ,
621643 [ Guid . NewGuid ( ) . ToString ( ) , "some desc" , "some make" ]
622644 ) ;
623- Thread . Sleep ( 200 ) ;
645+ await tcs . Task ;
646+ tcs . Reset ( ) ;
624647 Assert . Equal ( 2 , callCount ) ;
625648
626649 customTokenSource . Cancel ( ) ;
650+
627651 await db . Execute (
628652 "insert into assets(id, description, make) values (?, ?, ?)" ,
629653 [ Guid . NewGuid ( ) . ToString ( ) , "some desc" , "some make" ]
630654 ) ;
631- Thread . Sleep ( 200 ) ;
655+ await Task . Delay ( 100 ) ;
632656 Assert . Equal ( 2 , callCount ) ; // Same value
633657 }
634658
635659 [ Fact ( Timeout = 2000 ) ]
636- public async void WatchMultipleCancelledTest ( )
660+ public async void WatchSingleCancelledTest ( )
637661 {
638662 int callCount = 0 ;
639- var watchHandlerFactory = ( ) => new WatchHandler < IdResult >
663+ object lockObject = new object ( ) ;
664+
665+ var watchHandlerFactory = ( WrappedTCS tcs ) => new WatchHandler < IdResult >
640666 {
641- OnResult = ( result ) => callCount ++ ,
642- OnError = ( ex ) => Assert . Fail ( "An exception occurred: " + ex . ToString ( ) ) ,
667+ OnResult = ( result ) =>
668+ {
669+ lock ( lockObject )
670+ {
671+ callCount ++ ;
672+ }
673+ tcs . Complete ( ) ;
674+ } ,
675+ OnError = ( ex ) => Assert . Fail ( ex . ToString ( ) ) ,
643676 } ;
644677
645- var query1 = await db . Watch ( "select id from assets" , null , watchHandlerFactory ( ) ) ;
646- var query2 = await db . Watch ( "select id from customers" , null , watchHandlerFactory ( ) ) ;
647- Thread . Sleep ( 200 ) ;
678+ var tcsAlwaysRunning = new WrappedTCS ( ) ;
679+ var tcsCancelled = new WrappedTCS ( ) ;
680+ // Make queryAlwaysRunning slightly more expensive
681+ using var queryAlwaysRunning = await db . Watch ( "select id from assets order by id" , null , watchHandlerFactory ( tcsAlwaysRunning ) ) ;
682+ using var queryCancelled = await db . Watch ( "select id from assets" , null , watchHandlerFactory ( tcsCancelled ) ) ;
683+
684+ await Task . WhenAll ( tcsAlwaysRunning . Task , tcsCancelled . Task ) ;
685+ tcsAlwaysRunning . Reset ( ) ;
686+ tcsCancelled . Reset ( ) ;
648687 Assert . Equal ( 2 , callCount ) ;
649688
650689 await db . Execute (
651690 "insert into assets(id, description, make) values (?, ?, ?)" ,
652691 [ Guid . NewGuid ( ) . ToString ( ) , "some desc" , "some make" ]
653692 ) ;
693+ await Task . WhenAll ( tcsAlwaysRunning . Task , tcsCancelled . Task ) ;
694+ tcsAlwaysRunning . Reset ( ) ;
695+ tcsCancelled . Reset ( ) ;
696+ Assert . Equal ( 4 , callCount ) ;
697+
698+ // Close one query
699+ queryCancelled . Dispose ( ) ;
700+
654701 await db . Execute (
655702 "insert into assets(id, description, make) values (?, ?, ?)" ,
656703 [ Guid . NewGuid ( ) . ToString ( ) , "some desc" , "some make" ]
657704 ) ;
658- Thread . Sleep ( 200 ) ;
659- Assert . Equal ( 4 , callCount ) ;
705+ // Because queryAlwaysRunning is more expensive than queryCancelled, we would expect
706+ // queryCancelled to complete before it if the cancellation failed, hence we can just
707+ // wait for queryAlwaysRunning to complete.
708+ await tcsAlwaysRunning . Task ;
660709
661- db . UnsubscribeAllQueries ( ) ;
710+ Assert . Equal ( 5 , callCount ) ;
711+ }
712+
713+ [ Fact ( Timeout = 2000 ) ]
714+ public async void WatchMultipleCancelledTest ( )
715+ {
716+ int callCount = 0 ;
717+ object lockObject = new object ( ) ;
718+
719+ var watchHandlerFactory = ( WrappedTCS tcs ) => new WatchHandler < IdResult >
720+ {
721+ OnResult = ( result ) =>
722+ {
723+ lock ( lockObject )
724+ {
725+ callCount ++ ;
726+ }
727+ tcs . Complete ( ) ;
728+ } ,
729+ OnError = ( ex ) => Assert . Fail ( ex . ToString ( ) ) ,
730+ } ;
731+
732+ var tcs1 = new WrappedTCS ( ) ;
733+ var tcs2 = new WrappedTCS ( ) ;
734+ using var query1 = await db . Watch ( "select id from assets" , null , watchHandlerFactory ( tcs1 ) ) ;
735+ using var query2 = await db . Watch ( "select id from assets" , null , watchHandlerFactory ( tcs2 ) ) ;
736+
737+ await Task . WhenAll ( tcs1 . Task , tcs2 . Task ) ;
738+ tcs1 . Reset ( ) ;
739+ tcs2 . Reset ( ) ;
740+ Assert . Equal ( 2 , callCount ) ;
662741
663742 await db . Execute (
664743 "insert into assets(id, description, make) values (?, ?, ?)" ,
665744 [ Guid . NewGuid ( ) . ToString ( ) , "some desc" , "some make" ]
666745 ) ;
746+ await Task . WhenAll ( tcs1 . Task , tcs2 . Task ) ;
747+ tcs1 . Reset ( ) ;
748+ tcs2 . Reset ( ) ;
749+ Assert . Equal ( 4 , callCount ) ;
750+
751+ db . UnsubscribeAllQueries ( ) ;
752+
667753 await db . Execute (
668754 "insert into assets(id, description, make) values (?, ?, ?)" ,
669755 [ Guid . NewGuid ( ) . ToString ( ) , "some desc" , "some make" ]
670756 ) ;
671- Thread . Sleep ( 200 ) ;
757+ // Short wait to ensure unsubscription was successful
758+ await Task . Delay ( 100 ) ;
672759 Assert . Equal ( 4 , callCount ) ;
673760 }
674761}
0 commit comments