OpenSRE 网关核心架构解析:gateway/core 的进程编排、容量闸门与 Agent 驱动边界
【免费下载链接】opensreBuild your own AI SRE agents. The open source toolkit for the AI era.项目地址: https://gitcode.com/GitHub_Trending/op/opensre
本指南以仓库 gateway/core/AGENTS.md 为骨架,结合
bootstrap/、infrastructure/turn_host、infrastructure/scheduling与gateway/tests/下的源码与测试,系统拆解 OpenSRE 网关进程中"进程与叶子基础设施"层的设计:哪些代码可以放在gateway/core、谁能真正驱动 Agent、容量槽如何防止进程被击穿、进程引导与生命周期如何分工,以及这些约定如何被 AST 测试钉死。读完你将理解 OpenSRE 多传输(Slack/Discord/Telegram/Web)共享同一套 Agent 运行时与并发闸门的底层机制,并能判断新增代码应当落在哪一层。
一、gateway/core 是什么:进程与叶子基础设施
OpenSRE 网关是面向多聊天传输(chat transports)与 Web 的常驻进程。gateway/core/是所有 surface 共享的核心机制,其 AGENTS.md 第一句话就划定了它的身份:process and leaf infrastructure(进程与叶子基础设施)。
两条硬性约束决定了这个包的"叶子"属性:
- 不得导入
gateway.transports.*或gateway.web; - 只有 gateway/core/lifecycle/controller.py 可以通过 gateway/startup.py 拉起这些 surface(传输层与 Web 层)。
违反边界的代码会被 gateway/tests/test_package_borders.py 直接拦截——这是一个用 Python AST 解析 import 关系做断言的测试文件,属于"只许收缩、不许放宽"的白名单式约束(详见第五节)。
1.1 包结构与职责一览
AGENTS.md 用一张表定义了gateway/core下各子包的职责,这是理解整个核心层的第一张地图:
| Package | Role |
|---|---|
lifecycle/ | Composition root(controller)、凭据注水(credential hydration)、网关错误 |
process/ | 守护进程、轮询线程、就绪状态 |
middleware/ | 每个传输都要执行的回合级步骤(入站决策、身份策略、审批、注意力、锁) |
storage/ | 会话绑定与调查、事件与反馈存储 |
billing/ | 积分(credits)客户端 |
attachments/ | 附件辅助 |
session/ | 网关聊天上下文辅助 |
config/ | 日志 / 网关配置辅助 |
infrastructure.turn_host | 回合处理器、会话-Agent 池、可绑定输出、取消控制台、容量 |
注意最后一行:infrastructure.turn_host并不位于gateway/目录下,它属于infrastructure/包,但被gateway/core视作自己职责的一部分——这正体现了 AGENTS.md 强调的"叶子基础设施":turn_host 是网关回合执行的真正宿主,而gateway/core中的聊天代码只做调度与收尾。
传输层与web/可以import 上述包,但聊天传输的业务代码不属于gateway/core。反过来,gateway/core绝不能 import 任何聊天传输——从 controller.py 的模块注释可以看到,GatewayController只负责"启动 web 与聊天传输"这一编排动作,具体 worker 由gateway.startup提供。
1.2 叶子之间的依赖纪律
gateway/core内部还有一层更细的依赖纪律。例如 gateway/core/process/liveness.py 的注释点明:
A leaf: both supervision and component status need it, and putting it in either would make them import each other.
也就是说,liveness(判断进程 ID 是否存活)被supervision(守护进程)与component_status(组件状态)同时依赖,所以它被单独抽成叶子,避免兄弟模块互相 import。这种"依赖只朝下"的纪律贯穿整个核心层。
二、谁可以驱动 Agent:turn_host 的单一入口
AGENTS.md 中"Who may drive the agent"一节回答了 OpenSRE 网关里最关键的问题:在一个多传输进程里,谁有权构造端口、绑定回合、运行/刷新会话、格式化目标进度?
答案只有一个:infrastructure.turn_host。
只有 infrastructure.turn_host 负责: - 构建 ports(端口) - 绑定一个 turn(回合) - 运行或刷新一个 session(会话) - 格式化 goal progress(目标进度)gateway/下其他代码只能import harness 的类型(SessionCore、OutputSink、SlashPortsFactory、SessionGoal)用于函数签名,例如 gateway/core/chat_agent_build.py 与 controller.py 中都以SlashPortsFactory作为注入类型。一个自己执行回合的传输,就等于出现了"第二个 handler",这是设计上明确禁止的。
2.1 为什么 web 在豁免名单上(一个已知例外)
该边界由 gateway/tests/test_harness_behaviour_border.py 钉死,且其 allowlist只许收缩、不许扩张。当前名单上有一个已知例外:web/。
原因在 controller.py 中写得很清楚:POST /investigate直接内嵌(embed)Agent,因此它拿不到 turn-host 提供的一系列保障——Agent 复用、审批钩子、取消控制台、能力策略(capability policy)。这是设计者主动接受的权衡:同步 HTTP 调查接口需要即时拿到推理结果,无法走聊天式的异步回合管道,但代价就是它游离于统一保障之外。
2.2 TurnRunner:唯一回合处理器
从源码看,infrastructure.turn_host.turn_runner.TurnRunner是"唯一turn runner"(turn_runner.py)。它的定位是:
- 传输调度器(Slack/Discord/Telegram)是 ingress 适配器:它们做授权、解析会话、构建输出,然后调用这个回调;
TurnRunner对传输无关,接收(text, session, output, logger),运行回合并在 output 上收尾文本;- 进程级容量是可选的
gate,挂在同一对象上——不是第二个 runner; - Agent 复用由
SessionAgentPool负责。
TurnRunner暴露两个入口、共享同一个回合(turn_runner.py):
__call__:TurnCallback,四个聊天传输使用,返回 None;run:面向进程内调用方,需要以值的形式拿到TurnResult,接受调用方的 console、confirm_fn、is_tty与on_progress。
两个入口的每个关键字参数默认值都指向传输路径,保证二者永不漂移。在run内部,容量检查、取消检查、准入钩子(如计费的admit_metered_turn)与真正的_run_turn严格顺序执行:准入钩子故意放在槽内执行,目的是"容量拒绝的工作不得计费"。
SessionAgentPool(session_agents.py)则保证每个逻辑会话一个HeadlessAgent:会话间并发、会话内串行——同一会话的回合不能重叠,否则一个回合的输出会串到另一个回合的输出上,因此每个 session_id 维护一把独立的threading.Lock,锁的跨度覆盖整个 dispatch(不只是发放 Agent),防止下一个回合重定向一个仍在流式输出的 Agent。
三、容量槽(Capacity Slot):进程级并发闸门
AGENTS.md 用一个专门小节强调:进程内的每一个回合——聊天、POST /investigate、调查 worker、定时运行——都要从同一个process_turn_gate()取一个许可(permit)。这是"一个进程、一个闸门"的强一致性设计。
3.1 闸门从哪里来
process_turn_gate()在 infrastructure/turn_host/concurrency.py 中实现为进程级单例,底层是threading.BoundedSemaphore,其默认上限由两个环境变量决定(定义见 config/constants/turn_concurrency.py):
| 环境变量 | 作用 |
|---|---|
OPENSRE_SIZE_PROFILE | 部署规模档位:SMALL/MEDIUM/LARGE,映射到默认并发回合上限 |
OPENSRE_MAX_CONCURRENT_TURNS | 显式并发上限,覆盖规模档位默认值;回合是 I/O 密集型,小任务也可以跑多个 |
规模档位的默认上限在 concurrency.py 中是硬编码的:
| SizeProfile | 并发回合上限 |
|---|---|
SMALL | 1 |
MEDIUM | 2 |
LARGE | 4 |
configured_turn_limit()的解析规则很严谨:OPENSRE_MAX_CONCURRENT_TURNS若非正整数或无法解析,会被忽略并告警,退回规模档位默认值——一个拼写错误不会把闸门打到 0(打 0 等于永久拒绝服务)。这个行为有测试佐证:gateway/tests/runtime/test_concurrency_gate.py 用参数化测试逐一验证了 SMALL/MEDIUM/LARGE 的上限。
3.2 两种容量策略:drop 与 wait
闸门满了之后怎么办?AGENTS.md 明确给出恰好两种策略(实现见 infrastructure/process/turn_capacity/slots.py),调用方必须二选一,禁止手工配对acquire/release:
turn_slot(gate)—— drop(丢弃)。适用于"可以告诉对方稍后再试"的请求:聊天消息、HTTP 回合。它非阻塞地try_acquire(),闸门满时立即 yieldFalse,由调用方决定如何应答:
- 聊天路径:finalize
AT_CAPACITY_MESSAGE; - Web 路径:作为 503 响应体返回。
AT_CAPACITY_MESSAGE在 concurrency.py 中定义为唯一一句话、唯一一个位置:"OpenSRE is at capacity. Please try again shortly.",聊天与 Web 引用同一个常量,保证两端文案永不漂移。这条链路有端到端测试覆盖:gateway/tests/test_multi_actor_concurrency.py 验证所有被拒绝的 actor 都逐字收到该消息,test_turn_runner.py 验证TurnRunner在闸门满时 finalize 的就是这条消息。
queued_turn_slot(gate)—— wait(等待)。适用于"已经从队列认领、无法叫它重试"的工作:定时运行、HTTP worker。丢弃等于丢失,所以它阻塞式acquire()排队等槽。
两者的核心风险点相同:缺失finally会泄漏许可,而一个泄漏的许可 = 一个永远回答 "at capacity" 的进程。turn_slot/queued_turn_slot用contextmanager+try/finally把释放内建在协议里(slots.py),这正是"不要手工配对"的原因。测试 test_concurrency_gate.py 专门构造了一个抛异常的 handler,验证finally释放后闸门能立刻再次try_acquire()成功。
3.3 定时任务也走同一个闸门
关键设计:定时运行同样消耗一个回合。在 controller.py 的start_scheduler中,控制器把scheduler_runners().gated(self.turn_gate)的结果安装进调度器——SchedulerRunners.gated(infrastructure/scheduling/scheduler/runners.py)返回的是新的 bundle而非原地改写:每个定时 run 在queued_turn_slot(gate)内执行。这个"值而非副作用"的设计消除了历史上"双重 gating 导致一次运行扣两个 permit"的幂等问题。此外OPENSRE_SCHEDULER_MAX_CONCURRENT_RUNS还提供了调度回调自身的并发上限(默认 2,见 turn_concurrency.py),与进程闸门构成第二道闸。
四、进程引导(Process Boot)vs 生命周期(Lifecycle)
AGENTS.md 的最后两节回答了一个工程上极易混淆的问题:进程启动的公共部分放哪里?网关自己的编排放哪里?
4.1 唯一入口:bootstrap.process.configure_process
进程引导有且仅有一个入口:bootstrap/process.py 的configure_process(profile),并且网关必须传GATEWAY_PROFILE。AGENTS.md 明确警告:"不要给它包一层网关本地的 wrapper"。
configure_process的设计是一张有序步骤表:BootStep枚举定义了env → sentry → harness_adapters → scheduler_runners → capability_warnings → preload_llm六步(process.py),_STEP_ORDER元组固定了执行顺序(process.py),而 profile 只决定成员资格、不能发明新顺序——"没有哪个 profile 能在环境加载前跑适配器"。
GATEWAY_PROFILE(process.py)选择四步:
- ENV:
bootstrap_opensre_env_once加载本地环境; - SENTRY:以
SentryEntrypoint.GATEWAY初始化错误上报; - HARNESS_ADAPTERS:安装 Agent 工具解析所用的适配器与 CLI 认证检查器;
- CAPABILITY_WARNINGS + PRELOAD_LLM:输出沙箱能力告警、预热 LLM 客户端。
configure_process是幂等的(process.py):同一 profile 只执行一次,测试通过reset_process_runtime_for_tests()清理。仓库还为 CLI、Web、调度 worker、定时命令与嵌入式宿主各预置了 profile(CLI_PROFILE、WEB_PROFILE、SCHEDULER_WORKER_PROFILE、SCHEDULED_COMMAND_PROFILE、EMBEDDED_PROFILE),它们只是步骤子集不同——这正是"一个顺序表、多个 profile"的价值。
4.2 GatewayController:组合根与生命周期所有者
生命周期由 GatewayController 承担。start_gateway()的编排顺序非常明确:
配置日志 → set_ready(False) → 凭据注水(_load_credentials) → configure_process(GATEWAY_PROFILE) → 构建【一个】TurnRunner(gate=turn_gate, admission_check=admit_metered_turn) → start_surfaces()(web + 聊天传输一起启动) → start_scheduler()(进程内托管 cron/loop) → 发布组件状态、写 gateway.pid、set_ready(True)、注册 SIGINT/SIGTERM几个关键约束与源码一一对应:
- "不要包第二个 turn runner":
TurnRunner构造时通过gate=传入容量闸门,容量就在这个对象上(controller.py)。测试 test_concurrency_gate.py 专门断言生产路径上容量挂在TurnRunner自身、无需第二层包装(历史上有过ConcurrencyLimitedTurnHandler这样的包装器,被 test_package_borders.py 列入"隔离名单",只允许出现在测试中)。 - 凭据先于一切:
_load_credentials在传输、调度器、worker 启动之前完成注水(controller.py)。注水有两条路由(credential_hydration.py):integrations secret(部署 silo 的主路由)与credentials API(暂存回退路由,需 bootstrap secret 携带 token);两条都未配置则视为"功能关闭"(笔记本本地运行),而"配了但不全"则视为部署损坏、直接抛GatewayConfigurationError。 - 调度器托管可选:
OPENSRE_GATEWAY_HOST_SCHEDULER环境变量设为 false 时,调度器作为独立服务运行(MODE=scheduler),避免两个进程同时触发定时任务(controller.py)。
4.3 调度器宿主:薄调用 + 轮询重载
网关托管调度器的方式是薄调用:
scheduler_runners().gated(turn_gate).install() → infrastructure.scheduling.scheduler.runner.start_background_scheduler(...)控制器不持有调度逻辑,只安装"带闸门的 runner"并启动后台调度器(controller.py)。重载(reload)是跨进程信号机制(infrastructure/scheduling/scheduler/reload_signal.py):
- 写方是交互式 shell / CLI(
request_scheduler_reload):在~/.opensre/scheduler_reload_requested写一个时间戳文件,尽力而为——写失败只告警、绝不阻塞任务存储的变更; - 读方是网关控制器:一个后台线程以 2 秒为周期
watch_and_reconcile(RELOAD_POLL_SECONDS = 2.0),同时监听信号文件与任务存储文件的(mtime_ns, size)签名——即使尽力而为的信号丢了,下一次轮询也会因存储签名变化而收敛;idle 存储只是一次廉价的stat。
控制器只"轮询与重同步",这正是 AGENTS.md 强调的职责切分:shell/CLI 是写方、控制器是读方,不要在gateway/内再塞一个调度器。相关约束也被 test_package_borders.py 钉死:infrastructure.scheduling.scheduler不得 importTurnRunner——调度器是"生产者(producer)",没有用户、没有回合输出,绝不能变成第五个聊天通道。
4.4 进程层的叶子组件
gateway/core/process/提供守护进程所需的全部叶子能力:
- supervision.py:守护进程的 pidfile(
~/.opensre/gateway/gateway.pid)与日志(~/.opensre/gateway/gateway.log)管理;它只负责监督子进程,不启动调度器、不命名surfaces.gateway_entry; - liveness.py:
process_is_alive(pid),先waitpid(WNOHANG)收割子进程再kill(pid, 0)探测; - readiness.py:进程内就绪状态(
set_ready/is_gateway_ready),start_gateway开头置 false、结尾置 true; - component_status.py:组件状态发布,配合
supervision写入;shutdown_budget.py 为stop()的优雅停机分配时间预算。
五、边界如何被测试钉死:AST 级包依赖防线
AGENTS.md 中反复出现的"Pinned by …"并非口号。gateway/tests/test_package_borders.py用ast解析源码的 import 语句,把包边界变成可执行的断言,其核心规则包括:
- 传输之间是平级:任何聊天传输不得 import 另一个聊天传输(test_package_borders.py);
gateway.web不得 import 聊天传输或gateway.startup(L92-L95);gateway.core不得 import 聊天传输、gateway.web或gateway.transports(L97-L101)——这就是"叶子"边界的机器保证;- 只有
lifecycle/controller.py能 importgateway.startup(L140-L148); gateway.core不得命名任何聊天厂商:vendor 知识属于跟它对话的传输层(L104-L123)。例如历史上gateway/core/billing曾 importintegrations.slack.webapp_auth,导致计费客户端看起来 Slack 专属,后被移到厂商中立的模块;- 每个磁盘上的传输包必须出现在注册表中,反之亦然(L397-L414),防止"包存在但传输从未启动"或"注册了却不存在的包"两类静默故障;
- 生产
gateway/不得继承core.agent.Agent、不得调用build_agent、不得使用precomputed_action_tools(L315-L394):聊天必须走SessionAgentPool → HeadlessAgent的单一路径,绝不允许在网关里长出"第二个大脑"。
这套 AST 防线还覆盖了字符串字面量形式的依赖(例如守护进程历史上用python -m surfaces.gateway_entry硬编码子进程 argv 造成的隐藏环,L265-L283),保证可执行代码中的模块引用也受约束。
六、回合级中间件:每个传输都要走的固定步骤
gateway/core/middleware/是 AGENTS.md 表格中"每个传输都要执行"的回合级步骤,它把四份几乎相同的传输代码收敛为一份(inbound_decision.py):
inbound_decision:入站授权决策与"决策后编排"的统一实现——持久化策略变更、发送决策回复、拒绝时停止回合、遇到ROTATE_SESSION哨兵时轮换会话;传输特有的行为(如何发送回复、拒绝是否显式回复、解析器是否需要锁)留在调用方;identity_policy:加载/持久化聊天传输的入站身份策略,存储在平台集成记录的credentials["identity_policy"]上;首次运行无记录时给一个宽松默认(入站启用、空 allowlist)(identity_policy.py);approvals:传输无关的写工具审批闸门——requires_approval=True的工具在 CLI 上是Proceed? [Y/n],在聊天网关上是 Approve/Deny 按钮,这里只放与传输无关的部分:把点击连接到等待中的工具调用的ApprovalBroker、按钮标识符(opensre_approval_approve/opensre_approval_deny)、harness 钩子;审批最长等待 180 秒、参数预览截断 400 字符(approvals.py)。每个传输提供自己的ApprovalPrompter:Slack 用 Block Kit、Discord 用消息组件、Telegram 用内联键盘;- 另有
attention、conversation_locks、active_turns、terminal_outcome等步骤共同构成回合管道。
七、存储层与计费:叶子基础设施的落地
- 会话绑定:
gateway/core/storage/session/binding_store.py采用JSON 文件 + 写临时文件再 rename的原子方案,明确不用数据库、不用 SQLite——因为上下文根目录是 NFS 挂载,SQLite 的 advisory locking 不可靠,而 rename 是原子的;单写者假设来自 Slack Socket Mode 的单消费者特性(binding_store.py); - 事件与反馈:
storage/events/提供调查事件仓储与建表,storage/feedback/jsonl.py用 JSONL 持久化反馈; - 计费:
billing/提供积分客户端与admit_metered_turn准入钩子——它作为TurnRunner的admission_check注入,在容量槽内执行,保证被容量拒绝的回合不会被计费。
结语:一张边界图记住 gateway/core
用一句话概括 OpenSRE 网关核心的设计哲学:叶子只管叶子的事,编排只留一个根,并发只有一个闸,驱动 Agent 只有一条路。
gateway/core是叶子与组合根:不碰传输、不碰 Web,只有lifecycle/controller.py能拉起 surface;- 回合执行只属于
infrastructure.turn_host,TurnRunner是唯一 runner,SessionAgentPool保证会话内串行、会话间并发; - 容量是进程级单例闸门,
turn_slot(drop)与queued_turn_slot(wait)二选一,绝无第三种; - 引导走唯一的
configure_process(GATEWAY_PROFILE),生命周期归GatewayController,调度器是"薄安装 + 轮询重载"; - 以上全部边界由 gateway/tests/test_package_borders.py 以 AST 断言固化,只许收缩、不许放宽。
对新增代码的读者而言,这份文档就是一张决策表:要启动一个传输 worker?去gateway/startup组合,而不是在 core 里 import 传输;要驱动一个 Agent 回合?调用TurnRunner,而不是自己构建 harness;要开一个并发的回合来源?从process_turn_gate()取槽,选好 drop 或 wait 策略,并永远让释放发生在finally里。
【免费下载链接】opensreBuild your own AI SRE agents. The open source toolkit for the AI era.项目地址: https://gitcode.com/GitHub_Trending/op/opensre
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考