gateway/,99 个文件,其中 run.py 单文件 1.55 MB —— 是整个项目最大的模块。这一章讲一个常驻进程如何让 22 个聊天软件都能触达同一个智能体。
先说清楚没有网关会怎样。假设你想让智能体支持 Slack 和 Telegram 两个平台:
网关就是把这些共性抽出来的那一层。它是一个长期不关闭的后台进程,同时连着所有平台,负责:
gateway/platforms/base.py,333 KB。这是那份「插座标准」—— 规定任何一个平台想接进来必须提供什么。
class BasePlatformAdapter(ABC):
# —— 消息长度限制 ——
def max_message_length_for_chat(self, chat_id: str) -> int
def message_len_fn(self) -> Callable[[str], int]
def message_len_fn_for_chat(self, chat_id: str) -> Callable[[str], int]
# —— 流式输出能力 ——
def supports_draft_streaming(self, ...) -> bool
def prefers_fresh_final_streaming(self, ...) -> bool
def streaming_overflow_limit(self) -> Optional[int]
async def send_draft(self, ...)
# —— 权限模型 ——
def enforces_own_access_policy(self) -> bool
def authorization_is_upstream(self) -> bool
# —— 渲染 ——
def render_message_event(self, event, sink) -> None
def format_tool_event(self, event, *, mode: str = "all", ...) -> str
def format_tool_preview(self, preview: "ToolPreview") -> str
def set_status_text(self, chat_id: str, text: Optional[str]) -> None
hermes-agent/gateway/platforms/base.py
对比两种写法:
# 写法 A:在主流程里判断平台
if platform == 'slack':
max_len = 40000
elif platform == 'discord':
max_len = 2000
elif platform == 'telegram':
max_len = 4096
elif platform == 'sms':
max_len = 160
# ... 22 个分支,而且这样的判断散落在几十处
# 写法 B:问适配器
max_len = adapter.max_message_length_for_chat(chat_id)
写法 A 的问题不是难看,是「新增平台要改几十个地方」。而且你不知道要改哪几个 —— 只能靠搜索和运气。
写法 B 下,新增一个平台 = 实现一组能力声明。主流程一行都不用动。而且抽象基类会强制你实现所有必需的方法,漏一个直接报错。
这就是这一层能撑住 22 个平台的根本原因。
| 能力 | 各平台的实际差异 |
|---|---|
能不能编辑已发消息supports_draft_streaming |
Slack / Telegram / Discord:能。所以可以流式更新同一条消息,用户看到文字逐渐生长。 短信:不能。只能发新消息 —— 流式输出根本没法做,只能等全部生成完再发一条。 |
偏好重发最终版吗prefers_fresh_final_streaming |
有些平台编辑消息会触发通知或者把消息顶到最新,体验很差。这类平台宁可「流式过程用草稿,最终结果发一条新的」。 |
单条消息长度上限max_message_length_for_chat |
Discord 2,000 / Telegram 4,096 / 短信 160 / Slack 约 40,000。 注意这个方法接收 chat_id 参数 —— 因为同一平台的不同频道可能有不同限制(比如企业版和免费版)。 |
长度怎么算message_len_fn |
更微妙:「长度」的定义各平台不同。有的按字符数,有的按 UTF-16 码元,有的把 emoji 算多个。所以返回的是一个计算函数而不是一个数字。 |
平台自己管权限吗enforces_own_access_policy |
企业微信 / 飞书:有完整的企业权限体系,能进到这个群的人就是被授权的。 IRC:完全没有权限概念。任何人都能发消息,必须由智能体自己做授权。 |
授权在上游吗authorization_is_upstream |
有些接入方式(比如通过企业网关代理)已经在上游做过身份认证,智能体这里不该再问一遍。 |
各平台的消息格式完全不同。网关把它们统一成一个事件对象:
class MessageType(Enum): ...
class ProcessingOutcome(Enum): ...
class MessageEvent:
...
def is_command(self) -> bool # 是不是斜杠命令
def get_command(self) -> Optional[str] # 命令名
def get_command_args(self) -> str # 命令参数
class CachedMedia:
def context_note(self) -> str # 附件的上下文说明
class TextDebounceState: ... # ★ 文本去抖动
class SendResult: ...
class EphemeralReply(str): # ★ 带过期时间的临时回复
def __new__(cls, text: str, ttl_seconds: Optional[int] = None): ...
def text(self) -> str: ...
TextDebounceState(文本去抖动状态)解决的是这个场景:
EphemeralReply(临时回复)是一个继承自 str 的类,带一个过期时间。用于那些「说完就该消失」的消息 —— 比如「正在思考…」这类状态提示。在支持的平台上,这类消息会在一段时间后自动删除,不污染聊天记录。
让它继承 str 是一个很实用的设计:所有原本处理字符串的代码不用改,照常工作;只有关心过期时间的代码才去检查它是不是 EphemeralReply。
run.py 里的函数名本身就是一份「长驻进程会遇到什么问题」的清单:
def _hygiene_cooldown_for_failure(...)
def _reset_hygiene_failure_streak(gateway, session_key: str) -> None
def hygiene_compaction_recovered(...)
def _record_hygiene_cooldown(...)
「卫生」(hygiene)在这里指的是后台自动做的上下文整理。如果某个会话的压缩连续失败,就给它一个冷却期,别再反复尝试 —— 否则会持续消耗资源而且持续失败。
_reset_hygiene_failure_streak(重置失败连击)这个命名说明它跟踪的是连续失败次数,成功一次就清零。这和第 3 章会讲的幂等锁是同一个思路。
def _is_transient_network_error(exc: BaseException) -> bool
长驻进程必须区分「网络抖了一下」和「真的出问题了」。前者应该静默重试,后者应该告警。如果不区分,用户会被无意义的网络波动告警淹没,最终屏蔽所有告警。
def _redact_gateway_user_facing_secrets(text: str) -> str
def _redact_approval_command(cmd: "str | None") -> str
网关会把一些内部信息发到聊天窗口(错误提示、审批请求)。这些文本里可能夹带 API 密钥、令牌、密码。而聊天窗口是多人可见的、会被搜索的、会被归档的。所以发出去之前必须脱敏。
注意有两个脱敏函数:一个给通用文本,一个专门给「待审批的命令」。因为命令里的密钥形态不一样(可能在环境变量赋值里、在参数里、在管道里)。
def _gateway_provider_error_reply(text: str) -> str
def _looks_like_gateway_provider_error(text: str) -> bool
def _sanitize_gateway_final_response(platform: Any, text: str) -> str
模型服务返回的错误信息通常是给开发者看的(含堆栈、内部错误码、请求 ID)。直接发到聊天窗口对用户毫无意义。所以要识别出来并翻译成人话。
gateway/restart.py
gateway/restart_loop_guard.py ← ★ 重启风暴防护
重启风暴是长驻进程的经典故障:进程崩溃 → 自动重启 → 启动时又崩 → 又重启 → 无限循环。每秒重启几十次,日志被刷爆,CPU 打满。
守卫的做法通常是:记录最近的重启次数和间隔,如果在短时间内重启太多次,就停止自动重启并保持崩溃状态 —— 让人来看一眼。
gateway/memory_monitor.py 内存占用监控
gateway/agent_cache_pressure.py 智能体缓存压力
gateway/drain_control.py ★ 优雅排空
gateway/disk_status.py 磁盘状态
智能体缓存压力:网关会缓存智能体实例(同一场会话复用同一个)。但每个实例都持有完整的消息历史 —— 几十场活跃会话就能吃掉几 GB 内存。所以需要监控压力并在必要时淘汰不活跃的实例。
优雅排空:要关闭网关时,不能直接杀掉 —— 手头正在处理的消息会丢。正确做法是「停止接受新消息,把手头的处理完,然后退出」。
gateway/delivery.py
gateway/delivery_ledger.py ← ★ 投递账本
gateway/rich_sent_store.py
gateway/message_timestamps.py
gateway/dead_targets.py ← 死目标(比如被删掉的频道)
投递账本记录「哪条回复已经发到哪里了」。它防的是重复投递:网络超时时你不知道消息发出去没有,重试可能导致用户收到两遍。有账本就能判断。
死目标处理的是:智能体要回复的那个频道被删了、机器人被踢出群了、用户拉黑了。这些投递会永久失败,必须识别出来并停止重试,否则重试队列会无限增长。
def _status_template_to_regex(template: str) -> str
def _gateway_compression_progress_notices_enabled() -> bool
def _prepare_gateway_status_message(platform, event_type: str, message: str) -> Optional[str]
async def _send_or_update_status_coro(adapter, chat_id, status_key, content, metadata)
def render_notice_line(notice) -> str
智能体在长任务中需要给用户反馈「我还在干活」。但在聊天软件里做这件事很微妙:
_status_template_to_regex 这个函数值得注意:它把状态消息的模板转成正则表达式。用途应该是「识别聊天记录里哪些消息是自己之前发的状态消息」,从而可以更新或删除它们 —— 因为平台 API 返回的消息 ID 可能已经丢失,只能靠内容匹配来找。
def _is_fresh_gateway_interruption(...)
def build_resume_recovery_note(...)
def _prepare_resume_pending_message(...)
def _build_replay_entry(...)
def _startup_restore_drain_timeout_secs() -> float
def _auto_continue_freshness_window() -> float
这一组函数处理的是:网关重启后,那些「进行到一半」的会话怎么办。
这一组函数是「长驻」和「命令行工具」的本质差别所在。
命令行工具崩了就崩了,用户重新跑一次。长驻进程崩了之后必须自己判断「刚才干到哪了、该不该继续、要不要告诉用户」 —— 而且判断依据只有磁盘上的状态。
这部分逻辑在架构图上完全看不见,但它占了网关模块相当大的比重。
gateway/hooks.py
gateway/builtin_hooks/
网关层有自己的钩子系统,让插件可以在消息进出的关键节点插入逻辑。和第 9 章讲的插件系统配合,构成了「不改核心代码就能扩展网关行为」的能力。
1.55 MB 的单文件是这一层最直观的代价。它不是设计缺陷,是「22 个平台 × 每个平台的边角情况」组合爆炸的必然结果。
从函数名可以看出,这个文件里塞了:冷却策略、错误分类、脱敏、状态渲染、时间戳处理、审批转发、进度线程解析、平台显示配置、Telegram 特有的提及格式转换(_telegramize_command_mentions)……
每一个都很小,但加起来就是 1.55 MB。而且它们大多无法被抽象掉 —— 因为它们本质上就是在处理外部世界的不规则性。