orchestration

package
v0.2.75 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 18, 2026 License: MIT Imports: 13 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func RunWithTransfers added in v0.2.72

func RunWithTransfers(
	ctx context.Context,
	registry *AgentRegistry,
	tool *TransferTool,
	startAgentID string,
	input string,
	maxTransfers int,
) (result string, chain []string, err error)

RunWithTransfers runs an agent and follows any transfers it requests.

maxTransfers bounds the chain. Without a bound two agents can transfer to each other indefinitely, each call costing a request -- the routing equivalent of an infinite loop, and expensive.

Returns the final answer and the chain of agent IDs that produced it, so a caller can log or display how a request was routed.

Types

type AgentPool

type AgentPool struct {
	// contains filtered or unexported fields
}

AgentPool represents a pool of specialized agents

func NewAgentPool

func NewAgentPool() *AgentPool

NewAgentPool creates a new agent pool

func (*AgentPool) Add

func (p *AgentPool) Add(id string, agent *agent.Agent, description string)

Add adds an agent to the pool

func (*AgentPool) Get

func (p *AgentPool) Get(id string) (*agent.Agent, bool)

Get retrieves an agent from the pool

func (*AgentPool) GetDescription

func (p *AgentPool) GetDescription(id string) (string, bool)

GetDescription retrieves an agent's description

func (*AgentPool) List

func (p *AgentPool) List() map[string]*agent.Agent

List returns all agents in the pool

func (*AgentPool) ListDescriptions

func (p *AgentPool) ListDescriptions() map[string]string

ListDescriptions returns all agent descriptions

type AgentRegistry

type AgentRegistry struct {
	// contains filtered or unexported fields
}

AgentRegistry maintains a registry of available agents

func NewAgentRegistry

func NewAgentRegistry() *AgentRegistry

NewAgentRegistry creates a new agent registry

func (*AgentRegistry) Get

func (r *AgentRegistry) Get(id string) (*agent.Agent, bool)

Get retrieves an agent from the registry

func (*AgentRegistry) List

func (r *AgentRegistry) List() map[string]*agent.Agent

List returns all registered agents

func (*AgentRegistry) Register

func (r *AgentRegistry) Register(id string, agent *agent.Agent)

Register registers an agent with the registry

type CodeOrchestrator

type CodeOrchestrator struct {
	// contains filtered or unexported fields
}

CodeOrchestrator orchestrates agents using code-defined workflows

func NewCodeOrchestrator

func NewCodeOrchestrator(registry *AgentRegistry) *CodeOrchestrator

NewCodeOrchestrator creates a new code orchestrator

func (*CodeOrchestrator) ExecuteWorkflow

func (o *CodeOrchestrator) ExecuteWorkflow(ctx context.Context, workflow *Workflow) (string, error)

ExecuteWorkflow executes a workflow, running each task as soon as its dependencies have completed.

Scheduling is dependency-driven: every task gets one goroutine that waits on its dependencies' completion channels. There is no central monitor.

The previous design had a monitor goroutine that received on an unbuffered completion channel and re-scanned for runnable tasks. It had four defects:

  • Deadlock. executeTask ended with an unbuffered `completionCh <- task.ID` and no select on ctx.Done(). Once the monitor saw every task finished it returned, so a second task finishing at the same moment blocked forever on that send, its deferred wg.Done() never ran, and wg.Wait() hung.
  • Duplicate execution. The monitor launched any task in TaskPending, but TaskRunning was set inside the spawned goroutine. Two completions arriving before that write both launched the same task -- running the agent twice and billing for it twice.
  • Data races. workflow.Results and workflow.Errors were written from N goroutines with no lock at all, and task.Status was written by workers while the monitor read it.
  • wg.Add from the monitor goroutine, potentially after wg.Wait had already observed a zero counter.

Every wg.Add now happens before wg.Wait, all shared state is mutex-guarded, and a cycle is rejected up front rather than deadlocking.

type DelegationAgent

type DelegationAgent struct {
	*agent.Agent
	// contains filtered or unexported fields
}

DelegationAgent is an agent that can delegate tasks to other agents

func NewDelegationAgent

func NewDelegationAgent(baseAgent *agent.Agent, registry *AgentRegistry) *DelegationAgent

NewDelegationAgent creates a new delegation agent

func (*DelegationAgent) Delegate

func (a *DelegationAgent) Delegate(ctx context.Context, targetAgentID string, query string, preserveMemory bool) (string, error)

Delegate delegates a task to another agent

type GraphNode added in v0.2.71

type GraphNode struct {
	// ID identifies this node; other nodes reference it in DependsOn.
	ID string

	// AgentID is the registered agent to run.
	AgentID string

	// Input is this node's own input, before dependency results are appended.
	Input string

	// DependsOn lists node IDs that must complete first.
	DependsOn []string
}

GraphNode is one agent invocation in a Graph workflow.

type HandoffRequest

type HandoffRequest struct {
	// TargetAgentID is the ID of the agent to hand off to
	TargetAgentID string

	// Reason explains why the handoff is happening
	Reason string

	// Context contains additional context for the target agent
	Context map[string]interface{}

	// Query is the query to send to the target agent
	Query string

	// PreserveMemory indicates whether to copy memory to the target agent
	PreserveMemory bool
}

HandoffRequest represents a request to hand off to another agent

type HandoffResult

type HandoffResult struct {
	// AgentID is the ID of the agent that handled the request
	AgentID string

	// Response is the response from the agent
	Response string

	// Completed indicates whether the task was completed
	Completed bool

	// NextHandoff is the next handoff request, if any
	NextHandoff *HandoffRequest
}

HandoffResult represents the result of a handoff

type LLMOrchestrator

type LLMOrchestrator struct {
	// contains filtered or unexported fields
}

LLMOrchestrator orchestrates the execution of a query using multiple agents

func NewLLMOrchestrator

func NewLLMOrchestrator(registry *AgentRegistry, planner interfaces.LLM) *LLMOrchestrator

NewLLMOrchestrator creates a new LLM orchestrator

func (*LLMOrchestrator) Execute

func (o *LLMOrchestrator) Execute(ctx context.Context, query string) (string, error)

Execute executes a query using the orchestrator

func (*LLMOrchestrator) WithLogger

func (o *LLMOrchestrator) WithLogger(logger logging.Logger) *LLMOrchestrator

WithLogger sets the logger for the orchestrator

type LLMRouter

type LLMRouter struct {
	// contains filtered or unexported fields
}

LLMRouter uses an LLM to determine which agent should handle a request

func NewLLMRouter

func NewLLMRouter(llm interfaces.LLM) *LLMRouter

NewLLMRouter creates a new LLM router

func (*LLMRouter) Route

func (r *LLMRouter) Route(ctx context.Context, query string, context map[string]interface{}) (string, error)

Route determines which agent should handle a request

func (*LLMRouter) WithLogger

func (r *LLMRouter) WithLogger(logger logging.Logger) *LLMRouter

WithLogger sets the logger for the router

type Orchestrator

type Orchestrator struct {
	// contains filtered or unexported fields
}

Orchestrator orchestrates handoffs between agents

func NewOrchestrator

func NewOrchestrator(registry *AgentRegistry, router Router) *Orchestrator

NewOrchestrator creates a new orchestrator

func (*Orchestrator) HandleRequest

func (o *Orchestrator) HandleRequest(ctx context.Context, query string, initialContext map[string]interface{}) (*HandoffResult, error)

HandleRequest handles a request, potentially routing it through multiple agents

func (*Orchestrator) WithLogger

func (o *Orchestrator) WithLogger(logger logging.Logger) *Orchestrator

WithLogger sets the logger for the orchestrator

type Pattern added in v0.2.71

type Pattern struct {
	// contains filtered or unexported fields
}

Pattern is a composed multi-agent workflow, ready to run.

func Graph added in v0.2.71

func Graph(registry *AgentRegistry, finalNodeID string, nodes ...GraphNode) *Pattern

Graph builds a workflow from an explicit dependency graph, for shapes the three named patterns do not cover.

Each node names the agent to run and the node IDs it depends on. Cycles and references to unknown nodes are rejected when the workflow runs.

func Loop added in v0.2.71

func Loop(registry *AgentRegistry, agentID string, iterations int) *Pattern

Loop builds a workflow that runs one agent repeatedly, feeding each iteration's result into the next.

iterations must be at least 1. Unlike ADK's LoopAgent there is no model-driven escape hatch yet: the count is fixed, which is why the parameter is required rather than defaulted -- an accidental unbounded loop against a paid API is not a failure mode worth offering.

func Parallel added in v0.2.71

func Parallel(registry *AgentRegistry, combinerID string, agentIDs ...string) *Pattern

Parallel builds a workflow that runs every agent concurrently on the same input, then combines their results with a final agent.

combinerID receives each branch's output appended to its input, in the order the branches were declared.

result, err := orchestration.Parallel(reg, "summarize",
    "legal-review", "security-review", "cost-review").
    Run(ctx, "assess this proposal")

func Sequential added in v0.2.71

func Sequential(registry *AgentRegistry, agentIDs ...string) *Pattern

Sequential builds a workflow that runs agents one after another, feeding each agent's result into the next.

These three patterns were always expressible with AddTask and dependency lists, but every caller had to hand-roll the wiring -- and the sequential case in particular is easy to get subtly wrong, since "sequential" is encoded as each task depending on the previous one rather than declared.

result, err := orchestration.Sequential(reg, "research", "draft", "edit").
    Run(ctx, "write about tide pools")

The result of the final agent is the result of the workflow.

func (*Pattern) Run added in v0.2.71

func (p *Pattern) Run(ctx context.Context, input string) (string, error)

Run executes the pattern. input is supplied to every task that has no input of its own, which for Sequential and Loop means the first step and for Parallel means every branch.

func (*Pattern) Workflow added in v0.2.71

func (p *Pattern) Workflow() *Workflow

Workflow exposes the underlying workflow, so a pattern can be inspected or extended before it runs.

type Plan

type Plan struct {
	// Steps is the list of steps in the plan
	Steps []Step `json:"steps"`

	// FinalAgentID is the ID of the agent that should provide the final response
	FinalAgentID string `json:"final_agent_id"`
}

Plan represents an orchestration plan

type Router

type Router interface {
	Route(ctx context.Context, query string, context map[string]interface{}) (string, error)
}

Router determines which agent should handle a request

type SimpleRouter

type SimpleRouter struct {
	// contains filtered or unexported fields
}

SimpleRouter routes requests based on a simple keyword matching

func NewSimpleRouter

func NewSimpleRouter() *SimpleRouter

NewSimpleRouter creates a new simple router

func (*SimpleRouter) AddRoute

func (r *SimpleRouter) AddRoute(keyword string, agentID string)

AddRoute adds a route to the router

func (*SimpleRouter) Route

func (r *SimpleRouter) Route(ctx context.Context, query string, context map[string]interface{}) (string, error)

Route determines which agent should handle a request

type Step

type Step struct {
	// AgentID is the ID of the agent to execute
	AgentID string `json:"agent_id"`

	// Input is the input to provide to the agent
	Input string `json:"input"`

	// Description explains the purpose of this step
	Description string `json:"description"`

	// DependsOn lists the IDs of steps that must complete before this one
	DependsOn []string `json:"depends_on,omitempty"`
}

Step represents a single step in an orchestration plan

type Task

type Task struct {
	// ID is the unique identifier for the task
	ID string

	// AgentID is the ID of the agent to execute the task
	AgentID string

	// Input is the input to provide to the agent
	Input string

	// Dependencies are the IDs of tasks that must complete before this one
	Dependencies []string

	// Status is the current status of the task
	Status TaskStatus

	// Result is the result of the task
	Result string

	// Error is any error that occurred during execution
	Error error
}

Task represents a task to be executed by an agent

type TaskStatus

type TaskStatus string

TaskStatus represents the status of a task

const (
	// TaskPending indicates the task is pending
	TaskPending TaskStatus = "pending"

	// TaskRunning indicates the task is running
	TaskRunning TaskStatus = "running"

	// TaskCompleted indicates the task is completed
	TaskCompleted TaskStatus = "completed"

	// TaskFailed indicates the task failed
	TaskFailed TaskStatus = "failed"
)

type TransferRequest added in v0.2.72

type TransferRequest struct {
	// TargetAgentID is the agent to transfer to.
	TargetAgentID string

	// Reason is why, in the model's words. Worth keeping: it is the only
	// explanation a human debugging a routing decision will have.
	Reason string

	// Query is what to ask the target agent.
	Query string
}

TransferRequest is a model's decision to hand off.

type TransferTool added in v0.2.72

type TransferTool struct {
	// contains filtered or unexported fields
}

TransferTool lets a model hand a conversation to another agent by calling a tool, rather than by emitting a marker in its prose.

The existing handoff mechanism parses free-form output:

regexp.MustCompile(`\[HANDOFF:([a-zA-Z0-9_-]+):([^\]]+)\]`)

which requires the prompt to teach a bespoke syntax, breaks when the model paraphrases or wraps the marker in a code fence, and cannot constrain the target -- a hallucinated agent name is only discovered after the fact. A tool call is structured, and the agent name can be constrained to an enum so the model cannot invent one.

Transfer is a request, not an execution. The tool records the decision and returns; the caller runs the target agent. That keeps the depth guard, the cycle check and usage accounting under the caller's control rather than buried in a tool.

func NewTransferTool added in v0.2.72

func NewTransferTool(registry *AgentRegistry, descriptions map[string]string) *TransferTool

NewTransferTool creates a transfer tool over a registry.

descriptions maps agent ID to a one-line summary of what that agent is for. They are what the model routes on, so an agent with no description is effectively unreachable.

func (*TransferTool) Description added in v0.2.72

func (t *TransferTool) Description() string

Description implements interfaces.Tool.

func (*TransferTool) Execute added in v0.2.72

func (t *TransferTool) Execute(_ context.Context, args string) (string, error)

Execute implements interfaces.Tool.

func (*TransferTool) Name added in v0.2.72

func (t *TransferTool) Name() string

Name implements interfaces.Tool.

func (*TransferTool) Parameters added in v0.2.72

func (t *TransferTool) Parameters() map[string]interfaces.ParameterSpec

Parameters implements interfaces.Tool.

func (*TransferTool) Requested added in v0.2.72

func (t *TransferTool) Requested() (*TransferRequest, bool)

Requested returns the transfer the model asked for, if any, and clears it.

Clearing on read means a handle reused across turns cannot replay a stale decision from a previous turn.

func (*TransferTool) Run added in v0.2.72

func (t *TransferTool) Run(ctx context.Context, input string) (string, error)

Run implements interfaces.Tool.

type Workflow

type Workflow struct {
	// Tasks is the list of tasks in the workflow
	Tasks []*Task

	// Results is a map of task IDs to results
	Results map[string]string

	// Errors is a map of task IDs to errors
	Errors map[string]error

	// FinalTaskID is the ID of the task that produces the final result
	FinalTaskID string
}

Workflow represents a workflow of tasks

func NewWorkflow

func NewWorkflow() *Workflow

NewWorkflow creates a new workflow

func (*Workflow) AddTask

func (w *Workflow) AddTask(id string, agentID string, input string, dependencies []string)

AddTask adds a task to the workflow

func (*Workflow) SetFinalTask

func (w *Workflow) SetFinalTask(id string)

SetFinalTask sets the final task

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL