Sub-Workflows

سير العمل الفرعي هو سير عمل كامل يتم تشغيله كمنفذ داخل سير عمل أصل. يمكنك هذا من إنشاء أنظمة معقدة من كتل إنشاء سير عمل أصغر وقابلة لإعادة الاستخدام - لكل منها سياق تنفيذ معزول وإدارة الحالة وتوجيه الرسائل.

نظرة عامة

تكون مهام سير العمل الفرعية مفيدة عندما تريد:

  • تحليل التعقيد — تقسيم سير عمل كبير إلى وحدات أصغر وقابلة للاختبار بشكل مستقل.
  • إعادة استخدام منطق سير العمل — تضمين نفس سير العمل الفرعي في مهام سير عمل أصل متعددة.
  • عزل الحالة — احتفظ بالحالة الداخلية لكل سير عمل فرعي منفصلة عن الأصل.
  • التحكم في تدفق البيانات — تدخل الرسائل سير العمل الفرعي وتتركه فقط من خلال حوافه، مع عدم وجود بث عبر المستويات.

عند إضافة سير عمل فرعي إلى سير عمل أصل، فإنه يتصرف مثل أي منفذ آخر: يتلقى رسائل الإدخال، ويشغل الرسم البياني الداخلي الخاص به حتى الاكتمال، وينتج رسائل الإخراج لمنفذي انتقال البيانات من الخادم.

إنشاء Sub-Workflow

في 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، يتم التفاف سير العمل الفرعي كعامل ويتبع عقود الإدخال/الإخراج للمنفذ العامل (string، ChatMessage، IEnumerable<ChatMessage>).

سلوك الإخراج

بشكل افتراضي، عندما ينتج سير عمل فرعي مخرجات (عبر YieldOutputAsync)، تتم إعادة توجيه هذه المخرجات كرسائل إلى المنفذين المتصلين في سير العمل الأصل. وهذا يمكن منفذي انتقال البيانات من الخادم من معالجة نتائج سير العمل الفرعي.

ExecutorOptions تتحكم الفئة في هذا السلوك:

خيار Default 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;
    }
}

Note

من منظور المتصل بسير العمل الأصل، لا يوجد فرق بين طلب من منفذ المستوى الأعلى وطلب من سير عمل فرعي. يتعامل إطار العمل مع التوجيه بشفافية.

كيفية عملها

عندما يوجه سير العمل الأصل رسالة إلى منفذ سير العمل الفرعي:

  1. تسليم الإدخال — تتم إعادة توجيه الرسالة إلى منفذ بدء سير العمل الداخلي. مع BindAsExecutor، يجب أن يتطابق نوع الرسالة مع الأنواع المتوقعة لمنفذ البدء. باستخدام AsAIAgent، يتم تسوية الرسائل إلى ChatMessage تنسيق.
  2. التنفيذ الداخلي — يقوم سير العمل الداخلي بتشغيل حلقة التراكب الخاصة به.
  3. مجموعة الإخراج — يتم تجميع أحداث إخراج سير العمل الداخلي. مع BindAsExecutor، تحتفظ المخرجات بالأنواع الأصلية الخاصة بها. باستخدام AsAIAgent، يتم تحويل المخرجات إلى رسائل استجابة العامل.
  4. إعادة توجيه الطلب — إذا كان سير العمل الداخلي يحتوي على طلبات معلقة، تتم إعادة توجيهها إلى سير العمل الأصل للتعامل معها (راجع الطلبات والاستجابات).
  5. إرسال انتقال البيانات من الخادم — يتم إرسال الرسائل الناتجة إلى المنفذ التالي في سير العمل الأصل.

نظرا لأن سير العمل الداخلي يحتفظ بسياق التنفيذ الخاص به، فإن حالته مستقلة عن سير العمل الأصل.

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

Note

يضيف كل مستوى تداخل حمل تنفيذ لأن سير العمل الداخلي يقوم بتشغيل حلقة التداخل الفائقة الخاصة به. حافظ على عمق التداخل معقولا للسيناريوهات الحساسة للأداء.

معالجة الأخطاء

عند فشل سير عمل فرعي، يتم نشر الخطأ إلى سير العمل الأصل ك 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);

إنشاء Sub-Workflow

في 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

المعلمه Type Default 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" تظهر الأحداث التي يتم إنتاجها داخل سير عمل تابع عبر دفق الحدث الأصل تلقائيا. وهي تنسب إلى WorkflowExecutorid الخاصية (وليس المنفذ الداخلي الذي انبعث منها في الأصل)، مما يحافظ على التغليف. والأهم من ذلك، تحتفظ هذه الأحداث بالتسمية "intermediate" بغض النظر عن كيفية تعيين WorkflowExecutor الأصل في قوائمه أو intermediate_output_from الخاصة 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:

  1. تسليم الإدخال — تتم إعادة توجيه الرسالة إلى منفذ بدء سير العمل الداخلي. يجب أن يتطابق نوع الرسالة مع أنواع الإدخال المتوقعة لمنفذ البدء.
  2. التنفيذ الداخلي — يقوم سير العمل الداخلي بتشغيل حلقة التراكب الخاصة به حتى الاكتمال، أو حتى يحتاج إلى إدخال خارجي.
  3. مجموعة الإخراج — يتم تجميع أحداث إخراج سير العمل الداخلي وإعادة توجيهها استنادا allow_direct_output إلى الإعداد.
  4. إعادة توجيه الطلب — إذا كان سير العمل الداخلي يحتوي على طلبات معلقة، تتم إعادة توجيهها استنادا propagate_request إلى الإعداد (راجع الطلبات والاستجابات).
  5. تراكم الاستجابة - WorkflowExecutor يجمع الاستجابات ويستأنف سير العمل الفرعي فقط عند تلقي جميع الاستجابات المتوقعة لتنفيذ معين.
  6. إرسال انتقال البيانات من الخادم — يتم إرسال المخرجات إلى المنفذ التالي في سير العمل الأصل.

يحافظ سير العمل الفرعي على حالته الداخلية الخاصة بشكل مستقل عن الأصل. يتم توجيه الرسائل فقط من خلال الحواف التي تربط 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()
)

Note

يضيف كل مستوى تداخل حمل تنفيذ لأن سير العمل الداخلي يقوم بتشغيل حلقة التداخل الفائقة الخاصة به. حافظ على عمق التداخل معقولا للسيناريوهات الحساسة للأداء.

تحذير

تشترك جميع عمليات التنفيذ المتزامنة 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)

إنشاء Sub-Workflow

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

تتم إعادة توجيه أحداث سير العمل المخصصة المنبعثة داخل سير العمل التابع إلى دفق الحدث الأصل. يتم الاحتفاظ بأحداث بدء سير العمل الفرعي ودورة حياة فائقة داخلية بحيث يظل الدفق الأصل مركزا على الأحداث ذات المعنى الخارجي.

الطلبات والاستجابات

تدعم مهام سير العمل الفرعية آلية الطلب والاستجابة . عندما يقوم منفذ داخل سير العمل الفرعي بنشر طلب خارجي، يؤهل ربط سير العمل الفرعي معرف منفذ الطلب عن طريق إلحاق معرف الربط. على سبيل المثال، يصبح ApprovalSubWorkflow.ApprovalPort منفذ الطلب التابع المسمى 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
}

كيفية عملها

عندما يوجه سير العمل الأصل رسالة إلى ربط سير العمل الفرعي:

  1. تسليم الإدخال - يقبل الربط الرسائل التي تطابق أنواع الإدخال المقبولة لسير العمل التابع ويلصقها في منفذ بدء سير العمل التابع.
  2. التنفيذ الداخلي — يتم تشغيل سير العمل التابع في نفس بيئة التنفيذ قيد المعالجة ويحافظ على حلقة التراكب الخاصة به.
  3. إعادة توجيه الإخراج — يتم إرسال القيم التابعة workflow.OutputEvent كرسائل من الربط إلى المنفذين الأصليين المتلقيين للمعلومات، كما يتم إرجاعها من الأصل إذا كان الربط مدرجا في WithOutputFrom.
  4. إعادة توجيه الطلب - تتم إعادة إرسال الطلبات التابعة workflow.RequestInfoEvent بمعرفات منفذ مؤهلة ويمكن توجيهها من خلال روابط الأصل RequestPort .
  5. إعادة توجيه الحدث — تتم إضافة أحداث سير عمل تابعة مخصصة إلى الدفق الأصل. تظهر الأخطاء كقيم أصل workflow.ErrorEvent مع تسجيل معرف سير العمل الفرعي.
  6. إرسال انتقال البيانات من الخادم — تستمر الرسائل الناتجة عبر حواف سير العمل الأصل.

يحافظ سير العمل التابع على الحالة وتوجيه الرسائل بشكل منفصل عن الأصل. تعبر الرسائل الحدود فقط من خلال الحواف المتصلة بربط سير العمل الفرعي.

تداخل متعدد المستويات

يمكن تداخل مهام سير العمل الفرعية مع العمق العشوائي. يتم ربط كل سير عمل تابع قبل إضافته إلى سير العمل الذي يحتوي عليه:

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

Note

يضيف كل مستوى تداخل حمل تنفيذ لأن سير العمل التابع يقوم بتشغيل حلقة التداخل الفائقة الخاصة به. حافظ على عمق التداخل معقولا للسيناريوهات الحساسة للأداء.

معالجة الأخطاء

عندما يصدر سير عمل تابع خطأ، يقوم ربط سير العمل الفرعي بإعادة توجيهه إلى سير العمل الأصل ك 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)

إذا تمت استعادة نقطة تحقق أثناء وجود طلب معلق لسير عمل تابع، يعيد التشغيل الأصل المستعادة نشر حدث معلومات الطلب المؤهل. يمكن للمتصل إنشاء استجابة من هذا الطلب الذي تمت إعادة نشره وإرساله مرة أخرى من خلال مقبض التشغيل الأصل.

الخطوات التالية