@@ -536,6 +536,57 @@ public async Task SendAndWaitAsync_DroppedIdle_Fallback_Rechecks_Activity_After_
536536 }
537537 }
538538
539+ [ Fact ]
540+ public async Task SendAndWaitAsync_DroppedIdle_Fallback_Flushes_Event_Enqueued_During_Final_Barrier ( )
541+ {
542+ await using var server = await FakeCopilotServer . StartAsync ( ) ;
543+ server . ConfigureEventEnqueueDuringFinalBarrier ( ) ;
544+ await using var client = new CopilotClient ( new CopilotClientOptions { Connection = RuntimeConnection . ForUri ( server . Url ) } ) ;
545+ await using var session = await client . CreateSessionAsync ( new SessionConfig
546+ {
547+ OnPermissionRequest = PermissionHandler . ApproveAll
548+ } ) ;
549+
550+ var queuedHandlerStarted = new TaskCompletionSource ( TaskCreationOptions . RunContinuationsAsynchronously ) ;
551+ var releaseQueuedHandler = new TaskCompletionSource ( TaskCreationOptions . RunContinuationsAsynchronously ) ;
552+ using var subscription = session . On < AssistantTurnEndEvent > ( evt =>
553+ {
554+ if ( evt . Data . TurnId == "activity-reactivation-barrier" )
555+ {
556+ DispatchEvent ( session , new AssistantTurnEndEvent
557+ {
558+ Data = new AssistantTurnEndData { TurnId = "queued-during-final-barrier" }
559+ } ) ;
560+ }
561+ else if ( evt . Data . TurnId == "queued-during-final-barrier" )
562+ {
563+ queuedHandlerStarted . TrySetResult ( ) ;
564+ releaseQueuedHandler . Task . GetAwaiter ( ) . GetResult ( ) ;
565+ }
566+ } ) ;
567+
568+ var completionTask = session . SendAndWaitAsync (
569+ new MessageOptions { Prompt = "flush an event queued during the final barrier" } ,
570+ timeout : TimeSpan . FromSeconds ( 5 ) ) ;
571+
572+ try
573+ {
574+ var first = await Task . WhenAny ( queuedHandlerStarted . Task , completionTask )
575+ . WaitAsync ( TimeSpan . FromSeconds ( 2 ) ) ;
576+ Assert . Same ( queuedHandlerStarted . Task , first ) ;
577+ await Task . Delay ( 200 ) ;
578+ Assert . False ( completionTask . IsCompleted , "Completion must wait for events enqueued after the stable snapshot." ) ;
579+
580+ releaseQueuedHandler . TrySetResult ( ) ;
581+ var response = await completionTask ;
582+ Assert . Equal ( "completed response" , response ? . Data . Content ) ;
583+ }
584+ finally
585+ {
586+ releaseQueuedHandler . TrySetResult ( ) ;
587+ }
588+ }
589+
539590 [ MethodImpl ( MethodImplOptions . NoInlining ) ]
540591 private static async Task < WeakReference < CopilotSession > > CreateDroppedSessionAsync ( CopilotClient client )
541592 {
@@ -647,6 +698,7 @@ private sealed class FakeCopilotServer : IAsyncDisposable
647698 private bool _emitDelayedIdle ;
648699 private bool _fallbackMethodsUnavailable ;
649700 private bool _reactivateDuringFinalBarrier ;
701+ private bool _enqueueDuringFinalBarrier ;
650702 private bool _hasActiveWork ;
651703 private int _activityRequestCount ;
652704
@@ -742,6 +794,12 @@ public void ConfigureActivityReactivationDuringFinalBarrier()
742794 _reactivateDuringFinalBarrier = true ;
743795 }
744796
797+ public void ConfigureEventEnqueueDuringFinalBarrier ( )
798+ {
799+ ConfigureDroppedIdleCompletion ( ) ;
800+ _enqueueDuringFinalBarrier = true ;
801+ }
802+
745803 public void CompleteReactivatedActivity ( )
746804 {
747805 Volatile . Write ( ref _hasActiveWork , false ) ;
@@ -835,7 +893,7 @@ private async Task HandleRequestAsync(Stream stream, JsonElement request, Cancel
835893 if ( method == "session.metadata.activity" )
836894 {
837895 activityRequestNumber = Interlocked . Increment ( ref _activityRequestCount ) ;
838- if ( _reactivateDuringFinalBarrier && activityRequestNumber == 2 )
896+ if ( ( _reactivateDuringFinalBarrier || _enqueueDuringFinalBarrier ) && activityRequestNumber == 2 )
839897 {
840898 await EmitActivityReactivationBarrierEventAsync ( stream , cancellationToken ) ;
841899 }
0 commit comments