OpencodeServeClient 设计解析:Onyx Craft 基于opencode serve传输层的沙箱 Agent 客户端
【免费下载链接】danswerOpen Source AI Platform - AI Chat with advanced features that works with every LLM项目地址: https://gitcode.com/GitHub_Trending/da/danswer
OpencodeServeClient是 Onyx Craft 中驱动沙箱内 Agent 回合的唯一传输层客户端:它以进程内 Python 客户端的形式封装单个 pod 中的opencode serve实例,将 opencode 的/eventSSE 事件流翻译成 Craft 前端可消费的 SandboxEvent(ACP 协议事件),并内置断线补洞、取消/中止、权限自动应答与成本观测能力。本文基于 设计文档 展开,结合 serve_client.py 与 event_bus.py 的真实实现,讲清它的公开接口、线程模型、事件翻译规则、gap-fill 重连算法、取消路径、认证配置与测试策略,让读者既能读懂设计意图,也能对照源码理解落地细节。
背景与定位:为什么要一个opencode serve客户端
该文档是opencode serve迁移方案的配套设计(迁移方案本身记录于迁移计划,后续 ACP 层清理见 drop-acp-layer.md)。迁移的根本动机是修复ACP terminator-drop bug:旧的 ACP 执行路径会在某些竞态下丢失终止信号,导致回合永远卡住。OpencodeServeClient作为 Phase-1 交付物,在SandboxManager.send_message之后替换掉ACPExecClient/DockerACPExecClient,但对外契约保持不变——send_message仍然返回一个产出 ACP 事件的生成器,因此session/manager.py、scheduled_tasks/executor.py、SSE 编码层与数据包日志等所有调用方都无需改动。
从部署演进看,迁移后的架构中opencode serve是 pod 内长驻进程,由入口点 supervisor 管理;OpencodeServeClient是"一个客户端对应一个 opencode HTTP 目标",每次调用在SandboxManager.send_message内部创建、随调用结束销毁。更完整的部署视角可参考 docker-opencode-serve.md,实际部署中的坑可参考 deploy-gotchas.md 与 event-stream-pitfalls.md。
范围界定(迁移计划覆盖、本文不展开):pod spec、入口点 supervisor、Dockerfile 改动;BuildSession.opencode_session_id持久化列;KubernetesSandboxManager.send_message/DockerSandboxManager.send_message的接线;Phase-2/3/4/5 的滚动上线机制。
公开接口:一个"只讲 HTTP"的薄客户端
设计文档中OpencodeServeClient的核心方法是:
class OpencodeServeClient: def __init__( self, base_url: str, # "http://10.0.0.42:4096" password: str | None, # None 表示开发环境;集群内必填 *, client_info: dict[str, Any] | None = None, timeouts: ClientTimeouts | None = None, ) -> None: ... # 会话生命周期 def health_check(self) -> bool: ... # GET /doc,200 即 True def ensure_session(self, opencode_session_id, *, directory, title=None) -> str: ... def delete_session(self, opencode_session_id, *, directory) -> bool: ... # 承重方法:发送提示词并流式返回事件 def send_message(self, opencode_session_id, message, *, timeout=...) \ -> Generator[ACPEvent, None, None]: ... # 生成器外部的取消入口 def abort(self, opencode_session_id) -> None: ... # 供测试断言 gap-fill 的重连辅助 def list_messages(self, opencode_session_id) -> list[Message]: ...ClientTimeouts是三个命名超时的数据类:
| 超时字段 | 默认值(设计) | 实现中的环境变量与默认值 |
|---|---|---|
connect_timeout | 5s | OPENCODE_SERVE_CONNECT_TIMEOUT,默认5.0 |
request_timeout | 30s | OPENCODE_SERVE_REQUEST_TIMEOUT,默认30.0 |
event_read_timeout | 60s | OPENCODE_SERVE_EVENT_READ_TIMEOUT,默认60.0 |
注意:实现中的默认值从 configs.py 拉取,部署方可通过环境变量调优而无需改动客户端代码。event_read_timeout是/eventSSE 流的空闲超时——超过该时长没有字节到达,读取器就触发重连。
方法语义细节
ensure_session:幂等、可从任意 API 副本安全调用。若传入opencode_session_id,先GET /session/{id}校验存在性:200 直接复用,404 则回退创建;否则POST /session新建。一个关键实现细节是directory参数:opencode-serve 通过每个路由上的?directory=查询参数(Instance.provide中间件)按目录隔离会话存储,body 里的目录字段会被静默忽略。不带directory会把会话创建到服务器启动目录(/workspace),破坏按会话的文件系统隔离。delete_session:尽力而为,失败返回 False 但不影响 Onyx 侧会话删除。send_message:内部流程见下文,对外契约是"产出PromptResponse(或Error)作为回合终结事件"。abort:POST /session/{id}/abort,可安全地与在途的send_message生成器并发——opencode 将入站 abort 视为会话状态翻转,生成器在/event上看到终结信号后产出合成的Error。
内部架构:为什么是"读者线程 + 队列",而不是 asyncio
send_message是一个同步 Python 生成器,调用方(sandbox manager → session manager)在单线程上同步迭代它;但 opencode 的/event是推送式流,后台读取器不可避免。设计文档给出的模型是每次调用一个守护读者线程:
┌─────────────────────────┐ ┌─────────────────────┐ │ caller thread │ │ /event reader │ │ for ev in send_msg(): │ <── ACPEvent ──│ thread (daemon) │ │ yield ev │ via Queue │ - httpx.stream │ │ │ │ - parse SSE │ │ │ │ - translate + │ │ │ │ enqueue │ └─────────────────────────┘ └─────────────────────┘- 每次
send_message调用对应一个queue.Queue[ACPEvent | _ReaderError | _ReaderEnded]; - 读者线程在
send_message内启动、在退出时销毁(成功、出错或GeneratorExit),不越过单次调用存活; - 读者入队前按
sessionID做事件关联——/event是实例级的,同一流上混合多个会话的事件; - 关键防挂起机制:读者在 SSE 连接关闭或看到终结符并干净退出时,向队列放入哨兵
_ReaderEnded(reason);调用方线程每次出队都检查该哨兵,因此读者线程死亡永远不会让调用方无限期挂起。这正是此前"数据包丢失调查"中 Bug A 的修复——在设计层面落地,且永不作为回归目标。
事件流(读者线程内部)
┌─────────────────────────────┐ GET /event ─── SSE chunks ──► │ buffer until "\n\n" │ └──────────────┬──────────────┘ │ one event ▼ ┌──────────────────────────────┐ │ json.loads(data line) │ └──────────────┬───────────────┘ │ evt.properties.info.sessionID ── filter ──┐ │ │ ▼ ▼ ┌──────────────────────────────┐ drop │ translate (see below) │ └──────────────┬──────────────┘ │ ▼ ┌──────────────────────────────┐ │ queue.put(ACPEvent) │ └──────────────────────────────┘为什么不用 asyncio:现有SandboxManager.send_message契约是同步生成器,调用方(FastAPI 同步端点、定时任务 worker)都是同步的。把 asyncio 引入这条路径意味着要改造所有调用方,为一个客户端不值得。httpx.stream+ 守护线程正是现有 ACP 客户端已经使用的模式。
实现演进:从 per-call 读者线程到共享PodEventBus
设计文档中"每次调用一个读者线程"的方案,在最终实现中升级为每 pod 一个长期存活的共享事件总线event_bus.py:PodEventBus维护一条GET /eventSSE 订阅,通过subscribe(session_id)扇出到每个会话的订阅队列(Queue(maxsize=500),满时计数丢弃并告警)。send_message内部改为self._event_bus.subscribe(opencode_session_id)→ 等待stream_ready→POST prompt_async→_consume_from_bus排空订阅队列并翻译事件。这带来几个设计文档没有的收益:
- 跨调用复用连接:同一 pod 内多个并发回合共享一条 SSE 流,不再为每个回合建立新连接;
- 统一重连:总线以 1s→2s→4s…(上限 30s)指数退避重连,连续失败 20 次后自关闭并向订阅者投递
BUS_CLOSED_SENTINEL; - 子代理(subagent)路由:总线维护
_child_to_parent/_parent_to_children映射,把session.created记录的父子关系用于将后代会话事件向上游转发,前端即可在父会话流中看到子代理事件; - 401 自愈:
/event收到 401 时通过reload_auth重新拉取凭据(处理 peer pod 轮换密码的场景)。
事件翻译:opencode /event→ SandboxEvent(ACP 事件)
翻译核心是一个纯函数(无 I/O、无self),因此极易单测:
def translate_opencode_event( raw: dict[str, Any], session_id: str, state: _TurnState, ) -> Iterable[ACPEvent]: """Translate one opencode /event payload into 0..N ACPEvents. Returns an iterable because a single opencode event can imply two ACP events (e.g. a `message.updated` with `time.completed` set both finalizes streaming AND emits PromptResponse). Pure — call it from tests with hand-rolled dicts."""返回迭代器而非单事件,是因为一条 opencode 事件可能蕴含多条 ACP 事件(例如最终的message.updated既要冲刷流又要产出PromptResponse)。实现中该函数位于 serve_client.py,签名进一步演化:新增可选的fetch_message(REST 消息水合回调,解决"delta 先于 message.updated 到达"的竞态)以及parent_resolver/children_resolver/fetch_message_by_session(子代理会话路由)。
映射表(函数的事实来源)
| opencode 事件类型 | 过滤条件 | 产出 |
|---|---|---|
server.connected | 总是 | 无——仅设置"流就绪"标记 |
session.created | 匹配 session_id | 无 |
session.next.agent.switched | 匹配 | 无 |
session.next.model.switched | 匹配 | 无 |
message.part.delta | 匹配,目标 part 角色=assistant,类型=text | AgentMessageChunk(content=TextContent(text=delta)) |
message.part.delta | 匹配,目标 part 角色=assistant,类型=reasoning | AgentThoughtChunk(content=TextContent(text=delta)) |
message.part.updated | 匹配,类型=tool,status=pending,part.id首次出现 | ToolCallStart(...) |
message.part.updated | 匹配,类型=tool,后续出现 | ToolCallProgress(... status=running\|completed) |
message.part.updated | 类型=text | 无(token 流已由 delta 发出) |
message.updated | 匹配,role=assistant,time.completed 非空 | 冲刷缓冲事件,然后PromptResponse(stopReason=...) |
session.idle | 匹配 | 兜底终结符:若尚未产出PromptResponse,立即产出 |
session.status | 匹配,status=idle | 兜底终结符:同上 |
session.error | 匹配 | Error(code=..., message=...) |
permission.asked | 匹配 | 自动允许:POST /session/.../permissions/{id}body{"response": "once"};不向消费者产出任何事件;以 WARN 记录权限与模式并上报指标opencode_unexpected_permission_ask |
permission.replied | 匹配 | 无(仅信息性) |
server.heartbeat | 总是 | 无(或作为SSEKeepalive透传到上层) |
session.diff、session.updated(终结后) | 匹配 | 无 |
| 其他 | — | DEBUG 日志,忽略 |
兜底终结符是对 ACP terminator-drop bug 在 serve 层的纵深防御:Phase 0 的经验数据显示三个终结信号都会可靠触发;代码对先到者产出PromptResponse并忽略其余。
重要实现差异:真正的回合终结信号
实现代码(serve_client.py)对终结语义做了修正:message.updated不是回合终结符——opencode 在每一个 step的 assistant 消息上都会发出带time.completed的message.updated(工具调用 step、文本 step 等),真正的 end-of-turn 信号是session.status: idle/session.idle。因此message.updated在实现中只负责:登记 assistant 消息 id、记录last_finish、产出上下文用量包(ContextUsagePacket),以及当消息携带 error 时产出终结事件。终结统一由_emit_terminator完成,它保证每回合至多产出一次PromptResponse/Error,并正确区分stopReason(opencode 的"stop"映射为 schema 的"end_turn",max_tokens/refusal/cancelled透传)。
工具调用 content 合成(翻译器逻辑)
前端 parsePacket.ts 从content[].type==="diff"读 diff 数据、从content[].type==="content"读文件内容。但 opencode serve 在工具 part 上不产出content数组——只有state.input/state.output/state.metadata。翻译器因此合成content数组,使前端保持零改动。字段名映射以测试报告锁定:
edit工具(state.status到达completed):
content = [ { "type": "diff", "path": state.input["filePath"], "oldText": state.input["oldString"], "newText": state.input["newString"], } ]read工具(state.status到达completed):
content = [ { "type": "content", "content": { "type": "text", "text": state.output, # opencode 返回带行号字符串 }, } ] # 前端的 extractFileContent 通过 /^\d+\| /gm 正则剥掉行号——开箱即用。bash与task工具:无需合成content,前端直接读rawOutput.output(见下方字段映射的raw_output行)。
raw_input/raw_output字段名映射(保证前端现有getRawInput/getRawOutput助手零改动):
| ACP 字段 | Opencode 来源 | 翻译器动作 |
|---|---|---|
raw_input | state.input(已是 camelCase,与前端的filePath/oldString等回退链一致) | 原样透传 |
raw_output | state.output(纯字符串或对象) | 字符串则包裹为{"output": state.output};dict 则原样透传 |
工具名 → ACPtitle和kind的推导,实现中用两个小查表_TOOL_KIND/_TOOL_TITLE完成(serve_client.py),与前端的NAME_MAP/TOOL_KIND_MAP对齐,覆盖bash、read、write、edit、patch、apply_patch、glob、grep、list、task、todowrite、webfetch、websearch及 opencode 1.15.x 新增的lsp、skill、question、invalid;未知工具回退为kind="other"、title "Running tool"。opencode 的工具状态值(pending/running/completed/error/cancelled)通过_TOOL_STATUS_MAP归一为 Onyx schema 的pending/in_progress/completed/failed。
为什么需要 per-turn 状态对象
三件事需要跨事件状态,由_TurnState(serve_client.py)承载:
ToolCallStart是"part.id 首次出现"——opencode 会对同一工具 part 发出多条message.part.updated(随state.status变迁)。用seen_tool_calls: set[str]记录已见过的工具 part id 即可。- 幂等终结符——一旦产出
PromptResponse,后续任何兜底终结信号都是 no-op。一个布尔terminator_yielded即可。 - 按文本 part 的累计器用于补洞——
local_text: dict[str, str]记录partID → 已产出的累计文本,供message.part.updated上的 gap-fill 对账使用。
实现中的_TurnState比设计文档更丰富:还包含assistant_message_ids/user_message_ids(角色缓存,避免对每条 delta 重复发 REST 水合请求)、part_types(part 类型缓存,因为 delta 事件本身不带类型,只有 part id)、task_child_by_call(task 工具 callID → 子会话映射)、child_states(每个子代理会话独立状态)等。凡是不需要跨事件关联的信息一律不进状态。
重连与补洞(Gap-fill):/event断开后不丢一个事件
/event不支持Last-Event-ID(opencode 上游 issue #25657),纯 TCP 重试会丢失断开窗口内的全部事件。原始方案想用GET /session/:id/message做快照,但经验测试证明:快照中的part.text在流式期间是空的,只有回合终结后才填充——对回合中途恢复毫无用处。
可靠的对账点是message.part.updated事件本身:对每个文本 part,opencode 在实时流上至少发出两条——part 创建时(空text)和 part 终结时(text= 完整累计内容),其间还有工具边界触发的中间更新;每条都携带累计的part.text。如果断开窗口内错过了 delta,该 part 的下一条message.part.updated就能通过比较累计长度找回缺失内容。
算法
读者线程内维护local_text: dict[str, str](partID → 已累计文本)。
每次message.part.delta(field == "text"):
local_text[partID] += delta- 产出
AgentMessageChunk(text=delta)
每次message.part.updated(type == "text"):
expected = properties.part.textlocal = local_text.get(partID, "")- 若
len(expected) > len(local):错过了 delta——补发AgentMessageChunk(text=expected[len(local):]),然后local_text[partID] = expected - 若
expected == local:no-op(稳态情形) - 若
len(expected) < len(local):告警并保留local(除非 opencode 回退,否则不应发生;按数据完整性问题处理)
实现将该逻辑抽象为共享函数_reconcile_part_text(serve_client.py),文本 part 与 reasoning part 各自封装为_reconcile_text_part/_reconcile_reasoning_part(分别产出AgentMessageChunk/AgentThoughtChunk)。注意实现中 delta 与对账共用local_text累计器——reasoning 与 text 的 partID 全局唯一,可共用同一字典而无碰撞。
httpx.stream抛错/连接在服务端未关闭时结束:
- 不要用
GET /session/{id}/message做快照——回合中途它帮不上忙; - 退避重连
/event(1s、2s、4s,最多 3 次); - 等待
server.connected; - 在途 part 的下一条
message.part.updated会自动通过上面的对账补洞——无需专门的"补洞模式"。
边界情况一:回合恰好在断开窗口内完成。新流上不会再有本回合事件。重连后静默超过MAX_GAP_WAIT_SECONDS=10,回退到GET /session/{id}/message(终结后它已被完整填充),找到 assistant 消息,把尚未拿到的文本合成一条AgentMessageChunk,再从快照的info.finish产出PromptResponse(stopReason=...)。实现在_post_disconnect_snapshot与list_messages中体现(serve_client.py)。
边界情况二:重连本身失败。3 次尝试后,向队列推入_ReaderError("event stream lost")并退出;调用方线程循环捕获哨兵后产出Error。
gap-fill 逻辑按设计拆为两个纯函数便于单测:_reconcile_text_part()(per-event 钩子)与_post_disconnect_snapshot()(罕见的终结后回退)。
取消路径:三种触发,一个机制
三种不同触发,统一走POST /session/{id}/abort:
- 调用方关闭 SSE 流 →
send_message内产生GeneratorExit:用try/except GeneratorExit:包裹主yield,先 abort 再 re-raise。 - 外部
/cancelAPI 端点(本次迁移新增):直接调用client.abort(session_id)。在途的send_message生成器看到session.status变化后产出Error(或在 opencode 1.15.7 的行为下只是合成的兜底终结符——需在 Phase-2 测试中验证)。 - 生成器内墙钟超时:同一代码路径——
abort,然后产出Error(code=-1),返回。
旧 ACP 路径依赖GeneratorExit传播进cancel()调用,这个职责位于 sandbox-manager 层;现在它移入send_message内部。集中化后定时任务不再需要自己的GeneratorExit管道——直接调用abort即可。
实现中send_message的取消语义进一步细化:GeneratorExit时仅在已成功 POST 过 prompt才发 abort(未发提示词则无需中止);_consume_from_bus还支持可选的should_interrupt回调(约 1 秒轮询一次),让调用方能确定性结束回合——先 abort 再自行产出PromptResponse(stopReason="cancelled"),而不是等待可能永不到来的session.idle,避免被中断且无事件的回合钉住槽位直到超时。
认证与配置
API 服务端所需环境变量:
OPENCODE_SERVE_PORT(默认4096)——与沙箱 Dockerfile 的 EXPOSE 指令对应;OPENCODE_SERVER_PASSWORD_SOURCE——每 pod 密码的读取方式,二选一:- 每个 pod 一个 Kubernetes Secret(与现有
ONYX_PAT模式一致); - 由集群级 Secret + sandbox-id 确定性派生(更省事;pod 环境本身就是 secret 存储,安全边界相同)。
- 每个 pod 一个 Kubernetes Secret(与现有
客户端构造函数接收password——由 sandbox manager 负责获取密码,OpencodeServeClient不关心密码来源。
HTTP 细节(以 configs.py 与 serve_client.py 的实现为准):
Authorization: Basic ${base64(username:password)}。设计文档称 username 默认"onyx",但实现已修正:opencode 的 serve 实现在仅设置OPENCODE_SERVER_PASSWORD时把用户名硬编码为"opencode",任何其他值都会得到 401(已对 opencode 1.15.7 验证)。因此OPENCODE_SERVER_USERNAME = "opencode";/event上Accept: text/event-stream,其余端点Accept: application/json;- POST/PATCH 使用
Content-Type: application/json。
实现中还有两层健壮性细节值得注意:冷 pod 重试(_http_with_cold_pod_retry)——沙箱 pod 已 K8s-Ready 但 opencode-serve 尚未绑定 4096 端口时,ConnectError总是可重试(TCP 拒绝证明服务端从未见到请求,重试不会产生重复副作用),而RemoteProtocolError仅在调用方声明idempotent=True时才重试(服务端可能已处理请求,对非幂等 POST 重试会制造孤立会话——这正是该传输层要避免的 bug);401 密码自愈(_request)——收到 401 时通过reload_password回调重新拉取密码并重建 httpx 客户端(peer pod 轮换密码导致本地缓存失效的场景)。
错误表面化
两层错误:
| 来源 | 客户端如何暴露 |
|---|---|
POST /session/{id}/prompt_async返回非 2xx | 终结符之前产出Error(code=http_status, message=body[:200]);读者线程关闭 |
/event上的session.error事件 | 产出Error(code=-2, message=event.properties.message);若事件同时携带info.time.completed,视为终结符 |
| 读者线程崩溃(httpx 异常、JSON 解析失败等) | 通过_ReaderError哨兵合成Error(code=-3, message="event stream error: {e}") |
| 墙钟超时 | Error(code=-1, message="Timeout waiting for response") |
| 调用方发起的 abort | 不产出——GeneratorExit在POST /abort后传播 |
所有错误事件还会把 opencode 的requestID(若事件中存在)追加进消息,用于与opencode serve日志交叉关联。
实现中的错误码改为语义常量(serve_client.py):TURN_ERROR_CODE_SESSION(会话错误)、TURN_ERROR_CODE_TIMEOUT(超时)、TURN_ERROR_CODE_TRANSPORT(传输层)。超时进一步细分:inactivity 超时(timeout,随回合活动续期)产出ActivityTimeoutError;绝对墙钟超时(absolute_timeout)产出Error(TURN_ERROR_CODE_TIMEOUT, "Turn exceeded maximum duration")。被中止的消息(MessageAbortedError)被正确映射为PromptResponse(stopReason="cancelled"),这样消费者不会把用户打断当成回合失败。
测试策略
外部依赖单元测试(针对真实opencode serve,subprocess 跑在临时目录,位于backend/tests/external_dependency_unit/craft/):
test_serve_client_basic.py——ensure_session、同一会话上连续三次提示词,断言事件有序且每回合恰好一个PromptResponse;test_serve_client_terminator_backstops.py——跑一个回合,用注入的代理删掉message.updated,断言客户端仍能通过session.idle终结且只产出一个PromptResponse(Phase 0 显示该竞态在 serve 上罕见,但兜底是承重设计,必须测);test_serve_client_reconnect.py——回合中途切断/event代理,验证重连 + 补洞产出的最终累计器与未切断的运行一致;test_serve_client_abort.py——发提示词后 100ms 中止,验证生成器产出Error(-1)(或GeneratorExit传播,视取消路径而定),且同一会话的下一次提示词干净启动;test_serve_client_tool_call.py——驱动 bash 工具提示词,断言每个工具 part 恰好一个ToolCallStart,且ToolCallProgress状态循环到completed。
纯函数单元测试(backend/tests/unit/):
test_translate_opencode_event.py——罐装 dict 进、ACPEvents 出,断言完整映射表,包含message.part.delta与message.part.updated的区别,防止未来贡献者回归;test_gap_fill_diff.py——罐装快照 + 罐装"已发出事件",断言合成事件与实时流产出一致。
单元测试是承重的线缆契约锁定;外部依赖单元测试是防止 opencode 升级改变行为的集成安全网。仓库中另有针对 401 密码自愈的单元测试 test_serve_client_401_reload.py。
关键设计决策(2026-05-22 定案)
1. 权限流——Path A(内部自动处理,线缆格式冻结)
OpencodeServeClient在内部处理permission.asked:不上抛给前端,也不向消费者产出RequestPermissionRequest。生产环境中 Onyx 生成的opencode.json已对所用每个工具类别固定*: allow(见 opencode_config.py 的build_opencode_config),权限询问理论上不应发生;一旦发生,说明 opencode 引入了尚未配置的新权限类别——按配置漂移 bug处理。行为:默认响应自动允许(POST /session/.../permissions/{id}body{"response": "once"}),与现 ACP 路径行为一致;遥测上以 WARN 记录权限类型与模式并递增 ERROR 级指标;内部方法_auto_respond_permission不属公开 API。
实现进一步区分了connect_app权限(对应 no-op 的connect_app工具):它走 connect-card 流程——把待决请求存入缓存并向前端公告,由决策端点带外应答,消费循环的超时回退在用户未决定时干净拒绝(_handle_connect_app_permission与_reject_expired_connect_app_permissions)。Path B(真实用户审批 UI)是产品功能而非迁移需求,推迟。
2.OPENCODE_SERVER_PASSWORD来源——每 pod K8s Secret
每个沙箱 pod 获得一个含新生成密码的独立 Secret,以OPENCODE_SERVER_PASSWORD环境变量挂载到sandbox容器。sandbox manager 在provision()期间与现有ONYX_PATSecret 一起生成密码并创建 Secret。不用集群级派生密码的理由:横向移动——若沙箱内 agent 能窃取集群级 secret,就能知道所有沙箱的密码;per-pod 隔离把爆炸半径限制在一个沙箱。且ONYX_PAT已是 per-pod Secret 供给,复用该模式保持 K8s manager 对称;运维开销约 10 行kubernetes.client.V1Secret创建代码,pod 删除级联清理 Secret。
3. 多副本并发——不加锁,客户端处理 409
真实并发路径罕见(双标签页用户;定时任务 vs 用户)。opencode 的session.status: busy状态强烈暗示其prompt_async按会话串行化——要么排队第二个提示词,要么以 409 拒绝。客户端把prompt_async的非 2xx 视为软信号:409 Conflict(会话忙)→ 等待/event流上的下一个session.idle(最多 30s),然后重试一次prompt_async;再失败则暴露为Error。其他非 2xx → 直接产出Error(code=status, message=body)并结束。配套计数器指标opencode_serve_busy_retries用于观察该路径在生产的触发频率;若频繁触发再升级为 Redis 锁——在经验信号出现之前不提前引入复杂度。
4. Token 用量/成本捕获——新的LLMFlow.OPENCODE_TURNspan
终结符message.updated载荷携带成本观测所需的全部数据:
"cost": 0.00107985, "tokens": {"total": ..., "input": ..., "output": ..., "reasoning": ..., "cache": {"read": ..., "write": ...}}, "modelID": "gpt-4o-mini", "providerID": "openai"实现步骤:给 flows.py 的LLMFlow枚举加OPENCODE_TURN;回合开始时经traced_llm_call(flow=LLMFlow.OPENCODE_TURN, model=…, provider=…)开启 generation span(model/provider 来自opencode.json配置,或由首个session.next.model.switched事件填充);终结时设置 span 属性cost、tokens.input、tokens.output、tokens.total、tokens.reasoning、tokens.cache.read、tokens.cache.write并关闭。不做 per-token 延迟的 span 字段——底层 LLM 调用是 opencode 发起的,不是我们;聚合成本/token 即我们拥有的可观测性。这是并行工作项:不阻塞客户端在ACP_TRANSPORT=serve标志后落地,但必须在生产翻转标志前落地,否则过渡期会丢失成本遥测。
实现中该能力以ContextUsagePacket的形式体现在翻译层:message.updated携带的info.tokens(含 cache.read/write)与info.cost被聚合为ContextUsagePacket(used_tokens=..., cost=...)直接产出。
从设计到代码:文档骨架与实现的对应
设计文档给出的代码骨架展示了核心结构(serve_client.py 的__init__与其一致):
# backend/onyx/server/features/build/sandbox/opencode/serve_client.py class OpencodeServeClient: def __init__( self, base_url, password, *, event_bus, client_info=None, timeouts=None ): self._base_url = base_url.rstrip("/") self._auth = httpx.BasicAuth("onyx", password) if password else None self._timeouts = timeouts or ClientTimeouts() # Unary-only client. ``request_timeout`` bounds GET/POST against /session, # /prompt_async, /abort, etc. The long-lived ``/event`` SSE stream lives on # the shared per-pod PodEventBus, which owns its own httpx.stream with # ``event_read_timeout`` — that way the bus's per-frame idle timeout is # not capped by this client's unary read timeout. self._http = httpx.Client( base_url=self._base_url, auth=self._auth, timeout=httpx.Timeout( connect=self._timeouts.connect_timeout, read=self._timeouts.request_timeout, write=self._timeouts.request_timeout, pool=self._timeouts.connect_timeout, ), )设计与实现的差异点(上文已逐一展开):用户名从onyx修正为opencode;per-call 读者线程演进为共享PodEventBus;message.updated从"主终结符"修正为"每 step 信息源",终结交由session.idle/session.status兜底;超时从单一值拆为 inactivity 超时 + 绝对墙钟超时;新增 401 自愈与冷 pod 重试。骨架中"_reader_loop与_consume_until_terminator共同实现死读者快速失败:_consume_until_terminator中每次q.get(timeout=1.0)都检查_ReaderEnded哨兵,若在终结符前看到它则合成Error并返回"——这一"杜绝 15 分钟挂起"的结构性修复,在实现中由PodEventBus的BUS_CLOSED_SENTINEL+_consume_from_bus的显式检查延续,并因unsubscribe的确定性清理而更加健壮。
至此,OpencodeServeClient的设计与实现闭环:对外是一个契约不变的同步生成器客户端,对内以共享事件总线 + 纯函数翻译器 + 累计对账补洞 + 统一中止机制,解决了 ACP 传输层的终止符丢失、断线丢包与跨会话事件路由问题。剩余未决问题(何时翻转传输标志、何时删除 ACP 代码)已不属于客户端库本身,详见 drop-acp-layer.md。
【免费下载链接】danswerOpen Source AI Platform - AI Chat with advanced features that works with every LLM项目地址: https://gitcode.com/GitHub_Trending/da/danswer
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考