From 631ec9f86be2ee445b60046b14db220f57d20478 Mon Sep 17 00:00:00 2001 From: quifox <289420841+quifox@users.noreply.github.com> Date: Wed, 30 Sep 2026 17:03:14 +0900 Subject: [PATCH] .NET: Publish synthetic terminal responses with their events --- .../Responses/InMemoryResponsesService.cs | 12 +- .../InMemoryResponsesServiceStreamingTests.cs | 391 ++++++++++++++++++ 2 files changed, 397 insertions(+), 6 deletions(-) create mode 100644 dotnet/tests/Microsoft.Agents.AI.Hosting.OpenAI.UnitTests/InMemoryResponsesServiceStreamingTests.cs diff --git a/dotnet/src/Microsoft.Agents.AI.Hosting.OpenAI/Responses/InMemoryResponsesService.cs b/dotnet/src/Microsoft.Agents.AI.Hosting.OpenAI/Responses/InMemoryResponsesService.cs index 365851508d..f9fb04fa2c 100644 --- a/dotnet/src/Microsoft.Agents.AI.Hosting.OpenAI/Responses/InMemoryResponsesService.cs +++ b/dotnet/src/Microsoft.Agents.AI.Hosting.OpenAI/Responses/InMemoryResponsesService.cs @@ -556,7 +556,7 @@ private async Task ExecuteResponseAsync(string responseId, ResponseState state, // Update response status to completed if not already in a terminal state if (!state.IsTerminal) { - state.Response = state.Response! with + var completedResponse = state.Response! with { Status = ResponseStatus.Completed }; @@ -565,7 +565,7 @@ private async Task ExecuteResponseAsync(string responseId, ResponseState state, var completedEvent = new StreamingResponseCompleted { SequenceNumber = sequenceNumber, - Response = state.Response + Response = completedResponse }; state.AddStreamingEvent(completedEvent); @@ -574,7 +574,7 @@ private async Task ExecuteResponseAsync(string responseId, ResponseState state, catch (OperationCanceledException) { // Update response status to cancelled - state.Response = state.Response! with + var cancelledResponse = state.Response! with { Status = ResponseStatus.Cancelled }; @@ -583,7 +583,7 @@ private async Task ExecuteResponseAsync(string responseId, ResponseState state, var cancelledEvent = new StreamingResponseCancelled { SequenceNumber = sequenceNumber, - Response = state.Response + Response = cancelledResponse }; state.AddStreamingEvent(cancelledEvent); @@ -591,7 +591,7 @@ private async Task ExecuteResponseAsync(string responseId, ResponseState state, catch (Exception ex) { // Update response status to failed - state.Response = state.Response! with + var failedResponse = state.Response! with { Status = ResponseStatus.Failed, Error = new ResponseError @@ -605,7 +605,7 @@ private async Task ExecuteResponseAsync(string responseId, ResponseState state, var failedEvent = new StreamingResponseFailed { SequenceNumber = sequenceNumber, - Response = state.Response + Response = failedResponse }; state.AddStreamingEvent(failedEvent); diff --git a/dotnet/tests/Microsoft.Agents.AI.Hosting.OpenAI.UnitTests/InMemoryResponsesServiceStreamingTests.cs b/dotnet/tests/Microsoft.Agents.AI.Hosting.OpenAI.UnitTests/InMemoryResponsesServiceStreamingTests.cs new file mode 100644 index 0000000000..720411484d --- /dev/null +++ b/dotnet/tests/Microsoft.Agents.AI.Hosting.OpenAI.UnitTests/InMemoryResponsesServiceStreamingTests.cs @@ -0,0 +1,391 @@ +// Copyright (c) Microsoft. All rights reserved. + +using System; +using System.Collections.Generic; +using System.Reflection; +using System.Threading; +using System.Threading.Tasks; +using System.Threading.Tasks.Sources; +using Microsoft.Agents.AI.Hosting.OpenAI.Responses; +using Microsoft.Agents.AI.Hosting.OpenAI.Responses.Models; +using Microsoft.Extensions.AI; +using Microsoft.Extensions.Caching.Memory; + +namespace Microsoft.Agents.AI.Hosting.OpenAI.UnitTests; + +/// +/// Tests publication and consumption of stored response events. +/// +public sealed class InMemoryResponsesServiceStreamingTests +{ + private static readonly TimeSpan s_timeout = TimeSpan.FromSeconds(10); + + [Theory] + [InlineData("completed")] + [InlineData("failed")] + [InlineData("cancelled")] + public async Task SynthesizedTerminalEvent_IsPublishedWithResponseAsync(string outcome) + { + // Arrange + var executor = new ControlledResponseExecutor(); + using var service = new InMemoryResponsesService(executor); + using var timeout = new CancellationTokenSource(s_timeout); + Response response = await service.CreateResponseAsync(new CreateResponse + { + Input = ResponseInput.FromText("hello"), + Background = true + }, timeout.Token); + await executor.ContinuationRegistered.Task.WaitAsync(timeout.Token); + + // Observe the private publication lock without introducing a production test hook. + var cache = (IMemoryCache)typeof(InMemoryResponsesService) + .GetField("_cache", BindingFlags.Instance | BindingFlags.NonPublic)!.GetValue(service)!; + Assert.True(cache.TryGetValue(response.Id, out object? state)); + Assert.NotNull(state); + object syncRoot = state.GetType().GetField("_lock", BindingFlags.Instance | BindingFlags.NonPublic)!.GetValue(state)!; + await using IAsyncEnumerator reader = service.GetResponseStreamingAsync(response.Id) + .GetAsyncEnumerator(timeout.Token); + var producer = new Thread(() => executor.Complete(outcome)) { IsBackground = true }; + Task? moveNext = null; + bool completedBeforePublication = false; + + // Act + try + { + lock (syncRoot) + { + producer.Start(); + + // The registered continuation runs inline on this dedicated thread. After disposing + // the executor, its only blocking operation is acquiring the publication lock. + Assert.True(SpinWait.SpinUntil( + () => executor.DisposalThreadId == producer.ManagedThreadId && + (producer.ThreadState & ThreadState.WaitSleepJoin) != 0, + s_timeout)); + Assert.Equal(producer.ManagedThreadId, executor.ContinuationThreadId); + + // Reenter the same monitor as a reader. Until the producer can append its event, + // a reader must remain pending rather than observe terminal state and stop. + moveNext = reader.MoveNextAsync().AsTask(); + completedBeforePublication = moveNext.IsCompleted; + } + } + finally + { + Assert.True(producer.Join(s_timeout)); + } + + // Assert + Assert.NotNull(moveNext); + Assert.True(await moveNext.WaitAsync(timeout.Token)); + Assert.False(completedBeforePublication); + Assert.Equal($"response.{outcome}", reader.Current.Type); + var terminal = Assert.IsAssignableFrom(reader.Current); + Assert.True(terminal.Response.IsTerminal); + Assert.Same(terminal.Response, await service.GetResponseAsync(response.Id, timeout.Token)); + Assert.False(await reader.MoveNextAsync()); + } + + [Theory] + [InlineData("completed")] + [InlineData("failed")] + [InlineData("cancelled")] + public async Task SynthesizedTerminalEvent_PreservesOutputForReadersAndReplayAsync(string outcome) + { + // Arrange + var output = new ResponsesAssistantMessageItemResource + { + Id = "msg_partial", + Status = ResponsesMessageItemResourceStatus.InProgress, + Content = [new ItemContentOutputText { Text = "Partial output", Annotations = [] }] + }; + var outputEvent = new StreamingOutputItemAdded { SequenceNumber = 1, OutputIndex = 0, Item = output }; + var executor = new ControlledResponseExecutor((_, _) => [outputEvent]); + using var service = new InMemoryResponsesService(executor); + using var timeout = new CancellationTokenSource(s_timeout); + Response response = await service.CreateResponseAsync(new CreateResponse + { + Input = ResponseInput.FromText("hello"), + Background = true + }, timeout.Token); + await executor.ContinuationRegistered.Task.WaitAsync(timeout.Token); + await using IAsyncEnumerator first = service.GetResponseStreamingAsync(response.Id) + .GetAsyncEnumerator(timeout.Token); + await using IAsyncEnumerator second = service.GetResponseStreamingAsync(response.Id) + .GetAsyncEnumerator(timeout.Token); + Task firstMove; + Task secondMove; + bool firstCompletedBeforePublication; + bool secondCompletedBeforePublication; + Task? cancellation = null; + bool executionCancellationRequested = false; + bool cancellationCompletedBeforePublication = false; + try + { + Assert.True(await first.MoveNextAsync()); + Assert.True(await second.MoveNextAsync()); + Assert.Same(outputEvent, first.Current); + Assert.Same(outputEvent, second.Current); + firstMove = first.MoveNextAsync().AsTask(); + secondMove = second.MoveNextAsync().AsTask(); + firstCompletedBeforePublication = firstMove.IsCompleted; + secondCompletedBeforePublication = secondMove.IsCompleted; + + // Act + if (outcome == "cancelled") + { + cancellation = service.CancelResponseAsync(response.Id, timeout.Token); + executionCancellationRequested = executor.ExecutionCancellationToken.IsCancellationRequested; + cancellationCompletedBeforePublication = cancellation.IsCompleted; + } + } + finally + { + executor.Complete(outcome); + } + + // Settle both moves before asserting their earlier state so a failed assertion cannot + // dispose an iterator while MoveNextAsync is still pending. + bool[] receivedTerminal = await Task.WhenAll(firstMove, secondMove).WaitAsync(timeout.Token); + + // Assert + Assert.False(firstCompletedBeforePublication); + Assert.False(secondCompletedBeforePublication); + Assert.All(receivedTerminal, Assert.True); + if (cancellation is not null) + { + Assert.True(executionCancellationRequested); + Assert.False(cancellationCompletedBeforePublication); + } + + StreamingResponseEvent terminalEvent = first.Current; + Assert.Same(terminalEvent, second.Current); + Assert.Equal($"response.{outcome}", terminalEvent.Type); + Assert.Equal(2, terminalEvent.SequenceNumber); + var terminal = Assert.IsAssignableFrom(terminalEvent); + Assert.Same(output, Assert.Single(terminal.Response.Output)); + Assert.Same(terminal.Response, await service.GetResponseAsync(response.Id, timeout.Token)); + if (outcome == "failed") + { + Assert.NotNull(terminal.Response.Error); + Assert.Equal("execution_error", terminal.Response.Error.Code); + Assert.Equal("Test executor failure.", terminal.Response.Error.Message); + } + + if (cancellation is not null) + { + Assert.Same(terminal.Response, await cancellation.WaitAsync(timeout.Token)); + } + + Assert.False(await first.MoveNextAsync()); + Assert.False(await second.MoveNextAsync()); + List replay = await ReadEventsAsync(service.GetResponseStreamingAsync(response.Id), timeout.Token); + Assert.Collection(replay, item => Assert.Same(outputEvent, item), item => Assert.Same(terminalEvent, item)); + List resumed = await ReadEventsAsync( + service.GetResponseStreamingAsync(response.Id, startingAfter: 1), timeout.Token); + Assert.Same(terminalEvent, Assert.Single(resumed)); + } + + [Fact] + public async Task CreateResponseStreamingAsync_CancelledReader_DoesNotCancelProducerOrOtherReaderAsync() + { + // Arrange + var executor = new ControlledResponseExecutor(); + using var service = new InMemoryResponsesService(executor); + using var timeout = new CancellationTokenSource(s_timeout); + using var readerCancellation = new CancellationTokenSource(); + await using IAsyncEnumerator cancelledReader = service.CreateResponseStreamingAsync(new CreateResponse + { + Input = ResponseInput.FromText("hello"), + Stream = true + }).GetAsyncEnumerator(readerCancellation.Token); + Task cancelledMove = cancelledReader.MoveNextAsync().AsTask(); + await executor.ContinuationRegistered.Task.WaitAsync(timeout.Token); + await using IAsyncEnumerator otherReader = service.GetResponseStreamingAsync(executor.ResponseId) + .GetAsyncEnumerator(timeout.Token); + Task otherMove = otherReader.MoveNextAsync().AsTask(); + + // Act + Exception? cancellationException; + bool executionCancellationRequested; + bool otherCompletedBeforePublication; + try + { + readerCancellation.Cancel(); + cancellationException = await Record.ExceptionAsync(() => cancelledMove.WaitAsync(timeout.Token)); + executionCancellationRequested = executor.ExecutionCancellationToken.IsCancellationRequested; + otherCompletedBeforePublication = otherMove.IsCompleted; + } + finally + { + executor.Complete("completed"); + } + + bool receivedTerminal = await otherMove.WaitAsync(timeout.Token); + + // Assert + Assert.IsAssignableFrom(cancellationException); + Assert.False(executionCancellationRequested); + Assert.False(otherCompletedBeforePublication); + Assert.True(receivedTerminal); + StreamingResponseEvent terminalEvent = otherReader.Current; + Assert.IsType(terminalEvent); + Assert.False(await otherReader.MoveNextAsync()); + List replay = await ReadEventsAsync(service.GetResponseStreamingAsync(executor.ResponseId), timeout.Token); + Assert.Same(terminalEvent, Assert.Single(replay)); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task ExecutorTerminalEvent_IsPreservedWithoutSynthesizedCompletionAsync(bool incomplete) + { + // Arrange + StreamingResponseEvent? suppliedEvent = null; + var executor = new ControlledResponseExecutor((context, request) => + { + var response = new Response + { + Id = context.ResponseId, + CreatedAt = 123, + Background = request.Background, + Status = incomplete ? ResponseStatus.Incomplete : ResponseStatus.Completed, + Output = [], + Usage = ResponseUsage.Zero, + Tools = [] + }; + suppliedEvent = incomplete + ? new StreamingResponseIncomplete { SequenceNumber = 42, Response = response } + : new StreamingResponseCompleted { SequenceNumber = 42, Response = response }; + return [suppliedEvent]; + }); + using var service = new InMemoryResponsesService(executor); + using var timeout = new CancellationTokenSource(s_timeout); + Response initialResponse = await service.CreateResponseAsync(new CreateResponse + { + Input = ResponseInput.FromText("hello"), + Background = true + }, timeout.Token); + await executor.ContinuationRegistered.Task.WaitAsync(timeout.Token); + + // Act + List initial; + try + { + initial = await ReadEventsAsync(service.GetResponseStreamingAsync(initialResponse.Id), timeout.Token); + } + finally + { + executor.Complete("completed"); + } + + List replay = await ReadEventsAsync(service.GetResponseStreamingAsync(initialResponse.Id), timeout.Token); + + // Assert + Assert.Same(suppliedEvent, Assert.Single(initial)); + Assert.Same(suppliedEvent, Assert.Single(replay)); + Assert.NotNull(suppliedEvent); + Assert.Equal(42, suppliedEvent.SequenceNumber); + var terminal = Assert.IsAssignableFrom(suppliedEvent); + Assert.Same(terminal.Response, await service.GetResponseAsync(initialResponse.Id, timeout.Token)); + } + + private static async Task> ReadEventsAsync( + IAsyncEnumerable events, CancellationToken cancellationToken) + { + List result = []; + await foreach (StreamingResponseEvent item in events.WithCancellation(cancellationToken)) + { + result.Add(item); + } + + return result; + } + + private sealed class ControlledResponseExecutor( + Func>? initialEventsFactory = null) : IResponseExecutor, + IAsyncEnumerable, IAsyncEnumerator, IValueTaskSource + { + private ManualResetValueTaskSourceCore _completion = new() { RunContinuationsAsynchronously = false }; + private int _continuationThreadId; + private int _disposalThreadId; + private IReadOnlyList _events = []; + private int _nextEvent; + + public TaskCompletionSource ContinuationRegistered { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously); + + public int ContinuationThreadId => Volatile.Read(ref this._continuationThreadId); + + public int DisposalThreadId => Volatile.Read(ref this._disposalThreadId); + + public StreamingResponseEvent Current { get; private set; } = null!; + + public string ResponseId { get; private set; } = string.Empty; + + public CancellationToken ExecutionCancellationToken { get; private set; } + + public ValueTask ValidateRequestAsync(CreateResponse request, CancellationToken cancellationToken = default) + => ValueTask.FromResult(null); + + public IAsyncEnumerable ExecuteAsync( + AgentInvocationContext context, + CreateResponse request, + IReadOnlyList? conversationHistory = null, + CancellationToken cancellationToken = default) + { + this._events = initialEventsFactory?.Invoke(context, request) ?? []; + this.ResponseId = context.ResponseId; + this.ExecutionCancellationToken = cancellationToken; + return this; + } + + public IAsyncEnumerator GetAsyncEnumerator(CancellationToken cancellationToken = default) => this; + + public ValueTask MoveNextAsync() + { + if (this._nextEvent < this._events.Count) + { + this.Current = this._events[this._nextEvent++]; + return ValueTask.FromResult(true); + } + + return new(this, this._completion.Version); + } + + public ValueTask DisposeAsync() + { + Volatile.Write(ref this._disposalThreadId, Environment.CurrentManagedThreadId); + return default; + } + + public void Complete(string outcome) + { + if (outcome == "failed") + { + this._completion.SetException(new InvalidOperationException("Test executor failure.")); + } + else if (outcome == "cancelled") + { + this._completion.SetException(new OperationCanceledException()); + } + else + { + this._completion.SetResult(false); + } + } + + public bool GetResult(short token) + { + Volatile.Write(ref this._continuationThreadId, Environment.CurrentManagedThreadId); + return this._completion.GetResult(token); + } + + public ValueTaskSourceStatus GetStatus(short token) => this._completion.GetStatus(token); + + public void OnCompleted(Action continuation, object? state, short token, ValueTaskSourceOnCompletedFlags flags) + { + this._completion.OnCompleted(continuation, state, token, flags); + this.ContinuationRegistered.SetResult(); + } + } +}