Skip to content

Agent 三面两线编排 #18

Description

@SATA260

Agent 三面两线

目标设计,尚未落地。只改运行时编排,不改 HTTP/SSE 契约、前端三层、业务目录或 Git。产品语义(Session / Turn / 审批整单 / 记忆 / Git HTTP)保持不变。


当前 Loop

一次 Run 是进程内长循环:

Start → Worker.Submit(runID) → Execute while
  ├─ Load / Compact / Stream
  ├─ Dispatch
  ├─ 无 tool → terminate
  ├─ 待批 → return waiting_approval
  └─ 事实先落库再 Bus → SSE

问题:

  • 审批后靠 RunStatus + checkpoint 猜该从哪继续。
  • RunStatus 既回答「停在哪」(waiting_approval / completed),又回答「正在干什么」(loading_context / running_llm / executing_tools)。
  • 推进靠 Worker 的 chan string,但粒度是整段 Run,不是一拍。
  • 无状态函数 Load / Stream / Dispatch 被 Runtime 包进一条调用栈,同时做决策、I/O、写库、发事件、循环。

现有问题

问题 表现
循环与阻塞绑在一起 一拍模型流占死 Worker goroutine,没有「一拍作业」单位,无法跨进程重试
恢复靠重跑整段 Execute 取消、Continue、进程启动补领都走同一条 while,用 checkpoint 猜进度而不是读 AgentState
细状态泄漏成产品状态 前端 thinking 认 loading_context / running_llm / executing_tools,Run 同时回答「停在哪」和「正在干什么」
观察与推进职责没写死 Worker 按 runID 重整段执行,没有 (run_id, step_index) 幂等键
决策散落在 if 里 Execute / runModelTurn / runTools 各自判断下一步,Brain 无法单测,闸门插不进 Dispatch 的「复检」与「执行」之间

What

用三面 + 两线替换长循环。

三面

  • 状态面:AgentState 是 Run 的唯一真相。粗状态 + checkpoint + step_index + 取消标志。只有协调器事务能写。
  • 执行面:一次推进单位是 step(state, phase) → StepResult。Brain 只产出指令;call_llmStreamcall_tools_batchDispatch
  • 观察面:沿用 AgentEvent + Bus + SSE。先落库再广播,可回放。

两线

  • 观察总线只通知。
  • 执行总线只投 StepJob(run_id, step_index, phase)
  • 两者逻辑上禁止共用订阅表。插件只挂观察面,不得投作业。
flowchart TB
  Handler[Handler]

  subgraph statePlane [状态面]
    Coordinator[Coordinator]
    AgentState[AgentState]
  end

  subgraph execBus [执行总线]
    StepJob[StepJob]
  end

  subgraph execPlane [执行面]
    Worker[Worker]
    Engine[Engine]
    Brain[Brain]
  end

  subgraph obsPlane [观察面]
    AgentEvent[AgentEvent]
    Bus[Bus]
    SSE[SSE]
  end

  Handler -->|"CreateAgentState ClaimSession Enqueue"| Coordinator
  Handler -->|RequestCancel| Coordinator
  Coordinator -->|Enqueue| StepJob
  StepJob --> Worker
  Worker -->|"TryClaimStep LoadAgentState"| Coordinator
  Worker --> Engine
  Engine --> Brain
  Engine -->|StepResult| Coordinator
  Coordinator -->|CommitStep| AgentState
  Coordinator -->|AppendFact| AgentEvent
  AgentEvent --> Bus
  Bus --> SSE
  Coordinator -->|Next| StepJob
Loading

Run 粗状态收成 queued | running | waiting_human | cancelling | 终态loading_context / running_llm / executing_tools 降为观察事件。审批仍是同一 run_id,不开新 Run。


Why

问题 不变量
循环与阻塞 恢复、审批、取消都变成「写 AgentState + 投作业」,不重入 while;进程可以死,AgentState 不能死
细状态泄漏 Run 只回答「停在哪」;进度只走事件,不在粗状态里画执行细节
观察与推进 观察可丢、执行必须幂等;同一 (run_id, step_index) 只允许一个执行者
决策散落 决策进 Brain,调度留 DispatchDispatch 只回答这批 Call 能不能跑、怎么跑、Result 是什么

How

对外结构与接口

type RunStatus string     // queued | running | waiting_human | cancelling | completed | failed | cancelled
type Phase string        // init | user_input | llm_result | tools_batch_result | human_approved | human_abort | compression_result | error
type InstructionType string // load_context | call_llm | call_tools_batch | request_human_approve | compress_context | finish

// AgentState 是一次 Agent 执行的可序列化状态。
type AgentState struct {
	SessionID       string
	RunID           string
	TurnID          *string
	Status          RunStatus
	StepIndex       int               // 已提交拍号;下一拍必须是 StepIndex+1
	Config          RunConfigSnapshot // 启动后冻结,只读
	CancelRequested bool              // 用户已请求取消
	StopReason      *StopReason       // 终态原因
	ForceFinish     bool              // 超 max turns 后剥 tools 收尾
	Checkpoint      ToolCheckpoint    // 工具恢复点
	PendingApproval *string           // 未裁定审批 id
	StartedAt       *time.Time
	FinishedAt      *time.Time
}

// ToolCheckpoint 记录同一批 tool_call 的执行状态。
type ToolCheckpoint struct {
	TurnID    string
	Completed []string      // 已执行完毕的 tool_call_id
	Approved  []string      // 已批准 tool_call_id
	Denied    []string      // 已拒绝 tool_call_id
	Pending   []tool.Call   // 待执行 tool_call(可能多个)
	Results   []tool.Result // 已产生结果,按原始顺序
}

type StepJob struct {
	RunID     string
	StepIndex int               // 幂等去重键
	Phase     Phase             // 为什么被叫醒
	Payload   json.RawMessage
	Attempt   int               // 本拍重试次数
}

type CallToolsBatchPayload struct {
	Calls       []tool.Call   // 本次全部 tool_call
	Mode        string        // "serial" | "parallel"
	MaxParallel int           // 并行上限
}

type ToolsBatchResultPayload struct {
	Results []tool.Result // 含成功 / 失败 / 拒绝 / 跳过,按原始顺序
}

type Instruction struct {
	Type    InstructionType
	Payload json.RawMessage
}

type Fact struct {
	Type    EventType
	TurnID  *string
	Payload json.RawMessage
}

type StepResult struct {
	State AgentState
	Facts    []Fact
	Next     *StepJob // 只有 running 时才非空
}

type Brain interface {
	Decide(phase Phase, payload json.RawMessage, state AgentState) ([]Instruction, error)
}

// Engine 执行一拍。不写库、不发事件、不调度下一拍。
type Engine interface {
	Step(ctx context.Context, state AgentState, job StepJob) (StepResult, error)
}

// Coordinator 是 AgentState 的唯一写入点,也是两总线唯一扇出点。
type Coordinator interface {
	CreateAgentState(sessionID, triggerMessageID string, mode AgentMode, config RunConfigSnapshot) (runID string, err error) // 创建 queued 的 AgentState
	ClaimSession(sessionID, runID string) error                          // 标为 Session 的 active Run
	Enqueue(job StepJob) error                                           // 向执行总线投递作业
	TryClaimStep(runID string, stepIndex int) (ok bool, err error)       // 互斥领取一拍
	LoadAgentState(runID string) (AgentState, error)                         // 加载 AgentState
	AppendFact(runID string, fact Fact) (AgentEvent, error)              // 拍内写事实并发布事件
	CommitStep(runID string, result StepResult) error                    // 提交 AgentState 并扇出后续作业
	RequestCancel(runID string) error                                    // 请求取消 Run
	RecoverActive() error                                                // 启动时补投可继续作业
	DequeueNext(sessionID, finishedRunID string) error                 // 终态后领取下一个 queued Run
}

// Worker 消费执行总线。
type Worker interface {
	Start(ctx context.Context) error // 启动消费者
	Submit(job StepJob) error        // 投递作业;队列满则失败,不丢 AgentState
	Cancel(runID string)             // 取消 Run 的飞行拍
}

推进一拍

Enqueue(StepJob)
  → TryClaimStep
  → LoadAgentState
  → Brain.Decide(phase, state) → Instruction
  → Executor

call_llm:流内每条 delta 走 AppendFact(不推进 step_index),流结束 CommitStep,若仍 runningEnqueue(phase=llm_result)

call_tools_batch:整批复检;有一条需审批则整批不执行,返回 WaitingApproval。否则按 Mode 串行或 MaxParallel 并行执行,结果按原始顺序落 ToolsBatchResultPayload

分支:

  • WaitingApproval → 写 Approval + checkpoint → waiting_human → 停投。
  • 否则 Results 写入 AgentState → 仍 runningEnqueue(phase=tools_batch_result)

CommitStep 原子校验粗状态边与 step_index,写 Run / Turn / Message / checkpoint,递增 last_event_seq;提交后再 Publish 观察事件,并视 Next 投执行作业。

Brain 决策表:

phase 指令
user_input call_llm
llm_result 无 tool finish
llm_result 有 tool call_tools_batch
tools_batch_result call_llm
human_approved 带 checkpoint 再 call_tools_batch
取消 finish(cancelled)

一拍一条主指令。禁止一拍里「模型 → 工具 → 模型」。Dispatch 仍是无状态入口;需要插 BeforeToolCall 时再拆成 InspectCall / RunPrepared。

用户动作流程

// 用户发起会话
Session.Create(...)

// 用户输入文本
Message.Insert(...)
Run.Start(...)
Session.ClaimActiveRun(...)
Event.Append(run.created)
Coordinator.Enqueue(StepJob{Phase: user_input})

// 执行:user_input 拍
TryClaimStep(...)
LoadAgentState(...)
Brain.Decide(user_input, state) → call_llm
Context.Load(...)
Context.CompactIfNeeded(...)
Context.Build(...)
Model.Stream(...)
  Event.Append(assistant.delta)
Message.Insert(...)
Memory.IndexMessage(...)
CommitStep(...) // 若还有 tool 则投下一拍

// 执行:llm_result 拍(有 tool)
Brain.Decide(llm_result, state) → call_tools_batch
Tool.Dispatch(...)
  Event.Append(tool.call_started)
  Event.Append(tool.execution_started)
  Event.Append(tool.execution_result)
Approval.Create(...)
Event.Append(tool.approval_required)
CommitStep(...) // 状态 waiting_human

// 用户审批
Approval.Decide(...)
Coordinator.RecordToolDecisions(...)
Event.Append(tool.approval_decided)
Coordinator.Enqueue(StepJob{Phase: human_approved})

// 执行:human_approved 拍
Tool.Dispatch(...)
Message.Insert(...)
Memory.IndexMessage(...)
Event.Append(tool.execution_result)
CommitStep(...) // 状态 running,投 tools_batch_result

// 执行:tools_batch_result 拍
Brain.Decide(tools_batch_result, state) → call_llm

// 用户取消
Run.Cancel(...)
CancelInFlight(...)
Event.Append(run.cancelled)
Session.ClearActiveRun(...)
Coordinator.DequeueNext(...)

// 继续 / 重试
Run.Continue(...)
Run.Retry(...)

// 打开会话或重连
Event.Subscribe(...)

关键约束

  • 同一 Session 至多一个非终态 active Run(waiting_human 仍占用)。
  • 同一 Run 至多一拍在飞。
  • RunConfigSnapshot 启动后只读。
  • 插件事件不是 AgentEvent
  • 观察总线不得 Submit

落地范围

分三步,每步可单独上线:

  1. Execute 收成一拍,拍末自投;waiting_approval 仍退出。观察总线不动。
  2. 细状态写入时映射为 runningCanTransition 换成粗状态边。
  3. Worker 从 chan string 改为 chan StepJob

不含 Redis 化总线、子 Agent、把记忆或 Git 改成指令。

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    moduleSingle-module objects, interfaces, and design

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions