Alfolyamatok

Az al-munkafolyamatok olyan teljes munkafolyamatok, amelyek végrehajtóként futnak egy szülő-munkafolyamaton belül. Ez lehetővé teszi összetett rendszerek írását kisebb, újrafelhasználható munkafolyamat-építőelemekből – mindegyik saját elkülönített végrehajtási környezettel, állapotkezeléssel és üzenet-útválasztással.

Overview

Az al-munkafolyamatok akkor hasznosak, ha:

  • Összetettség felbontása – egy nagy munkafolyamat kisebb, egymástól függetlenül tesztelhető egységekre bontható.
  • Munkafolyamat-logika újrafelhasználása – ugyanazt az al-munkafolyamatot ágyazza be több szülő-munkafolyamatba.
  • Az állapot elkülönítése – az egyes al-munkafolyamatok belső állapotának elkülönítése a szülőtől.
  • Az adatfolyam szabályozása – az üzenetek csak az élein keresztül lépnek be és hagyják el az al-munkafolyamatot, és nem történik közvetítés több szinten keresztül.

Amikor egy al-munkafolyamatot hozzáad egy szülő munkafolyamathoz, ugyanúgy viselkedik, mint bármely más végrehajtó: bemeneti üzeneteket fogad, belső gráfját a befejezésig futtatja, és kimeneti üzeneteket hoz létre az alsóbb rétegbeli végrehajtók számára.

Sub-Workflow létrehozása

A C#-ban az al-munkafolyamatok kétféleképpen írhatóak:

  • Közvetlen kötés – munkafolyamat BindAsExecutor() közvetlen beágyazása végrehajtóként a szülő munkafolyamatba. Ez megőrzi az al-munkafolyamat natív bemeneti/kimeneti típusait.
  • Ügynökburkolás – használja a AsAIAgent()-t a munkafolyamat ügynökké alakításához, majd adja hozzá az ügynököt a szülő munkafolyamathoz. Ez akkor hasznos, ha a szülő munkafolyamat ügynökalapú végrehajtókat használ.

Közvetlen kötés a BindAsExecutorral

A BindAsExecutor() bővítménymetódus átalakít egy munkafolyamatot olyan ExecutorBinding munkafolyamattá, amely közvetlenül hozzáadható egy szülő munkafolyamathoz:

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();

Ezzel BindAsExecutoraz al-munkafolyamat beírt bemeneti és kimeneti típusai megmaradnak – a szülő munkafolyamat az al-munkafolyamat által várt és előállított tényleges típusok alapján irányítja át az üzeneteket.

Ügynökburkolás az AsAIAgenttel

Amikor a szülő munkafolyamat ügynökalapú végrehajtókat használ, konvertálja a belső munkafolyamatot ügynökhasználatra a következő módon: AsAIAgent(). Az WorkflowBuilder ügynök automatikusan egy végrehajtóba burkolódik:

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();

A belső munkafolyamat egyetlen lépésben fut a szülő munkafolyamat szemszögéből. A koordinátor üzeneteket küld az elemzési folyamatnak, amely belsőleg fut specialist1 → specialist2, majd továbbítja az eredményt a véleményezőnek.

Tip

Használja a BindAsExecutor() elemet gépelt végrehajtókkal és a AsAIAgent() elemet ügynökalapú munkafolyamatokkal való munka során. A munkafolyamatok ügynökké alakításának konfigurálásáról további információt a Munkafolyamatok ügynökként című témakörben talál.

Bemeneti és kimeneti típusok

Ha egy munkafolyamatot al-munkafolyamatként használ, megőrzi a belső végrehajtók típusszerződéseit.

Ezzel BindAsExecutoregyütt az al-munkafolyamat-végrehajtó ugyanazokat a bemeneti típusokat fogadja el, mint a belső munkafolyamat kezdő végrehajtója, és ugyanazokat a kimeneti típusokat küldi el, amelyeket a belső munkafolyamat előállít. A szülő munkafolyamat éleinek össze kell kapcsolniuk azokat a végrehajtókat, amelyek kimeneti típusai megegyeznek az al-munkafolyamat várt bemeneti típusaival, és az al-munkafolyamat kimeneti típusainak meg kell egyeznie az alsóbb rétegbeli végrehajtók várt bemenetével.

A AsAIAgent segítségével az al-munkafolyamat ügynökként van becsomagolva, és követi az Ügynök Végrehajtó bemeneti/kimeneti szerződéseit (string, ChatMessage, IEnumerable<ChatMessage>).

Kimeneti viselkedés

Alapértelmezés szerint, amikor egy al-munkafolyamat kimeneteket állít elő (a YieldOutputAsync segítségével), azokat üzenetekként továbbítják a szülő munkafolyamat csatlakoztatott végrehajtóinak. Ez lehetővé teszi az alsóbb rétegbeli végrehajtók számára a munkafolyamat-részeredmények feldolgozását.

Az ExecutorOptions osztály a következő viselkedést vezérli:

Option Alapértelmezett Description
AutoSendMessageHandlerResultObject true Az al-munkafolyamat kimeneteinek továbbítása üzenetként a csatlakoztatott végrehajtóknak a szülőgráfban.
AutoYieldOutputHandlerResultObject false Az al-munkafolyamat kimeneteit közvetlenül a szülő munkafolyamat kimeneti eseménystreamjének adja vissza.

Ha AutoYieldOutputHandlerResultObject engedélyezve van, az al-munkafolyamat kimenetei megkerülik a szülő belső útválasztását, és közvetlenül a szülő munkafolyamat hívójának lesznek kézbesítve.

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

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

Kérelmek és válaszok

Az al-munkafolyamatok teljes mértékben támogatják a kérés- és válaszmechanizmust . Amikor az al-munkafolyamatban egy végrehajtó kérést küld (például emberi bemenet igénylésére), a WorkflowHostExecutor rendszer egy RequestInfoEvent továbbítja a szülő munkafolyamatnak – az al-munkafolyamat végrehajtójának azonosítója elő van állítva a portazonosítóra (példáulSubWorkflow.GuessNumber).

Ez a minősítés biztosítja, hogy amikor a szülő munkafolyamat választ kap, vissza tudja irányítani a választ a megfelelő al-munkafolyamat-példányra. A szülő munkafolyamat az al-munkafolyamat-kérelmeket ugyanazzal a válaszmechanizmussal kezeli, mint bármely más kérés:

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;
    }
}

Megjegyzés:

A fölérendelt munkafolyamat-hívó szempontjából nincs különbség a legfelső szintű végrehajtótól érkező kérések és az al-munkafolyamatból érkező kérések között. A keretrendszer transzparensen kezeli az útválasztást.

Hogyan működik?

Amikor a szülő munkafolyamat átirányít egy üzenetet az al-munkafolyamat végrehajtójának:

  1. Bemeneti kézbesítés – az üzenet a belső munkafolyamat kezdő végrehajtójának továbbítja. Ezzel BindAsExecutoraz üzenettípusnak meg kell egyeznie a kezdő végrehajtó várt típusával. Az AsAIAgent használatával az üzenetek ChatMessage formátumra normalizálódnak.
  2. Belső végrehajtás – a belső munkafolyamat saját superstep ciklust futtat.
  3. Kimeneti gyűjtemény – a rendszer összegyűjti a belső munkafolyamat kimeneti eseményeit. A BindAsExecutorkimenetek megtartják az eredeti típusukat. A AsAIAgentkimenetek ügynöki válaszüzenetekké alakulnak.
  4. Kérelemtovábbítás – ha a belső munkafolyamat függőben lévő kérésekkel rendelkezik, a rendszer továbbítja őket a szülő munkafolyamatnak kezelés céljából (lásd : Kérések és válaszok).
  5. Alsóbb rétegbeli küldés – az így kapott üzenetek a szülő munkafolyamat következő végrehajtójának lesznek elküldve.

Mivel a belső munkafolyamat saját végrehajtási környezetet tart fenn, az állapota független a szülő-munkafolyamattól.

Tip

A munkafolyamatok ügynökké alakításának konfigurálásáról, beleértve a streamelési viselkedést és a kivételkezelést, tekintse meg a Munkafolyamatok ügynökként című témakört.

Többszintű beágyazás

Az al-munkafolyamatok tetszőleges mélységbe ágyazhatók. Minden szint saját végrehajtási környezetet tart fenn:

// 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();

Megjegyzés:

Minden beágyazási szint végrehajtási többletterhelést eredményez, mivel a belső munkafolyamat saját szupersztep-ciklust futtat. Tartsa ésszerű határok között a fészkelés mélységét a teljesítményérzékeny forgatókönyvekben.

Hibakezelés

Ha egy al-munkafolyamat meghiúsul, a hiba propagálása a szülő munkafolyamatba SubworkflowErrorEventtörténik. A szülő munkafolyamat az eseménystreamen keresztül figyelheti meg ezeket a hibákat:

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

Ha az al-munkafolyamat kezeletlen kivételt tapasztal, a szülő munkafolyamat végrehajtása folytatódik, de az al-munkafolyamat végrehajtója leállítja a további üzenetek feldolgozását.

Ellenőrző pontok használata

Amikor ellenőrzőpontot hoz létre a szülő-munkafolyamaton, a rendszer szerializálja az al-munkafolyamat-ügynök munkamenet-állapotát a szülő végrehajtó ellenőrzőpont-adatainak részeként. A visszaállítás során a munkamenet állapota deszerializálódik, lehetővé téve a szülő munkafolyamat számára, hogy a változatlan alfolyamat állapotával folytatódjon.

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);

Sub-Workflow létrehozása

Pythonban egy al-munkafolyamatot úgy hozhat létre, hogy egy Workflow beágyaz egy WorkflowExecutor-be, és hozzáadja azt egy szülő munkafolyamathoz.

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()
)

A belső munkafolyamat egyetlen lépésben fut a szülő munkafolyamat szemszögéből. A koordinátor üzeneteket küld az elemzési folyamatnak, amely belsőleg fut specialist1 → specialist2, majd továbbítja az eredményt a véleményezőnek.

WorkflowExecutor paraméterek

Paraméter Típus Alapértelmezett Description
workflow Workflow A munkafolyamat-példány, amely végrehajtóként van körbefuttatva.
id str A végrehajtó egyedi azonosítója.
allow_direct_output bool False Amikor Trueaz al-munkafolyamat kimenetei közvetlenül a szülő munkafolyamat eseményfolyamára kerülnek, ahelyett, hogy üzenetként küldenék őket a csatlakoztatott végrehajtóknak.
propagate_request bool False Amikor az al-munkafolyamatból érkező kérések a szülő munkafolyamat eseményfolyamára kerülnek propagálásra, rendszeres kérésinformációs eseményekként jelennek meg. Amikor Falsea rendszer becsomagolta SubWorkflowRequestMessage a kérelmeket a fölérendelt végrehajtók elfogására.

Almunkafolyamatok beágyazása

A Workflow példányait kifejezetten foglalja egy WorkflowExecutor elembe, mielőtt hozzáadja őket egy szülő munkafolyamathoz. Az ügynökök közvetlenül átadhatók a(z) WorkflowBuilder számára, de a nyers Workflow példányokhoz ez a burkoló szükséges.

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()
)

Az explicit körbefuttatás lehetővé teszi a következőket:

  • Rendeljen hozzá egy adott végrehajtóazonosítót a hivatkozáshoz több élen.
  • Használja ugyanazt a WorkflowExecutor példányt újra a gráfon.
# 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()
)

Bemeneti és kimeneti típusok

A becsomagolt munkafolyamat örökli a típusjelölését:

  • A bemeneti típusok megfelelnek a becsomagolt munkafolyamat kezdő végrehajtójának bemeneti típusainak (valamint SubWorkflowResponseMessage a továbbított kérelmekre adott válaszok kezeléséhez).
  • A kimeneti típusok megegyeznek a becsomagolt munkafolyamat kimeneti típusaival. Ha az al-munkafolyamat bármely végrehajtója képes a kérés-válaszra, SubWorkflowRequestMessage akkor kimeneti típusként is szerepel.

Ez azt jelenti, hogy a szülő munkafolyamat éleinek össze kell kapcsolniuk azokat a végrehajtókat, amelyek kimeneti típusai megegyeznek az al-munkafolyamat várt bemeneti típusaival. Hasonlóképpen, az alsóbb rétegbeli végrehajtóknak el kell fogadniuk az al-munkafolyamat által előállított típusokat:

# 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}")

Kimeneti viselkedés

Alapértelmezés szerint (allow_direct_output=False) amikor egy al-munkafolyamat kimeneteket hoz létre, yield_outputa rendszer üzenetekként továbbítja ezeket a kimeneteket a szülő munkafolyamatban send_messagelévő csatlakoztatott végrehajtóknak. Ez lehetővé teszi az alsóbb rétegbeli végrehajtók számára, hogy a fölérendelt gráf részeként dolgozzák fel az al-munkafolyamat eredményeit.

Amikor allow_direct_output=True, az al-munkafolyamat kimenetei közvetlenül a szülő munkafolyamat eseményfolyamához kerülnek. Az al-munkafolyamat kimenetei a szülő munkafolyamat kimenetei lesznek, megkerülve a szülő belső végrehajtói útválasztását:

# 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)

Köztes kibocsátások gyermek-munkafolyamatokból

"intermediate" a gyermek-munkafolyamatban létrehozott események automatikusan továbbítódnak a szülő munkafolyamat eseményfolyamába. Ezek a WorkflowExecutor saját id elemének tulajdoníthatók (nem annak a belső végrehajtónak, amely eredetileg kibocsátotta őket), ami megőrzi az enkapszulációt. Döntő fontosságú, hogy ezek az események megtartják a "intermediate" címkét , függetlenül attól, hogy a szülő hogyan jelöli meg a WorkflowExecutor saját output_from vagy intermediate_output_from listáiban.

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}")

Kérelmek és válaszok

Az al-munkafolyamatok teljes mértékben támogatják a kérés- és válaszmechanizmust . Amikor egy végrehajtó egy al-munkafolyamatban hív ctx.request_info(), a WorkflowExecutor rendszer elfogja a kérést, és a propagate_request beállítás alapján kezeli azt.

Kérelmek lehallgatása a szülői munkafolyamatban (alapértelmezett)

Alapértelmezés szerint propagate_request=False, az al-munkafolyamatból érkező kérések be vannak csomagolva SubWorkflowRequestMessage-be, és a szülő munkafolyamat csatlakoztatott végrehajtóinak lesznek elküldve. Ez lehetővé teszi, hogy a szülő végrehajtók helyileg kezeljék a kérést:

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)

A create_response() metódus ellenőrzi, hogy a válasz adattípusa megegyezik-e az eredeti kérés várt típusával. Ha a típusok nem egyeznek, akkor TypeError kivétel keletkezik.

Important

Amikor visszaküldi a választ, használja a target_id=request.executor_id-t a SubWorkflowResponseMessage megfelelő WorkflowExecutor példányra való irányításához.

Kérelmek propagálása külső hívóknak

Ezzel propagate_request=Trueaz al-munkafolyamatból érkező kéréseket a rendszer a standard request_info mechanizmus használatával propagálja a szülő munkafolyamat eseményfolyamára. A szülő munkafolyamat hívója ugyanúgy kezeli ezeket a kéréseket, mint bármely más, emberi beavatkozást igénylő kérést.

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)

Hogyan működik?

Amikor a szülő munkafolyamat átirányít egy üzenetet a WorkflowExecutorkövetkezőhöz:

  1. Bemeneti kézbesítés – az üzenet a belső munkafolyamat kezdő végrehajtójának továbbítja. Az üzenettípusnak meg kell egyeznie a kezdő végrehajtó várt bemeneti típusával.
  2. Belső végrehajtás – a belső munkafolyamat a saját superstep ciklusát futtatja a befejezésig, vagy amíg külső bemenetre nem van szüksége.
  3. Kimeneti gyűjtemény – a belső munkafolyamat kimeneti eseményei a beállítás alapján allow_direct_output lesznek összegyűjtve és továbbítva.
  4. Kérelmek továbbítása – ha a belső munkafolyamat függőben lévő kérésekkel rendelkezik, a rendszer a propagate_request beállítás alapján továbbítja őket (lásd : Kérések és válaszok).
  5. Válaszfelhalmozódás – a WorkflowExecutor válaszok összegyűjtése és az al-munkafolyamat folytatása csak akkor történik meg, ha egy adott végrehajtásra vonatkozó összes várt válasz érkezett.
  6. Alsóbb rétegbeli küldés – a kimenetek a szülő munkafolyamat következő végrehajtójának lesznek elküldve.

Az al-munkafolyamat a szülőtől függetlenül fenntartja saját belső állapotát. Az üzenetek csak a szülőgráf többi részéhez csatlakozó WorkflowExecutor éleken keresztül vannak irányítva – a beágyazott szinteken nem jelennek meg üzenetek.

Többszintű beágyazás

Az al-munkafolyamatok tetszőleges mélységbe ágyazhatók. Minden szint saját végrehajtási környezetet tart fenn:

# 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()
)

Megjegyzés:

Minden beágyazási szint végrehajtási többletterhelést eredményez, mivel a belső munkafolyamat saját szupersztep-ciklust futtat. Tartsa ésszerű határok között a fészkelés mélységét a teljesítményérzékeny forgatókönyvekben.

Warning

WorkflowExecutor összes egyidejű végrehajtása ugyanazt a mögöttes munkafolyamat-példányt osztja meg. Az al-munkafolyamat végrehajtóinak állapot nélkülinek kell lenniük az egyidejű végrehajtások közötti interferencia elkerülése érdekében.

Hibakezelés

Ha egy al-munkafolyamat meghiúsul, a hiba propagálása a szülő munkafolyamatba történik. A WorkflowExecutor rendszer rögzíti a sikertelen eseményt az al-munkafolyamatból, és a szülőkörnyezetben hibaeseménysé alakítja:

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)

Ha az al-munkafolyamat kezeletlen kivételt tapasztal, a szülő munkafolyamat hibaüzenetet kap a kivétel részleteivel, beleértve az al-munkafolyamat azonosítóját is.

Ellenőrző pontok használata

Az al-munkafolyamatok támogatják az ellenőrzőpont-készítést. Ha ellenőrzőpontot hoznak létre a szülő munkafolyamaton, a WorkflowExecutor szerializálja annak belső állapotát, beleértve a belső munkafolyamat végrehajtási állapotát és a gyorsítótárazott üzeneteket. A visszaállítás során ez az állapot deszerializálva van, lehetővé téve, hogy a szülő munkafolyamat a meglévő al-munkafolyamattal együtt folytatódjon.

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)

Sub-Workflow létrehozása

A Go nyelvben úgy hozhat létre al-munkafolyamatot, hogy létrehoz egy *workflow.Workflow elemet, majd a inproc.BindSubworkflowAsExecutor használatával a szülő munkafolyamathoz köti.

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
}

A kötött gyermek munkafolyamat egyetlen végrehajtóként fut a szülő-munkafolyamat szemszögéből. Az üzenetek a kötésen keresztül lépnek be, a gyermek munkafolyamat saját belső gráfot futtat, a gyermek munkafolyamat kimenetei pedig vissza lesznek irányítva a szülőgráfba.

Bemeneti és kimeneti típusok

A kötés örökli a burkolt munkafolyamat protokollját. A szülő munkafolyamat küldhet olyan üzeneteket, amelyek futtatókörnyezeti típusai megegyeznek a gyermek munkafolyamat elfogadott bemeneti típusaival, a kötés pedig a gyermek munkafolyamat által előállított kimeneti típusokat teszi elérhetővé üzenettípusként és kimeneti típusként is.

Ez azt jelenti, hogy a szülő munkafolyamat éleinek össze kell kapcsolniuk azokat a végrehajtókat, akiknek a kimeneti típusai megegyeznek az al-munkafolyamat által elfogadott bemenetekkel, az alsóbb rétegbeli végrehajtóknak pedig a gyermek munkafolyamat által létrehozott típusokat kell kezelnie:

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()

Kimeneti viselkedés

Amikor egy gyermek-munkafolyamat kimenetet állít elő, az almunkafolyamat-kötés ezt a kimenetet üzenetként küldi el a kötéstől a szülő munkafolyamat kapcsolódó végrehajtóinak. Ha a szülő munkafolyamat az al-munkafolyamat kötését WithOutputFromis megjelöli, ugyanaz az érték lesz kibocsátva, mint egy szülő workflow.OutputEvent , akinek ExecutorID az al-munkafolyamat kötésazonosítója.

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()

A gyermek munkafolyamatban kibocsátott egyéni munkafolyamat-eseményeket a rendszer a szülőesemény-adatfolyamba továbbítja. A részmunkafolyamat saját indítási és szuperlépés-életciklushoz kapcsolódó eseményei belső használatúak maradnak, így a szülőfolyam a kívülről is releváns eseményekre összpontosíthat.

Kérelmek és válaszok

Az al-munkafolyamatok támogatják a kérés- és válaszmechanizmust . Amikor egy végrehajtó az al-munkafolyamatban közzétesz egy külső kérést, az al-munkafolyamat kötése a kötésazonosító előerősítésével minősíti a kérelemport azonosítóját. A gyermekkérelem-port ApprovalPort neve például a szülő munkafolyamatban lesz ApprovalSubWorkflow.ApprovalPort .

Ha a gyermekkérelmet a szülő-munkafolyamaton keresztül szeretné felszínre helyezni, adjon hozzá egy szülőt RequestPort , amely rendelkezik a munkafolyamat-rész-kötés és a port közötti minősített azonosítóval és útvonalkérésekkel és válaszokkal:

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()

A hívó kezeli a szülő munkafolyamat-adatfolyamból érkező kérést, és ugyanazon a lefuttatási leírón keresztül küldi vissza a választ. Az al-munkafolyamat-kötés eltávolítja a minősített névtérelőtagot, mielőtt továbbítja a választ a gyermek-munkafolyamatnak.

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)
    }
}

Használjon predikátumokat a kérelem és a válasz élének szűkítéséhez:

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

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

Hogyan működik?

Amikor a szülő munkafolyamat átirányít egy üzenetet az al-munkafolyamat kötéséhez:

  1. Bemeneti kézbesítés – a kötés fogadja a gyermek munkafolyamat elfogadott bemeneti típusainak megfelelő üzeneteket, és a gyermek munkafolyamat kezdő végrehajtójába alakítja őket.
  2. Belső végrehajtás – az alárendelt munkafolyamat ugyanabban az azonos folyamaton belüli végrehajtási környezetben fut, és saját szuperlépés-ciklust tart fenn.
  3. Kimenettovábbítás — a gyermek workflow.OutputEvent értékek üzenetekként kerülnek elküldésre a kötésből a lejjebb lévő szülő-végrehajtóknak, és a szülőből is kiadódnak, ha a kötés szerepel a WithOutputFrom listában.
  4. Kérelemtovábbítás – a gyermekkérelmek workflow.RequestInfoEvent újra kibocsáthatók minősített portazonosítókkal, és szülőkötéseken RequestPort keresztül irányíthatók.
  5. Eseménytovábbítás – a rendszer hozzáadja az egyéni gyermek munkafolyamat-eseményeket a szülőstreamhez. A hibák szülő workflow.ErrorEvent értékekként jelennek meg, miközben az al-munkafolyamat azonosítója rögzítésre kerül.
  6. Alsóbb rétegbeli küldés – az eredményként kapott üzenetek a fölérendelt munkafolyamat élein haladnak tovább.

A gyermek munkafolyamat a szülőtől elkülönítve tartja az állapotát és az üzenetek útválasztását. Az üzenetek csak a munkafolyamat-részkötéshez csatlakoztatott éleken keresztül lépik át a határt.

Többszintű beágyazás

Az al-munkafolyamatok tetszőleges mélységbe ágyazhatók. Minden egyes alárendelt munkafolyamat kötésre kerül, mielőtt hozzáadják az azt tartalmazó munkafolyamathoz:

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()

Megjegyzés:

Minden egyes beágyazási szint végrehajtási többletterheléssel jár, mivel a gyermekmunkafolyamat a saját szupersztep-ciklusát futtatja. Tartsa ésszerű határok között a fészkelés mélységét a teljesítményérzékeny forgatókönyvekben.

Hibakezelés

Ha egy gyermek-munkafolyamat hibát ad vissza, az al-munkafolyamat kötése workflow.ErrorEvent formájában továbbítja azt a szülő munkafolyamathoz, és a SubWorkflowID értékét a kötésazonosítóra állítja. A szülő munkafolyamat ugyanazokat a hibákat figyelheti meg, mint a legfelső szintű munkafolyamat-hibákhoz használt eseménystream:

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)
    }
}

A gyermekesemények továbbítása során felmerülő hibák a hozzájuk csatolt al-munkafolyamat-azonosítóval a szülő workflow.ErrorEvent értékeivé is konvertálódnak.

Ellenőrző pontok használata

Az al-munkafolyamatok támogatják az ellenőrzőpont-készítést. Amikor a szülő munkafolyamat ellenőrzőpontot hoz létre, az almunkafolyamat-kötés a gyermek-munkafolyamat ellenőrzőpont-kezelőjét és a függőben lévő minősített válaszportok leképezéseit a szülő végrehajtójának állapotában tárolja. A visszaállítás során a gyermek munkafolyamat a beágyazott végrehajtási állapottal folytatható, beleértve a függőben lévő kéréseket is.

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)

Ha egy ellenőrzőpont vissza lett állítva, miközben a gyermek munkafolyamat függőben lévő kéréssel rendelkezik, a visszaállított szülőfuttatás újra közzéteheti a minősített kéreleminformációs eseményt. A hívó létrehozhat egy választ abból az újraközölt kérésből, majd visszaküldheti azt a szülő futtatás kezelőjén keresztül.

Következő lépések