mirror of
https://github.com/microsoft/agent-framework.git
synced 2026-06-16 21:04:09 +08:00
.NET: Mitigation for issue 1315 (#1432)
* Adding mitigation for issue 1315, will follow up with long-term plan. * Apply suggestion from @lokitoth Co-authored-by: Jacob Alber <jaalber@microsoft.com> * Apply suggestion from @Copilot Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Extracting common updates to a helper method. * Fixing spelling mistake. * Update per PR feedback. --------- Co-authored-by: Jacob Alber <jaalber@microsoft.com> Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
This commit is contained in:
committed by
GitHub
Unverified
parent
d1ac2d9331
commit
64b57639b8
@@ -126,10 +126,11 @@ public sealed class InProcessExecutionEnvironment : IWorkflowExecutionEnvironmen
|
||||
string? runId = null,
|
||||
CancellationToken cancellationToken = default) where TInput : notnull
|
||||
{
|
||||
AsyncRunHandle runHandle = await this.BeginRunAsync(workflow, checkpointManager: null, runId: runId, [], cancellationToken)
|
||||
.ConfigureAwait(false);
|
||||
var runHandle = await this.GetRunHandleWithTurnTokenAsync(workflow: workflow, input: input, checkpointManager: null, runId: runId, cancellationToken).ConfigureAwait(false);
|
||||
|
||||
return await runHandle.EnqueueAndRunAsync(input, cancellationToken).ConfigureAwait(false);
|
||||
Run run = new(runHandle);
|
||||
await run.RunToNextHaltAsync(cancellationToken).ConfigureAwait(false);
|
||||
return run;
|
||||
}
|
||||
|
||||
/// <inheritdoc/>
|
||||
@@ -139,10 +140,11 @@ public sealed class InProcessExecutionEnvironment : IWorkflowExecutionEnvironmen
|
||||
string? runId = null,
|
||||
CancellationToken cancellationToken = default) where TInput : notnull
|
||||
{
|
||||
AsyncRunHandle runHandle = await this.BeginRunAsync(workflow, checkpointManager: null, runId: runId, [typeof(TInput)], cancellationToken)
|
||||
.ConfigureAwait(false);
|
||||
var runHandle = await this.GetRunHandleWithTurnTokenAsync(workflow: workflow, input: input, checkpointManager: null, runId: runId, cancellationToken).ConfigureAwait(false);
|
||||
|
||||
return await runHandle.EnqueueAndRunAsync(input, cancellationToken).ConfigureAwait(false);
|
||||
Run run = new(runHandle);
|
||||
await run.RunToNextHaltAsync(cancellationToken).ConfigureAwait(false);
|
||||
return run;
|
||||
}
|
||||
|
||||
/// <inheritdoc/>
|
||||
@@ -153,10 +155,11 @@ public sealed class InProcessExecutionEnvironment : IWorkflowExecutionEnvironmen
|
||||
string? runId = null,
|
||||
CancellationToken cancellationToken = default) where TInput : notnull
|
||||
{
|
||||
AsyncRunHandle runHandle = await this.BeginRunAsync(workflow, checkpointManager, runId: runId, [], cancellationToken)
|
||||
.ConfigureAwait(false);
|
||||
var runHandle = await this.GetRunHandleWithTurnTokenAsync(workflow: workflow, input: input, checkpointManager: checkpointManager, runId: runId, cancellationToken).ConfigureAwait(false);
|
||||
|
||||
return await runHandle.WithCheckpointingAsync(() => runHandle.EnqueueAndRunAsync(input, cancellationToken))
|
||||
Run run = new(runHandle);
|
||||
await run.RunToNextHaltAsync(cancellationToken).ConfigureAwait(false);
|
||||
return await runHandle.WithCheckpointingAsync(() => new ValueTask<Run>(run))
|
||||
.ConfigureAwait(false);
|
||||
}
|
||||
|
||||
@@ -168,10 +171,11 @@ public sealed class InProcessExecutionEnvironment : IWorkflowExecutionEnvironmen
|
||||
string? runId = null,
|
||||
CancellationToken cancellationToken = default) where TInput : notnull
|
||||
{
|
||||
AsyncRunHandle runHandle = await this.BeginRunAsync(workflow, checkpointManager, runId: runId, [typeof(TInput)], cancellationToken)
|
||||
.ConfigureAwait(false);
|
||||
var runHandle = await this.GetRunHandleWithTurnTokenAsync(workflow: workflow, input: input, checkpointManager: checkpointManager, runId: runId, cancellationToken).ConfigureAwait(false);
|
||||
|
||||
return await runHandle.WithCheckpointingAsync(() => runHandle.EnqueueAndRunAsync(input, cancellationToken))
|
||||
Run run = new(runHandle);
|
||||
await run.RunToNextHaltAsync(cancellationToken).ConfigureAwait(false);
|
||||
return await runHandle.WithCheckpointingAsync(() => new ValueTask<Run>(run))
|
||||
.ConfigureAwait(false);
|
||||
}
|
||||
|
||||
@@ -204,4 +208,48 @@ public sealed class InProcessExecutionEnvironment : IWorkflowExecutionEnvironmen
|
||||
return await runHandle.WithCheckpointingAsync<Run>(() => new(new Run(runHandle)))
|
||||
.ConfigureAwait(false);
|
||||
}
|
||||
|
||||
// Helper to construct a RunHandle with the provided input enqueued. If the starting executor supports it, a TurnToken will be enqueued also.
|
||||
private async ValueTask<AsyncRunHandle> GetRunHandleWithTurnTokenAsync<TInput>(
|
||||
Workflow workflow,
|
||||
TInput input,
|
||||
CheckpointManager? checkpointManager,
|
||||
string? runId,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
var knownTypes = new List<Type>() { typeof(TInput) };
|
||||
var needsTurnToken = await StartingExecutorHandlesTurnTokenAsync<TInput>(workflow).ConfigureAwait(false);
|
||||
if (needsTurnToken)
|
||||
{
|
||||
knownTypes.Add(typeof(TurnToken));
|
||||
}
|
||||
|
||||
AsyncRunHandle runHandle = await this.BeginRunAsync(workflow, checkpointManager: checkpointManager, runId: runId, knownTypes, cancellationToken)
|
||||
.ConfigureAwait(false);
|
||||
|
||||
await runHandle.EnqueueMessageAsync(input, cancellationToken).ConfigureAwait(false);
|
||||
|
||||
if (needsTurnToken)
|
||||
{
|
||||
await runHandle.EnqueueMessageAsync(new TurnToken(emitEvents: true), cancellationToken).ConfigureAwait(false);
|
||||
}
|
||||
|
||||
return runHandle;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Helper method to detect if the starting executor of a given workflow accepts the provided input type as well as a TurnToken.
|
||||
/// </summary>
|
||||
private static async ValueTask<bool> StartingExecutorHandlesTurnTokenAsync<TInput>(Workflow workflow)
|
||||
{
|
||||
if (workflow.Registrations.TryGetValue(workflow.StartExecutorId, out var registration))
|
||||
{
|
||||
// Create instance to check type
|
||||
Executor startExecutor = await registration.CreateInstanceAsync(string.Empty)
|
||||
.ConfigureAwait(false);
|
||||
return startExecutor.CanHandle(typeof(TInput)) && startExecutor.CanHandle(typeof(TurnToken));
|
||||
}
|
||||
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,200 @@
|
||||
// Copyright (c) Microsoft. All rights reserved.
|
||||
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Linq;
|
||||
using System.Runtime.CompilerServices;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using FluentAssertions;
|
||||
using Microsoft.Extensions.AI;
|
||||
|
||||
namespace Microsoft.Agents.AI.Workflows.UnitTests;
|
||||
|
||||
/// <summary>
|
||||
/// Tests for InProcessExecution to verify streaming and non-streaming execution behavior.
|
||||
/// </summary>
|
||||
public class InProcessExecutionTests
|
||||
{
|
||||
/// <summary>
|
||||
/// The non-streaming version (RunAsync) should execute the workflow and produce events,
|
||||
/// similar to the streaming version (StreamAsync + TrySendMessageAsync).
|
||||
/// </summary>
|
||||
[Fact]
|
||||
public async Task RunAsyncShouldExecuteWorkflowAsync()
|
||||
{
|
||||
// Arrange: Create a simple agent that responds to messages
|
||||
var agent = new SimpleTestAgent("test-agent");
|
||||
var workflow = AgentWorkflowBuilder.BuildSequential(agent);
|
||||
var inputMessage = new ChatMessage(ChatRole.User, "Hello");
|
||||
|
||||
// Act: Execute using non-streaming RunAsync
|
||||
Run run = await InProcessExecution.RunAsync(workflow, new List<ChatMessage> { inputMessage });
|
||||
|
||||
// Assert: The workflow should have executed and produced events
|
||||
RunStatus status = await run.GetStatusAsync();
|
||||
status.Should().Be(RunStatus.Idle, "workflow should complete execution");
|
||||
|
||||
// The run should have events (at minimum, a WorkflowOutputEvent)
|
||||
run.OutgoingEvents.Should().NotBeEmpty("workflow should produce events during execution");
|
||||
|
||||
// Check that we have an agent execution event
|
||||
var agentEvents = run.OutgoingEvents.OfType<AgentRunUpdateEvent>().ToList();
|
||||
agentEvents.Should().NotBeEmpty("agent should have executed and produced update events");
|
||||
|
||||
// Check that we have output events
|
||||
var outputEvents = run.OutgoingEvents.OfType<WorkflowOutputEvent>().ToList();
|
||||
outputEvents.Should().NotBeEmpty("workflow should produce output events");
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// This test shows that the streaming version works correctly when TurnToken is sent following a message.
|
||||
/// </summary>
|
||||
[Fact]
|
||||
public async Task StreamAsyncWithTurnTokenShouldExecuteWorkflowAsync()
|
||||
{
|
||||
// Arrange: Create a simple agent that responds to messages
|
||||
var agent = new SimpleTestAgent("test-agent");
|
||||
var workflow = AgentWorkflowBuilder.BuildSequential(agent);
|
||||
var inputMessage = new ChatMessage(ChatRole.User, "Hello");
|
||||
|
||||
// Act: Execute using streaming version with TurnToken
|
||||
await using StreamingRun run = await InProcessExecution.StreamAsync(workflow, new List<ChatMessage> { inputMessage });
|
||||
|
||||
// Send TurnToken to actually trigger execution (this is the key step)
|
||||
bool messageSent = await run.TrySendMessageAsync(new TurnToken(emitEvents: true));
|
||||
messageSent.Should().BeTrue("TurnToken should be accepted");
|
||||
|
||||
// Collect events
|
||||
List<WorkflowEvent> events = new();
|
||||
await foreach (WorkflowEvent evt in run.WatchStreamAsync())
|
||||
{
|
||||
events.Add(evt);
|
||||
}
|
||||
|
||||
// Assert: The workflow should have executed and produced events
|
||||
RunStatus status = await run.GetStatusAsync();
|
||||
status.Should().Be(RunStatus.Idle, "workflow should complete execution");
|
||||
|
||||
events.Should().NotBeEmpty("workflow should produce events during execution");
|
||||
|
||||
// Check that we have agent execution events
|
||||
var agentEvents = events.OfType<AgentRunUpdateEvent>().ToList();
|
||||
agentEvents.Should().NotBeEmpty("agent should have executed and produced update events");
|
||||
|
||||
// Check that we have output events
|
||||
var outputEvents = events.OfType<WorkflowOutputEvent>().ToList();
|
||||
outputEvents.Should().NotBeEmpty("workflow should produce output events");
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// This test compares the behavior of RunAsync vs StreamAsync to highlight the difference.
|
||||
/// Both should produce similar results, but as of issue #1315, RunAsync fails to execute.
|
||||
/// </summary>
|
||||
[Fact]
|
||||
public async Task RunAsyncAndStreamAsyncShouldProduceSimilarResultsAsync()
|
||||
{
|
||||
// Arrange: Create the same workflow for both tests
|
||||
var agent1 = new SimpleTestAgent("test-agent-1");
|
||||
var workflow1 = AgentWorkflowBuilder.BuildSequential(agent1);
|
||||
|
||||
var agent2 = new SimpleTestAgent("test-agent-2");
|
||||
var workflow2 = AgentWorkflowBuilder.BuildSequential(agent2);
|
||||
|
||||
var inputMessage = new ChatMessage(ChatRole.User, "Test message");
|
||||
|
||||
// Act 1: Execute using RunAsync (non-streaming)
|
||||
Run nonStreamingRun = await InProcessExecution.RunAsync(workflow1, new List<ChatMessage> { inputMessage });
|
||||
var nonStreamingEvents = nonStreamingRun.OutgoingEvents.ToList();
|
||||
|
||||
// Act 2: Execute using StreamAsync (streaming) with TurnToken
|
||||
await using StreamingRun streamingRun = await InProcessExecution.StreamAsync(workflow2, new List<ChatMessage> { inputMessage });
|
||||
await streamingRun.TrySendMessageAsync(new TurnToken(emitEvents: true));
|
||||
|
||||
List<WorkflowEvent> streamingEvents = new();
|
||||
await foreach (WorkflowEvent evt in streamingRun.WatchStreamAsync())
|
||||
{
|
||||
streamingEvents.Add(evt);
|
||||
}
|
||||
|
||||
// Assert: Both should have produced events
|
||||
// The streaming version works (we know this from the issue report)
|
||||
streamingEvents.Should().NotBeEmpty("streaming version should produce events");
|
||||
|
||||
// The non-streaming version should also produce events (this is the bug being tested)
|
||||
nonStreamingEvents.Should().NotBeEmpty("non-streaming version should also produce events");
|
||||
|
||||
// Both should have similar types of events
|
||||
var streamingAgentEvents = streamingEvents.OfType<AgentRunUpdateEvent>().Count();
|
||||
var nonStreamingAgentEvents = nonStreamingEvents.OfType<AgentRunUpdateEvent>().Count();
|
||||
|
||||
nonStreamingAgentEvents.Should().Be(streamingAgentEvents,
|
||||
"both versions should produce the same number of agent events");
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Simple test agent that echoes back the input message.
|
||||
/// </summary>
|
||||
private sealed class SimpleTestAgent : AIAgent
|
||||
{
|
||||
private readonly string _name;
|
||||
|
||||
public SimpleTestAgent(string name)
|
||||
{
|
||||
this._name = name;
|
||||
}
|
||||
|
||||
public override string Name => this._name;
|
||||
|
||||
public override AgentThread GetNewThread() => new SimpleTestAgentThread();
|
||||
|
||||
public override AgentThread DeserializeThread(System.Text.Json.JsonElement serializedThread,
|
||||
System.Text.Json.JsonSerializerOptions? jsonSerializerOptions = null) => new SimpleTestAgentThread();
|
||||
|
||||
public override Task<AgentRunResponse> RunAsync(
|
||||
IEnumerable<ChatMessage> messages,
|
||||
AgentThread? thread = null,
|
||||
AgentRunOptions? options = null,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
var lastMessage = messages.LastOrDefault();
|
||||
var responseMessage = new ChatMessage(ChatRole.Assistant, $"Echo: {lastMessage?.Text ?? "no message"}");
|
||||
return Task.FromResult(new AgentRunResponse(responseMessage));
|
||||
}
|
||||
|
||||
public override async IAsyncEnumerable<AgentRunResponseUpdate> RunStreamingAsync(
|
||||
IEnumerable<ChatMessage> messages,
|
||||
AgentThread? thread = null,
|
||||
AgentRunOptions? options = null,
|
||||
[EnumeratorCancellation] CancellationToken cancellationToken = default)
|
||||
{
|
||||
await Task.Yield();
|
||||
|
||||
var lastMessage = messages.LastOrDefault();
|
||||
var responseText = $"Echo: {lastMessage?.Text ?? "no message"}";
|
||||
|
||||
string messageId = Guid.NewGuid().ToString("N");
|
||||
|
||||
// Yield role first
|
||||
yield return new AgentRunResponseUpdate(ChatRole.Assistant, this._name)
|
||||
{
|
||||
AuthorName = this._name,
|
||||
MessageId = messageId
|
||||
};
|
||||
|
||||
// Then yield content
|
||||
yield return new AgentRunResponseUpdate(ChatRole.Assistant, responseText)
|
||||
{
|
||||
AuthorName = this._name,
|
||||
MessageId = messageId
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Simple thread implementation for SimpleTestAgent.
|
||||
/// </summary>
|
||||
private sealed class SimpleTestAgentThread : InMemoryAgentThread
|
||||
{
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user