A .NET Agent Orchestration Framework - The Orchestrator
Introduction
This is part two of my series on the agent orchestration framework I'm building in .NET. You really should read the first part, Introduction, before this one, so if you haven't already, take the time to do so now!
In this post I'm going to talk about the orchestrator, the central piece of the framework. I'm going to describe its contract and how we should use it. Some related parts will be covered in other posts.
The IOrchestrator Interface
The orchestrator is main part of the framework, as I said before. It is the one that coordinates - with some assistance - the different agents, distributes the work to be done, enforces reviews and revision processes, then gathers all responses and summarises them. The interface that describes the orchestrator is IOrchestrator (surprise, surprise), and this is what it looks like:
public interface IOrchestrator : IDisposable
{
event EventHandler<OrchestrationEventArgs>? OrchestrationStarting;
event EventHandler<OrchestrationEventArgs>? OrchestrationStopped;
event EventHandler<AgentTrainingEventArgs>? AgentTrainingStarting;
event EventHandler<AgentTrainingEventArgs>? AgentTrainingStopped;
event EventHandler<AgentExecutionEventArgs>? AgentExecutionStarting;
event EventHandler<AgentExecutionEventArgs>? AgentExecutionStopped;
event EventHandler<AgentReviewEventArgs>? AgentReviewStarting;
event EventHandler<AgentReviewEventArgs>? AgentReviewStopped;
event EventHandler<AgentReviseEventArgs>? AgentReviseStarting;
event EventHandler<AgentReviseEventArgs>? AgentReviseStopped;
Task<OrchestrationResult> OrchestrateAsync(string prompt, CancellationToken cancellationToken);
bool IsProcessing { get; }
}
I took out all comments and some less-important stuff. In a nutshell, we have:
- The IOrchestrator interface implements IDisposable, this is essentially to clear the event handlers, as the orchestrator itself is pretty much stateless - it does contain references to a number of services, which may or may not be stateless (all of the built-in services described in this post are stateless)
- It exposes a number of events, for the different stages of the orchestration process (more on this in a moment)
- The most important method, and the only one shown here, is OrchestrateAsync. This is the method that starts the orchestration process, it is asynchronous and returns an OperationResult
- A property that indicates whether or not the orchestrator is currently processing (IsProcessing); it is enabled while the OrchestrateAsync method is running
The only class that implements IOrchestrator is Orchestrator, and it is sealed, meaning, not meant to be derived from.
There is a configuration class, OrchestratorOptions, which can be provided to the Orchestrator constructor. This is what it contains:
public sealed class OrchestratorOptions
{
public bool EnableEvents { get; set; }
public bool EnableInternetSearch { get; set; }
public bool EnablePeerReview { get; init; }
public int MaxReviseAttempts { get; init; }
public int? MaxTokens { get; init; }
}The properties are:
- EnableEvents: controls whether or not the orchestrator raises events for each stage of the processing; the default is true
- EnableInternetSearch: controls whether Internet searches are allowed; the default is false
- EnablePeerReview: whether or not to enforce review of agent responses; the default is true
- MaxReviseAttempts: the maximum number of revise attempts before the agent response is considered final; the default is 5
- MaxTokens: the maximum number of tokens that can be consumed for the processing of an orchestration, including all agents that are involved; null, meaning, no limitation will be enforced
The Orchestration Process
The orchestrator deals with agents (IAgent instances, of which I won't talk much here), so it must receive a non-empty collection of agents, with distinct names. It also receives a prompt and, somehow, distributes that prompt - or, actually, a specific prompt that is derived from the original one - to one or some of the agents it knows.
The process starts with the OrchestrateAsync method, which is responsible for creating a session id, and starting the flow for a given prompt. This prompt is first sanitised before usage, so as to prevent unwanted contents. There are many stages in the agent processing, and some events are raised. This flow consists of:
- Filter the original prompt before using it (IPromptFilter)
- Plan what to do, meaning, what are the sub-tasks and the required specialisations (IPlanner)
- Assign each plan assignment to a concrete agent (IAgentSelector)
- Execute each agent with its own assignment (IAgentScheduler)
- For each finished agent (author), select another agent to review its work (IAgentSelector)
- Ask the review agent to review and produce feedback (*)
- Ask the author to revise based on the feedback (*)
- When all reviews are finished, create a summary of all the final responses (ISynthesizer)
(*) these steps may occur multiple times
When the Orchestrator is processing, its IsProcessing property will be true, and no other orchestration operation can be started. The orchestration execution can be cancelled at any time through the CancellationToken parameter.
A session is essentially a call to OrchestrateAsync, an orchestration. Its id should be unique and is present as a property in all classes that are produced or consumed by the orchestration process, its purpose is to add a correlation identifier.
The following events are raised during the orchestration process, which should be self-explanatory:
- OrchestrationStarting: raised once when the orchestration starts (OrchestrateAsync)
- OrchestrationStopped: raised when OrchestrateAsync finishes
- AgentTrainingStarting: raised for each agent when it is being trained
- AgentTrainingStopped: raised when an agent has finished training
- AgentExecutionStarting: raised when an agent's execution is starting
- AgentExecutionStopped: raised when the agent execution finished
- AgentReviewStarting: raised when a review process is starting
- AgentReviewStopped: raised when a review process finished
- AgentReviseStarting: raised when an agent is beginning to revise its work based on a review
- AgentReviseStopping: raised when an agent's revise process has finished
One note: events are only raised if OrchestratorOptions.EnableEvents is set to true, which is the default.
Some events - AgentTraining*, AgentExecution*, AgentReview*, AgentRevise* - can, of course, be raised multiple times, one for each agent - AgentTraining*, AgentExecution* - or maybe multiple times for a single agent - AgentReview*, AgentRevise*. The Orchestration* events only fire once per orchestration. All these events are raised from another thread, so they do not block the working of the orchestrator, and throwing an exception from an handler will be ignored silently. Each event has an argument that is specific to the event and adds context properties including a SessionId property.
The orchestrator also produces many logs, using the Information, Warning, Debug, or Error, severity levels. These can be enabled or disabled by the usual ways. All operations and their results are logged. Traces are also produced, including the following activities:
- orchestrator.orchestrate: when the orchestrator starts
- orchestrator.plan: when the orchestrator is planning
- orchestrator.synthesize: when the orchestrator is synthesising the responses
The following tags are included:
- orchestrator.session_id: the session id
- orchestrator.prompt_length: the prompt length
- orchestrator.max_tokens: the maximum configured number of tokens to consume
- orchestrator.enable_peer_review: whether peer review is enabled
- orchestrator.max_revise_attempts: the maximum configured revise attempts
- orchestrator.assignments_count: the number of assignments after the plan
- orchestrator.tokens_consumed: the total number of tokens consumed at the end of the orchestration
At the end of the OrchestrateAsync method the orchestrator returns an OrchestrationResult instance. It contains these properties (simplified):
public sealed record OrchestrationResult
{ public string SessionId { get; }
public string FinalResponse { get; }
public IReadOnlyList<AgentResponse> AgentResponses { get; }
public IReadOnlyList<ReviewResult> Reviews { get; }
public bool WasShortCircuited { get; }
public bool Cancelled { get; }
public long TotalTokens { get; }
public bool Success { get; }
public string? Error { get; } }
Where:
- AgentResponses contains the final responses, one for each agent that was involved in the orchestration
- Cancelled is set to true if the operation was cancelled through the CancellationToken parameter
- Error contains an error message if the orchestration was not finished successfully or was cancelled
- FinalResponse property contains the synthesised responses from all the agents
- Reviews contains all the final reviews
- SessionId is the id of the current session
- Success means that the orchestration finished successfully
- TotalTokens contains the ever-increasing number of tokens used by the processing
- WasShortCircuited is flagged if the work was handled by a single agent
So, it is safe to look at FinalResponse if Success is true, or to Error otherwise.
Service Dependencies
The Orchestrator class, besides the IAgent collection, makes use of the following services, which actually implement the important parts of the flow:
| Service | Purpose | Implementations |
| IAgentScheduler | Schedules agents for execution | ParallelAgentScheduler (default) SequentialAgentScheduler |
| IAgentSelector | Selects the appropriate agent the execution or review according to some algorithm | DefaultAgentSelector |
| ISharedStateManager | Manages the shared state which may be used to pass information between agents | InMemorySharedStateManager |
| IPlanner | Creates the execution plan | DefaultPlanner |
| IPromptFilter | Filters the prompt before using it | DefaultPromptFilter |
| ISynthesizer | Summarises all the final agent responses | DefaultSynthesizer |
All of these services are injected through the Orchestrator constructor, and are all optional, meaning, the default implementation will be used if one is not provided.
There are two included implementations of IAgentScheduler: one that processes each agent sequentially (SequentialAgentScheduler), and another one that executes all in "parallel" (ParallelAgentScheduler), which is the default. I will dwell in this in a future post.
DefaultAgentSelector uses one algorithm for selecting the executing agent and another for the reviewer, I will also cover them shortly.
Both DefaultPlanner and DefaultSynthesizer implementations depend on the existence of an IAgent with a specialization of Planning and Synthesis, respectively. If one does not exist, the implementation will be sub-optimal.
The InMemorySharedStateManager implementation of ISharedStateManager manages state in memory, essentially a collection of key-value pairs. It is provided as a sample, more robust implementations might include out-of-process storage.
The included implementation of IPromptFilter, DefaultPromptFilter, just returns the original prompt, makes no attempt to filter or change it. It is provided just as a sample.
All of these services are required for the orchestration to work, but different implementations can lead to very different results. It is interesting to notice that the Orchestrator class itself knows nothing about AI, LLMs, or Microsoft.Extensions.AI: this knowledge exists in some of these services.
Conclusion
We are just starting to get into the technical aspects of my orchestration framework, here I just covered one of the most important parts of it; I certainly did not cover all of the orchestrator, but I hope you can get a good picture of it. In the next post I will cover the agents, which are also very important, as they wrap the connection to an LLM provider and model.
I hope you find this interesting, feel free to ask or comment whatever you like!
Comments
Post a Comment