Microsoft 代理程式框架工作流程 - 人力介入(HITL)

本頁概述了 Microsoft 代理框架工作流程系統中的 人工參與(HITL) 互動。 HITL 透過工作流程中的 請求與回應 處理機制實現,允許執行者向外部系統(如人工操作員)發送請求,並等待回應後再執行工作流程。

概觀

工作流程中的執行程式可以將要求傳送至工作流程外部,並等待回應。 這對於執行器需要與外部系統互動的場景非常有用,例如人機迴圈互動或任何其他非同步操作。

我們來建立一個工作流程,請人工操作員猜數字,並用執行器判斷猜測是否正確。

在工作流程中啟用請求和回應處理

請求和響應通過稱為 RequestPort的特殊類型處理。

A RequestPort 是一個通訊通道,允許執行者發送請求並接收回應。 當執行者向 發送 RequestPort訊息時,請求埠會發出包含請求細節的 a RequestInfoEvent 。 外部系統可以監聽這些事件,處理請求,並將回應回傳給工作流程。 框架會根據原始請求自動將回應路由回適當的執行者。

// Create a request port that receives requests of type NumberSignal and responses of type int.
var numberRequestPort = RequestPort.Create<NumberSignal, int>("GuessNumber");

將輸入埠新增至工作流程。

JudgeExecutor judgeExecutor = new(42);
var workflow = new WorkflowBuilder(numberRequestPort)
    .AddEdge(numberRequestPort, judgeExecutor)
    .AddEdge(judgeExecutor, numberRequestPort)
    .WithOutputFrom(judgeExecutor)
    .Build();

定義 JudgeExecutor 需要一個目標數字,並能判斷猜測是否正確。 若不正確,系統會再次發送請求,透過RequestPort要求新的猜測。

internal enum NumberSignal
{
    Init,
    Above,
    Below,
}

internal sealed class JudgeExecutor() : Executor<int>("Judge")
{
    private readonly int _targetNumber;
    private int _tries;

    public JudgeExecutor(int targetNumber) : this()
    {
        this._targetNumber = targetNumber;
    }

    public override async ValueTask HandleAsync(int message, IWorkflowContext context, CancellationToken cancellationToken = default)
    {
        this._tries++;
        if (message == this._targetNumber)
        {
            await context.YieldOutputAsync($"{this._targetNumber} found in {this._tries} tries!", cancellationToken);
        }
        else if (message < this._targetNumber)
        {
            await context.SendMessageAsync(NumberSignal.Below, cancellationToken: cancellationToken);
        }
        else
        {
            await context.SendMessageAsync(NumberSignal.Above, cancellationToken: cancellationToken);
        }
    }
}

在 Python 中,執行者使用 ctx.request_info() 發送請求,並使用 @response_handler 裝飾器處理回應。

我們來建立一個工作流程,請人工操作員猜數字,並用執行器判斷猜測是否正確。

在工作流程中啟用請求和回應處理

from dataclasses import dataclass

from agent_framework import (
    Executor,
    WorkflowBuilder,
    WorkflowContext,
    handler,
    response_handler,
)


@dataclass
class NumberSignal:
    hint: str  # "init", "above", or "below"


class JudgeExecutor(Executor):
    def __init__(self, target_number: int):
        super().__init__(id="judge")
        self._target_number = target_number
        self._tries = 0

    @handler
    async def handle_guess(self, guess: int, ctx: WorkflowContext[int, str]) -> None:
        self._tries += 1
        if guess == self._target_number:
            await ctx.yield_output(f"{self._target_number} found in {self._tries} tries!")
        elif guess < self._target_number:
            await ctx.request_info(request_data=NumberSignal(hint="below"), response_type=int)
        else:
            await ctx.request_info(request_data=NumberSignal(hint="above"), response_type=int)

    @response_handler
    async def on_human_response(
        self,
        original_request: NumberSignal,
        response: int,
        ctx: WorkflowContext[int, str],
    ) -> None:
        await self.handle_guess(response, ctx)


judge = JudgeExecutor(target_number=42)
workflow = WorkflowBuilder(start_executor=judge).build()

@response_handler裝飾器會自動註冊方法,以處理指定要求和回應類型的回應。 該框架會根據 和 original_request 參數的response類型註解,將收到的回應匹配到正確的處理程序。

工作流程支援透過 RequestPort的人機參與循環模式,該模式會暫停執行並等待外部輸入。

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

approval := approvalPort.Bind()
finalize := workflow.NewExecutor("FinalizeExecutor", func(approved bool) string {
    if approved {
        return "Request approved by the human reviewer"
    }
    return "Request rejected by the human reviewer"
}).Bind()

wf, err := workflow.NewBuilder(approval).
    AddEdge(approval, finalize).
    WithOutputFrom(finalize).
    Build()

A RequestPort 定義了工作流程與外部世界之間的類型請求/回應通道。 當執行器抵達請求埠時,工作流程會暫停並發出外部請求事件。 當外部回應出現時,工作流程才會恢復。

處理請求和回應

在收到要求時,一個 RequestPort 會發出一個 RequestInfoEvent。 您可以訂閱這些事件,以處理來自工作流程的傳入要求。 當您收到來自外部系統的回應時,請使用回應機制將其傳回工作流程。 框架會自動將回應路由到發送原始請求的執行器。

await using StreamingRun handle = await InProcessExecution.RunStreamingAsync(workflow, NumberSignal.Init);
await foreach (WorkflowEvent evt in handle.WatchStreamAsync())
{
    switch (evt)
    {
        case RequestInfoEvent requestInputEvt:
            // Handle `RequestInfoEvent` from the workflow
            int guess = ...; // Get the guess from the human operator or any external system
            await handle.SendResponseAsync(requestInputEvt.Request.CreateResponse(guess));
            break;

        case WorkflowOutputEvent outputEvt:
            // The workflow has yielded output
            Console.WriteLine($"Workflow completed with result: {outputEvt.Data}");
            return;
    }
}

小提示

完整可執行專案請參閱 完整範例

執行器可以直接發送請求,而不需要單獨的元件。 當執行者呼叫ctx.request_info()時,工作流程會伴隨WorkflowEvent一起發出type == "request_info"。 您可以訂閱這些事件,以處理來自工作流程的傳入要求。 當您收到來自外部系統的回應時,請使用回應機制將其傳回工作流程。 框架會自動將回應路由到執行器的方法 @response_handler

from collections.abc import AsyncIterable

from agent_framework import WorkflowEvent


async def process_event_stream(stream: AsyncIterable[WorkflowEvent]) -> dict[str, int] | None:
    """Process events from the workflow stream to capture requests."""
    requests: list[tuple[str, NumberSignal]] = []
    async for event in stream:
        if event.type == "request_info":
            requests.append((event.request_id, event.data))

    # Handle any pending human feedback requests.
    if requests:
        responses: dict[str, int] = {}
        for request_id, request in requests:
            guess = ...  # Get the guess from the human operator or any external system.
            responses[request_id] = guess
        return responses

    return None

# Initiate the first run of the workflow with an initial guess.
# Runs are not isolated; state is preserved across multiple calls to run.
stream = workflow.run(25, stream=True)

pending_responses = await process_event_stream(stream)
while pending_responses is not None:
    # Run the workflow until there is no more human feedback to provide,
    # in which case this workflow completes.
    stream = workflow.run(stream=True, responses=pending_responses)
    pending_responses = await process_event_stream(stream)

小提示

請參閱此 完整範例 以取得完整的可執行檔案。

等待 workflow.RequestInfoEvent、根據請求建立回應,然後使用該回應繼續執行流程:

run, err := inproc.Default.Run(ctx, wf, "Approve deployment to production?")
if err != nil {
    return err
}

var request *workflow.ExternalRequest
for evt := range run.NewEvents() {
    if requestEvent, ok := evt.(workflow.RequestInfoEvent); ok {
        request = requestEvent.Request
        break
    }
}

response, err := request.CreateResponse(true)
if err != nil {
    return err
}

if _, err := run.Resume(ctx, response); err != nil {
    return err
}

for evt := range run.NewEvents() {
    if output, ok := evt.(workflow.OutputEvent); ok {
        fmt.Println(output.Output)
    }
}

小提示

請參閱人在迴路中的範例,以查看完整且可執行的檔案。

人機回路與代理協作

RequestPort上述所描述的模式適用於自訂執行者和WorkflowBuilder。 在使用 代理協調 (如序列、並行或群組聊天工作流程)時, 工具審核 是透過人工在迴圈中的請求/回應機制來完成。

代理人可以使用需要人工批准才能執行的工具。 當代理嘗試呼叫需要核准的工具時,工作流程會暫停,並像 RequestInfoEvent 模式一樣發出 RequestPort,但事件酬載中包含的是 ToolApprovalRequestContent(C# 和 Go)或含有 Contenttype == "function_approval_request"(Python),而不是自訂請求型別。

對於代理程式需要先從使用者蒐集更多資訊並反覆確認後才能繼續的互動情境;而非僅僅核准或拒絕工具呼叫;請使用 交接協調流程。 交接預設為互動式:當某個代理在回應後未將交接轉給其他代理時,控制權會返回使用者以進行下一次輸入,從而可在協同流程中進行多輪來回互動。 順序、並行及群組聊天編排本身不會暫停以等待自由形式的使用者輸入;當您需要在步驟之間進行這類控制時,請在自訂 RequestPort 工作流程中將它們與 WorkflowBuilder 搭配使用。

檢查點和請求

欲了解更多檢查點資訊,請參閱檢查點。

建立檢查點時,擱置中的請求也會儲存為檢查點狀態的一部分。 當您從檢查點還原時,任何擱置中的請求都會重新發出為 RequestInfoEvent 物件,讓您能夠擷取並回應它們。 你也可以從檢查點恢復,並透過將 checkpoint_idresponses 一併傳遞給 workflow.run(...),在同一次呼叫中提供回應。

還原後,請監聽重新觸發的請求事件,並透過先前針對你的語言所示範的相同回應機制來回應。

後續步驟