Delarbetsflöden

Ett underarbetsflöde är ett fullständigt arbetsflöde som körs som en exekutor inom ett överordnat arbetsflöde. På så sätt kan du skapa komplexa system från mindre, återanvändbara arbetsflödesbyggstenar – var och en med sin egen isolerade körningskontext, tillståndshantering och meddelanderoutning.

Overview

Underarbetsflöden är användbara när du vill:

  • Dela upp komplexitet – dela upp ett stort arbetsflöde i mindre, oberoende testbara enheter.
  • Återanvänd arbetsflödeslogik – bädda in samma underarbetsflöde i flera överordnade arbetsflöden.
  • Isolerat tillstånd – håll varje underarbetsflödes interna tillstånd separat från det överordnade.
  • Kontrollera dataflödet – meddelanden anger och lämnar underarbetsflödet endast via dess kanter, utan sändning över nivåer.

När ett underarbetsflöde läggs till i ett överordnat arbetsflöde fungerar det som alla andra köre: det tar emot indatameddelanden, kör sin interna graf till slutförande och genererar utdatameddelanden för underordnade utförare.

Skapa ett Delarbetsflöde

I C# skapar du underarbetsflöden på två sätt:

  • Direktbindning – används BindAsExecutor() för att bädda in ett arbetsflöde direkt som köre i det överordnade arbetsflödet. Detta bevarar underarbetsflödets interna indata-/utdatatyper.
  • Agentomslutning – använd AsAIAgent() för att konvertera ett arbetsflöde till en agent och lägg sedan till agenten i det överordnade arbetsflödet. Detta är användbart när det överordnade arbetsflödet använder agentbaserade utförare.

Direktbindning med BindAsExecutor

Tilläggsmetoden BindAsExecutor() konverterar ett arbetsflöde till ett ExecutorBinding som kan läggas till direkt i ett överordnat arbetsflöde:

using Microsoft.Agents.AI.Workflows;

// Create executors for the inner workflow
UppercaseExecutor uppercase = new();
ReverseExecutor reverse = new();
AppendSuffixExecutor append = new(" [PROCESSED]");

// Build the inner workflow
var innerWorkflow = new WorkflowBuilder(uppercase)
    .AddEdge(uppercase, reverse)
    .AddEdge(reverse, append)
    .WithOutputFrom(append)
    .Build();

// Bind the inner workflow as an executor
ExecutorBinding subWorkflowExecutor = innerWorkflow.BindAsExecutor("TextProcessingSubWorkflow");

// Build the parent workflow using the sub-workflow executor
PrefixExecutor prefix = new("INPUT: ");
PostProcessExecutor postProcess = new();

var parentWorkflow = new WorkflowBuilder(prefix)
    .AddEdge(prefix, subWorkflowExecutor)
    .AddEdge(subWorkflowExecutor, postProcess)
    .WithOutputFrom(postProcess)
    .Build();

Med BindAsExecutorbevaras underarbetsflödets typinmatnings- och utdatatyper – det överordnade arbetsflödet dirigerar meddelanden baserat på de faktiska typer som underarbetsflödet förväntar sig och producerar.

Agentinkapsling med AsAIAgent

När det överordnade arbetsflödet använder agentbaserade utförare konverterar du det inre arbetsflödet till en agent med .AsAIAgent() WorkflowBuilder omsluter automatiskt agenten i en exekutor:

using Microsoft.Agents.AI;
using Microsoft.Agents.AI.Workflows;

// Create agents for the inner workflow
AIAgent specialist1 = chatClient.AsAIAgent("You are specialist 1. Analyze the data.");
AIAgent specialist2 = chatClient.AsAIAgent("You are specialist 2. Validate the analysis.");

// Build the inner workflow
var innerWorkflow = new WorkflowBuilder(specialist1)
    .AddEdge(specialist1, specialist2)
    .Build();

// Convert the inner workflow to an agent
AIAgent innerWorkflowAgent = innerWorkflow.AsAIAgent(
    id: "analysis-pipeline",
    name: "Analysis Pipeline",
    description: "A sub-workflow that analyzes and validates data"
);

// Create agents for the parent workflow
AIAgent coordinator = chatClient.AsAIAgent("You are a coordinator. Delegate tasks to the team.");
AIAgent reviewer = chatClient.AsAIAgent("You are a reviewer. Review the final output.");

// Build the parent workflow with the sub-workflow
var parentWorkflow = new WorkflowBuilder(coordinator)
    .AddEdge(coordinator, innerWorkflowAgent)
    .AddEdge(innerWorkflowAgent, reviewer)
    .Build();

Det inre arbetsflödet körs som ett enda steg från det överordnade arbetsflödets perspektiv. Koordinatorn skickar meddelanden till analyspipelinen, som internt kör specialist1 → specialist2och vidarebefordrar sedan resultatet till granskaren.

Tips/Råd

Använd BindAsExecutor() när du arbetar med inskrivna utförare och AsAIAgent() när du arbetar med agentbaserade arbetsflöden. Mer information om hur du konfigurerar konverteringen från arbetsflöde till agent finns i Arbetsflöden som agenter.

Indata- och utdatatyper

När ett arbetsflöde används som ett underarbetsflöde bevaras typkontrakten för dess interna utförare.

Med BindAsExecutoraccepterar underarbetsflödeskörningen samma indatatyper som det inre arbetsflödets startexekutor och skickar samma utdatatyper som det inre arbetsflödet skapar. Det överordnade arbetsflödets kanter måste ansluta utförare vars utdatatyper matchar underarbetsflödets förväntade indatatyper, och underarbetsflödets utdatatyper måste matcha de förväntade indata som förväntas av underordnade utförare.

Med AsAIAgent omsluts underarbetsflödet som en agent och följer Agent Executor-indata-/utdatakontrakt (string, ChatMessage, IEnumerable<ChatMessage>).

Utdatabeteende

Som standard vidarebefordras dessa utdata som meddelanden till anslutna utförare i det överordnade arbetsflödet när ett underarbetsflöde skapar utdata (via YieldOutputAsync). Detta gör det möjligt för underordnade utförare att bearbeta underarbetsflödesresultat.

Klassen ExecutorOptions styr det här beteendet:

Option Standardinställning Description
AutoSendMessageHandlerResultObject true Vidarebefordra utdata från underarbetsflödet som meddelanden till anslutna utförare i det överordnade diagrammet.
AutoYieldOutputHandlerResultObject false Ge utdata från underarbetsflödet direkt till det överordnade arbetsflödets utdatahändelseström.

När AutoYieldOutputHandlerResultObject är aktiverat kringgår utdata från underarbetsflödet den överordnades interna routning och levereras direkt till anroparen för det överordnade arbetsflödet.

var options = new ExecutorOptions
{
    AutoYieldOutputHandlerResultObject = true,
};

ExecutorBinding subWorkflowExecutor = innerWorkflow.BindAsExecutor("SubWorkflow", options);

Begäranden och svar

Underarbetsflöden har fullt stöd för mekanismen för begäran och svar . När en exekverare i underarbetsflödet skickar en begäran (till exempel för att begära mänsklig indata) WorkflowHostExecutor vidarebefordras RequestInfoEvent till det överordnade arbetsflödet med ett kvalificerat port-ID – delarbetsflödeskörningens ID prefixas med port-ID:t (till exempel SubWorkflow.GuessNumber).

Den här kvalificeringen säkerställer att när det överordnade arbetsflödet får ett svar kan svaret dirigeras tillbaka till rätt instans av underarbetsflödet. Det överordnade arbetsflödet hanterar begäranden under arbetsflödet med samma svarsmekanism som andra begäranden:

await using StreamingRun handle = await InProcessExecution.RunStreamingAsync(parentWorkflow, input);
await foreach (WorkflowEvent evt in handle.WatchStreamAsync())
{
    switch (evt)
    {
        case RequestInfoEvent requestInfoEvt:
            // The request may originate from the sub-workflow
            // Handle it and send the response back
            var response = requestInfoEvt.Request.CreateResponse(myResponseData);
            await handle.SendResponseAsync(response);
            break;

        case WorkflowOutputEvent outputEvt:
            Console.WriteLine($"Output: {outputEvt.Data}");
            break;
    }
}

Anmärkning

Från den överordnade arbetsflödesuppringarens perspektiv finns det ingen skillnad mellan en begäran från en utförare på den översta nivån och en begäran från ett underarbetsflöde. Ramverket hanterar routningen transparent.

Så här fungerar det

När det överordnade arbetsflödet dirigerar ett meddelande till exekveraren av underarbetsflödet:

  1. Indataleverans – meddelandet vidarebefordras till det inre arbetsflödets startexekutor. Med BindAsExecutor måste meddelandetypen matcha de förväntade typerna för startutföraren. Med AsAIAgentnormaliseras meddelanden till ChatMessage format.
  2. Inre körning – det inre arbetsflödet har en egen supersteploop.
  3. Utdatasamling – det inre arbetsflödets utdatahändelser samlas in. Med BindAsExecutorbehåller utdata sina ursprungliga typer. Med AsAIAgentkonverteras utdata till agentsvarsmeddelanden.
  4. Vidarebefordran av begäranden – om det inre arbetsflödet har väntande begäranden vidarebefordras de till det överordnade arbetsflödet för hantering (se Begäranden och svar).
  5. Distribution nedströms – de resulterande meddelandena skickas till nästa exekutor i det överordnade arbetsflödet.

Eftersom det inre arbetsflödet har en egen körningskontext är dess tillstånd oberoende av det överordnade arbetsflödet.

Tips/Råd

Mer information om hur du konfigurerar konverteringen från arbetsflöde till agent, inklusive strömningsbeteende och undantagshantering, finns i Arbetsflöden som agenter.

Kapsling på flera nivåer

Underarbetsflöden kan kapslas till godtyckligt djup. Varje nivå behåller sin egen exekveringskontext:

// Level 1: Data preparation pipeline
var dataPipeline = new WorkflowBuilder(fetcher)
    .AddEdge(fetcher, cleaner)
    .Build();

AIAgent dataPipelineAgent = dataPipeline.AsAIAgent(
    id: "data-pipeline",
    name: "Data Pipeline"
);

// Level 2: Analysis pipeline (contains the data pipeline)
var analysisPipeline = new WorkflowBuilder(dataPipelineAgent)
    .AddEdge(dataPipelineAgent, analyzer)
    .Build();

AIAgent analysisPipelineAgent = analysisPipeline.AsAIAgent(
    id: "analysis-pipeline",
    name: "Analysis Pipeline"
);

// Level 3: Top-level orchestration
var topWorkflow = new WorkflowBuilder(coordinator)
    .AddEdge(coordinator, analysisPipelineAgent)
    .AddEdge(analysisPipelineAgent, reporter)
    .Build();

Anmärkning

Varje kapslingsnivå lägger till exekveringsbelastning eftersom det inre arbetsflödet kör en egen superstegs-slinga. Håll kapslingsdjupet rimligt för prestandakänsliga scenarier.

Felhantering

När ett underarbetsflöde misslyckas sprids felet till det överordnade arbetsflödet som en SubworkflowErrorEvent. Det överordnade arbetsflödet kan observera dessa fel via händelseströmmen:

await foreach (WorkflowEvent evt in handle.WatchStreamAsync())
{
    if (evt is SubworkflowErrorEvent subError)
    {
        Console.WriteLine($"Sub-workflow '{subError.ExecutorId}' failed: {subError.Data}");
    }
}

Om underarbetsflödet stöter på ett ohanterat undantag fortsätter körningen av det överordnade arbetsflödet, men underarbetsflödeskörningen slutar bearbeta ytterligare meddelanden.

Kontrollpunkter

När en kontrollpunkt tas i föräldraarbetsflödet, serialiseras underarbetsflödesagentens sessionstillstånd som en del av föräldraexekverarens kontrollpunktsdata. Vid återställningen deserialiseras sessionstillståndet, vilket gör att det överordnade arbetsflödet kan återupptas med underarbetsflödets tillstånd intakt.

CheckpointManager checkpointManager = CheckpointManager.CreateInMemory();

// Run the parent workflow with checkpointing
StreamingRun run = await InProcessExecution
    .RunStreamingAsync(parentWorkflow, input, checkpointManager);

await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
    // Process events, including those from sub-workflows
}

// Resume from a checkpoint
CheckpointInfo checkpoint = run.Checkpoints[^1];
StreamingRun resumedRun = await InProcessExecution
    .ResumeStreamingAsync(parentWorkflow, checkpoint, checkpointManager);

Skapa ett Delarbetsflöde

I Python skapar du ett underarbetsflöde genom att omsluta ett Workflow i ett WorkflowExecutor och lägga till det i ett överordnat arbetsflöde.

from agent_framework import WorkflowBuilder, WorkflowExecutor

# Create agents for the inner workflow
specialist1 = client.as_agent(name="Specialist1", instructions="Analyze the data.")
specialist2 = client.as_agent(name="Specialist2", instructions="Validate the analysis.")

# Build the inner workflow
inner_workflow = (
    WorkflowBuilder(start_executor=specialist1)
    .add_edge(specialist1, specialist2)
    .build()
)

# Wrap as an executor
inner_workflow_executor = WorkflowExecutor(
    workflow=inner_workflow,
    id="analysis-pipeline",
)

# Create agents for the parent workflow
coordinator = client.as_agent(name="Coordinator", instructions="Delegate tasks to the team.")
reviewer = client.as_agent(name="Reviewer", instructions="Review the final output.")

# Build the parent workflow with the sub-workflow
parent_workflow = (
    WorkflowBuilder(start_executor=coordinator)
    .add_edge(coordinator, inner_workflow_executor)
    .add_edge(inner_workflow_executor, reviewer)
    .build()
)

Det inre arbetsflödet körs som ett enda steg från det överordnade arbetsflödets perspektiv. Koordinatorn skickar meddelanden till analyspipelinen, som internt kör specialist1 → specialist2och vidarebefordrar sedan resultatet till granskaren.

WorkflowExecutor-Parametrar

Parameter Type Standardinställning Description
workflow Workflow Arbetsflödesinstansen som ska omslutas som en köre.
id str Unikt identifierare för den här köraren.
allow_direct_output bool False När Truereturneras utdata från underarbetsflödet direkt till det överordnade arbetsflödets händelseström i stället för att skickas som meddelanden till anslutna utförare.
propagate_request bool False När Truesprids begäranden från underarbetsflödet till det överordnade arbetsflödets händelseström som vanliga informationshändelser för begäranden. När False slutförs, omsluts begäranden i SubWorkflowRequestMessage för avlyssning av överordnade utförare.

Kapsla in underarbetsflöden

Omslut uttryckligen Workflow-instanser i ett WorkflowExecutor innan du lägger till dem i ett överliggande arbetsflöde. Agenter kan skickas direkt till WorkflowBuilder, men råa Workflow instanser kräver den här omslutningen.

from agent_framework import WorkflowExecutor

inner_workflow_executor = WorkflowExecutor(inner_workflow, id="analysis_pipeline")

parent_workflow = (
    WorkflowBuilder(start_executor=coordinator)
    .add_edge(coordinator, inner_workflow_executor)
    .add_edge(inner_workflow_executor, reviewer)
    .build()
)

Med explicit omslutning kan du:

  • Tilldela ett specifikt exekutor-ID att användas som referens i flera kanter.
  • Återanvänd samma WorkflowExecutor instans i diagrammet.
# Explicit wrapping — create the WorkflowExecutor yourself
inner_workflow_executor = WorkflowExecutor(
    workflow=inner_workflow,
    id="analysis-pipeline",
)

parent_workflow = (
    WorkflowBuilder(start_executor=coordinator)
    .add_edge(coordinator, inner_workflow_executor)
    .add_edge(inner_workflow_executor, reviewer)
    .build()
)

Indata- och utdatatyper

Ärver WorkflowExecutor sin typsignatur från det omslutna arbetsflödet:

  • Indatatyper matchar det omslutna arbetsflödets startexekutorindatatyper (plus SubWorkflowResponseMessage för hantering av svar på vidarebefordrade begäranden).
  • Utdatatyperna matchar utdatatyperna för det omslutna arbetsflödet. Om någon exekverare i underarbetsflödet är begär-svar-kompatibel, SubWorkflowRequestMessage inkluderas även som en utdatatyp.

Det innebär att det överordnade arbetsflödets kanter måste ansluta utförare vars utdatatyper matchar underarbetsflödets förväntade indatatyper. På samma sätt måste underordnade utförare acceptera de typer som underarbetsflödet skapar:

# The sub-workflow's start executor accepts TextProcessingRequest
# So the parent executor must send TextProcessingRequest
class Orchestrator(Executor):
    @handler
    async def start(self, texts: list[str], ctx: WorkflowContext[TextProcessingRequest]) -> None:
        for text in texts:
            await ctx.send_message(TextProcessingRequest(text=text))

# The sub-workflow yields TextProcessingResult
# So the downstream executor must handle TextProcessingResult
class ResultCollector(Executor):
    @handler
    async def collect(self, result: TextProcessingResult, ctx: WorkflowContext) -> None:
        print(f"Received: {result}")

Utdatabeteende

Som standard (allow_direct_output=False) vidarebefordras dessa utdata som meddelanden till anslutna utförare i det överordnade arbetsflödet med hjälp av yield_outputnär ett underarbetsflöde skapar utdata via send_message. Detta gör det möjligt för underordnade utförare att bearbeta delarbetsflödesresultat som en del av det överordnade diagrammet.

När allow_direct_output=Truereturneras utdata från underarbetsflödet direkt till det överordnade arbetsflödets händelseström. Utdata från underarbetsflödet blir utdata från det överordnade arbetsflödet och kringgår den överordnades interna körningsroutning:

# Outputs go directly to parent's event stream
sub_workflow_executor = WorkflowExecutor(
    workflow=inner_workflow,
    id="analysis-pipeline",
    allow_direct_output=True,
)

# The caller receives sub-workflow outputs directly
async for event in parent_workflow.run(input_data, stream=True):
    if event.type == "output":
        # This output came from the sub-workflow
        print(event.data)

Mellanliggande utsläpp från underordnade arbetsflöden

"intermediate" händelser som skapas i ett underordnat arbetsflöde bubblar upp automatiskt via den överordnade händelseströmmen. De tillskrivs WorkflowExecutors egna id (inte den inre exekverare som ursprungligen genererade dem), vilket bevarar inkapslingen. Avgörande är att dessa händelser behåller "intermediate" etiketten oavsett hur den överordnade anger WorkflowExecutor i sina egna output_from eller intermediate_output_from listor.

async for event in parent_workflow.run(input_data, stream=True):
    if event.type == "intermediate":
        # Attributed to the WorkflowExecutor id, e.g. "analysis-pipeline"
        print(f"[{event.executor_id}] intermediate: {event.data}")
    elif event.type == "output":
        print(f"Terminal output: {event.data}")

Begäranden och svar

Underarbetsflöden har fullt stöd för mekanismen för begäran och svar . När en verkställare inom ett underarbetsflöde anropar ctx.request_info() så fångar WorkflowExecutor begäran och hanterar den baserat på inställningen propagate_request.

Avlyssna Begäranden i det Överordnade Arbetsflödet (Standard)

Med propagate_request=False (standardinställningen) omsluts begäranden från underarbetsflödet i en SubWorkflowRequestMessage och skickas till anslutna utförare i det överordnade arbetsflödet. På så sätt kan överordnade utförare hantera begäran lokalt:

from agent_framework import (
    SubWorkflowRequestMessage,
    SubWorkflowResponseMessage,
)


class ParentHandler(Executor):
    @handler
    async def handle_request(
        self,
        request: SubWorkflowRequestMessage,
        ctx: WorkflowContext[SubWorkflowResponseMessage],
    ) -> None:
        # Inspect the original request from the sub-workflow
        original_data = request.source_event.data

        # Create and send a response back to the sub-workflow
        response = request.create_response(my_response_data)
        await ctx.send_message(response, target_id=request.executor_id)

Metoden create_response() verifierar att svarsdatatypen matchar den förväntade typen från den ursprungliga begäran. Om typerna inte matchar genereras en TypeError .

Important

När du skickar tillbaka svaret använder du target_id=request.executor_id för att dirigera SubWorkflowResponseMessage till rätt WorkflowExecutor instans.

Sprida begäranden till externa anropare

Med propagate_request=Truesprids begäranden från underarbetsflödet till det överordnade arbetsflödets händelseström med hjälp av standardmekanismen request_info . Det överordnade arbetsflödets anropare hanterar dessa begäranden på samma sätt som andra mänskliga begäranden i loopen:

sub_workflow_executor = WorkflowExecutor(
    workflow=inner_workflow,
    id="analysis-pipeline",
    propagate_request=True,
)

# Run the parent workflow and handle propagated requests
result = await parent_workflow.run(input_data)
request_info_events = result.get_request_info_events()
if request_info_events:
    responses = {}
    for event in request_info_events:
        # Handle each request (e.g., ask a human)
        responses[event.request_id] = get_human_response(event.data)
    result = await parent_workflow.run(responses=responses)

Så här fungerar det

När det överordnade arbetsflödet dirigerar ett meddelande till WorkflowExecutor:

  1. Indataleverans – meddelandet vidarebefordras till det inre arbetsflödets startexekutor. Meddelandetypen måste matcha startkörarens förväntade indatatyper.
  2. Inre körning – det inre arbetsflödet kör en egen superstegsloop till fullbordan eller tills den behöver externa indata.
  3. Utdatasamling – det inre arbetsflödets utdatahändelser samlas in och vidarebefordras baserat på inställningen allow_direct_output .
  4. Vidarebefordran av begäranden – om det inre arbetsflödet har väntande begäranden vidarebefordras de baserat på propagate_request inställningen (se Begäranden och svar).
  5. Svarackumulering — den WorkflowExecutor samlar in svar och återupptar underarbetsflödet först när alla förväntade svar för en viss körning har tagits emot.
  6. Nedströmssändning – utdata skickas till nästa köre i det överordnade arbetsflödet.

Underarbetsflödet upprätthåller sitt eget interna tillstånd oberoende av det överordnade. Meddelanden dirigeras endast genom kanterna som ansluter WorkflowExecutor till resten av det överordnade diagrammet – det finns inget meddelande som sänder över kapslingsnivåer.

Kapsling på flera nivåer

Underarbetsflöden kan kapslas till godtyckligt djup. Varje nivå behåller sin egen exekveringskontext:

# Level 1: Data preparation pipeline
data_pipeline = (
    WorkflowBuilder(start_executor=fetcher)
    .add_edge(fetcher, cleaner)
    .build()
)

data_pipeline_executor = WorkflowExecutor(data_pipeline, id="data_pipeline")

# Level 2: Analysis pipeline (contains the data pipeline)
analysis_pipeline = (
    WorkflowBuilder(start_executor=data_pipeline_executor)
    .add_edge(data_pipeline_executor, analyzer)
    .build()
)

analysis_pipeline_executor = WorkflowExecutor(analysis_pipeline, id="analysis_pipeline")

# Level 3: Top-level orchestration
top_workflow = (
    WorkflowBuilder(start_executor=coordinator)
    .add_edge(coordinator, analysis_pipeline_executor)
    .add_edge(analysis_pipeline_executor, reporter)
    .build()
)

Anmärkning

Varje kapslingsnivå lägger till exekveringsbelastning eftersom det inre arbetsflödet kör en egen superstegs-slinga. Håll kapslingsdjupet rimligt för prestandakänsliga scenarier.

Varning

Alla samtidiga körningar av en WorkflowExecutor delar samma underliggande arbetsflödesinstans. Utförare i underarbetsflödet bör vara tillståndslösa för att undvika störning mellan samtida körningar.

Felhantering

När ett underarbetsflöde misslyckas sprids felet till det överordnade arbetsflödet. Samlar WorkflowExecutor in den misslyckade händelsen från underarbetsflödet och konverterar den till en felhändelse i den överordnade kontexten:

async for event in parent_workflow.run(input_data, stream=True):
    if event.type == "error":
        print(f"Sub-workflow failed: {event.details.message}")
    elif event.type == "output":
        print(event.data)

Om underarbetsflödet stöter på ett ohanterat undantag får det överordnade arbetsflödet en felhändelse med undantagsinformationen, inklusive underarbetsflödets ID.

Kontrollpunkter

Underarbetsflöden stöder kontrollpunkter. När en kontrollpunkt tas i det överordnade arbetsflödet, serialiserar WorkflowExecutor sitt interna tillstånd, vilket omfattar det interna arbetsflödets körningsförlopp och eventuella cachelagrade meddelanden. Vid återställning är det här tillståndet deserialiserat, vilket gör att det överordnade arbetsflödet kan återupptas med underarbetsflödet intakt.

from agent_framework import FileCheckpointStorage, WorkflowBuilder

checkpoint_storage = FileCheckpointStorage(storage_path="./checkpoints")

# Build the parent workflow with checkpointing
parent_workflow = (
    WorkflowBuilder(
        start_executor=coordinator,
        checkpoint_storage=checkpoint_storage,
    )
    .add_edge(coordinator, inner_workflow_executor)
    .add_edge(inner_workflow_executor, reviewer)
    .build()
)

# Run with automatic checkpointing
async for event in parent_workflow.run("Analyze the dataset", stream=True):
    if event.type == "output":
        print(event.data)

# Resume from a checkpoint
checkpoints = await checkpoint_storage.list_checkpoints(workflow_name=parent_workflow.name)
async for event in parent_workflow.run(
    checkpoint_id=checkpoints[-1].checkpoint_id,
    checkpoint_storage=checkpoint_storage,
    stream=True,
):
    if event.type == "output":
        print(event.data)

Skapa ett Delarbetsflöde

I Go skapar du ett underarbetsflöde genom att skapa ett *workflow.Workflow och binda det till det överordnade arbetsflödet med inproc.BindSubworkflowAsExecutor.

package main

import (
    "context"
    "fmt"
    "slices"
    "strings"

    "github.com/microsoft/agent-framework-go/workflow"
    "github.com/microsoft/agent-framework-go/workflow/inproc"
)

func buildParentWorkflow() (*workflow.Workflow, error) {
    uppercase := workflow.NewExecutor("UppercaseExecutor", strings.ToUpper).Bind()
    reverse := workflow.NewExecutor("ReverseExecutor", reverseString).Bind()
    appendSuffix := workflow.NewExecutor("AppendSuffixExecutor", func(input string) string {
        return input + " [PROCESSED]"
    }).Bind()

    textProcessing, err := workflow.NewBuilder(uppercase).
        AddEdge(uppercase, reverse).
        AddEdge(reverse, appendSuffix).
        WithOutputFrom(appendSuffix).
        Build()
    if err != nil {
        return nil, err
    }

    textProcessingExecutor := inproc.BindSubworkflowAsExecutor(
        textProcessing,
        "TextProcessingSubWorkflow",
    )

    prefix := workflow.NewExecutor("PrefixExecutor", func(input string) string {
        return "INPUT: " + input
    }).Bind()
    postProcess := workflow.NewExecutor("PostProcessExecutor", func(input string) string {
        return "[FINAL] " + input + " [END]"
    }).Bind()

    return workflow.NewBuilder(prefix).
        AddEdge(prefix, textProcessingExecutor).
        AddEdge(textProcessingExecutor, postProcess).
        WithOutputFrom(postProcess).
        Build()
}

func reverseString(input string) string {
    runes := []rune(input)
    slices.Reverse(runes)
    return string(runes)
}

func runWorkflow(ctx context.Context, parentWorkflow *workflow.Workflow) error {
    run, err := inproc.Default.RunStreaming(ctx, parentWorkflow, "hello")
    if err != nil {
        return err
    }
    defer run.Close(ctx)

    for event, err := range run.WatchStream(ctx) {
        if err != nil {
            return err
        }
        if output, ok := event.(workflow.OutputEvent); ok {
            fmt.Println(output.Output)
        }
    }
    return nil
}

Det bundna underordnade arbetsflödet körs som en enda exekverare ur det överordnade arbetsflödets perspektiv. Meddelanden som matas in via bindningen, det underordnade arbetsflödet kör ett eget internt diagram och det underordnade arbetsflödets utdata dirigeras tillbaka till det överordnade diagrammet.

Indata- och utdatatyper

Bindningen ärver protokollet från det inbäddade arbetsflödet. Det överordnade arbetsflödet kan skicka meddelanden vars körningstyper matchar det underordnade arbetsflödets godkända indatatyper, och bindningen exponerar det underordnade arbetsflödets utdatatyper som både meddelandetyper och utdatatyper.

Det innebär att det överordnade arbetsflödets kanter måste ansluta utförare vars utdatatyper matchar underarbetsflödets godkända indata, och underordnade utförare måste hantera de typer som det underordnade arbetsflödet ger:

type TextProcessingRequest struct {
    Text string
}

type TextProcessingResult struct {
    Text string
}

orchestrator := workflow.NewExecutor("Orchestrator", func(ctx *workflow.Context, texts []string) error {
    for _, text := range texts {
        if err := ctx.SendMessage("", TextProcessingRequest{Text: text}); err != nil {
            return err
        }
    }
    return nil
}).Bind()

collector := workflow.NewExecutor("Collector", func(result TextProcessingResult) {
    fmt.Println(result.Text)
}).Bind()

Utdatabeteende

När ett underarbetsflöde genererar utdata skickar underarbetsflödets bindning dessa utdata som ett meddelande från bindningen till anslutna exekverare i det överordnade arbetsflödet. Om det överordnade arbetsflödet också markerar underarbetsflödesbindningen med WithOutputFrom, genereras samma värde som ett överordnat workflow.OutputEvent, vars ExecutorID är underarbetsflödesbindningens ID.

subWorkflowExecutor := inproc.BindSubworkflowAsExecutor(textProcessing, "TextProcessingSubWorkflow")
postProcess := workflow.NewExecutor("PostProcessExecutor", func(input string) string {
    return "[FINAL] " + input
}).Bind()

parentWorkflow, err := workflow.NewBuilder(subWorkflowExecutor).
    AddEdge(subWorkflowExecutor, postProcess).
    WithOutputFrom(subWorkflowExecutor).
    WithOutputFrom(postProcess).
    Build()

Anpassade arbetsflödeshändelser som utlöses i det underordnade arbetsflödet vidarebefordras till den överordnade händelseströmmen. Delarbetsflödets egna livscykelhändelser för start och supersteg hålls internt så att huvudströmmen förblir fokuserad på händelser som är meningsfulla externt.

Begäranden och svar

Underarbetsflöden stöder mekanismen för begäran och svar . När en exekverare inuti underarbetsflödet skickar en extern begäran kvalificerar underarbetsflödets bindning ID:t för begärandeporten genom att sätta bindnings-ID:t framför det. Till exempel blir ApprovalPort en port för underordnad begäran med namnet ApprovalSubWorkflow.ApprovalPort i det överordnade arbetsflödet.

Om du vill visa den underordnade begäran via det överordnade arbetsflödet lägger du till en överordnad RequestPort med det kvalificerade ID:t och dirigerar begäranden och svar mellan bindningen under arbetsflödet och den porten:

import "reflect"

approvalPort := workflow.RequestPort{
    ID:       "ApprovalPort",
    Request:  reflect.TypeFor[string](),
    Response: reflect.TypeFor[bool](),
}

approvalWorkflow, err := workflow.NewBuilder(approvalPort.Bind()).
    Build()
if err != nil {
    return err
}

approvalSubWorkflow := inproc.BindSubworkflowAsExecutor(
    approvalWorkflow,
    "ApprovalSubWorkflow",
)

qualifiedApprovalPort := workflow.RequestPort{
    ID:       "ApprovalSubWorkflow.ApprovalPort",
    Request:  approvalPort.Request,
    Response: approvalPort.Response,
}
qualifiedApproval := qualifiedApprovalPort.Bind()

parentWorkflow, err := workflow.NewBuilder(approvalSubWorkflow).
    AddDirectEdge(approvalSubWorkflow, qualifiedApproval, false, externalRequestOnly).
    AddDirectEdge(qualifiedApproval, approvalSubWorkflow, false, externalResponseOnly).
    Build()

Anroparen hanterar begäran från den överordnade arbetsflödesströmmen och skickar tillbaka svaret via samma körningshandtag. Bindningen för underarbetsflödet tar bort det kvalificerade prefixet innan svaret levereras till det underordnade arbetsflödet.

run, err := inproc.Default.RunStreaming(ctx, parentWorkflow, "Approve deployment?")
if err != nil {
    return err
}
defer run.Close(ctx)

for event, err := range run.WatchStream(ctx) {
    if err != nil {
        return err
    }
    switch event := event.(type) {
    case workflow.RequestInfoEvent:
        response, err := event.Request.CreateResponse(true)
        if err != nil {
            return err
        }
        if err := run.SendResponse(ctx, response); err != nil {
            return err
        }
    case workflow.OutputEvent:
        fmt.Println(event.Output)
    }
}

Använd predikat för att hålla begärande- och svarskanterna smala:

func externalRequestOnly(msg any) bool {
    _, ok := msg.(*workflow.ExternalRequest)
    return ok
}

func externalResponseOnly(msg any) bool {
    _, ok := msg.(*workflow.ExternalResponse)
    return ok
}

Så här fungerar det

När huvudarbetsflödet dirigerar ett meddelande till underarbetsflödets bindning:

  1. Indataleverans – bindningen accepterar meddelanden som matchar det underordnade arbetsflödets godkända indatatyper och lägger till dem i det underordnade arbetsflödets startexekutor.
  2. Inre körning – det underordnade arbetsflödet körs i samma processkörningsmiljö och underhåller en egen superstegsloop.
  3. Vidarebefordran av utdata — underordnade värden för workflow.OutputEvent skickas som meddelanden från bindningen till efterföljande överordnade exekverare och ges också som utdata från den överordnade om bindningen anges i WithOutputFrom.
  4. Vidarebefordran av begäran – underordnade workflow.RequestInfoEvent begäranden skickas på nytt med kvalificerade port-ID:n och kan dirigeras via överordnade RequestPort bindningar.
  5. Vidarebefordran av händelser — anpassade händelser från underordnade arbetsflöden läggs till i föräldraströmmen. Fel visas som överordnade workflow.ErrorEvent värden med underarbetsflödes-ID registrerat.
  6. Nedströmssändning – resulterande meddelanden fortsätter genom de överordnade arbetsflödeskanterna.

Det underordnade arbetsflödet håller sitt tillstånd och sin meddelanderoutning åtskilda från det överordnade. Meddelanden passerar endast gränsen genom kanterna som är anslutna till bindningen under arbetsflödet.

Kapsling på flera nivåer

Underarbetsflöden kan kapslas till godtyckligt djup. Varje underarbetsflöde binds innan det läggs till i arbetsflödet som innehåller det:

fraudCheck, err := workflow.NewBuilder(analyzePatterns).
    AddEdge(analyzePatterns, calculateRiskScore).
    WithOutputFrom(calculateRiskScore).
    Build()
if err != nil {
    return err
}

fraudCheckExecutor := inproc.BindSubworkflowAsExecutor(fraudCheck, "FraudCheck")

payment, err := workflow.NewBuilder(validatePayment).
    AddEdge(validatePayment, fraudCheckExecutor).
    AddEdge(fraudCheckExecutor, chargePayment).
    WithOutputFrom(chargePayment).
    Build()
if err != nil {
    return err
}

paymentExecutor := inproc.BindSubworkflowAsExecutor(payment, "Payment")
shippingExecutor := inproc.BindSubworkflowAsExecutor(shipping, "Shipping")

orderWorkflow, err := workflow.NewBuilder(orderReceived).
    AddEdge(orderReceived, paymentExecutor).
    AddEdge(paymentExecutor, shippingExecutor).
    AddEdge(shippingExecutor, orderCompleted).
    WithOutputFrom(orderCompleted).
    Build()

Anmärkning

Varje kapslingsnivå medför exekveringsomkostnader eftersom underarbetsflödet kör en egen loop för supersteg. Håll kapslingsdjupet rimligt för prestandakänsliga scenarier.

Felhantering

När ett underordnat arbetsflöde genererar ett fel vidarebefordrar bindningen under arbetsflödet den till det överordnade arbetsflödet som ett workflow.ErrorEvent och anger SubWorkflowID bindnings-ID:t. Det överordnade arbetsflödet kan observera dessa fel via samma händelseström som används för arbetsflödesfel på den översta nivån:

for event, err := range run.WatchStream(ctx) {
    if err != nil {
        return err
    }
    switch event := event.(type) {
    case workflow.ErrorEvent:
        if event.SubWorkflowID != "" {
            return fmt.Errorf("sub-workflow %q failed: %w", event.SubWorkflowID, event.Error)
        }
        return event.Error
    case workflow.ExecutorFailedEvent:
        return fmt.Errorf("executor %q failed: %w", event.ExecutorID, event.Error)
    }
}

Fel som uppstår vid vidarebefordran av underordnade händelser konverteras också till överordnade workflow.ErrorEvent värden med det kopplade underarbetsflödes-ID:t.

Kontrollpunkter

Underarbetsflöden stöder kontrollpunkter. När det överordnade arbetsflödet skapar en kontrollpunkt lagrar underarbetsflödesbindningen det underordnade arbetsflödets kontrollpunktshanterare och eventuella väntande mappningar för kvalificerade svarsportar i det överordnade exekveringstillståndet. Vid återställning kan underarbetsflödet återupptas med sitt kapslade exekveringstillstånd intakt, inklusive väntande förfrågningar.

checkpointManager := checkpoint.NewInMemoryManager()
environment := inproc.Default.WithCheckpointing(checkpointManager)

var checkpoints []workflow.CheckpointInfo
run, err := environment.RunStreaming(ctx, parentWorkflow, "hello")
if err != nil {
    return err
}
defer run.Close(ctx)

for event, err := range run.WatchStream(ctx) {
    if err != nil {
        return err
    }
    if completed, ok := event.(workflow.SuperStepCompletedEvent); ok {
        if completed.CompletionInfo != nil && completed.CompletionInfo.CheckpointInfo != nil {
            checkpoints = append(checkpoints, *completed.CompletionInfo.CheckpointInfo)
        }
    }
}

if len(checkpoints) == 0 {
    return fmt.Errorf("no checkpoints were created")
}

resumedRun, err := environment.ResumeStreaming(ctx, parentWorkflow, checkpoints[len(checkpoints)-1])
if err != nil {
    return err
}
defer resumedRun.Close(ctx)

Om en checkpunkt återställs medan ett underarbetsflöde har en väntande begäran, återpublicerar den återställda överordnade körningen händelsen med information om den kvalificerade begäran. Anroparen kan skapa ett svar från den ompublicerade begäran och skicka tillbaka den via det överordnade körningshandtaget.

Nästa steg