Agent 工程体系 · 第 94/98 篇。内容以 2026 年 9 月可验证的公开规范和稳定接口为基线;框架版本敏感能力会明确标注,不把实验行为写成通用保证。

Go 实现 Agent Runtime:状态、工具、流式、Context、并发和持久化

Agent Runtime 不是“调用一次大模型的函数”,而是一个负责驱动多步执行的运行时。它需要回答一组相互关联的问题:

  • 当前任务处于什么状态?
  • 大模型这一步观察到了什么?
  • 它决定调用哪个工具?
  • 工具调用是否允许并发?
  • 用户是否已经看到了部分输出?
  • 请求超时后,工具和后台任务是否会停止?
  • 进程崩溃后,如何从上一次安全位置恢复?
  • 重试工具时,如何避免重复扣款、重复发货或重复写入?
  • 一个运行中的 Agent,何时算真正结束?

OpenAI 的 Agents 文档将 Agent 描述为能够规划、调用工具、协作并维护足够状态以完成多步工作的应用;其运行循环的核心也是“调用模型、检查输出、执行工具或切换 Agent、直到得到最终答案”。(developers.openai.com) Anthropic 则强调,Agent 的关键不是预先写死所有步骤,而是让模型依据每一步获得的环境事实动态决定后续动作;同时必须设置最大迭代次数等停止条件。(anthropic.com)

本文以 Go 标准库为基础,构造一个不依赖具体模型厂商的 Agent Runtime。模型供应商、工具协议、事件格式和存储都通过接口隔离。这样做的目的不是重新实现某个 SDK,而是把运行时的控制边界、状态语义、并发模型和恢复机制讲清楚。


一、先定义 Agent Runtime 的执行语义

1. Agent 是一个闭环,而不是一次函数调用

将一次 Agent 运行抽象为:

St+1=δ(St,Ot,At,Rt)S_{t+1} = \delta(S_t, O_t, A_t, R_t)

其中:

  • StS_t:第 tt 步开始时的运行状态;
  • OtO_t:本轮模型看到的观察结果;
  • AtA_t:模型选择的动作,例如调用工具、请求人工确认或直接回答;
  • RtR_t:动作执行后的反馈;
  • δ\delta:运行时定义的状态转移函数。

典型执行顺序是:

用户输入
   ↓
准备模型输入
   ↓
调用模型
   ↓
解析模型输出
   ├── 最终答案 → 完成
   ├── 工具调用 → 执行工具 → 写入工具结果 → 继续
   ├── 人工确认 → 暂停 → 等待恢复
   └── 非法输出/错误 → 失败或重试

模型负责“提出下一步动作”,Runtime 负责“判断动作是否合法、执行动作、记录结果以及决定是否继续”。

这一区分很重要。如果把所有控制权都交给模型,模型可能:

  • 重复调用同一个工具;
  • 在工具失败后无限重试;
  • 生成不存在的工具名;
  • 在已经超时的请求中继续执行后台任务;
  • 在用户没有确认的情况下执行有副作用的动作。

因此,Runtime 必须是状态机的拥有者,模型只是状态机中的一个决策参与者。

2. 运行状态不等于对话历史

很多实现把所有内容都塞进一个 []Message。这能完成简单聊天,但不足以表达生产运行时的状态。

一个较完整的状态至少包括:

type RunStatus string

const (
	StatusRunning   RunStatus = "running"
	StatusWaiting   RunStatus = "waiting_approval"
	StatusCompleted RunStatus = "completed"
	StatusFailed    RunStatus = "failed"
	StatusCanceled  RunStatus = "canceled"
)

type RunState struct {
	RunID       string
	Status      RunStatus
	Step        int
	Version     int64
	Messages    []Message
	Variables   map[string]any
	Pending     []ToolCall
	LastError   string
	LeaseOwner  string
	LeaseExpire time.Time
}

其中:

  • Messages 是模型可见的交互历史;
  • Variables 是运行时变量,例如订单号、用户身份、检索结果摘要;
  • Pending 是尚未完成的动作;
  • Step 是逻辑步号,不应简单等同于网络请求次数;
  • Version 用于乐观并发控制;
  • Status 表示整个运行的生命周期;
  • LeaseOwnerLeaseExpire 用于防止多个 Worker 同时处理同一运行。

对话历史回答“模型说过什么”,运行状态回答“系统现在允许做什么”。

例如:

消息历史:
模型:我将为订单 A1001 退款。
工具调用:refund(order_id=A1001)
工具结果:需要人工确认。

运行状态:
status = waiting_approval
pending = [refund(A1001)]
allowed_next_action = approve | reject

如果只保存消息历史,恢复进程时还需要重新推断“退款调用是否已经发出”“是否等待确认”“是否可以重试”。这是不安全的。


二、状态、观察、决策、行动和反馈

1. 观察是事实,决策是意图

Runtime 应将模型输出拆成两类:

type ModelDecision struct {
	Text       string
	ToolCalls  []ToolCall
	Finish     bool
	NeedHuman  bool
}

type ToolCall struct {
	ID        string
	Name      string
	Arguments json.RawMessage
}

ToolCall 只是模型产生的意图,不代表工具已经执行。

例如模型输出:

{
  "tool_calls": [
    {
      "id": "call-1",
      "name": "get_weather",
      "arguments": {"city": "杭州"}
    }
  ]
}

此时系统只能说:

模型决定调用 get_weather

不能说:

天气已经查询成功

只有工具实际返回:

{
  "temperature": 28,
  "condition": "晴"
}

之后,Runtime 才能把这个结果作为新的观察写入上下文。

2. 工具结果必须成为下一轮模型输入

一次完整工具循环如下:

第 0 步:
用户消息 → 模型
模型决策:调用 get_weather

第 1 步:
执行 get_weather
工具反馈:杭州,28℃,晴

第 2 步:
用户消息 + 工具调用 + 工具结果 → 模型
模型决策:生成最终答案

如果工具结果只写日志、不写入下一轮输入,模型就无法根据环境事实修正计划。

形式化地说,模型决策应当满足:

At=π(St,Ot)A_t = \pi(S_t, O_t)

其中 π\pi 是模型策略。如果工具反馈没有进入 Ot+1O_{t+1},那么下一步决策实际使用的是旧观察:

At+1=π(St,Ot)A_{t+1} = \pi(S_t, O_t)

这会导致模型重复执行已经完成的动作,或者基于过期信息继续决策。

3. 终止条件不是“模型说结束”这么简单

Runtime 至少需要检查以下终止条件:

1. 模型给出最终答案,且没有未完成工具调用;
2. 达到最大步骤数;
3. 超过总截止时间;
4. 用户主动取消;
5. 工具或模型返回不可恢复错误;
6. 进入人工确认状态;
7. 运行被系统策略拒绝。

可以定义:

stop(S)=final(S)failed(S)canceled(S)waiting(S)step(S)N\text{stop}(S) = \text{final}(S) \lor \text{failed}(S) \lor \text{canceled}(S) \lor \text{waiting}(S) \lor \text{step}(S) \geq N

其中 NN 是最大逻辑步数。

反例是只要模型返回文本就结束:

模型:我需要查询订单状态。
Runtime:把这句话当作最终答案返回。

这会让“计划”被误当成“结果”。正确做法是:只要输出中存在工具调用,就不能进入完成态。


三、工具系统:从函数调用升级为可治理动作

1. 工具接口必须包含 Context

Go 的 context.Context 用于传递截止时间、取消信号和请求范围内的值;官方文档要求 Context 通常作为函数第一个参数传递,并且派生出的取消函数应在所有控制流路径上调用。(pkg.go.dev)

工具接口可以定义为:

type Tool interface {
	Name() string
	Description() string
	InputSchema() json.RawMessage
	Execute(ctx context.Context, args json.RawMessage) (json.RawMessage, error)
}

一个天气工具:

type WeatherTool struct{}

func (WeatherTool) Name() string {
	return "get_weather"
}

func (WeatherTool) Description() string {
	return "查询城市当前天气"
}

func (WeatherTool) InputSchema() json.RawMessage {
	return json.RawMessage(`{
		"type": "object",
		"required": ["city"],
		"properties": {
			"city": {"type": "string"}
		}
	}`)
}

func (WeatherTool) Execute(ctx context.Context, args json.RawMessage) (json.RawMessage, error) {
	var input struct {
		City string `json:"city"`
	}
	if err := json.Unmarshal(args, &input); err != nil {
		return nil, fmt.Errorf("invalid arguments: %w", err)
	}
	if input.City == "" {
		return nil, errors.New("city is required")
	}

	select {
	case <-time.After(100 * time.Millisecond):
		return json.Marshal(map[string]any{
			"city":        input.City,
			"temperature": 28,
			"condition":   "晴",
		})
	case <-ctx.Done():
		return nil, ctx.Err()
	}
}

这里有三个关键点:

  1. 参数解析失败属于输入错误,不应直接 panic;
  2. 业务调用使用 ctx,否则 Runtime 取消后工具仍可能阻塞;
  3. 工具结果使用结构化 JSON,便于模型、日志和恢复逻辑共同消费。

2. 工具描述是模型接口的一部分

工具不仅是 Go 函数,还包含:

  • 名称;
  • 功能描述;
  • 参数约束;
  • 返回结构;
  • 错误语义;
  • 副作用说明;
  • 是否允许并发;
  • 是否需要人工确认;
  • 是否支持幂等重试。

例如退款工具不能只描述为:

退款

更准确的契约应包括:

根据订单号发起退款。
该操作会产生外部资金副作用。
同一 idempotency_key 重复调用不会重复退款。
金额由服务端根据订单状态计算,调用方不能自行指定最终金额。
执行前需要人工确认。

模型可以决定“想退款”,但不能通过工具参数绕过服务端权限和业务状态检查。

3. 工具错误必须区分类型

至少区分四类错误:

type ToolErrorKind string

const (
	ErrInvalidInput ToolErrorKind = "invalid_input"
	ErrUnauthorized ToolErrorKind = "unauthorized"
	ErrTransient    ToolErrorKind = "transient"
	ErrPermanent    ToolErrorKind = "permanent"
)

错误类型决定后续动作:

错误类型 Runtime 行为
参数非法 将结构化错误反馈给模型,通常不自动重试
未授权 立即暂停或失败
临时错误 按策略重试,但受次数和截止时间限制
永久错误 反馈给模型或结束运行
Context 取消 停止当前动作,不应继续后台执行

不要把所有错误都变成:

{"error": "tool failed"}

因为模型无法区分“参数写错”和“服务暂时不可用”,Runtime 也无法选择正确的恢复策略。


四、一个可运行的 Go Runtime 骨架

下面的示例使用一个假的模型实现,重点展示运行循环、工具执行、事件流和 Context。它不依赖外部 API,保存为 main.go 后可以直接运行。

package main

import (
	"context"
	"encoding/json"
	"errors"
	"fmt"
	"sync"
	"time"
)

type Message struct {
	Role    string `json:"role"`
	Content string `json:"content"`
}

type ToolCall struct {
	ID        string
	Name      string
	Arguments json.RawMessage
}

type ModelDecision struct {
	Text      string
	ToolCalls []ToolCall
	Finish    bool
}

type Model interface {
	Decide(ctx context.Context, messages []Message) (ModelDecision, error)
}

type Tool interface {
	Name() string
	Execute(ctx context.Context, args json.RawMessage) (json.RawMessage, error)
}

type Event struct {
	Type string
	Data any
}

type Runtime struct {
	Model Model
	Tools map[string]Tool
	MaxSteps int
	OnEvent func(Event)
}

func (r *Runtime) emit(e Event) {
	if r.OnEvent != nil {
		r.OnEvent(e)
	}
}

func (r *Runtime) Run(ctx context.Context, input string) (string, error) {
	messages := []Message{
		{Role: "user", Content: input},
	}

	for step := 1; step <= r.MaxSteps; step++ {
		if err := ctx.Err(); err != nil {
			return "", err
		}

		r.emit(Event{
			Type: "step.started",
			Data: map[string]any{"step": step},
		})

		decision, err := r.Model.Decide(ctx, messages)
		if err != nil {
			return "", fmt.Errorf("model decision at step %d: %w", step, err)
		}

		if len(decision.ToolCalls) == 0 && decision.Finish {
			r.emit(Event{
				Type: "run.completed",
				Data: decision.Text,
			})
			return decision.Text, nil
		}

		if len(decision.ToolCalls) == 0 {
			return "", errors.New("model returned neither final answer nor tool call")
		}

		for _, call := range decision.ToolCalls {
			tool, ok := r.Tools[call.Name]
			if !ok {
				return "", fmt.Errorf("unknown tool: %s", call.Name)
			}

			r.emit(Event{
				Type: "tool.started",
				Data: call,
			})

			result, err := tool.Execute(ctx, call.Arguments)
			if err != nil {
				r.emit(Event{
					Type: "tool.failed",
					Data: map[string]any{
						"id":    call.ID,
						"name":  call.Name,
						"error": err.Error(),
					},
				})

				messages = append(messages, Message{
					Role:    "tool",
					Content: fmt.Sprintf(`{"tool_call_id":%q,"error":%q}`, call.ID, err.Error()),
				})
				continue
			}

			r.emit(Event{
				Type: "tool.completed",
				Data: map[string]any{
					"id":     call.ID,
					"name":   call.Name,
					"result": string(result),
				},
			})

			messages = append(messages, Message{
				Role:    "tool",
				Content: fmt.Sprintf(`{"tool_call_id":%q,"result":%s}`, call.ID, result),
			})
		}
	}

	return "", fmt.Errorf("maximum steps exceeded: %d", r.MaxSteps)
}

type DemoModel struct {
	mu    sync.Mutex
	calls int
}

func (m *DemoModel) Decide(ctx context.Context, messages []Message) (ModelDecision, error) {
	m.mu.Lock()
	defer m.mu.Unlock()

	m.calls++
	if m.calls == 1 {
		return ModelDecision{
			ToolCalls: []ToolCall{{
				ID:        "call-1",
				Name:      "get_weather",
				Arguments: json.RawMessage(`{"city":"杭州"}`),
			}},
		}, nil
	}

	return ModelDecision{
		Text:   "杭州当前 28℃,天气晴。",
		Finish: true,
	}, nil
}

type WeatherTool struct{}

func (WeatherTool) Name() string {
	return "get_weather"
}

func (WeatherTool) Execute(ctx context.Context, args json.RawMessage) (json.RawMessage, error) {
	var input struct {
		City string `json:"city"`
	}
	if err := json.Unmarshal(args, &input); err != nil {
		return nil, err
	}
	if input.City == "" {
		return nil, errors.New("city is required")
	}

	timer := time.NewTimer(50 * time.Millisecond)
	defer timer.Stop()

	select {
	case <-timer.C:
		return json.Marshal(map[string]any{
			"city":        input.City,
			"temperature": 28,
			"condition":   "晴",
		})
	case <-ctx.Done():
		return nil, ctx.Err()
	}
}

func main() {
	rt := &Runtime{
		Model: &DemoModel{},
		Tools: map[string]Tool{
			"get_weather": WeatherTool{},
		},
		MaxSteps: 5,
		OnEvent: func(e Event) {
			fmt.Printf("[%s] %+v\n", e.Type, e.Data)
		},
	}

	ctx, cancel := context.WithTimeout(context.Background(), time.Second)
	defer cancel()

	answer, err := rt.Run(ctx, "杭州今天的天气怎么样?")
	if err != nil {
		panic(err)
	}

	fmt.Println("answer:", answer)
}

预期输出类似:

[step.started] map[step:1]
[tool.started] {call-1 get_weather {"city":"杭州"}}
[tool.completed] map[id:call-1 name:get_weather result:{"city":"杭州","condition":"晴","temperature":28}]
[step.started] map[step:2]
[run.completed] 杭州当前 28℃,天气晴。
answer: 杭州当前 28℃,天气晴。

这个例子有意保持简单,但已经展示了 Runtime 的最小闭环:

  1. 用户输入进入消息历史;
  2. 模型生成工具调用;
  3. Runtime 校验工具是否存在;
  4. 工具接收同一个取消链;
  5. 工具结果进入下一轮消息;
  6. 模型根据结果生成最终答案;
  7. 生命周期事件同步发出。

生产实现还需要加入参数 Schema 校验、权限检查、重试、幂等键、持久化和人工确认。


五、流式输出:传输机制不是运行状态

1. 流式的含义

流式输出是指 Runtime 不等待完整最终结果,而是在执行过程中持续向客户端发送事件。

这些事件可以包括:

run.started
step.started
model.text.delta
tool.started
tool.completed
model.text.completed
run.completed
run.failed

OpenAI 的 Responses API 使用有类型的语义事件传输流式结果,例如文本增量、工具调用参数增量、完成事件和错误事件。(developers.openai.com)

在 Go 中,Runtime 可以通过 channel 暴露事件:

func (r *Runtime) Stream(ctx context.Context, input string) <-chan Event {
	out := make(chan Event, 16)

	go func() {
		defer close(out)

		r.OnEvent = func(e Event) {
			select {
			case out <- e:
			case <-ctx.Done():
			}
		}

		_, _ = r.Run(ctx, input)
	}()

	return out
}

客户端消费:

for event := range rt.Stream(ctx, "查询杭州天气") {
	switch event.Type {
	case "model.text.delta":
		fmt.Print(event.Data)
	case "tool.started":
		fmt.Println("\n正在调用工具...")
	case "run.completed":
		fmt.Println("\n运行完成")
	}
}

2. 增量文本不能作为唯一事实源

流式文本适合展示,但不适合直接作为持久化状态。

例如客户端已经收到:

退款将

此时进程崩溃。数据库中如果只保存这段文本,恢复后无法判断:

  • 这是最终答案的一部分,还是模型尚未完成;
  • 是否还有工具调用;
  • 是否已经产生了副作用;
  • 是否应该从模型重新请求。

因此需要区分:

展示事件:model.text.delta
事实事件:model.output.completed
状态事件:run.checkpointed
动作事件:tool.started / tool.completed

增量事件可以丢失或重新发送;决定运行状态的事件则必须可靠记录。

3. 背压和断开连接

流式通道存在消费者速度小于生产者速度的情况:

Runtime 生成事件:1000/s
客户端消费事件:100/s

如果无限制地向 channel 写入,会造成内存增长;如果无缓冲直接发送,模型和工具执行又会被前端网络阻塞。

常见做法是:

  • 使用有限缓冲;
  • 将“状态持久化”和“用户展示”分成两条路径;
  • 客户端断开时取消展示 Context;
  • 对关键状态事件先落盘,再尝试发送;
  • 对非关键增量事件允许丢弃或合并。

客户端断开不一定意味着 Agent 任务应该取消。交互请求和后台运行需要使用不同的 Context:

requestCtx:浏览器连接生命周期
runCtx:Agent 任务生命周期

如果直接把 requestCtx 传给整个后台任务,用户刷新页面可能导致重要任务被取消。反过来,如果所有任务都脱离请求 Context,又会导致取消信号无法传播。


六、Context:取消链、截止时间和请求范围

1. Context 不是全局配置

正确使用方式是:

func Run(ctx context.Context, input string) error
func Execute(ctx context.Context, args json.RawMessage) error
func CallModel(ctx context.Context, req Request) error

不应把 Context 存在 Runtime 结构体中:

// 不推荐
type Runtime struct {
	ctx context.Context
}

原因是 Runtime 通常会被多个运行共享,而 Context 属于单次请求或单次运行。官方 Go 文档明确建议不要把 Context 存入结构体,而是显式作为参数传递。(pkg.go.dev)

2. 多层截止时间

Agent 通常有三层时间限制:

请求截止时间:HTTP 请求还能等待多久
运行截止时间:整个 Agent 最多执行多久
工具截止时间:某一次外部调用最多执行多久

可以这样派生:

requestCtx := r.Context()

runCtx, cancelRun := context.WithTimeout(requestCtx, 2*time.Minute)
defer cancelRun()

toolCtx, cancelTool := context.WithTimeout(runCtx, 10*time.Second)
defer cancelTool()

如果请求 Context 在 30 秒后取消,所有派生 Context 都会被取消;如果某个工具超过 10 秒,只有该工具调用超时,Runtime 仍可根据策略决定是否重试或进入失败态。

3. 取消必须被实际消费

下面的代码虽然接收了 Context,但没有使用它:

func badTool(ctx context.Context) error {
	time.Sleep(30 * time.Second)
	return nil
}

这会使工具在 Runtime 已经取消后继续运行。

正确做法是将阻塞操作改为可取消:

func goodTool(ctx context.Context) error {
	timer := time.NewTimer(30 * time.Second)
	defer timer.Stop()

	select {
	case <-timer.C:
		return nil
	case <-ctx.Done():
		return ctx.Err()
	}
}

对于 HTTP 请求、数据库查询、RPC 调用,应使用支持 Context 的 API,例如:

req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
if err != nil {
	return err
}

resp, err := http.DefaultClient.Do(req)

4. 取消不等于回滚

Context 取消只能表示“调用方不再等待或要求停止”,不能自动撤销已经完成的外部副作用。

例如:

1. Runtime 调用支付服务;
2. 支付服务已经扣款;
3. Runtime 在读取响应前超时;
4. Runtime 只看到 context deadline exceeded。

此时不能简单重试扣款。系统必须通过支付流水号或幂等键查询原始请求状态,再决定是否继续。


七、并发:工具调用可以并发,状态提交不能失序

1. 哪些动作可以并发

如果模型同时产生两个互不依赖的只读调用:

get_weather(杭州)
get_exchange_rate(CNY/USD)

可以并发执行,以降低总延迟:

Tparallelmax(T1,T2)T_{\text{parallel}} \approx \max(T_1, T_2)

而不是:

Tserial=T1+T2T_{\text{serial}} = T_1 + T_2

但以下调用通常不能直接并发:

创建订单
扣款
发货
取消订单

因为它们之间可能存在业务依赖或顺序要求。

2. 并发工具执行示例

type toolResult struct {
	call   ToolCall
	result json.RawMessage
	err    error
}

func executeTools(
	ctx context.Context,
	tools map[string]Tool,
	calls []ToolCall,
) []toolResult {
	results := make([]toolResult, len(calls))
	var wg sync.WaitGroup

	for i, call := range calls {
		i, call := i, call
		wg.Add(1)

		go func() {
			defer wg.Done()

			tool, ok := tools[call.Name]
			if !ok {
				results[i] = toolResult{
					call: call,
					err:  fmt.Errorf("unknown tool: %s", call.Name),
				}
				return
			}

			result, err := tool.Execute(ctx, call.Arguments)
			results[i] = toolResult{
				call:   call,
				result: result,
				err:    err,
			}
		}()
	}

	wg.Wait()
	return results
}

这里用 results[i] 保持了输入顺序,因此即使工具完成顺序不同,传给模型的结果仍然稳定。

如果使用“谁先完成谁先写入”的方式,可能出现:

本次运行:
tool-B 先完成
tool-A 后完成

重试运行:
tool-A 先完成
tool-B 后完成

模型收到的消息顺序不同,后续决策也可能不同。这是确定性被破坏的一个来源。

3. 并发限制不能只靠 Goroutine 数量

生产环境还需要限制:

  • 单个运行最多并发工具数;
  • 单个租户最多并发运行数;
  • 某个工具的全局并发数;
  • 外部服务的速率;
  • 单次运行的累计成本。

可以使用信号量:

type Semaphore chan struct{}

func (s Semaphore) Acquire(ctx context.Context) error {
	select {
	case s <- struct{}{}:
		return nil
	case <-ctx.Done():
		return ctx.Err()
	}
}

func (s Semaphore) Release() {
	<-s
}

但并发控制应位于工具调度层,而不是让每个工具自己决定。否则 Runtime 无法知道当前运行是否已经超过资源上限。

4. 共享状态必须串行提交

工具可以并发执行,但状态更新最好通过单一提交点完成:

并发工具执行
    ↓
按 call index 收集结果
    ↓
一次性追加 tool result 事件
    ↓
更新内存状态
    ↓
写 checkpoint

否则两个工具同时修改 MessagesStepVariables,会产生数据竞争和事件顺序不一致。Go 的 -race 可以帮助发现内存级数据竞争,但它无法发现“事件已经持久化却没有更新状态版本”这类业务竞态。


八、持久化:事件日志、Checkpoint 和恢复

1. 为什么只保存最终状态不够

如果只保存一份最新快照:

run_id = R1
status = running
step = 3

发生故障后,你不知道:

  • 第 3 步模型是否已经返回;
  • 工具调用是否已经发出;
  • 工具是否成功;
  • 结果是否写入消息;
  • checkpoint 是工具前还是工具后生成的。

事件日志提供了过程证据:

1. run.created
2. model.requested
3. model.completed
4. tool.started
5. tool.completed
6. state.updated
7. checkpoint.created

Checkpoint 则是某个事件位置的可恢复快照:

checkpoint = {
    run_id: "R1",
    event_seq: 7,
    state: {...}
}

恢复时不必从第一条事件重新计算,而是:

加载 checkpoint
→ 从 event_seq + 1 继续读取事件
→ 重建待执行动作
→ 继续 Runtime

2. 事件和状态的最小模型

type EventRecord struct {
	Seq       int64
	RunID     string
	Type      string
	Payload   json.RawMessage
	CreatedAt time.Time
}

type Checkpoint struct {
	RunID     string
	EventSeq  int64
	StateJSON []byte
	CreatedAt time.Time
}

存储接口:

type Store interface {
	AppendEvent(ctx context.Context, event EventRecord) error
	LoadCheckpoint(ctx context.Context, runID string) (Checkpoint, error)
	SaveCheckpoint(ctx context.Context, cp Checkpoint) error
	AcquireLease(ctx context.Context, runID, owner string, ttl time.Duration) (bool, error)
	ReleaseLease(ctx context.Context, runID, owner string) error
}

真实数据库中,AppendEventSaveCheckpoint 最好在同一事务中完成,或者至少使用严格的事件序号和版本条件。

3. Checkpoint 的安全位置

不是每个时刻都适合生成 checkpoint。

安全位置通常是:

A. 新用户输入已经持久化;
B. 模型输出已经持久化;
C. 工具结果已经持久化;
D. 工具副作用的提交状态已经可查询;
E. 运行进入 waiting_approval;
F. 运行进入 completed 或 failed。

不安全的恢复点:

工具请求已经发送,但 tool.started 尚未落盘;
工具已经扣款,但 tool.completed 尚未落盘;
内存状态已经改变,但 checkpoint 尚未保存。

这些窗口会造成“重放时以为没有执行”和“实际上已经执行”的矛盾。

4. 事件日志与 Checkpoint 的关系

事件日志是事实序列,Checkpoint 是加速恢复的缓存。

如果:

事件日志:不可修改、按序追加
Checkpoint:可以重建、可以过期

那么即使 Checkpoint 损坏,也可以从事件日志重新构造状态。

反过来,如果只保存 Checkpoint,不保存事件日志,系统无法回答审计问题:

谁触发了退款?
模型何时决定退款?
用户是否确认?
工具返回了什么?
第一次调用是否已经超时?

九、租约:防止多个 Worker 同时接管一个运行

1. 租约的语义

租约是带过期时间的所有权:

Worker A 获得 run R1 的租约
租约有效期:2026-09-01 10:00:00 至 10:01:00

只有租约持有者可以修改运行状态。

Worker 需要周期性续租:

当前时间 + 续租窗口 < lease_expire

如果 Worker 崩溃,租约过期后,其他 Worker 才能接管。

2. 为什么不能只用分布式锁

锁通常表达:

现在谁持有?

租约还表达:

持有者多久不续租就失效?

对于可能崩溃的 Agent Worker,永久锁会导致任务永久卡住。租约必须有 TTL。

但是租约也不是绝对安全的。网络分区时可能出现:

Worker A 认为自己仍持有租约;
Worker B 认为租约已经过期并接管。

因此,数据库更新还需要 fencing token 或版本号:

UPDATE runs
SET state = ?, version = version + 1
WHERE run_id = ?
  AND lease_owner = ?
  AND version = ?;

如果更新影响行数为 0,说明租约已经失效或状态版本发生变化,Worker 必须停止提交。

3. 租约失效后的正确行为

Worker 检测到租约失效后:

1. 停止启动新的工具调用;
2. 取消当前运行 Context;
3. 不再写入状态;
4. 保留本地日志;
5. 让新 Worker 从最新 checkpoint 恢复。

不能继续“先把当前动作做完再说”,因为当前动作可能已经属于新的所有者。


十、恢复与幂等:至少一次执行是常态

1. 为什么很难保证恰好一次

考虑以下时序:

t1:Runtime 发送 refund 请求
t2:支付服务完成退款
t3:Runtime 进程崩溃
t4:新 Worker 恢复,看到没有 tool.completed
t5:Runtime 再次发送 refund 请求

如果支付服务不支持幂等,可能发生两次退款。

因此,分布式 Agent Runtime 通常只能把工具执行设计为:

至少一次尝试 + 幂等业务语义

而不是简单声称“exactly once”。

2. 幂等键如何生成

幂等键必须稳定地表示同一个逻辑动作:

idempotency_key = hash(run_id + step + tool_call_id)

例如:

func idempotencyKey(runID string, step int, callID string) string {
	sum := sha256.Sum256(
		[]byte(fmt.Sprintf("%s:%d:%s", runID, step, callID)),
	)
	return hex.EncodeToString(sum[:])
}

同一个 run_id、逻辑步和工具调用 ID 重试时,得到同一个键。

但要注意:如果模型重新生成了一个新的 tool_call_id,它可能被误认为是新的动作。对于有副作用的工具,Runtime 应额外根据业务主键判断是否已经存在相同意图。

3. 工具必须提供查询接口

对于不可逆副作用,最好设计成:

submit_action(idempotency_key, request)
query_action(idempotency_key)

恢复时:

如果只有 tool.started,没有 tool.completed:
    先 query_action
    如果已完成:写入 completed
    如果不存在:才允许重新 submit
    如果状态未知:进入人工确认或补偿流程

这比盲目重试安全得多。


十一、人工确认和暂停恢复

有些工具不能由模型直接执行:

退款
删除数据
发布代码
发送外部邮件
修改生产配置

运行状态需要进入:

waiting_approval

并持久化待确认动作:

{
  "run_id": "R1",
  "status": "waiting_approval",
  "pending": [
    {
      "tool": "refund",
      "arguments": {"order_id": "A1001"},
      "reason": "该操作会产生资金副作用"
    }
  ]
}

人工确认后,恢复请求应包含:

{
  "run_id": "R1",
  "decision": "approve",
  "approved_by": "operator-17",
  "expected_version": 12
}

Runtime 必须检查:

  1. 运行仍然处于 waiting_approval
  2. 版本仍然是预期版本;
  3. 待确认工具参数没有被篡改;
  4. 操作者有权限;
  5. 确认是否已经处理过。

OpenAI Agents 文档也将中断运行视为一种可恢复状态:运行可能没有最终输出,但会返回待处理的中断项和可继续使用的状态快照。(developers.openai.com)


十二、确定性:不是让模型每次都输出一样

“确定性恢复”不是要求模型输出完全相同,而是要求 Runtime 对已经发生的事实保持一致。

需要固定的部分包括:

1. 事件序号;
2. 工具调用 ID;
3. 工具结果;
4. 消息排序;
5. 状态版本;
6. 幂等键;
7. 重试次数;
8. 随机种子或采样配置;
9. 工具返回的时间和外部版本信息。

需要重新决策的部分则要明确标记:

模型调用失败且没有响应
工具临时失败
人工拒绝后要求模型重新规划

一个合理的恢复策略是:

如果模型响应已经持久化:
    不重新调用模型,直接重放响应。

如果工具结果已经持久化:
    不重新执行工具,直接恢复状态。

如果只持久化了 tool.started:
    查询幂等状态,再决定恢复或重试。

如果模型请求尚未产生可验证响应:
    可以重新调用模型,但要记录 retry_attempt。

反例是每次恢复都从完整消息历史重新请求模型。即使输入相同,模型也可能因为采样、服务端版本、工具结果时间或上下文压缩不同而产生不同动作。


十三、完整组件关系和故障路径

flowchart TD
    U[用户请求] --> R[Runtime]
    R --> C[Context 截止时间/取消]
    R --> S[加载 State 或 Checkpoint]
    R --> M[调用 Model]
    M --> D[Model Decision]
    D --> V[策略/权限/Schema 校验]
    V -->|最终答案| E1[写入 completed 事件]
    V -->|工具调用| Q[工具调度器]
    V -->|人工确认| E2[写入 waiting_approval]
    Q --> T1[工具 A]
    Q --> T2[工具 B]
    T1 --> F[工具反馈]
    T2 --> F
    F --> L[追加 Event Log]
    L --> CP[保存 Checkpoint]
    CP --> R
    E1 --> OUT[流式事件/最终响应]
    E2 --> OUT
    R --> LEASE[租约续期]
    LEASE --> STORE[(持久化存储)]

故障路径应明确:

模型超时:
    取消模型请求 → 判断是否可重试 → 保留运行状态

工具超时:
    取消工具 Context → 查询幂等状态 → 重试或暂停

Worker 崩溃:
    租约过期 → 新 Worker 获取租约 → 加载 checkpoint → 重放未完成事件

客户端断开:
    取消展示流 → 根据运行类型决定是否取消 Agent

数据库写入失败:
    不宣布工具完成 → 保留本地错误 → 由恢复机制继续

状态版本冲突:
    放弃当前提交 → 重新加载最新状态 → 判断自己是否已失去租约

每条路径都要避免一个错误假设:网络调用返回错误,不等于远端没有执行。


十四、生产诊断:先看状态,再看模型

Agent 失败时,不应只打印最终错误。至少要记录以下字段:

run_id
tenant_id
trace_id
step
state_version
event_seq
model_request_id
tool_call_id
tool_name
idempotency_key
attempt
lease_owner
deadline
error_kind

诊断顺序建议是:

1. 判断运行是否还活着

status = running       是否长时间不变?
lease_expire           是否已经过期?
last_event_at          最后事件是什么?

2. 判断卡在哪一层

等待模型:
    检查模型请求、网络和截止时间

等待工具:
    检查 tool.started 后是否有 completed/failed

等待人工:
    检查 approval 是否已经写入

等待持久化:
    检查事务、连接池和版本冲突

3. 判断是否可能重复执行

如果存在:

tool.started
但没有
tool.completed

不能直接重试有副作用工具。必须先查询幂等键对应的外部状态。

4. 判断是模型错误还是 Runtime 错误

例如:

模型选择了不存在的工具

可能是模型输出或工具描述问题。

而:

工具已成功执行,但 Runtime 恢复后重复执行

则是幂等或持久化设计错误,不能通过修改提示词解决。


十五、常见误解和边界

误解一:Agent Runtime 就是一个 for 循环

for 循环只能表达控制流程,不能自动解决:

  • 状态持久化;
  • 取消传播;
  • 并发安全;
  • 工具副作用;
  • 崩溃恢复;
  • 租约;
  • 审计;
  • 流式背压。

循环是核心,但不是完整 Runtime。

误解二:流式发送出去的内容就是最终结果

流式增量只是传输事件。最终结果必须由完整模型输出或完成事件确认。否则客户端看到“退款成功”时,系统可能只是输出了模型的计划。

误解三:Context 超时后,远端动作自动取消

Context 只能取消本地调用链,是否能取消远端动作取决于协议和服务端实现。支付、消息发送、订单创建等动作必须通过幂等键和状态查询处理不确定结果。

误解四:并发越高,Agent 越快

并发只能减少相互独立等待的时间。它会增加:

  • 外部服务限流概率;
  • 结果排序复杂度;
  • 状态提交冲突;
  • 成本峰值;
  • 资源隔离难度。

读操作可以并发,写操作必须根据业务依赖和幂等能力决定。

误解五:保存完整消息历史就能恢复

消息历史不能替代:

  • 工具执行记录;
  • 状态版本;
  • 待确认动作;
  • 租约信息;
  • 幂等键;
  • Checkpoint 位置。

它只是恢复所需数据的一部分。


十六、一个可落地的最小生产边界

一个小型但可靠的 Agent Runtime,至少应具备以下结构:

Runtime
├── Model Adapter
├── Tool Registry
├── Policy Validator
├── Context / Deadline Manager
├── Sequential State Reducer
├── Concurrent Tool Executor
├── Event Store
├── Checkpoint Store
├── Lease Manager
├── Stream Publisher
└── Recovery Worker

执行时遵循以下不变量:

不变量 1:
没有持久化的工具调用,不能被认为已经执行。

不变量 2:
没有工具结果,不能把工具调用视为完成。

不变量 3:
没有满足版本和租约条件,不能提交状态。

不变量 4:
有副作用的工具重试前,必须查询或使用幂等语义。

不变量 5:
Context 取消后,不再启动新的动作。

不变量 6:
流式展示失败,不应自动等同于运行失败。

不变量 7:
达到最大步骤数、截止时间或策略限制时,必须进入明确终态。

Agent 的灵活性来自模型能够动态选择下一步;系统的可靠性来自 Runtime 对状态、工具、Context、并发和持久化施加边界。前者决定任务能否处理开放问题,后者决定任务能否在超时、断连、崩溃和重复投递之后仍然保持可解释、可恢复和不重复产生副作用。


系列导航与关联阅读

官方资料

本文依据 Agent、模型、协议与框架官方资料重新梳理;正文、示例与生产清单由 WR BLOG 编写。