Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Подпроцесс — это полный рабочий процесс, который исполняется в родительском рабочем процессе. Это позволяет создавать сложные системы из небольших, многократно используемых стандартных блоков рабочего процесса — каждый из них имеет собственный изолированный контекст выполнения, управление состоянием и маршрутизацию сообщений.
Overview
Подпроцессы полезны, когда вам нужно:
- Разложение сложности — разорвать большой рабочий процесс на меньшие, независимо тестируемые единицы.
- Повторное использование логики рабочего процесса — внедрение одного и того же подчиненного рабочего процесса в несколько родительских рабочих процессов.
- Изолирование состояния — сохраняйте внутреннее состояние каждого подпроцесса отдельно от родительского процесса.
- Управление потоком данных — сообщения входят и покидают подпроцесс только через его границы, без передачи данных на другие уровни.
Когда подпроцесс добавляется в родительский рабочий процесс, он ведет себя как любой другой исполнитель: он получает входные сообщения, запускает внутренний граф до завершения и создает выходные сообщения для последующих исполнителей.
Создание Подпроцесса
В C# вы создаете вложенные рабочие процессы двумя способами:
-
Прямая привязка — используется
BindAsExecutor()для внедрения рабочего процесса непосредственно в качестве исполнителя в родительский рабочий процесс. Это сохраняет собственные типы входных и выходных данных подпроцесса. -
Оболочка агента — используется
AsAIAgent()для преобразования рабочего процесса в агент, а затем добавьте агент в родительский рабочий процесс. Это полезно, если родительский рабочий процесс использует исполнителей на основе агента.
Прямая привязка с bindAsExecutor
Метод BindAsExecutor() расширения преобразует рабочий процесс в ExecutorBinding рабочий процесс, который можно добавить непосредственно в родительский рабочий процесс:
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();
С BindAsExecutor, типизированные входные и выходные типы вложенных рабочих процессов сохраняются — родительский рабочий процесс направляет сообщения на основе фактических типов, которые ожидаются и создаются вложенным рабочим процессом.
Упаковка агента с помощью AsAIAgent
Если родительский рабочий процесс использует исполнителей на основе агента, преобразуйте внутренний рабочий процесс в агент с помощью AsAIAgent().
WorkflowBuilder упаковывает агент в исполнительный механизм автоматически:
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();
Внутренний рабочий процесс работает в виде одного шага с точки зрения родительского процесса. Координатор отправляет сообщения в конвейер анализа, который внутренне выполняется specialist1 → specialist2, а затем результат пересылается рецензенту.
Tip
Используйте BindAsExecutor() при работе с типизированными исполнителями и AsAIAgent() при работе с рабочими процессами на основе агента. Дополнительные сведения о настройке преобразования рабочих процессов в агенты см. в разделе "Рабочие процессы в качестве агентов".
Типы входных и выходных данных
Если рабочий процесс используется в качестве вложенного рабочего процесса, он сохраняет контракты типов своих внутренних исполнителей.
При этом BindAsExecutorисполнитель вложенного рабочего процесса принимает те же типы входных данных, что и начальный исполнитель внутреннего рабочего процесса, и отправляет те же типы выходных данных, что и внутренний рабочий процесс. Края родительского рабочего процесса должны подключать исполнителей, типы выходных данных которых соответствуют ожидаемым типам входных данных вложенного рабочего процесса, а типы выходных данных вложенного рабочего процесса должны соответствовать ожидаемым входным данным подчиненных исполнителей.
При помощи AsAIAgent подрабочий процесс упаковывается как агент и следует контрактам ввода/вывода Agent Executor (string, ChatMessage, IEnumerable<ChatMessage>).
Поведение выходных данных
По умолчанию, когда вложенный рабочий процесс создает выходные данные (через YieldOutputAsync), эти выходные данные перенаправляются как сообщения подключенным исполнителям в родительском рабочем процессе. Это позволяет вторичным исполнителям обрабатывать результаты суб-рабочих процессов.
Класс ExecutorOptions управляет этим поведением:
| Опция | По умолчанию | Description |
|---|---|---|
AutoSendMessageHandlerResultObject |
true |
Перенаправьте выходные данные вложенного подпроцесса в виде сообщений подключенным исполнителям в родительском графе. |
AutoYieldOutputHandlerResultObject |
false |
Выводить выходные данные вложенного рабочего процесса непосредственно в поток событий выходных данных родительского рабочего процесса. |
Если параметр AutoYieldOutputHandlerResultObject включен, выходные данные вложенных рабочих процессов обходят внутреннюю маршрутизацию родительского процесса и передаются непосредственно вызывающему процессу родительского рабочего процесса.
var options = new ExecutorOptions
{
AutoYieldOutputHandlerResultObject = true,
};
ExecutorBinding subWorkflowExecutor = innerWorkflow.BindAsExecutor("SubWorkflow", options);
Запросы и ответы
Вложенные рабочие процессы полностью поддерживают механизм запроса и ответа . Когда исполнитель внутри дочернего рабочего процесса отправляет запрос (например, запрашивать данные человека), WorkflowHostExecutor перенаправляет RequestInfoEvent его в родительский рабочий процесс с соответствующим идентификатором порта — идентификатор исполнителя вложенного рабочего процесса добавляется к идентификатору порта (например, SubWorkflow.GuessNumber).
Эта квалификация гарантирует, что, когда родительский рабочий процесс получает ответ, он может направлять его обратно в правильный экземпляр подпроцесса. Родительский рабочий процесс обрабатывает запросы вложенного рабочего процесса с помощью того же механизма ответа, что и любой другой запрос:
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;
}
}
Замечание
С точки зрения вызывающего родительского рабочего процесса нет разницы между запросом от исполнителя верхнего уровня и запросом от дочернего рабочего процесса. Платформа прозрачно обрабатывает маршрутизацию.
Принцип работы
Когда родительский рабочий процесс направляет сообщение в обработчик вложенных рабочих процессов:
-
Доставка входных данных — сообщение пересылается в начальный исполнитель внутреннего рабочего процесса. При этом
BindAsExecutorтип сообщения должен соответствовать ожидаемым типам начального исполнителя. При этомAsAIAgentсообщения нормализуются вChatMessageформате. - Внутреннее выполнение — внутренний процесс выполняет собственный цикл супершага.
-
Сбор выходных данных — собираются выходные события внутреннего рабочего процесса. При этом
BindAsExecutorвыходные данные сохраняют исходные типы. При этомAsAIAgentвыходные данные преобразуются в сообщения ответа агента. - Переадресация запросов — если внутренний рабочий процесс имеет ожидающие запросы, они перенаправляются в родительский рабочий процесс для обработки (см. раздел "Запросы и ответы").
- Переадресация сообщений — полученные сообщения отправляются следующему исполнителю в родительском рабочем процессе.
Так как внутренний рабочий процесс поддерживает собственный контекст выполнения, его состояние не зависит от родительского рабочего процесса.
Tip
Дополнительные сведения о настройке преобразования рабочего процесса в агент, включая потоковую передачу и обработку исключений, см. в разделе "Рабочие процессы как агенты".
Многоуровневая вложенность
Подпроцессы могут быть вложены на произвольную глубину. Каждый уровень поддерживает собственный контекст выполнения:
// 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();
Замечание
Каждый уровень вложения увеличивает нагрузку на выполнение, так как внутренний рабочий процесс выполняет свой собственный цикл суперстепа. Сохраняйте глубину вложения разумно для сценариев с учетом производительности.
Обработка ошибок
При сбое вложенного рабочего процесса ошибка передается родительскому рабочему процессу как SubworkflowErrorEvent. Родительский рабочий процесс может наблюдать за этими ошибками через поток событий:
await foreach (WorkflowEvent evt in handle.WatchStreamAsync())
{
if (evt is SubworkflowErrorEvent subError)
{
Console.WriteLine($"Sub-workflow '{subError.ExecutorId}' failed: {subError.Data}");
}
}
Если вложенный рабочий процесс сталкивается с необработанным исключением, выполнение родительского рабочего процесса продолжается, но исполнитель вложенных рабочих процессов перестает обрабатывать дальнейшие сообщения.
Создание контрольных точек
При выполнении точки восстановления в родительском рабочем процессе состояние сеанса агента подпроцесса сериализуется как часть данных точки восстановления родительского процесса. При восстановлении состояние сеанса десериализируется, что позволяет родительскому рабочему процессу возобновить работу с состоянием дочернего рабочего процесса без изменений.
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);
Создание Подпроцесса
В Python вы создаете под-рабочий процесс, обернув Workflow в WorkflowExecutor и добавив его в родительскую рабочий процесс.
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()
)
Внутренний рабочий процесс работает в виде одного шага с точки зрения родительского процесса. Координатор отправляет сообщения в конвейер анализа, который внутренне выполняется specialist1 → specialist2, а затем результат пересылается рецензенту.
Параметры WorkflowExecutor
| Parameter | Тип | По умолчанию | Description |
|---|---|---|---|
workflow |
Workflow |
— | Экземпляр рабочего процесса, чтобы обернуть в качестве исполнителя. |
id |
str |
— | Уникальный идентификатор для этого исполнителя. |
allow_direct_output |
bool |
False |
Когда True выходные данные из подчиненных рабочих процессов передаются непосредственно в поток событий родительского рабочего процесса, а не отправляются в виде сообщений подключенным исполнителям. |
propagate_request |
bool |
False |
Когда True запросы из дочернего рабочего процесса передаются в поток событий родительского рабочего процесса, они рассматриваются как обычные информационные события о запросах. Когда False запросы упаковываются в SubWorkflowRequestMessage для перехвата родительскими исполнителями. |
Обёртывание вложенных рабочих процессов
Явно оборачивайте экземпляры Workflow в WorkflowExecutor перед добавлением их в родительский рабочий процесс. Агенты могут передаваться напрямую в WorkflowBuilder, но для необработанных экземпляров Workflow требуется эта обёртка.
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()
)
Явная оболочка позволяет:
- Назначьте определенный идентификатор исполнителя для ссылки в нескольких узлах.
- Повторно используйте один и тот же экземпляр
WorkflowExecutorпо всему графу.
# 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()
)
Типы входных и выходных данных
Элемент WorkflowExecutor наследует подпись типа от обернутого рабочего процесса.
-
Типы входных данных соответствуют типам входных данных начального исполнителя рабочего процесса (а также
SubWorkflowResponseMessageдля обработки ответов на перенаправленные запросы). -
Типы выходных данных соответствуют типам выходных данных в оболочке рабочего процесса. Если любой исполнитель в дочернем рабочем процессе имеет возможность ответа на запрос,
SubWorkflowRequestMessageтакже включается в качестве выходного типа.
Это означает, что края родительского рабочего процесса должны соединять исполнителей, чьи типы выходных данных соответствуют ожидаемым типам входных данных в подчинённом рабочем процессе. Аналогичным образом, подчиненные исполнители должны принимать типы данных, генерируемые подпроцессом.
# 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}")
Поведение выходных данных
По умолчанию (allow_direct_output=Falseесли вложенный рабочий процесс создает выходные данные через yield_output), эти выходные данные перенаправляются как сообщения подключенным исполнителям в родительском рабочем процессе с помощью send_message. Это позволяет последующим исполнителям обрабатывать результаты подпроцесса в рамках родительского графа.
Когда allow_direct_output=Trueвыходные данные вложенных рабочих процессов предоставляются непосредственно потоку событий родительского рабочего процесса. Выходные данные подпроцесса становятся выходными данными родительского потока работ, минуя внутреннюю маршрутизацию исполнителя родительского потока.
# 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)
Промежуточные выбросы дочерних рабочих процессов
"intermediate" События, созданные внутри дочернего рабочего процесса, автоматически всплывают в потоке событий родительского рабочего процесса. Они приписываются собственному WorkflowExecutor объекта id (а не внутреннему исполнителю, который первоначально их сгенерировал), что сохраняет инкапсуляцию. Что особенно важно, эти события сохраняют метку "intermediate" независимо от того, как родительский элемент обозначает WorkflowExecutor в своих списках output_from или intermediate_output_from.
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}")
Запросы и ответы
Вложенные рабочие процессы полностью поддерживают механизм запроса и ответа . При вызове ctx.request_info() внутри подчинённого рабочего процесса WorkflowExecutor перехватывает запрос и обрабатывает его на основе настройки propagate_request.
Перехват запросов в родительском рабочем процессе (по умолчанию)
При использовании параметра propagate_request=False (по умолчанию) запросы из вложенного рабочего процесса оборачиваются в SubWorkflowRequestMessage и отправляются подключенным исполнителям в родительском рабочем процессе. Это позволяет родительским исполнителям локально обрабатывать запрос:
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)
Метод create_response() проверяет, соответствует ли тип данных ответа ожидаемому типу из исходного запроса. Если типы не соответствуют, TypeError вызывается.
Important
При отправке ответа назад используйте target_id=request.executor_id для маршрутизации SubWorkflowResponseMessage в корректный экземпляр WorkflowExecutor.
Распространение запросов внешним вызывающим абонентам
При этом propagate_request=Trueзапросы из подпроцесса распространяются в поток событий родительского рабочего процесса с помощью стандартного request_info механизма. Вызывающий объект родительского рабочего процесса обрабатывает эти запросы так же, как и любой другой запрос человека в цикле:
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)
Принцип работы
Когда родительский рабочий процесс перенаправляет сообщение в WorkflowExecutor:
- Доставка входных данных — сообщение пересылается в начальный исполнитель внутреннего рабочего процесса. Тип сообщения должен соответствовать ожидаемым типам входных данных начального исполнителя.
- Внутреннее выполнение — внутренний процесс использует свой собственный цикл суперстепа до завершения или пока не требуется внешний ввод.
-
Сбор выходных данных — выходные события внутреннего рабочего процесса собираются и пересылаются в зависимости от настройки
allow_direct_output. -
Переадресация запросов — если внутренний рабочий процесс имеет ожидающие запросы, они пересылаются на
propagate_requestоснове параметра (см. раздел "Запросы и ответы"). -
Накопление ответов —
WorkflowExecutorсобирает ответы и возобновляет вложенный рабочий процесс только в том случае, если все ожидаемые ответы для заданного выполнения были получены. - Нисходящая диспетчеризация — выходные данные отправляются следующему исполнителю в родительском потоке работы.
Подпроцесс сохраняет собственное внутреннее состояние независимо от родительского. Сообщения перенаправляются только через края, соединяющие WorkflowExecutor со всеми остальными частями родительского графа, и отсутствует трансляция сообщений на разных уровнях вложенности.
Многоуровневая вложенность
Подпроцессы могут быть вложены на произвольную глубину. Каждый уровень поддерживает собственный контекст выполнения:
# 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()
)
Замечание
Каждый уровень вложения увеличивает нагрузку на выполнение, так как внутренний рабочий процесс выполняет свой собственный цикл суперстепа. Сохраняйте глубину вложения разумно для сценариев с учетом производительности.
Предупреждение
Все параллельные выполнения WorkflowExecutor используют одну и ту же базовую копию рабочего процесса. Исполнители внутри подпроцесса должны быть без отслеживания состояния, чтобы избежать помех между одновременными выполнениями.
Обработка ошибок
При сбое дочернего рабочего процесса ошибка распространяется на родительский рабочий процесс. Записывает WorkflowExecutor неудачное событие из вложенного рабочего процесса и преобразует его в событие ошибки в родительском контексте:
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)
Если подпроцесс сталкивается с необработанным исключением, основной процесс получает событие ошибки с подробными сведениями об исключении, включая идентификатор подпроцесса.
Создание контрольных точек
Подпроцессы поддерживают контрольные точки. При выполнении контрольной точки в родительском рабочем процессе сериализуется его внутреннее состояние, включая ход выполнения внутреннего рабочего процесса WorkflowExecutor и все кэшированные сообщения. При восстановлении это состояние десериализируется, позволяя родительскому рабочему процессу возобновить работу с вложенным рабочим процессом без изменений.
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)
Создание Подпроцесса
В Go создайте вложенный рабочий процесс, создав *workflow.Workflow и привязав его к родительскому рабочему процессу 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
}
Связанный дочерний рабочий процесс выполняется как один исполнитель с точки зрения родительского рабочего процесса. Сообщения вводятся через привязку, дочерний рабочий процесс запускает собственный внутренний граф, а выходные данные дочернего рабочего процесса направляются обратно в родительский граф.
Типы входных и выходных данных
Привязка наследует протокол обёрнутого рабочего процесса. Родительский рабочий процесс может отправлять сообщения, типы среды выполнения которых соответствуют принятым типам входных данных дочернего рабочего процесса, а привязка предоставляет выходные типы дочерних рабочих процессов как типы сообщений, так и типы выходных данных.
Это означает, что связи родительского рабочего процесса должны соединять исполнителей, чьи типы выходных данных соответствуют входным данным, которые принимает дочерний рабочий процесс, а последующие исполнители должны обрабатывать типы, выдаваемые дочерним рабочим процессом:
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()
Поведение выходных данных
Когда дочерний рабочий процесс формирует выходные данные, привязка подпроцесса отправляет их в виде сообщения подключенным исполнителям в родительском рабочем процессе. Если родительский рабочий процесс также помечает привязку вложенного рабочего процесса с помощью WithOutputFrom, то выводится то же значение как родительский workflow.OutputEvent, у которого ExecutorID — это идентификатор привязки вложенного рабочего процесса.
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()
Пользовательские события рабочего процесса, создаваемые внутри дочернего рабочего процесса, перенаправляются в родительский поток событий. Собственные события жизненного цикла запуска и супершага подпроцесса остаются внутренними, чтобы родительский поток оставался сосредоточенным на внешне значимых событиях.
Запросы и ответы
Вложенные рабочие процессы поддерживают механизм запроса и ответа . Когда исполнитель внутри дочернего рабочего процесса отправляет внешний запрос, привязка вложенного рабочего процесса квалифициирует идентификатор порта запроса путем подготовки идентификатора привязки. Например, порт запроса дочернего рабочего процесса с именем ApprovalPort становится ApprovalSubWorkflow.ApprovalPort в родительском рабочем процессе.
Чтобы сделать дочерний запрос доступным через родительский рабочий процесс, добавьте родительский порт RequestPort с полным идентификатором и направьте запросы и ответы между привязкой подчинённого рабочего процесса и этим портом:
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()
Вызывающий объект обрабатывает запрос из родительского потока рабочего процесса и отправляет ответ обратно через тот же дескриптор выполнения. Привязка вложенного рабочего процесса удаляет квалифицированный префикс перед передачей ответа дочернему рабочему процессу.
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)
}
}
Используйте предикаты, чтобы сузить границы запроса и ответа:
func externalRequestOnly(msg any) bool {
_, ok := msg.(*workflow.ExternalRequest)
return ok
}
func externalResponseOnly(msg any) bool {
_, ok := msg.(*workflow.ExternalResponse)
return ok
}
Принцип работы
Когда родительский рабочий процесс маршрутизирует сообщение в привязку к вложенному рабочему процессу:
- Доставка входных данных — привязка принимает сообщения, соответствующие поддерживаемым типам входных данных дочернего рабочего процесса, и помещает их в очередь исполнителя запуска дочернего рабочего процесса.
- Внутреннее выполнение — дочерний рабочий процесс выполняется в той же внутрипроцессной среде выполнения и сохраняет собственный цикл супершага.
-
Перенаправление выходных данных — дочерние значения
workflow.OutputEventотправляются в виде сообщений из привязки в родительские исполнительные модули ниже по потоку, а также возвращаются родителем, если привязка указана вWithOutputFrom. -
Переадресация запросов — дочерние
workflow.RequestInfoEventзапросы повторно отправляются с полными идентификаторами портов и могут маршрутизироваться через привязки родительскогоRequestPort. -
Переадресация событий — в родительский поток добавляются пользовательские дочерние события рабочего процесса. Ошибки отображаются в виде родительских
workflow.ErrorEventзначений с записанным идентификатором дочернего рабочего процесса. - Нисходящая отправка — итоговые сообщения продолжают проходить по связям родительского рабочего процесса.
Дочерний рабочий процесс сохраняет состояние и маршрутизацию сообщений отдельно от родительского. Сообщения пересекают эту границу только через связи, связанные с привязкой подпроцесса.
Многоуровневая вложенность
Подпроцессы могут быть вложены на произвольную глубину. Каждый дочерний рабочий процесс привязан перед добавлением в рабочий процесс, содержащий его:
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()
Замечание
Каждый уровень вложенности добавляет дополнительные накладные расходы при выполнении, поскольку дочерний рабочий процесс выполняет собственный цикл супершагов. Сохраняйте глубину вложения разумно для сценариев с учетом производительности.
Обработка ошибок
Когда дочерний рабочий процесс выдает ошибку, привязка подпроцесса передает ее в родительский рабочий процесс как workflow.ErrorEvent и устанавливает SubWorkflowID в значение идентификатора привязки. Родительский рабочий процесс может наблюдать за этими ошибками с помощью одного потока событий, который он использует для ошибок рабочего процесса верхнего уровня:
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)
}
}
Ошибки, возникающие при переадресации дочерних событий, также преобразуются в родительские workflow.ErrorEvent значения с вложенным идентификатором рабочего процесса.
Создание контрольных точек
Подпроцессы поддерживают контрольные точки. Когда родительский рабочий процесс создаёт точку сохранения, привязка подпроцесса сохраняет менеджер точек сохранения дочернего процесса и все ожидающие квалифицированные сопоставления портов ответов в состоянии исполнителя родительского процесса. При восстановлении дочерний рабочий процесс может продолжить выполнение с сохранением состояния вложенного выполнения, включая ожидающие запросы.
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)
Если выполняется восстановление из контрольной точки, пока в дочернем рабочем процессе есть незавершённый запрос, восстановленный родительский запуск повторно публикует событие с квалифицированной информацией о запросе. Вызывающая сторона может создать ответ на основе этого повторно опубликованного запроса и отправить его обратно через дескриптор родительского запуска.