从零搭建 Agent Harness 系列(十六)Protocol、WebSocket 与多渠道连接层
上一篇我们让一个 Go 进程可以同时承载多个 TCP 对话:每个连接拥有自己的 ChannelSession、Runtime 和 Session,不同 Session 可以并行运行,而同一个 Runtime 内部仍然只允许一个 Task。
但系列十五结尾还留下了几个关键问题:
1 | TCP 客户端和 Server 之间到底传什么格式? |
这一篇继续实现这些能力。本文对应当前 go-tiny-claw 中已经落地的连接层:
1 | JSON Line Protocol |
一、从 TCP 多连接到浏览器控制台
系列十五中的 TCP 连接可以这样使用:
1 | TCP Client |
TCP 是一个字节流。Server 能够读取到字节,但它并不知道这些字节代表什么业务动作。
如果客户端直接发送:
1 | 帮我读取 README |
Server 无法可靠判断这是一条 Prompt,还是一次中断、审批响应或者关闭请求。
因此需要在传输层之上定义应用协议:
1 | TCP / WebSocket |
同时,浏览器还有一个限制:浏览器 JavaScript 不能直接创建原始 TCP 连接。浏览器可以使用 HTTP、WebSocket 等标准 Web 协议,但不能像 Go 客户端一样连接 net.Conn。
所以最终的连接层变成:
1 | TCP Client ─────── TCP :8080 ───────┐ |
TCP 和 WebSocket 使用不同的传输适配器,但进入后面的 ChannelSession、Runtime 和 AgentEngine。
二、先定义连接协议
协议实现位于:
1 | internal/channel/protocol.go |
当前协议采用 JSON Line 形式。每条消息是一个 JSON 对象,并以换行结束。
1. 消息类型
1 | type MessageType string |
客户端发送 Prompt:
1 | {"type":"prompt","content":"请读取 README"} |
客户端请求中断:
1 | {"type":"interrupt"} |
客户端响应审批:
1 | {"type":"approval_response","request_id":"abc123","decision":"allow_once"} |
这里的 type 是连接层协议类型。它和 Reporter 发出的 Agent 事件类型不是同一个概念。
1 | MessageType |
2. Message 结构
1 | type Message struct { |
目前的消息结构仍然比较小,但已经覆盖了连续对话所需的控制面:
1 | Content Prompt 内容 |
未来如果需要协议版本、客户端 ID、任务 ID 和请求追踪,可以继续扩展这个结构,或者引入统一的 Envelope。
3. MessageReader 为什么接收 io.Reader
MessageReader 的构造函数不再要求调用方传入 *bufio.Reader:
1 | func NewMessageReader(input io.Reader) (*MessageReader, error) { |
这个边界很重要:
1 | 外部依赖:io.Reader |
因此以下输入都可以被协议层使用:
1 | net.Conn |
调用方不需要知道协议层是否使用缓冲。
4. Reader 负责协议校验
1 | func (r *MessageReader) Read() (Message, error) { |
协议层至少应该负责:
1 | 限制单条消息大小 |
它不应该负责启动 Agent 或执行工具。协议层只回答“收到了一条什么消息”。
三、MessageWriter 解决并发输出
Server 中一个 Session 可能同时产生多种输出:
1 | 模型文本 Delta |
这些事件可能来自不同 Goroutine。因此不能让多个 Goroutine 直接对同一个连接调用 json.Encoder。
当前 Writer 使用一个互斥锁:
1 | type MessageWriter struct { |
这里有两个设计点。
第一,Writer 接收 io.Writer,不和 TCP、WebSocket 绑定。
第二,Write 接收 any,因为它既可能写入连接控制消息:
1 | Message{Type: MessagePong} |
也可能写入 Agent 事件:
1 | reporter.Event{Type: reporter.EventToolCall} |
它的真实职责不是“只写 Message”,而是“串行写入 JSON 对象”。
同一个连接内,Writer 保证 JSON 字节不会交错;不同连接之间则因为底层连接不同,天然隔离。
同一个连接内,Writer 保证 JSON 字节不会交错;不同连接之间则因为底层连接不同,天然隔离。
四、从 Reporter 回调到结构化事件
AgentEngine 只依赖 Reporter 接口:
1 | AgentEngine |
它不会直接调用 fmt.Printf,也不会直接写 TCP。
1. Event 定义
文件:
1 | internal/reporter/event.go |
当前事件类型包括:
1 | const ( |
事件结构:
1 | type Event struct { |
2. JSONReporter
文件:
1 | internal/reporter/json_reporter.go |
JSONReporter 把 Engine 的回调转换成 Event:
1 | func (r *JSONReporter) OnToolCall( |
最终发送给浏览器的是:
1 | { |
同一个 Engine 可以继续使用 TerminalReporter:
1 | CLI → TerminalReporter → stdout |
Reporter 的抽象让输出渠道可以替换,而不需要修改 Agent Loop。
3. EventSink
文件:
1 | internal/channel/event_sink.go |
1 | type EventSink interface { |
JSONReporter 只依赖 EventSink:
1 | JSONReporter |
这样 Reporter 不需要知道底层是字节流还是 WebSocket Frame。
五、ChannelSession:连接和 Runtime 的绑定层
文件:
1 | internal/channel/channel_session.go |
一个 ChannelSession 的结构是:
1 | 外部连接 |
构造时,所有输出共享同一个 MessageWriter:
1 | writer, err := NewMessageWriter(conn) |
这里不能分别创建两个 Writer:
1 | 错误做法: |
两个 Writer 各自持有锁,无法保护彼此,可能导致输出交错。
正确做法是:
1 | Reporter ─┐ |
1. 消息循环
1 | func (s *ChannelSession) Run(ctx context.Context) error { |
2. Prompt 必须异步启动
如果 MessagePrompt 直接同步等待:
1 | task, _ := s.runtime.Start(ctx, prompt, s.reporter) |
那么消息循环会被阻塞。在 Agent 运行期间,下面这些消息都无法处理:
1 | interrupt |
当前实现启动 Task 后立即返回消息循环:
1 | task, err := s.runtime.Start(ctx, prompt, s.reporter) |
这就是 Server 能够在 Agent 执行期间响应中断和审批的原因。
这就是 Server 能够在 Agent 执行期间响应中断和审批的原因。
六、通道审批如何工作
终端审批可以直接读取 stdin,但 TCP 或 WebSocket 审批不能再启动一个 Reader 去读取同一连接。
否则会变成:
1 | ChannelSession.Reader ──┐ |
两个 Reader 会竞争输入,导致 Prompt 或审批响应被错误消费。
当前 ChannelApprovalHandler 使用一个 pending Map:
1 | type ChannelApprovalHandler struct { |
审批开始时:
1 | Engine |
前端收到:
1 | { |
点击允许一次后发送:
1 | { |
ChannelSession 的唯一输入循环收到响应后,调用:
1 | s.approval.Respond(message.RequestID, decision) |
pending Channel 被唤醒,ApprovalGate 才会把工具调用交给 Engine 后续执行。
审批响应和 Prompt 使用同一个输入循环,这是多渠道审批能够正确工作的关键。
七、为什么需要 WebSocket Adapter
浏览器使用 WebSocket Frame,而当前 MessageReader 使用 JSON Line:
1 | MessageReader 期待:{"type":"ping"}\n |
因此不能直接把 *websocket.Conn 传给 MessageReader。需要一个适配器,把 Frame 转换成 Reader/Writer 认识的字节流。
文件:
1 | internal/server/websocket_conn.go |
读取时,适配器从下一条 WebSocket 消息读取内容,并补充换行:
1 | func (c *websocketConn) Read(p []byte) (int, error) { |
写入时,适配器把一条 JSON Line 转换成一个 Text Frame:
1 | func (c *websocketConn) Write(p []byte) (int, error) { |
适配完成后,ChannelSession 不需要知道自己面对的是 TCP 还是 WebSocket:
1 | TCP Conn ────────────────┐ |
八、WebSocket Server 如何承载多个 Session
文件:
1 | internal/server/websocket_server.go |
TCP Server 和 WebSocket Server 的职责类似:
1 | 监听连接 |
WebSocket Session ID 使用独立前缀:
1 | ws-channel-session-1 |
TCP 使用:
1 | tcp-channel-session-1 |
两个 Server 共享同一个 RuntimeManager,但每个连接仍然拥有独立 Runtime:
1 | RuntimeManager |
启动时,当前 Server 使用两个端口:
1 | TCP :8080 |
这里没有强行把 TCP 和 WebSocket 复用到同一个端口,因为二者的连接握手不同。生产环境可以使用反向代理统一域名,也可以使用连接复用器,但学习阶段使用两个端口更清晰。
九、前端多 Session 控制台
前端目录:
1 | web-console |
技术栈:
1 | Vite |
1. 一个 Session 一条 WebSocket
前端维护一个 Socket Map:
1 | const sockets = useRef<Record<string, WebSocket>>({}) |
创建 Session 时:
1 | const id = `session-${sequence.current++}` |
每个连接的消息事件都带着自己的 Session ID 回到状态更新函数:
1 | socket.onmessage = (message) => { |
因此连接 A 的事件只会更新 Session A 的消息列表。
2. 流式文本聚合
Server 会连续发送:
1 | {"type":"text_delta","content":"第一段"} |
前端在 Session 内记录当前的 streamItemId:
1 | if (session.streamItemId) { |
这样每个 Delta 都会追加到当前 Agent 消息,而不是生成很多气泡。
3. 审批按钮
收到 approval_request 后,前端生成审批卡片:
1 | 需要审批 · write_file |
按钮最终发送:
1 | socket.send(JSON.stringify({ |
前端不直接执行工具,也不决定工具权限。它只是把用户选择传回 Server,真正的 Policy、GrantStore 和 Gate 仍然在 Runtime 内部。
4. 固定底部输入框
连续对话的输入框不能随着消息增长被推到页面之外。这里不是简单地给输入框设置 position: fixed,而是让工作区成为一个高度受控的 Flex 容器:消息区滚动,输入区不参与滚动。
1 | html, |
min-height: 0 很重要。Flex 子元素默认可能按照内容的最小高度撑开父容器,如果省略它,消息列表会把整个页面撑长,输入框仍然会被推走。
十、一次请求的完整时序
把各层连起来后,一条用户消息的生命周期如下:
1 | React Session A |
这里有两个容易混淆的异步边界。
第一,ChannelSession 的读循环不能直接同步执行 Agent。否则 Agent 调用模型、执行工具时,连接就无法继续读取 interrupt 或 approval_response。因此收到 prompt 后要启动任务,读循环继续工作。
第二,Agent 的输出不能直接写 WebSocket。Agent 只依赖 reporter.Reporter,由 JSONReporter 把领域事件转换成协议消息,再由 MessageWriter 串行写入连接。这样 Runtime 不知道当前连接是 TCP、WebSocket 还是未来的飞书。
审批的时序则是:
1 | Agent -> Policy/Gate -> approval_request -> Client |
request_id 是审批闭环的关联键。不能只依赖工具名,因为同一个 Session 中可能同时存在多个工具调用,多个连接也可能并行审批。
十一、如何验证这一阶段
先启动 Server:
1 | cd /Users/smsun/Documents/github/go-tiny-claw |
启动 Web 控制台:
1 | cd /Users/smsun/Documents/github/go-tiny-claw/web-console |
浏览器打开 http://localhost:5173。开发服务器会把 WebSocket 请求转发到 ws://127.0.0.1:8081/ws。打开两个会话后分别发送消息,可以验证:
- 两个会话能够同时运行任务。
- 会话 A 的流式事件不会显示到会话 B。
- 会话 A 按下中断时,只取消会话 A 的 Task。
- 会话 A 等待审批时,会话 B 仍然可以继续对话。
- 审批响应只唤醒对应的
request_id。 - Task 结束后会收到且只收到一个终态事件。
已有的 Go 测试和前端构建可以这样运行:
1 | cd /Users/smsun/Documents/github/go-tiny-claw |
测试重点不是只验证“能不能返回文本”,而是验证连接边界和并发边界:协议读写、审批等待、Session 隔离、Task 终止、事件顺序和竞态安全。
十二、距离生产级还缺什么
这一阶段解决的是“多连接可用”和“浏览器可接入”,还不能直接当作公网生产服务。下一步至少要补齐以下能力。
1. 身份认证与授权
当前连接建立后还没有可靠的用户身份。生产环境需要在握手或首条消息中完成认证,并把 UserID、TenantID、ProjectID 绑定到 Session。Policy 和 GrantStore 也要按用户、项目或租户隔离,不能让一个连接查询到另一个用户的授权。
2. TLS、Origin 和网络保护
WebSocket 生产部署应使用 TLS。服务端不能永久允许任意 Origin,要配置允许的前端域名。同时还需要限制连接数、单条消息大小、并发 Task 数和请求频率,避免一个客户端耗尽进程资源。
3. 超时、心跳和连接清理
需要分别设置握手超时、读超时、写超时、模型调用超时和工具执行超时。ping/pong 不能只作为协议示例,还要用来发现断开的连接。连接断开后必须取消当前 Task,并从 Manager 中清理 Session,避免 goroutine 和状态泄漏。
4. 协议版本和错误模型
协议消息应增加版本字段,并定义稳定的错误结构,例如错误码、可重试标记、关联的 request_id 和 task_id。不能只返回一段无法被前端稳定解析的中文错误字符串。
5. 事件关联与恢复
生产事件至少需要 session_id、task_id、event_id、时间戳和序号。客户端重连后才能根据序号补事件,而不是只能重新发起任务。长任务还需要持久化 Task 状态、审批状态和必要的事件日志。
6. 持久化和水平扩展
现在 Manager、Session 和 GrantStore 主要是进程内内存实现。单进程可以支持多个 TCP/WebSocket 连接,但服务重启后状态会丢失,也无法让多个实例共同管理同一个 Session。生产化需要把 Session、Task、授权和事件分别抽象出持久化接口,必要时使用 Redis、数据库和消息总线。
7. 可观测性和可靠性
日志要携带结构化的 session_id、task_id 和 request_id,指标要覆盖连接数、活跃 Task、模型延迟、工具延迟、审批等待时间、错误率和取消率。还需要补充 TCP/WebSocket 集成测试、断线测试、重复审批测试、并发压力测试和故障恢复测试。
8. 前端生产能力
控制台还需要登录、自动重连、历史消息加载、断线提示、重复发送保护、审批超时提示和权限展示。开发环境的 Vite 代理只解决本地调试问题,生产环境应由反向代理统一提供 HTTPS、WebSocket 转发和静态资源服务。
十三、这一阶段的架构总结
到这里,连接层可以整理成下面的依赖关系:
1 | TCP / WebSocket |
这套分层的核心价值是:
- Transport 只负责连接,不负责 Agent 业务。
- Protocol 只负责消息格式,不负责工具权限。
- Channel 只负责把外部消息翻译成 Runtime 调用,把内部事件翻译成外部消息。
- Runtime 负责 Session、Task、审批和 Agent 生命周期。
- Reporter 负责不同渠道的输出表达。
- Web Console 只是一个客户端,不会反向侵入核心引擎。
因此,增加 WebSocket 并不是复制一套 Agent 逻辑,而是增加一个新的 Transport 和一个新的入口适配器。未来接入飞书、HTTP API、MCP 或 A2A 时,也应沿着同样的方向扩展,而不是把渠道判断散落到 Engine、Tool 和 Task 中。
结语
从最初的单次命令行调用,到现在的 Task、Runtime、Manager、TCP 多连接、JSON Protocol、结构化事件、审批闭环和 WebSocket 控制台,go-tiny-claw 已经从一个 Agent Demo 变成了一个具备连接层和运行时边界的 Harness 原型。
但“能连接”不等于“能生产”。后续工作应优先补齐身份、协议安全、超时和恢复,再做 MCP 工具生态与 A2A 多 Agent 协作。只有先把连接、状态、事件和授权的边界稳定下来,外部协议越多,系统才不会越混乱。