.NET Workflows - Support agent level function invocation for declarative workflow (#1442)

* Checkpoint

* Checkpoint

* Checkpoint

* Good

* Namespace

* Namespace

* Dun

* Async Test

* AgentId

* Portable pattern

* Portable2

* Portable3

* Respond to comments

* Namespace

* Function call selection

* ToHashSet

* ToHashSet

* Updated

* Parameter name

* Final

* Tests
This commit is contained in:
Chris
2025-10-14 20:42:15 +00:00
committed by GitHub
parent a4039134de
commit a6b6937b94
25 changed files with 902 additions and 80 deletions
@@ -1,16 +1,20 @@
// Copyright (c) Microsoft. All rights reserved.
// Uncomment this to enable JSON checkpointing to the local file system.
#define CHECKPOINT_JSON
//#define CHECKPOINT_JSON
using System.Diagnostics;
using System.Reflection;
using System.Text.Json;
using Azure.AI.Agents.Persistent;
using Azure.Identity;
using Microsoft.Agents.AI.Workflows;
#if CHECKPOINT_JSON
using Microsoft.Agents.AI.Workflows.Checkpointing;
#endif
using Microsoft.Agents.AI.Workflows.Declarative;
using Microsoft.Agents.AI.Workflows.Declarative.Events;
using Microsoft.Agents.AI.Workflows.Declarative.Kit;
using Microsoft.Extensions.AI;
using Microsoft.Extensions.Configuration;
@@ -63,20 +67,21 @@ internal sealed class Program
#if CHECKPOINT_JSON
// Use a file-system based JSON checkpoint store to persist checkpoints to disk.
DirectoryInfo checkpointFolder = Directory.CreateDirectory(Path.Combine(".", $"chk-{DateTime.Now:YYmmdd-hhMMss-ff}"));
DirectoryInfo checkpointFolder = Directory.CreateDirectory(Path.Combine(".", $"chk-{DateTime.Now:yyMMdd-hhmmss-ff}"));
CheckpointManager checkpointManager = CheckpointManager.CreateJson(new FileSystemJsonCheckpointStore(checkpointFolder));
Checkpointed<StreamingRun> run = await InProcessExecution.StreamAsync(workflow, input, checkpointManager);
#else
// Use an in-memory checkpoint store that will not persist checkpoints beyond the lifetime of the process.
CheckpointManager checkpointManager = CheckpointManager.CreateInMemory();
#endif
Checkpointed<StreamingRun> run = await InProcessExecution.StreamAsync(workflow, input, checkpointManager);
bool isComplete = false;
InputResponse? response = null;
object? response = null;
do
{
ExternalRequest? inputRequest = await this.MonitorAndDisposeWorkflowRunAsync(run, response);
if (inputRequest is not null)
ExternalRequest? externalRequest = await this.MonitorAndDisposeWorkflowRunAsync(run, response);
if (externalRequest is not null)
{
Notify("\nWORKFLOW: Yield");
@@ -86,7 +91,7 @@ internal sealed class Program
}
// Process the external request.
response = HandleExternalRequest(inputRequest);
response = await this.HandleExternalRequestAsync(externalRequest);
// Let's resume on an entirely new workflow instance to demonstrate checkpoint portability.
workflow = this.CreateWorkflow();
@@ -107,11 +112,25 @@ internal sealed class Program
Notify("\nWORKFLOW: Done!\n");
}
/// <summary>
/// Create the workflow from the declarative YAML. Includes definition of the
/// <see cref="DeclarativeWorkflowOptions" /> and the associated <see cref="WorkflowAgentProvider"/>.
/// </summary>
/// <remarks>
/// The value assigned to <see cref="IncludeFunctions" /> controls on whether the function
/// tools (<see cref="AIFunction"/>) initialized in the constructor are included for auto-invocation.
/// </remarks>
private Workflow CreateWorkflow()
{
// Use DeclarativeWorkflowBuilder to build a workflow based on a YAML file.
AzureAgentProvider agentProvider = new(this.FoundryEndpoint, new AzureCliCredential())
{
// Functions included here will be auto-executed by the framework.
Functions = IncludeFunctions ? this.FunctionMap.Values : null,
};
DeclarativeWorkflowOptions options =
new(new AzureAgentProvider(this.FoundryEndpoint, new AzureCliCredential()))
new(agentProvider)
{
Configuration = this.Configuration,
//ConversationId = null, // Assign to continue a conversation
@@ -121,8 +140,18 @@ internal sealed class Program
return DeclarativeWorkflowBuilder.Build<string>(this.WorkflowFile, options);
}
/// <summary>
/// Configuration key used to identify the Foundry project endpoint.
/// </summary>
private const string ConfigKeyFoundryEndpoint = "FOUNDRY_PROJECT_ENDPOINT";
/// <summary>
/// Controls on whether the function tools (<see cref="AIFunction"/>) initialized
/// in the constructor are included for auto-invocation.
/// NOTE: By default, no functions exist as part of this sample.
/// </summary>
private const bool IncludeFunctions = true;
private static Dictionary<string, string> NameCache { get; } = [];
private static HashSet<string> FileCache { get; } = [];
@@ -132,6 +161,7 @@ internal sealed class Program
private PersistentAgentsClient FoundryClient { get; }
private IConfiguration Configuration { get; }
private CheckpointInfo? LastCheckpoint { get; set; }
private Dictionary<string, AIFunction> FunctionMap { get; }
private Program(string workflowFile, string? workflowInput)
{
@@ -142,12 +172,21 @@ internal sealed class Program
this.FoundryEndpoint = this.Configuration[ConfigKeyFoundryEndpoint] ?? throw new InvalidOperationException($"Undefined configuration setting: {ConfigKeyFoundryEndpoint}");
this.FoundryClient = new PersistentAgentsClient(this.FoundryEndpoint, new AzureCliCredential());
List<AIFunction> functions =
[
// Manually define any custom functions that may be required by agents within the workflow.
// By default, this sample does not include any functions.
//AIFunctionFactory.Create(),
];
this.FunctionMap = functions.ToDictionary(f => f.Name);
}
private async Task<ExternalRequest?> MonitorAndDisposeWorkflowRunAsync(Checkpointed<StreamingRun> run, InputResponse? response = null)
private async Task<ExternalRequest?> MonitorAndDisposeWorkflowRunAsync(Checkpointed<StreamingRun> run, object? response = null)
{
await using IAsyncDisposable disposeRun = run;
bool hasStreamed = false;
string? messageId = null;
await foreach (WorkflowEvent workflowEvent in run.Run.WatchStreamAsync())
@@ -211,11 +250,12 @@ internal sealed class Program
case AgentRunUpdateEvent streamEvent:
if (!string.Equals(messageId, streamEvent.Update.MessageId, StringComparison.Ordinal))
{
hasStreamed = false;
messageId = streamEvent.Update.MessageId;
if (messageId is not null)
{
string? agentId = streamEvent.Update.AuthorName;
string? agentId = streamEvent.Update.AgentId;
if (agentId is not null)
{
if (!NameCache.TryGetValue(agentId, out string? realName))
@@ -245,11 +285,18 @@ internal sealed class Program
await DownloadFileContentAsync(Path.GetFileName(messageUpdate.TextAnnotation?.TextToReplace ?? "response.png"), content);
}
break;
case RequiredActionUpdate actionUpdate:
Console.ForegroundColor = ConsoleColor.White;
Console.Write($"Calling tool: {actionUpdate.FunctionName}");
Console.ForegroundColor = ConsoleColor.DarkGray;
Console.WriteLine($" [{actionUpdate.ToolCallId}]");
break;
}
try
{
Console.ResetColor();
Console.Write(streamEvent.Data);
Console.Write(streamEvent.Update.Text);
hasStreamed |= !string.IsNullOrEmpty(streamEvent.Update.Text);
}
finally
{
@@ -260,7 +307,11 @@ internal sealed class Program
case AgentRunResponseEvent messageEvent:
try
{
Console.WriteLine();
if (hasStreamed)
{
Console.WriteLine();
}
if (messageEvent.Response.Usage is not null)
{
Console.ForegroundColor = ConsoleColor.DarkGray;
@@ -277,14 +328,31 @@ internal sealed class Program
return default;
}
private static InputResponse HandleExternalRequest(ExternalRequest request)
/// <summary>
/// Handle request for external input, either from a human or a function tool invocation.
/// </summary>
private async ValueTask<object> HandleExternalRequestAsync(ExternalRequest request) =>
request.Data.TypeId.TypeName switch
{
// Request for human input
_ when request.Data.TypeId.IsMatch<InputRequest>() => HandleInputRequest(request.DataAs<InputRequest>()!),
// Request for function tool invocation. (Only active when functions are defined and IncludeFunctions is true.)
_ when request.Data.TypeId.IsMatch<AgentToolRequest>() => await this.HandleToolRequestAsync(request.DataAs<AgentToolRequest>()!),
// Unknown request type.
_ => throw new InvalidOperationException($"Unsupported external request type: {request.GetType().Name}."),
};
/// <summary>
/// Handle request for human input.
/// </summary>
private static InputResponse HandleInputRequest(InputRequest request)
{
InputRequest? message = request.Data.As<InputRequest>();
string? userInput;
do
{
Console.ForegroundColor = ConsoleColor.DarkGreen;
Console.Write($"\n{message?.Prompt ?? "INPUT:"} ");
Console.Write($"\n{request.Prompt ?? "INPUT:"} ");
Console.ForegroundColor = ConsoleColor.White;
userInput = Console.ReadLine();
}
@@ -293,6 +361,30 @@ internal sealed class Program
return new InputResponse(userInput);
}
/// <summary>
/// Handle a function tool request by invoking the specified tools and returning the results.
/// </summary>
/// <remarks>
/// This handler is only active when <see cref="IncludeFunctions"/> is set to true and
/// one or more <see cref="AIFunction"/> instances are defined in the constructor.
/// </remarks>
private async ValueTask<AgentToolResponse> HandleToolRequestAsync(AgentToolRequest request)
{
Task<FunctionResultContent>[] functionTasks = request.FunctionCalls.Select(functionCall => InvokesToolAsync(functionCall)).ToArray();
await Task.WhenAll(functionTasks);
return AgentToolResponse.Create(request, functionTasks.Select(task => task.Result));
async Task<FunctionResultContent> InvokesToolAsync(FunctionCallContent functionCall)
{
AIFunction functionTool = this.FunctionMap[functionCall.Name];
AIFunctionArguments? functionArguments = functionCall.Arguments is null ? null : new(functionCall.Arguments.NormalizePortableValues());
object? result = await functionTool.InvokeAsync(functionArguments);
return new FunctionResultContent(functionCall.CallId, JsonSerializer.Serialize(result));
}
}
private static string? ParseWorkflowFile(string[] args)
{
string? workflowFile = args.FirstOrDefault();