using Aryx.AgentHost.Contracts; using Microsoft.Agents.AI; using Microsoft.Agents.AI.Workflows; namespace Aryx.AgentHost.Services; internal sealed class WorkflowRunner { public Workflow BuildWorkflow( WorkflowDefinitionDto workflowDefinition, IReadOnlyList agents, IReadOnlyList? workflowLibrary = null) { ArgumentNullException.ThrowIfNull(workflowDefinition); ArgumentNullException.ThrowIfNull(agents); Dictionary workflowLibraryMap = workflowLibrary? .Where(candidate => !string.IsNullOrWhiteSpace(candidate.Id)) .GroupBy(candidate => candidate.Id, StringComparer.Ordinal) .ToDictionary(group => group.Key, group => group.Last(), StringComparer.Ordinal) ?? new Dictionary(StringComparer.Ordinal); List agentIds = ResolveAgentIds(workflowDefinition, workflowLibraryMap); Dictionary agentMap = agentIds .Zip(agents, (agentId, agent) => (agentId, agent)) .ToDictionary(pair => pair.agentId, pair => pair.agent, StringComparer.Ordinal); return BuildWorkflow(workflowDefinition, agentMap, workflowLibraryMap); } private Workflow BuildWorkflow( WorkflowDefinitionDto workflowDefinition, IReadOnlyDictionary agentMap, IReadOnlyDictionary workflowLibrary) { WorkflowNodeDto startNode = workflowDefinition.Graph.Nodes.Single(node => string.Equals(node.Kind, "start", StringComparison.OrdinalIgnoreCase)); WorkflowNodeDto endNode = workflowDefinition.Graph.Nodes.Single(node => string.Equals(node.Kind, "end", StringComparison.OrdinalIgnoreCase)); WorkflowStateScopeCatalog stateCatalog = new(workflowDefinition.Settings.StateScopes); Dictionary routes = new(StringComparer.Ordinal); foreach (WorkflowNodeDto node in workflowDefinition.Graph.Nodes) { routes[node.Id] = CreateNodeRoute(node, agentMap, workflowLibrary, stateCatalog); } WorkflowBuilder builder = new(routes[startNode.Id].Entry); foreach (WorkflowNodeRoute route in routes.Values) { foreach ((ExecutorBinding source, ExecutorBinding target) in route.InternalEdges) { builder.AddEdge(source, target); } } foreach (WorkflowEdgeDto edge in workflowDefinition.Graph.Edges.Where(edge => string.Equals(edge.Kind, "direct", StringComparison.OrdinalIgnoreCase))) { Func? condition = WorkflowConditionEvaluator.Compile(edge); ExecutorBinding source = routes[edge.Source].Exit; ExecutorBinding target = routes[edge.Target].Entry; if (condition is null) { builder.AddEdge(source, target); } else { builder.AddEdge(source, target, condition); } } foreach (IGrouping fanOutGroup in workflowDefinition.Graph.Edges .Where(edge => string.Equals(edge.Kind, "fan-out", StringComparison.OrdinalIgnoreCase)) .GroupBy(edge => edge.Source, StringComparer.Ordinal)) { WorkflowEdgeDto[] fanOutEdges = fanOutGroup.ToArray(); ExecutorBinding source = routes[fanOutGroup.Key].Exit; ExecutorBinding[] targets = fanOutEdges.Select(edge => routes[edge.Target].Entry).ToArray(); Func?[] compiledConditions = fanOutEdges .Select(WorkflowConditionEvaluator.Compile) .ToArray(); bool hasConditionalRouting = fanOutEdges.Any(edge => edge.Condition is not null); if (!hasConditionalRouting) { builder.AddFanOutEdge(source, targets); continue; } builder.AddFanOutEdge( source, targets, (payload, _) => fanOutEdges .Select((edge, index) => (edge, index)) .Where(pair => compiledConditions[pair.index]?.Invoke(payload) ?? true) .Select(pair => pair.index) .ToArray()); } foreach (IGrouping fanInGroup in workflowDefinition.Graph.Edges .Where(edge => string.Equals(edge.Kind, "fan-in", StringComparison.OrdinalIgnoreCase)) .GroupBy(edge => edge.Target, StringComparer.Ordinal)) { builder.AddFanInBarrierEdge( fanInGroup.Select(edge => routes[edge.Source].Exit).ToArray(), routes[fanInGroup.Key].Entry); } if (!string.IsNullOrWhiteSpace(workflowDefinition.Name)) { builder = builder.WithName(workflowDefinition.Name); } return builder.WithOutputFrom(routes[endNode.Id].Exit).Build(); } private WorkflowNodeRoute CreateNodeRoute( WorkflowNodeDto node, IReadOnlyDictionary agentMap, IReadOnlyDictionary workflowLibrary, WorkflowStateScopeCatalog stateCatalog) { if (string.Equals(node.Kind, "start", StringComparison.OrdinalIgnoreCase)) { ExecutorBinding binding = new ChatForwardingExecutor(node.Id).BindExecutor(); return new WorkflowNodeRoute(binding); } if (string.Equals(node.Kind, "end", StringComparison.OrdinalIgnoreCase)) { ExecutorBinding binding = new WorkflowOutputMessagesExecutor(node.Id).BindExecutor(); return new WorkflowNodeRoute(binding); } if (string.Equals(node.Kind, "agent", StringComparison.OrdinalIgnoreCase)) { string agentId = !string.IsNullOrWhiteSpace(node.Config.Id) ? node.Config.Id : node.Id; if (!agentMap.TryGetValue(agentId, out AIAgent? agent)) { throw new InvalidOperationException($"Workflow node \"{node.Id}\" references unknown agent \"{agentId}\"."); } return new WorkflowNodeRoute(agent.BindAsExecutor(CopilotAgentBundle.CreateAgentHostOptions())); } if (string.Equals(node.Kind, "code-executor", StringComparison.OrdinalIgnoreCase)) { string implementation = NormalizeRequired(node.Config.Implementation, $"Workflow code executor \"{node.Id}\" requires an implementation."); ExecutorBinding binding = new WorkflowCodeExecutor(node.Id, implementation, stateCatalog).BindExecutor(); return new WorkflowNodeRoute(binding); } if (string.Equals(node.Kind, "function-executor", StringComparison.OrdinalIgnoreCase)) { string functionRef = NormalizeRequired(node.Config.FunctionRef, $"Workflow function executor \"{node.Id}\" requires a functionRef."); if (!WorkflowFunctionRegistry.IsSupported(functionRef)) { throw new InvalidOperationException( $"Workflow function executor \"{node.Id}\" references unsupported functionRef \"{functionRef}\"."); } ExecutorBinding binding = new WorkflowFunctionExecutor(node.Id, functionRef, node.Config.Parameters, stateCatalog).BindExecutor(); return new WorkflowNodeRoute(binding); } if (string.Equals(node.Kind, "request-port", StringComparison.OrdinalIgnoreCase)) { return CreateRequestPortRoute(node); } if (string.Equals(node.Kind, "sub-workflow", StringComparison.OrdinalIgnoreCase)) { WorkflowDefinitionDto subWorkflowDefinition = node.ResolveSubWorkflowDefinition(workflowLibrary); Workflow subWorkflow = BuildWorkflow(subWorkflowDefinition, agentMap, workflowLibrary); return new WorkflowNodeRoute(subWorkflow.BindAsExecutor(node.Id)); } throw new NotSupportedException($"Workflow node kind \"{node.Kind}\" is not executable yet."); } private static WorkflowNodeRoute CreateRequestPortRoute(WorkflowNodeDto node) { WorkflowRequestPortNodeDefinition definition = new( node.Id, NormalizeOptionalString(node.Label) ?? node.Id, NormalizeRequired(node.Config.PortId, $"Workflow request port \"{node.Id}\" requires a portId."), NormalizeRequired(node.Config.RequestType, $"Workflow request port \"{node.Id}\" requires a requestType."), NormalizeRequired(node.Config.ResponseType, $"Workflow request port \"{node.Id}\" requires a responseType."), NormalizeOptionalString(node.Config.Prompt)); RequestPort port = new(definition.PortId, typeof(WorkflowRequestPortPromptRequest), typeof(object)); ExecutorBinding entry = new WorkflowRequestPortIngressExecutor(definition, port).BindExecutor(); ExecutorBinding portBinding = new RequestPortBinding(port, false); ExecutorBinding exit = new WorkflowRequestPortResponseExecutor(node.Id).BindExecutor(); return new WorkflowNodeRoute(entry, exit, [(entry, portBinding), (portBinding, exit)]); } private static string NormalizeRequired(string? value, string errorMessage) => NormalizeOptionalString(value) ?? throw new InvalidOperationException(errorMessage); private static string? NormalizeOptionalString(string? value) => string.IsNullOrWhiteSpace(value) ? null : value.Trim(); private static List ResolveAgentIds( WorkflowDefinitionDto workflowDefinition, IReadOnlyDictionary workflowLibrary) { return workflowDefinition.GetAllAgentNodes(workflowLibrary) .Select(node => node.GetAgentId()) .ToList(); } private sealed record WorkflowNodeRoute( ExecutorBinding Entry, ExecutorBinding Exit, IReadOnlyList<(ExecutorBinding Source, ExecutorBinding Target)> InternalEdges) { public WorkflowNodeRoute(ExecutorBinding binding) : this(binding, binding, []) { } } }