opensre 网关聊天传输层(gateway/transports)架构指南:从入站适配器到统一启停循环
【免费下载链接】opensreBuild your own AI SRE agents. The open source toolkit for the AI era.项目地址: https://gitcode.com/GitHub_Trending/op/opensre
本文围绕 opensre 项目中
gateway/transports/目录的架构约定展开:它如何把 Slack、Discord、Telegram、Buzz 四种聊天平台统一成"一个平台一个包、彼此零依赖"的入站适配器,如何通过单一注册表与统一启停循环(start_transports/stop_transports)接入 AI SRE Agent 的 turn 运行器,以及测试如何以"关注点清单 + known-gaps 台账"的方式钉死每个传输的职责完整性。读完本文,你将掌握该网关在聊天接入侧的分层边界、契约接口与新增传输的标准路径,可直接对照 gateway/transports/AGENTS.md 与下方源码深入排查和二次开发。
一、传输层是什么:聊天平台侧的 ingress 适配器
在 opensre 的网关(gateway)架构里,gateway/transports/是所有入站聊天平台的归属地。用 API 框架的术语讲,每个传输包就是一个ingress 适配器:它负责完成平台侧消息的接收、鉴权、会话解析、turn 输出组装,并把最终的处理交给 turn 运行器(turn runner)。
根据 gateway/transports/AGENTS.md 的定位说明,每个传输包需要实现四条主链路:
- 授权(authorizes):校验消息来源是否可信(如 Slack 的签名校验、Telegram 的 token 校验);
- 解析会话(resolves a session):把平台消息映射到既有的 Agent 会话;
- 构建 turn 输出(builds turn output):把 Agent 的处理结果组装成该平台可发送的回复格式;
- 调用 turn runner:通过统一的 turn 契约把消息交给 Agent 执行。
而这几条链路的"公共零件"全部来自网关的共享层,传输包本身不持有业务分发逻辑:
- 传输实现的 turn 契约是
infrastructure.turn_host.turn_callback与infrastructure.turn_host.turn_output; - 传输注册的注册表是
gateway.transports.startup; - 传输复用的每-turn 共享步骤是
gateway.core.middleware; - 启动所有传输的门面(facade)是
gateway/startup.py。
一句话理解分层:传输包只负责"平台差异",所有"Agent 语义"都收敛到共享层,这样新增一个聊天平台时不会触碰任何 Agent 逻辑。
二、四个传输包与注册表映射
当前网关内置四个聊天传输包,每个包持有自己的 settings、入站 worker、安全逻辑、turn 输出与startup.py:
| 包 | 启动入口 | 运行形态 |
|---|---|---|
slack/ | startup.start_slack_worker | Socket Mode WebSocket 或 Events API HTTP |
discord/ | startup.start_discord_worker | Gateway WebSocket 长连接 |
telegram/ | startup.start_telegram_worker | 长轮询(long-poll) |
buzz/ | startup.start_buzz_worker | 提及轮询(mention-poll) |
这份映射并不是手写的文档,而是注册表源码本身。gateway/transports/startup.py 中定义了唯一的TRANSPORTS元组:
TRANSPORTS: tuple[TransportRegistration, ...] = ( TransportRegistration(TransportName.TELEGRAM, start_telegram_worker, "polling for messages"), TransportRegistration(TransportName.SLACK, start_slack_worker, "inbound connected"), TransportRegistration(TransportName.DISCORD, start_discord_worker, "connected via gateway"), TransportRegistration(TransportName.BUZZ, start_buzz_worker, "polling for messages"), )其中TransportRegistration是 gateway/transports/registration.py 中定义的冻结数据类,一行即一个注册条目:name(传输名)、start(启动函数)、running_status(启动成功后上报的状态文案)。
值得注意的架构细节:Web 不在这个注册表里。按 gateway/transports/names.py 的注释,Web 是一个 channel 但不是 chat transport,因此TransportName枚举只有四个成员:
class TransportName(StrEnum): TELEGRAM = "telegram" SLACK = "slack" DISCORD = "discord" BUZZ = "buzz"TransportName同时充当StartedGateway.transports字典与组件状态映射的键(见 gateway/startup.py),从而保证状态键与查询不会漂移。
三、统一启停循环:start_transports / stop_transports
注册表与 worker 的启停循环都住在 gateway/transports/startup.py——它是这个包里唯一允许导入 peer 模块的文件(且只允许导入各 peer 的startup子模块),这样做保证了"导入一个平台包绝不会连带加载另外三个平台的 SDK 栈"。
3.1 启动:一次遍历、三类结局
start_transports(gateway/transports/startup.py#L64-L97)对注册表做单次遍历,每个传输只会有三种结局:
- 正常启动:把
TransportHandle(name, worker, status)加入handles,状态记为注册表中的running_status; - 未配置:抛出
GatewayConfigurationError(如缺少凭据),记录为not configured (…)并跳过; - 启动失败:抛出
GatewayTransportFailedError(如 Discord readiness 超时),记录为failed (…)并跳过。
for registration in TRANSPORTS: try: worker, _settings = registration.start(logger=logger, handler=handler) except GatewayConfigurationError as exc: logger.warning("%s chat disabled: %s", registration.name.capitalize(), exc) statuses[registration.name] = f"not configured ({exc})" continue except GatewayTransportFailedError as exc: logger.warning("%s chat failed: %s", registration.name.capitalize(), exc) statuses[registration.name] = f"failed ({exc})" continue ...两个关键语义:
- 所有传输在同一趟遍历中启动,没有先后顺序依赖;某个传输失败/未配置时,网关继续服务其余成功启动的传输;
- 没有凭据 = 跳过而非报错,这是刻意的设计,让同一份网关代码既能在只配了 Telegram 的环境运行,也能在四个平台全配齐的环境运行。
ChatStartup返回值同时携带handles(已启动的 worker 列表)与statuses(所有尝试过的传输的状态),这样调用方无需深入状态映射就能上报"哪些没配置"。
3.2 停止:共享关闭预算
stop_transports(gateway/transports/startup.py#L100-L117)用ShutdownBudget对超时做统一管理:
- 逐个请求停止,即使某个失败也继续尝试其余——一个卡住的传输不能拖住其他传输;
- 超时预算按顺序共享:停止第一个 worker 花费的时间会从剩余预算中扣除,
budget.mark()/budget.consume(started)精确记账,避免累计超时。
budget = ShutdownBudget(timeout) stopped = True for handle in handles: started = budget.mark() stopped = handle.worker.stop(timeout=budget.remaining) and stopped budget.consume(started) return stopped默认停止超时取自config.constants.gateway.DEFAULT_STOP_TIMEOUT_SECONDS。
worker 的契约定义在 gateway/transports/worker.py:TransportWorker是一个仅含stop(*, timeout) -> bool的 Protocol;TransportStarter则是"入参不定、返回(worker, settings)二元组"的 Callable。这个"每个传输的 startup 返回 worker 加自有 settings 对象"的约定,让组合根能够持有解析后的配置而不必了解平台细节。
四、每个传输的 startup:平台差异收口处
AGENTS.md 明确规定"传输特定工作(settings 加载、Discord readiness 等待)留在各包的startup.py"。四个包的实现正好展示了三种典型形态。
4.1 Slack:双入站模式 + 单副本去重护栏
gateway/transports/slack/startup.py 支持两种入站传输,由SlackInboundTransport枚举选择,_TRANSPORT_STARTERS映射表"加一行即加一种模式,而不是加一个分支":
- Socket Mode(
_start_socket_mode):持有一条 WebSocket,适合无法开放公网端口的部署; - Events API HTTP(
_start_events_api_http):在自有端口上提供 HTTP 服务并接收带签名的 POST 请求。
其中_submit_turn值得注意:Slack 要求路由在 3 秒内应答,因此 turn 不在请求线程内执行,而是提交到共享执行器(stack.executor.submit(stack.dispatcher.dispatch, message)),请求快速返回,Agent 在后台处理。
Events API HTTP 模式还有一个强制的去重护栏(_build_handled_event_repository):Slack 的投递是"至少一次",重试落在另一副本上会被重复受理,因此该模式要求共享事件存储:
- 配置了
DATABASE_URL→ 使用共享的HandledSlackEventRepository; - 未配置 → 默认抛出
GatewayConfigurationError,提示设置DATABASE_URL或显式设置SLACK_GATEWAY_ALLOW_LOCAL_DEDUP=1接受"仅单副本安全的进程内去重"; - 即使选择了进程内回退,也会输出 warning 日志,保证这不是静默决定。
4.2 Discord:readiness 等待后统一返回
gateway/transports/discord/startup.py 的差异点是启动后要等待 Gateway 就绪:worker.wait_until_ready(timeout=settings.startup_timeout_seconds),超时则停止 worker 并抛GatewayTransportFailedError("startup timeout")。这样统一注册表可以把 Discord 与 Telegram/Slack 同等对待——start 返回时必然是活的 worker。
4.3 Telegram 与 Buzz:长轮询/提及轮询
gateway/transports/telegram/startup.py 与 gateway/transports/buzz/startup.py 结构几乎一致:加载平台 settings,创建PollingBackgroundworker,注入initialize_*_polling_runtime/shutdown_*_polling_runtime与 turn 回调。二者的共同点是把"轮询运行时"作为参数注入,传输包本身不持有 Agent/分发逻辑。
五、共享契约:TurnCallback 与 turn 输出
所有传输最终都汇入同一个回调契约。infrastructure/turn_host/turn_callback.py 定义了:
TurnCallback = Callable[[str, SessionCore, TurnOutput, logging.Logger], None]签名固定为(text, session, output, logger):每条入站聊天消息最终都归结为对这一契约的一次调用。契约的反向依赖约束同样严格——TurnCallback所在模块不得 import 任何传输,传输依赖契约、契约不依赖传输,形成单向依赖环。
turn 输出的合作式取消、单消息输出基类等约定(如SingleMessageTurnOutput声明turn_cancel)由infrastructure.turn_host.turn_output提供;而会话解析、入站决策(apply_inbound_decision)、审批(approvals)、注意力(attention)、对话锁(conversation locks)等每-turn 共享步骤全部沉淀在 gateway/core/middleware(含active_turns.py、approvals.py、attention.py、conversation_locks.py、identity_policy.py、inbound_decision.py、terminal_outcome.py)。任何两个传输都需要的能力,按规定必须提升到gateway.core,而不是在某个传输里复制一份。
六、组合根:gateway/startup.py 与 StartedGateway
Web 与聊天传输的组合发生在 gateway/startup.py,它把start_transports的产物与 Web 服务并成一个StartedGateway:
start_gateway先启动 Web 服务器,再启动所有聊天传输,把web与各聊天状态合并进同一个statuses字典;StartedGateway.transports以TransportName为键保存TransportHandle,stop()时先停 Web(独立超时WEB_STOP_TIMEOUT_SECONDS),再用共享预算停所有聊天 worker;GatewayController只持有这个不透明的StartedGateway,不知道各平台的启动细节——传输特定工作因此不会泄漏到控制器。
导入方向被严格钉死:只有GatewayController导入gateway.startup,也只有gateway.startup导入gateway.transports.startup;传输 peer 之间不得 importgateway.startup或gateway.web。
七、边界纪律与测试保障
AGENTS.md 指出,仅靠"边界测试"(禁止做什么)不足以保证每个传输"必须包含什么",为此仓库用两层测试把纪律固化下来。
7.1 传输契约测试:关注点清单 + known-gaps 台账
gateway/tests/test_transport_contract.py 用包源码 AST/文本检测(而非 import)验证每个传输是否实现了全部共享关注点,避免把平台 SDK 拉进测试。必须实现的关注点清单(_REQUIRED_CONCERNS)包括:
| 关注点 | 源码标记 |
|---|---|
| principal 作用域解析(turn 数据归属,任何存储访问前绑定) | def resolve_ |
| 围绕 turn 的存储作用域绑定 | bound_storage_scope |
| 共享入站决策步骤(会话生命周期走共享实现) | apply_inbound_decision |
| turn 超时设置(挂起的 turn 不能永久挂起会话) | turn_timeout_seconds |
| 停止命令处理(用户可终止运行中的 turn) | is_stop_command |
| 信用计量绑定(turn 在容量判定前绑定计费) | bound_turn_metering |
turn 输出声明turn_cancel(配合宿主侧取消) | self.turn_cancel/SingleMessageTurnOutput |
台账(_KNOWN_GAPS)当前为空——即每个传输都实现了全部关注点。台账语义是"账本而非白名单":断言是精确相等,补上一个缺口必须删除对应条目(台账只能缩小),新传输自动被发现、不得手写加入,因此新传输要么完全合规要么测试全红。
测试还禁止三类"被提升后又在本地复制"的实现(_FORBIDDEN_LOCAL_COPIES):本地重新实现 identity-policy 存储(def _load_policy)、容量判定前消费信用(consume_credits)、轮换会话的哨兵字面量("__ROTATE_SESSION__")。
对信用计量,测试还做了 AST 级检查:每次bound_turn_metering(...)调用的idempotency_key必须是含插值的 f-string(ast.JoinedStr且含ast.FormattedValue),防止常量或空 key 让客户端省略幂等头、从而重新引入双扣费问题。
7.2 包边界测试
- gateway/tests/test_package_borders.py:钉死 peer 之间的 import DAG,保证"一个平台 ≠ 四套 SDK 栈"的隔离;
- gateway/tests/discord/test_transport_borders.py:Discord ↔ Slack 的额外隔离验证。
八、新增一个聊天传输的标准路径
综合 AGENTS.md 与源码约定,为网关新增一个聊天平台的完整清单如下:
- 新建
gateway/transports/<name>/包,包含:settings.py(平台配置加载)、入站 worker、security.py等安全逻辑、turn 输出组装、以及startup.py; - 在
startup.py中实现start_<name>_worker(*, logger, handler) -> tuple[worker, settings]:加载 settings、启动 worker;未配置时抛GatewayConfigurationError,启动失败抛GatewayTransportFailedError; - 在 gateway/transports/names.py 的
TransportName中加入成员,并在 gateway/transports/startup.py 的TRANSPORTS元组中注册; - 实现全部七个共享关注点(见 7.1 清单)——契约测试会自动发现新包并强制全绿;
- 不得:import peer 包、import
gateway.startup/gateway.web、在本地复制已提升到gateway.core的实现。
按此路径,新增平台只需要交付"平台差异",Agent 分发、会话、审批、计量、停止命令等能力全部继承自共享层。
参考路径速查
- 架构约定: gateway/transports/AGENTS.md
- 注册表与启停循环: gateway/transports/startup.py、gateway/transports/registration.py、gateway/transports/worker.py、gateway/transports/names.py
- 各平台 startup: gateway/transports/slack/startup.py、gateway/transports/discord/startup.py、gateway/transports/telegram/startup.py、gateway/transports/buzz/startup.py
- 共享契约与组合根: infrastructure/turn_host/turn_callback.py、gateway/core/middleware、gateway/startup.py
- 契约与边界测试: gateway/tests/test_transport_contract.py、gateway/tests/test_package_borders.py、gateway/tests/discord/test_transport_borders.py
【免费下载链接】opensreBuild your own AI SRE agents. The open source toolkit for the AI era.项目地址: https://gitcode.com/GitHub_Trending/op/opensre
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考