从零搭建 Agent Harness 系列(十五)Task、Runtime 与多 Session 并行

前一篇我们把 Reporter 和 Runtime 从 Agent Engine 中拆了出来。Engine 负责执行 Agent Loop,Runtime 负责承载一次会话,Reporter 负责向外部渠道发布事件。

但如果只停留在这一层,Runtime 仍然只是一个可以被调用的对象。一个真正可用的 Agent Harness 还需要回答三个问题:

1
2
3
一次任务如何被启动、等待、取消和查询?
一个 Session 如何避免多个任务同时修改上下文?
一个 Go 进程如何同时承载多个独立对话?

这篇文章记录今天在 go-tiny-claw 中完成的三个方向:

1
2
3
Task:一次 Agent Run 的生命周期
RuntimeManager:多个 Runtime 的注册表
Multi-Session:同一个进程承载多个独立会话

一、为什么 Runtime 还需要 Task

最简单的 Runtime 调用方式是:

1
err := runtime.Run(ctx, prompt, reporter)

调用者只能同步等待结果。如果需要支持任务运行时查询状态、用户按 Ctrl-C 取消任务、多个调用者等待同一个任务,以及渠道连接关闭时取消任务,就需要把一次运行显式建模为 Task。

Task 表示一次具体的 Agent 执行,而不是一个 Session:

1
2
3
4
5
6
7
8
Session
└── 保存长期对话历史

Runtime
└── 管理一个 Session 的运行边界

Task
└── 表示一次 Prompt 的执行过程

因此 Runtime 的调用方式变成:

1
2
3
4
5
Runtime.Start
└── 立即返回 Task

Task.Wait
└── 等待 Agent Engine 执行完成

二、Task 的状态模型

当前 internal/runtime/task.go 中定义了四种状态:

1
2
3
4
5
6
7
8
type TaskStatus string

const (
TaskRunning TaskStatus = "running"
TaskCompleted TaskStatus = "completed"
TaskCanceled TaskStatus = "canceled"
TaskFailed TaskStatus = "failed"
)

Task 的核心字段是:

1
2
3
4
5
6
7
8
9
type Task struct {
id string
done chan struct{}
cancel context.CancelFunc

mu sync.RWMutex
status TaskStatus
err error
}

这里的 done 不是传递错误结果的 Channel,而是一个完成通知:

1
2
任务未完成:done 保持打开
任务完成:close(done)

多个调用者因此可以安全等待同一个任务:

1
<-task.Done()

因为读取的是关闭事件,而不是从 Channel 中消费一个错误值。Channel 关闭后,所有等待者都可以被唤醒。

Wait 的实现是:

1
2
3
4
5
6
7
8
func (t *Task) Wait() error {
<-t.done

t.mu.RLock()
defer t.mu.RUnlock()

return t.err
}

状态和错误由读写锁保护:

1
2
3
4
5
6
7
8
9
10
11
12
13
func (t *Task) Status() TaskStatus {
t.mu.RLock()
defer t.mu.RUnlock()

return t.status
}

func (t *Task) Err() error {
t.mu.RLock()
defer t.mu.RUnlock()

return t.err
}

任务结束时只允许从 running 转换一次:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
func (t *Task) finish(err error) {
t.mu.Lock()
defer t.mu.Unlock()

if t.status != TaskRunning {
return
}

t.err = err

switch {
case err == nil:
t.status = TaskCompleted
case errors.Is(err, context.Canceled):
t.status = TaskCanceled
default:
t.status = TaskFailed
}

close(t.done)
}

这里必须先写入状态和错误,最后关闭 done。这样所有被唤醒的等待者都能读取完整状态。

三、Runtime 保证一个 Session 内部不并发

Task 解决了单次任务的生命周期,但 Runtime 仍然需要限制同一 Session 的任务数量。

Runtime 保存当前活动任务:

1
2
3
4
5
6
7
type Runtime struct {
runner AgentRunner
session *ctxpkg.Session

mu sync.Mutex
active *Task
}

启动任务时先检查 active:

1
2
3
if r.active != nil {
return nil, ErrTaskRunning
}

任务结束后清理活动任务:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
func (r *Runtime) execute(
ctx context.Context,
task *Task,
rep reporter.Reporter,
) {
err := r.runner.Run(ctx, r.session, rep)
task.finish(err)

r.mu.Lock()
if r.active == task {
r.active = nil
}
r.mu.Unlock()
}

这产生了一个重要的并发边界:

1
2
同一个 Runtime:同一时间只能有一个 Task
不同 Runtime:可以并行运行不同 Task

这不是为了禁止整个系统并发,而是为了保护一个 Session 的消息历史不被两个 Agent Loop 同时写入。

四、为什么需要 RuntimeManager

一个 Runtime 只能管理一个 Session。如果系统只使用一个 Runtime,就只能支持一个对话。

因此增加一个进程级 Manager:

1
2
3
4
5
6
type Manager struct {
mu sync.RWMutex
factory *RuntimeFactory
runtimes map[string]*Runtime
creating map[string]struct{}
}

Manager 的职责不是运行 Agent,而是管理 Runtime 实例:

1
2
3
4
5
6
Create:创建并注册 Runtime
Add:手动注册 Runtime
Get:查询 Runtime
List:列出 Session ID
Count:查询 Runtime 数量
Destroy:取消任务并删除 Runtime

创建流程是:

1
2
3
4
5
6
7
连接建立

Manager.Create(sessionID)

RuntimeFactory 创建依赖

Manager.Add(sessionID, runtime)

creating 用来防止同一个 Session ID 被多个 Goroutine 同时创建:

1
2
3
4
5
6
7
8
9
if _, exists := m.runtimes[sessionID]; exists {
return nil, ErrRuntimeExists
}

if _, exists := m.creating[sessionID]; exists {
return nil, ErrRuntimeCreating
}

m.creating[sessionID] = struct{}{}

Factory 创建结束后释放占位:

1
2
3
4
5
defer func() {
m.mu.Lock()
delete(m.creating, sessionID)
m.mu.Unlock()
}()

销毁 Runtime 时,Manager 先从 Map 中删除,再取消活动任务:

1
2
3
4
delete(m.runtimes, sessionID)
m.mu.Unlock()

agentRuntime.Cancel()

不在持有 Manager 锁时调用 Cancel,可以避免 Manager 锁和 Runtime 锁形成嵌套等待。

五、RuntimeFactory 解决什么问题

如果每次建立连接时都在 Server 中手动创建 Registry、Approval、Engine、Session 和 Reporter,Server 很快会变成依赖组装代码的集合。

RuntimeFactory 统一创建每个 Runtime 所需的依赖:

1
2
3
4
5
6
7
8
RuntimeFactory
├── Session
├── Tool Registry
├── Approval Handler
├── GrantStore
├── Approval Gate
├── AgentEngine
└── Reporter

每个独立连接都应该拥有自己的 Runtime、Session、Registry、Approval Handler、GrantStore 和 Reporter。

Provider 可以由多个 Runtime 共享,但 Session、Reporter 和审批状态不能错误共享。

六、一个 Go 进程承载多个 Session

单机 CLI 的入口仍然是:

1
2
3
cmd/claw/main.go
└── os.Stdin + os.Stdout
└── 一个 REPL

多连接 Server 的入口则是:

1
2
3
cmd/claw_server/main.go
└── TCPServer
└── Accept 循环

TCP Server 每接收一个连接,就启动一个 Goroutine:

1
2
3
4
5
6
7
8
9
10
11
12
13
for {
conn, err := s.listener.Accept()
if err != nil {
return err
}

sessionID := fmt.Sprintf(
"tcp-channel-%d",
atomic.AddUint64(&s.sequence, 1),
)

go s.handleConnection(ctx, sessionID, conn)
}

一个连接对应一个 ChannelSession:

1
2
3
4
5
6
7
8
9
10
type ChannelSession struct {
id string
conn io.ReadWriteCloser
manager *runtimepkg.Manager
runtime *runtimepkg.Runtime
repl *cli.REPL

closeOnce sync.Once
closeErr error
}

它负责把网络连接和 Agent Runtime 绑定起来:

1
2
3
4
5
6
7
TCP Connection
├── bufio.Reader
├── io.Writer
├── ChannelSession
├── Runtime
├── Session
└── REPL

最终一个 Go 进程中的关系是:

1
2
3
4
5
6
7
8
9
10
11
12
Go Process
├── TCP Connection A
│ ├── Runtime A
│ └── Session A

├── TCP Connection B
│ ├── Runtime B
│ └── Session B

└── TCP Connection C
├── Runtime C
└── Session C

连接 A 的任务被取消,不应该影响连接 B:

1
2
3
4
5
6
Task A.Cancel()
└── Runtime A
└── Session A

Runtime B 继续运行
Session B 历史不变

七、连接生命周期

ChannelSession 的关闭必须同时处理两类资源:

1
2
3
4
5
Runtime
└── 取消当前 Task,并从 Manager 删除

Network Connection
└── 关闭 Reader 的底层连接,解除阻塞

因此 Server 退出时不能只关闭 Listener,还需要关闭所有活动 ChannelSession:

1
2
3
4
5
Server Shutdown
├── 关闭 Listener,阻止新连接
├── 遍历活动 ChannelSession
├── Destroy Runtime
└── Close TCP Connection

这也是多连接服务和一次性 CLI 最大的区别:除了任务本身,还需要管理连接资源和进程级退出。

八、测试应该验证什么

Task 和 Manager 的单元测试覆盖了:

1
2
3
4
5
6
任务完成
任务失败
任务取消
同一 Runtime 拒绝并发 Task
多个 Runtime 并发添加
Runtime 创建和销毁

多 Session 还需要补充 Channel 和 Server 集成测试:

1
2
3
4
5
6
ChannelSession 创建后 Manager 数量增加
ChannelSession Close 后 Runtime 被删除
重复 Close 不报错
两个 TCP 客户端可以同时连接
关闭客户端 A 不影响客户端 B
Server 关闭后所有连接和 Runtime 都退出

测试不应该调用真实 LLM Provider,而应该使用 fake Provider:

1
2
3
4
5
6
7
8
9
type fakeProvider struct{}

func (p *fakeProvider) Generate(
ctx context.Context,
messages []schema.Message,
tools []schema.ToolDefinition,
) (*schema.Message, error) {
return &schema.Message{}, nil
}

并发测试使用竞态检测:

1
GOCACHE=/tmp/go-tiny-claw-gocache go test -race ./...

九、这一阶段之后

完成 Task、RuntimeManager 和 TCP 多连接后,Agent Harness 已经从“一次性 CLI”进入“进程级多会话运行时”。

但这还不是完整的生产级网络服务,后续还需要:

1
2
3
4
5
6
7
8
9
10
11
TCP 消息协议
Prompt / Interrupt / Close / Heartbeat 消息类型
连接认证
TLS
读取和写入超时
最大消息长度
连接数限制
空闲连接清理
Session 持久化
Runtime 状态查询接口
WebSocket Channel

当前几个对象的边界是:

1
2
3
4
5
6
Task:一次执行
Session:一段长期对话
Runtime:一个 Session 的运行时边界
Manager:进程内多个 Runtime 的管理者
ChannelSession:一个外部连接与 Runtime 的绑定
TCPServer:连接接入和生命周期管理

把这些对象分开之后,Agent Engine 不需要知道用户来自本地 Terminal、TCP、WebSocket 还是其他渠道,只需要接收 Context、Session 和 Reporter,完成自己的 Agent Loop。