ScienceDiscovery
English GitHub

外部数据源限流机制

本文档描述外部数据源出站请求的统一限流底座:资源键抽象、排队与超时语义、错误码、默认参数,以及 LLM API 的挂接边界。

1. 背景与目标

论文/文献/数据库类公共 API(arXiv、NCBI eutils、UniProt 等)对客户端有 QPS 或并发限制。突发并发加上模型侧盲目重试会打穿上游配额,产生成串的 429 Too Many Requests,并以无结构错误回流给模型与用户。

限流底座在 Node 治理面(services/api)为每个上游资源提供可独立启用的最小请求间隔、最大并发、FIFO 等待队列深度、排队超时,以及取消释槽、429 冷却与排队时长审计。目标是把「盲目重试 + 裸 429」替换为「按源配置的有序排队 + 结构化可读错误」。

2. 与 MCP 源的关系与调用流程

2.1 关系边界

限流在 Node 控制面统一实现,不在各 Python MCP 源内部:

各源通过 manifest 的 governance 字段向底座提供参数(schema 见 packages/schema/src/mcp-source.ts;内建源构造见 packages/mcp-sources/src/public-biomed.ts 的 manifest()):

字段 作用
rateLimitGroup 资源键。内建源默认取 hosts[0](如 arXiv 为 export.arxiv.org);同 rateLimitGroup 的多个源共享同一配额,例如共享 eutils.ncbi.nlm.nih.gov 的 PubMed/ClinVar/GEO 共占 NCBI 的并发与 pacing。
maxConcurrentRequests 最大在途请求数。
minIntervalMs 最小请求间隔;未提供时 broker 在 rateLimitPerSecond 存在时推导 ceil(1000 / rps)。
rateLimitPerSecond 每秒放行上限,用于推导 minIntervalMs。
maxQueueDepth / queueTimeoutMs 等待队列深度与排队超时。

资源键是任意字符串。MCP 源之外预留命名空间:llm:<provider-host>(LLM,见第 6 节)、web:<provider>(Web Search/Fetch 尚未接入,见第 7 节)。

2.2 端到端调用链

Agent 工具 mcp__<source>__<tool>
  → createMcpWorkspaceTools.execute            packages/artifact-manager/src/mcp-workspace-tools.ts
  → McpGovernanceBroker.invoke                 packages/data-source/src/broker.ts
      会话可写检查 → 启用过滤 → 输入校验 → 权限 authorize
      → 结果缓存:命中则直接返回规范化结果(不再 acquire、不再出站)
      → ResourceRateLimiter.acquire(rateLimitGroup, …)   packages/data-source/src/resource-rate-limiter.ts
          队列满   → RATE_LIMIT_QUEUE_FULL
          排队超时 → RATE_LIMIT_QUEUE_TIMEOUT
      → McpNodeClient.invoke → MCP server → 上游 provider
          Gateway 内按 retryPolicy 对 429/5xx/transport 做有限重试(尊重 Retry-After)
          attempt 出现 rate-limited → ResourceRateLimiter.reportUpstreamRateLimit(rateLimitGroup, retryAfterMs)
      → lease.release()   finally 释放并发槽并 pump 队尾

工具名由 createMcpWorkspaceTools 包装为 mcp__<sourceId>__<toolId>,execute 统一走 McpGovernanceBroker.invoke;只有会话启用且 catalog 判定可用的工具才会被注入(注入链路详见 science-connectors.md 第 3 节)。

2.3 代码锚点

模块 路径 职责
限流底座 packages/data-source/src/resource-rate-limiter.ts acquire / release / reportUpstreamRateLimit:进程内并发、间隔、队列、冷却
治理挂接 packages/data-source/src/broker.ts McpGovernanceBroker.invoke:权限 → 缓存 → acquire → Gateway → 429 反馈 → release
治理参数 packages/mcp-sources/src/public-biomed.ts manifest() 构造 governance(rateLimitGroup、并发、间隔/QPS、队列等)
工具注入 packages/artifact-manager/src/mcp-workspace-tools.ts createMcpWorkspaceTools 把 MCP 工具包成 mcp__<source>__<tool>,execute 调 broker.invoke

2.4 Broker 逻辑要点

2.5 Gateway 重试与 Node 限流的分工

两层不重复,互补:

衔接点:Gateway 重试的 attempts 中出现 rate-limited 时,broker 调用 reportUpstreamRateLimit(rateLimitGroup, retryAfterMs),把上游指示的冷却写回底座,使队列中等待的后续调用推迟到冷却结束再放行,而不是放行后立刻再撞同一个 429。

2.6 不经此链的出站

以下出站不经 McpGovernanceBroker.invoke,不受本底座约束(与第 7 节「覆盖面与已知缺口」一致):

3. 语义

每个资源键的一次 acquire 按以下规则放行:

四个底座维度均为 optional;undefined 或省略表示不启用该维度,而不是采用隐式默认:

参数 undefined / 省略 显式配置
minIntervalMs 无 pacing;broker 仅在 manifest 显式配置 rateLimitPerSecond 时推导 ceil(1000/rps) ≥ 0;0 也表示无间隔;与 rps 同时配置时本字段优先
maxConcurrent(manifest 为 maxConcurrentRequests) 不限制同时在途请求数 正整数 ≥ 1
maxQueueDepth 无限队列;调用方接受对应的内存增长风险 整数 ≥ 0;0 表示不允许排队
queueTimeoutMs 排队不超时,只能由放行或 Abort/取消结束等待 > 0;到期返回结构化排队超时错误
signal 不提供主动取消 提供 AbortSignal;排队中取消即出队,执行中取消由调用方释放槽位

数值校验只针对已提供的字段。Broker 保留 governance 的缺失状态,不会把缺失字段回填为并发 1、队列 8 或排队 20 秒。

超时分层:

429 冷却:Gateway 重试明细中出现 rate-limited attempt 时,broker 调用 reportUpstreamRateLimit(key, retryAfterMs),该资源键的后续放行推迟到冷却结束。无 Retry-After 时至少冷却 1 秒;若最近一次 acquire 或等待队列存在更长 pacing,则取更长值。即使该键没有 pacing 配置,冷却仍安全生效。

4. 错误语义(对模型与用户可读)

场景 McpError.code retryable 说明
等待队列已满 RATE_LIMIT_QUEUE_FULL true 建议减少并行调用、稍后重试
排队超时 RATE_LIMIT_QUEUE_TIMEOUT true 附实际排队时长(审计 queueWaitMs)
上游 429(重试耗尽) RATE_LIMITED true 尽量携带真实 retryAfterMs
执行超时 TIMEOUT true 既有语义不变
用户取消 CANCELLED false 排队/执行中取消都会释放资源

Retry-After 传播:Python MCP server 的 HTTP 封装不再使用裸 raise_for_status()(它会丢弃响应头),而是抛出形如 HTTP 429 Too Many Requests from export.arxiv.org (retry-after: 5) 的可解析错误;Gateway 错误分类器锚定 retry-after token 取值,重试与冷却都能拿到上游指示。

5. 配置面与内建源显式参数

per-source 配置在 source manifest 的 governance 字段(packages/schema/src/mcp-source.ts):

参数 内建源显式值 arXiv 显式值 依据
minIntervalMs 由 rateLimitPerSecond 推导 3000 arXiv API 使用条款:约每 3 秒 1 请求
rateLimitPerSecond NCBI 3 / 其他 5 5(存在但被 minIntervalMs 覆盖) NCBI 无 key 上限 3/s
maxConcurrentRequests 2 1 arXiv 条款:单连接
maxQueueDepth 8(manifest 显式引用 DEFAULT_MCP_RATE_LIMIT_QUEUE) 8 一轮并行工具调用典型 ≤ 6
queueTimeoutMs 20000(manifest 显式引用同一推荐值) 20000 突发 6 个 arXiv 请求可在窗口内消化
retryPolicy.initialDelayMs 500 3000 重试间隔不低于源 pacing

这些值不是 broker 的隐式默认。当前 public-biomed 与 UniProt manifest 显式选择它们,因此现有源行为保持不变;自定义或未来 source 可省略任一治理字段,以关闭对应限制维度。

审计:每次 McpInvocation 记录 queueWaitMs(排队时长);429/重试明细在 attempts[];均可经 GET /api/sessions/:sessionId/mcp/invocations 查询。不记录 query 全文之外的新增敏感信息。

6. LLM API 边界(独立于数据源 admission)

LLM 请求不经过本文的数据源队列。当前模型 transport 位于 packages/model,使用 Node/undici 实现请求和 pre-stream 重试。

Native executor 直接使用产品 model layer;JiuwenSwarm executor 的模型请求通过 adapter 的 per-run LLM proxy 回到 ScienceDiscovery model gateway,因此仍复用产品 provider、proxy、retry 和 usage 语义。

当前 packages/model/src/client.ts 支持:

相关配置:

环境变量 默认 说明
SCIENCE_AGENT_LLM_TIMEOUT_SECONDS 600 请求 header timeout
SCIENCE_AGENT_LLM_MAX_RETRIES 2 pre-stream connect/429/5xx 最大重试次数

LLM provider 的 TPM/RPM admission 目前没有接入 ResourceRateLimiter。若新增全局模型限流,应放在 model/control-plane 语义中,而不是恢复旧 Python Gateway 限流路径。

7. 覆盖面与已知缺口

8. 测试