我们如何保持长程智能体存活

构建一个持久的事件响应代理,它能在工作进程崩溃时存活、等待人类批准、安全地重试并恢复而不重复已完成的工作。

我们如何保持长程智能体存活
AI模型价格对比 | AI工具导航 | ONNX模型库 | Vibe Coding教程 | PLC在线仿真器 | Tripo 3D | Meshy AI | ElevenLabs | KlingAI | ArtSpace | Phot.AI | InVideo

一个事件响应代理已经收集了日志、关联了指标并提出了修复方案。它正在等待工程师批准更改。

然后工作进程Pod在部署期间被替换。

对话可能仍然存在于数据库中。提议的操作可能仍然在跟踪中可见。但有几个更难的问题仍然存在:

  • 哪个步骤实际上已经完成?
  • 修复请求是否已经发送?
  • 失败的模型调用是否应该重试?
  • 运行是否可以等待批准直到明天?
  • 当工作进程返回时,执行应该从哪里恢复?

这些主要不是模型问题。它们是持久执行问题。

Temporal为我们提供了一种将编排编写为代码的方法,同时在崩溃、重启、超时和长时间等待中保留其执行状态。

它不会让代理更智能。它让代理的运行足够持久,能够在其周围的基础设施中存活。

1、核心思想一图概览

正常的代理框架通常帮助我们定义提示、工具、交接和推理循环。Temporal位于该逻辑之下,保持运行可恢复。

警报到达
    ↓
收集日志和指标
    ↓
要求模型进行诊断
    ↓
等待人类批准
    ↓
应用更改
    ↓
验证恢复
    ↓
关闭或升级事件

Temporal使用两个主要构建块对此建模:

Workflow(工作流)
确定性编排和持久状态

Activity(活动)
可能失败、超时或重试的外部工作

Workflow拥有顺序、分支、等待状态和关于下一步发生什么的决定。

Activity执行外部工作,例如调用LLM、查询监控API、写入数据库或重启服务。

这种分离很重要,因为Temporal可能会重放Workflow代码以重建其状态。因此Workflow代码必须是确定性的。网络调用、模型调用、数据库查询和现实世界的副作用应该放在Activity中。

2、Temporal实际持久化了什么

Temporal为每个Workflow执行记录仅追加的事件历史。历史捕获重要进度,如Workflow开始、Activity调度、Activity完成、信号、计时器、失败和重试。

当工作进程消失时,另一个工作进程可以重放该历史并重建Workflow的逻辑状态。已完成的Activity不会在重放期间简单地重复,因为它们记录的结果会返回给Workflow。

有一个重要的边界:

Temporal恢复Workflow的逻辑进度。它不会从进程停止的确切指令恢复任意Activity代码。

如果Activity在中途崩溃,该Activity可能会从头开始运行。这就是为什么重试安全的副作用仍然是我们的责任。

3、一个实用的事件修复工作流

以下示例保持核心机制可见。模型调用、监控调用和基础设施更改是Activity。Workflow通过信号协调它们并等待批准。

from dataclasses import dataclass
from datetime import timedelta

from temporalio import activity, workflow
from temporalio.common import RetryPolicy

@dataclass
class Incident:
    incident_id: str
    service: str
    summary: str

@dataclass
class RemediationPlan:
    action: str
    reason: str

@activity.defn
async def collect_evidence(incident: Incident) -> dict:
    # 在此调用日志和指标系统
    return {
        "service": incident.service,
        "error_rate": 0.18,
        "recent_deployment": True,
    }

@activity.defn
async def propose_remediation(evidence: dict) -> RemediationPlan:
    # 在此调用模型,提供证据和严格的输出模式
    return RemediationPlan(
        action="roll_back_latest_deployment",
        reason="Error rate increased immediately after deployment",
    )

@activity.defn
async def apply_remediation(plan: RemediationPlan) -> str:
    info = activity.info()
    idempotency_key = f"{info.workflow_run_id}-{info.activity_id}"
    # 将密钥传递给部署服务,这样重试不会创建第二个回滚请求
    return await deployment_api.apply(
        action=plan.action,
        idempotency_key=idempotency_key,
    )

@activity.defn
async def verify_recovery(service: str) -> bool:
    # 在此查询健康检查和监控
    return True

@workflow.defn
class IncidentRemediationWorkflow:
    def __init__(self) -> None:
        self.approved = False
        self.approved_by: str | None = None
    @workflow.signal
    def approve(self, approver: str) -> None:
        self.approved = True
        self.approved_by = approver
    @workflow.run
    async def run(self, incident: Incident) -> dict:
        retry = RetryPolicy(
            initial_interval=timedelta(seconds=2),
            backoff_coefficient=2.0,
            maximum_attempts=3,
        )
        evidence = await workflow.execute_activity(
            collect_evidence,
            incident,
            start_to_close_timeout=timedelta(seconds=30),
            retry_policy=retry,
        )
        plan = await workflow.execute_activity(
            propose_remediation,
            evidence,
            start_to_close_timeout=timedelta(seconds=60),
            retry_policy=retry,
        )
        await workflow.wait_condition(lambda: self.approved)
        change_id = await workflow.execute_activity(
            apply_remediation,
            plan,
            start_to_close_timeout=timedelta(minutes=2),
            retry_policy=retry,
        )
        recovered = await workflow.execute_activity(
            verify_recovery,
            incident.service,
            start_to_close_timeout=timedelta(seconds=30),
            retry_policy=retry,
        )
        return {
            "incident_id": incident.incident_id,
            "approved_by": self.approved_by,
            "change_id": change_id,
            "recovered": recovered,
        }

工作进程注册很小且故意无聊:

from temporalio.client import Client
from temporalio.worker import Worker

client = await Client.connect("localhost:7233")
worker = Worker(
    client,
    task_queue="incident-agent",
    workflows=[IncidentRemediationWorkflow],
    activities=[
        collect_evidence,
        propose_remediation,
        apply_remediation,
        verify_recovery,
    ],
)
await worker.run()

一个单独的客户端启动运行并稍后发送批准:

handle = await client.start_workflow(
    IncidentRemediationWorkflow.run,
    Incident("INC-1042", "checkout-api", "Error rate increased"),
    id="incident-INC-1042",
    task_queue="incident-agent",
)

await handle.signal(
    IncidentRemediationWorkflow.approve,
    "on-call-engineer",
)
result = await handle.result()

此示例故意比生产实现小,但责任边界是现实的:

  • 编排和等待状态存在于Workflow中;
  • 模型和工具调用存在于Activity中;
  • 重试按Activity配置;
  • 批准作为信号到达;
  • 外部更改使用幂等键。

4、使价值显而易见的失败测试

可以使用Temporal开发服务器运行本地演示:

temporal server start-dev

启动工作进程和Workflow后,我们可以在collect_evidence完成但批准到达之前停止工作进程。

Workflow在Temporal中保持打开。事件历史已包含完成的Activity结果和当前等待状态。

当工作进程重新启动时,Workflow被重放,返回到批准等待,并在批准信号到达时继续。

这是保存代理消息和保留执行之间的实际区别。消息持久性告诉我们说了什么。持久执行告诉系统什么已完成、什么待处理以及接下来允许运行什么。

5、重试不意味着恰好一次

Temporal Activity使用至少一次执行模型。

考虑这种失败:

回滚请求到达部署API
        ↓
API接受请求
        ↓
工作进程丢失响应
        ↓
Temporal观察到超时
        ↓
Activity被重试

Temporal无法知道外部系统是否完成了第一个请求。如果没有幂等键或协调检查,重试可能会重复副作用。

这是最重要的生产经验之一:

Temporal使重试持久。我们仍然必须使副作用对重试安全。

对于支付、部署、配置或消息传递系统,这通常意味着幂等键、去重记录或写前读协调步骤。

6、人类批准成为正常的工作流状态

人类批准步骤不应该要求工作进程保持存活数小时。

Temporal信号允许外部客户端异步更改Workflow状态。Workflow可以使用workflow.wait_condition等待该状态。

这使长时间等待变得普通:

计划提议
    ↓
Workflow持久等待
    ↓
工作进程可能停止或重新部署
    ↓
批准信号稍后到达
    ↓
Workflow恢复

对于可以更改基础设施、发送客户通信、批准支付或修改生产数据的代理系统,这比将批准保持在一个长HTTP请求或内存代理会话中更清晰的控制边界。

7、权衡和Temporal未解决的问题

Temporal解决执行持久性,而不是语义正确性。

它不决定:

  • 模型的诊断是否正确;
  • 证据是否充分;
  • 另一个推理回合是否值得其成本;
  • 提议的操作是否被授权;
  • 工具结果是否值得信任;
  • 代理是否在没有进展的情况下重复自己。

我们仍然需要评估、策略检查、范围凭据、工具沙箱、可观察性和人类升级。Temporal可以可靠地保留和执行这些决策,但它不会替代它们。

它还引入了操作和设计权衡。Workflow代码必须保持确定性和与重放兼容。长时间运行的Workflow可能积累大量事件历史,因此非常活跃的代理循环可能需要继续作为新实例。

团队必须在运行仍然打开时仔细版本控制Workflow代码。Temporal服务器是MIT许可的,因此我们可以自托管它,但这增加了服务、持久性、升级、容量规划和值班所有权。

8、Temporal与LangGraph和代理框架

Temporal和代理框架解决系统的不同部分。

代理框架
提示、工具、交接、模型循环

图框架
状态转换、分支、连接

Temporal
持久执行、重试、计时器、等待、恢复

我们并不总是需要所有三个。一个简短的同步助手可能不需要Temporal的基础设施。一个长时间运行的编码、研究、事件、入职或批准代理可能立即受益于它。

9、我们会在哪里使用它

当代理运行满足以下条件时,Temporal非常合适:

  • 持续时间超过一个请求或工作进程生命周期;
  • 等待人类或外部事件;
  • 跨越多个不可靠的服务;
  • 执行不应该不必要地重复的昂贵步骤;
  • 需要显式重试、截止日期和恢复;
  • 必须在事件后可检查;
  • 可能产生重要的现实世界副作用。

对于简单的聊天完成、简短的无状态工具调用或可以从头安全重启的后台作业,它可能是不必要的。

10、实际要点

生产代理不仅仅是模型加工具。它也是一个分布式过程,可以在网络失败、工作进程重启、批准需要数小时和外部系统返回模糊结果时保持活动。

Temporal为我们提供了一个持久的地方来保持该过程存活。

模型仍然可能推理不正确。工具仍然可能不安全。如果我们设计不当,重试仍然可能重复副作用。但一旦这些决策和边界被定义,Temporal可以确保运行不会在底层基础设施更改时简单地消失。

这就是真正的价值:

Temporal不会让我们的代理更智能。它让代理的执行足够持久以在现实世界中存活。

原文链接: How We Keep Long-Running AI Agents Alive with Temporal

汇智网翻译整理,转载请标明出处