working e2e orchestration completion.

This commit is contained in:
Shyju Krishnankutty
2026-01-14 16:18:21 -08:00
Unverified
parent 7f22a87a24
commit defb533b95
7 changed files with 47 additions and 66 deletions
@@ -68,7 +68,7 @@ internal sealed class BuiltInFunctionExecutor : IFunctionExecutor
}
}
if (durableTaskClient is null && context.FunctionDefinition.EntryPoint != BuiltInFunctions.RunWorkflowOrechstrtationFunctionEntryPoint)
if (durableTaskClient is null)
{
// This is not expected to happen since all built-in functions are
// expected to have a Durable Task client binding.
@@ -117,11 +117,6 @@ internal sealed class BuiltInFunctionExecutor : IFunctionExecutor
if (context.FunctionDefinition.EntryPoint == BuiltInFunctions.RunWorkflowOrechstrtationHttpFunctionEntryPoint)
{
//if (httpRequestData == null)
//{
// throw new InvalidOperationException($"HTTP request data binding is missing for the invocation {context.InvocationId}.");
//}
if (httpRequestData == null)
{
throw new InvalidOperationException($"HTTP request data binding is missing for the invocation {context.InvocationId}.");
@@ -23,7 +23,6 @@ internal static class BuiltInFunctions
internal static readonly string RunAgentHttpFunctionEntryPoint = $"{typeof(BuiltInFunctions).FullName!}.{nameof(RunAgentHttpAsync)}";
internal static readonly string RunAgentEntityFunctionEntryPoint = $"{typeof(BuiltInFunctions).FullName!}.{nameof(InvokeAgentAsync)}";
internal static readonly string RunWorkflowOrechstrtationHttpFunctionEntryPoint = $"{typeof(BuiltInFunctions).FullName!}.{nameof(RunWorkflowOrechstrtationHttpTriggerAsync)}";
internal static readonly string RunWorkflowOrechstrtationFunctionEntryPoint = $"{typeof(BuiltInFunctions).FullName!}.{nameof(RunWorkflowOrchestratorAsync)}";
internal static readonly string InvokeWorkflowActivityFunctionEntryPoint = $"{typeof(BuiltInFunctions).FullName!}.{nameof(InvokeWorkflowActivityAsync)}";
internal static readonly string RunAgentMcpToolFunctionEntryPoint = $"{typeof(BuiltInFunctions).FullName!}.{nameof(RunMcpToolAsync)}";
@@ -42,21 +41,6 @@ internal static class BuiltInFunctions
return runner.ExecuteActivityAsync(activityFunctionName, input, functionContext);
}
public static async Task<List<string>> RunWorkflowOrchestratorAsync(string taskOrchestrationContext, FunctionContext functionsContext)
{
var logger = functionsContext.GetLogger("BuiltInFunctions");
var outputs = new List<string>();
const string WorkflowName = "MyTestWorkflow"; // to do: get from TaskOrchestrationContext
if (logger.IsEnabled(LogLevel.Information))
{
logger.LogInformation("Orchestrator {WorkflowName} is executing. Input: {Input}", WorkflowName, taskOrchestrationContext);
}
//var runner = functionsContext.InstanceServices.GetService<DurableWorkflowRunner>();
//await runner!.RunAsync(null, WorkflowName);
return outputs;
}
// Exposed as an entity trigger via AgentFunctionsProvider
public static Task<string> InvokeAgentAsync(
[DurableClient] DurableTaskClient client,
@@ -86,14 +70,12 @@ internal static class BuiltInFunctions
[DurableClient] DurableTaskClient client,
FunctionContext context)
{
// to do: Retrieve the workflow and execute it.
var workflowName = context.FunctionDefinition.Name.Replace("http-", "");
var inputMessage = await req.ReadAsStringAsync();
//string instanceId = await client.ScheduleNewOrchestrationInstanceAsync("dafx-MyTestWorkflow");
string instanceId = await client.ScheduleNewOrchestrationInstanceAsync("WorkflowRunnerOrchestration", new DuableWorkflowRunRequest { WorkflowName = workflowName, Input = inputMessage! }); //OrchFunction"); // dafx-MyTestWorkflow");
HttpResponseData response = req.CreateResponse(HttpStatusCode.OK);
await response.WriteStringAsync($"InvokeWorkflowOrechstrtationAsync is invoked for {workflowName}.{instanceId}");
HttpResponseData response = req.CreateResponse(HttpStatusCode.Accepted);
await response.WriteStringAsync($"InvokeWorkflowOrechstrtationAsync is invoked for {workflowName}. Orchestration instanceId: {instanceId}");
return response;
}
@@ -9,6 +9,7 @@ using Microsoft.DurableTask.Worker;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.DependencyInjection.Extensions;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
namespace Microsoft.Agents.AI.Hosting.AzureFunctions;
@@ -68,7 +69,6 @@ public static class DurableOptionsExtensions
builder.UseWhen<BuiltInFunctionExecutionMiddleware>(static context =>
string.Equals(context.FunctionDefinition.EntryPoint, BuiltInFunctions.RunAgentHttpFunctionEntryPoint, StringComparison.Ordinal) ||
string.Equals(context.FunctionDefinition.EntryPoint, BuiltInFunctions.RunAgentMcpToolFunctionEntryPoint, StringComparison.Ordinal) ||
string.Equals(context.FunctionDefinition.EntryPoint, BuiltInFunctions.RunWorkflowOrechstrtationFunctionEntryPoint, StringComparison.Ordinal) ||
string.Equals(context.FunctionDefinition.EntryPoint, BuiltInFunctions.RunWorkflowOrechstrtationHttpFunctionEntryPoint, StringComparison.Ordinal) ||
string.Equals(context.FunctionDefinition.EntryPoint, BuiltInFunctions.InvokeWorkflowActivityFunctionEntryPoint, StringComparison.Ordinal) ||
string.Equals(context.FunctionDefinition.EntryPoint, BuiltInFunctions.RunAgentEntityFunctionEntryPoint, StringComparison.Ordinal));
@@ -85,9 +85,15 @@ public static class DurableOptionsExtensions
{
throw new InvalidOperationException("FunctionContext is not available in the orchestration context.");
}
var logger = tc.CreateReplaySafeLogger("WorkflowRunnerOrchestration");
DurableWorkflowRunner runner = functionContext.InstanceServices.GetRequiredService<DurableWorkflowRunner>();
return await runner.RunWorkflowOrchestrationAsync(tc, inputBindingData).ConfigureAwait(false);
var orchestrationResult = await runner.RunWorkflowOrchestrationAsync(tc, inputBindingData, logger);
if (logger.IsEnabled(LogLevel.Information))
{
logger.LogInformation("Durable workflow orchestration completed. Result:{Result}", string.Join(",", orchestrationResult));
}
return orchestrationResult;
}));
builder.Services.AddSingleton<IFunctionMetadataTransformer, DurableWorkflowFunctionMetadataTransformer>();
@@ -21,23 +21,23 @@ internal sealed class DurableWorkflowRunner
this._options = durableOptions.Workflows;
}
internal async Task<List<string>> RunWorkflowOrchestrationAsync(TaskOrchestrationContext context, DuableWorkflowRunRequest input)
internal async Task<List<string>> RunWorkflowOrchestrationAsync(TaskOrchestrationContext taskOrchestrationContext, DuableWorkflowRunRequest input, ILogger logger)
{
ArgumentNullException.ThrowIfNull(context);
ArgumentNullException.ThrowIfNull(taskOrchestrationContext);
ArgumentNullException.ThrowIfNull(input);
string workflowName = input.WorkflowName;
this._logger.LogAttemptingToRunWorkflow(workflowName);
logger.LogAttemptingToRunWorkflow(workflowName);
if (!this._options.Workflows.TryGetValue(workflowName, out Workflow? wf))
{
throw new InvalidOperationException($"Workflow '{workflowName}' not found.");
}
this._logger.LogRunningWorkflow(wf.Name);
logger.LogRunningWorkflow(wf.Name);
var result = await this.RunExecutorsInWorkFlowAsync(context, wf, input.Input).ConfigureAwait(false);
var result = await this.RunExecutorsInWorkFlowAsync(taskOrchestrationContext, wf, input.Input, logger);
return [result];
}
@@ -45,16 +45,17 @@ internal sealed class DurableWorkflowRunner
private async Task<string> RunExecutorsInWorkFlowAsync(
TaskOrchestrationContext taskOrchestrationContext,
Workflow workflow,
string initialInput)
string initialInput,
ILogger logger)
{
List<string> executorResult = [];
if (this._logger.IsEnabled(LogLevel.Information))
if (logger.IsEnabled(LogLevel.Information))
{
foreach (WorkflowExecutorInfo executorInfo in WorkflowHelper.GetExecutorsFromWorkflowInOrder(workflow))
{
string triggerName = this.BuildTriggerName(workflow.Name!, executorInfo.ExecutorId);
this._logger.LogInformation(
logger.LogInformation(
" Scheduling executor '{ExecutorId}' (IsAgentic: {IsAgentic}) with trigger name '{TriggerName}'",
executorInfo.ExecutorId,
executorInfo.IsAgenticExecutor,
@@ -69,7 +70,7 @@ internal sealed class DurableWorkflowRunner
else
{
string AgentName = this.GetAgentNameFromExecutorId(workflow.Name!, executorInfo.ExecutorId);
this._logger.LogInformation(
logger.LogInformation(
" Invoking agentic executor '{ExecutorId}'",
AgentName);
DurableAIAgent agent = taskOrchestrationContext.GetAgent(AgentName);
@@ -78,7 +79,7 @@ internal sealed class DurableWorkflowRunner
AgentThread destinationThread = agent.GetNewThread();
var agentResponse = await agent.RunAsync(input, destinationThread);
executorResult.Add(agentResponse.Text);
this._logger.LogInformation(
logger.LogInformation(
"Agentic executor '{ExecutorId}' completed with response: {AgentResponse}",
AgentName,
agentResponse);
@@ -162,7 +163,7 @@ internal sealed class DurableWorkflowRunner
"activity-run",
CancellationToken.None).ConfigureAwait(false);
// Create a minimal workflow context for the executor
// Create a minimal workflow taskOrchestrationContext for the executor
MinimalActivityContext context = new(executorPair.Key);
// Execute the executor with the input
@@ -1,4 +1,4 @@
// Copyright (c) Microsoft. All rights reserved.
// Copyright (c) Microsoft. All rights reserved.
using Microsoft.Agents.AI.DurableTask;
using Microsoft.Azure.Functions.Worker.Builder;
@@ -38,7 +38,6 @@ public static class FunctionsApplicationBuilderExtensions
builder.UseWhen<BuiltInFunctionExecutionMiddleware>(static context =>
string.Equals(context.FunctionDefinition.EntryPoint, BuiltInFunctions.RunAgentHttpFunctionEntryPoint, StringComparison.Ordinal) ||
string.Equals(context.FunctionDefinition.EntryPoint, BuiltInFunctions.RunAgentMcpToolFunctionEntryPoint, StringComparison.Ordinal) ||
string.Equals(context.FunctionDefinition.EntryPoint, BuiltInFunctions.RunWorkflowOrechstrtationFunctionEntryPoint, StringComparison.Ordinal) ||
string.Equals(context.FunctionDefinition.EntryPoint, BuiltInFunctions.RunWorkflowOrechstrtationHttpFunctionEntryPoint, StringComparison.Ordinal) ||
string.Equals(context.FunctionDefinition.EntryPoint, BuiltInFunctions.InvokeWorkflowActivityFunctionEntryPoint, StringComparison.Ordinal) ||
string.Equals(context.FunctionDefinition.EntryPoint, BuiltInFunctions.RunAgentEntityFunctionEntryPoint, StringComparison.Ordinal));