Python 交接文档:切换至模型中转站 RPC v2

适用对象:Python Agent 服务团队(当前直接调用上游大模型的 Python 服务)

文档版本:003-model-relay-station(基于提交 28ad927 的实现代码)

最后更新:2026-06-11


⚙️ ai-boss 实测补注(2026-06-15)

接入侧实测结论,补充本文档的内网视角:

  • 可达地址:本文 §2.1 的内网 Nacos 名 http://neargo-token:8080 仅 neargo 集群内可解析; ai-boss 部署环境经网关访问:https://gateway-pos-dev.neargo.ai/neargo-token (RPC 端点 …/internal/relay/chat,dev;生产经 NEARGO_RELAY_BASE_URL 覆盖)。
  • 可用对话模型GET …/internal/relay/models 实测):当前deepseek-chat(默认,别名 gpt-4o)。 不传 model 的默认模型即 deepseek-chat示例中的 qwen-plus);qwen-plus / 各 gemini-* 对话名实测 11002 model unavailablegemini-image图片生成(非对话)。
  • 已验证:默认 / 指定 deepseek-chat 实发正常(meta → token → usage → done);缺 braId11009

1. 概述

模型中转站(003-model-relay-station) 是 neargo 平台唯一的大模型出口。 在此之前,Python 服务持有上游 API Key 并直接调用大模型; 切换后,Python 服务不再持有任何上游明文 Key, 所有大模型请求均通过 neargo-token 服务提供的内部 RPC 端点统一转发。

为什么必须切换

现状(直连)目标(经中转站)
Python 持有上游 API Key 明文Python 持有中转站服务 Key(ngr_ 前缀),不接触上游 Key
计量在网关,Python 计量逻辑分散中转站统一计量、准入、风控、failover
换厂商/模型需改 Python 代码管理端切换模型配置,Python 零改动
无法限速/限额调用方RPM/TPM/金额上限统一治理

端点位置

POST /internal/relay/chat

服务名(Nacos):neargo-token 内网直连示例:http://neargo-token:8080/internal/relay/chat

仅内网可达,网关层不暴露此端点。


2. 认证:服务 Key

Key 格式与获取

  • 格式:ngr_ 前缀的不透明字符串,例如 ngr_xxxxxxxxxxxxxxxxxxxxxxxxxxxx
  • 由平台运维在管理端创建「调用方(relay_caller)」记录后,颁发服务 Key
  • Key 仅在创建时回显一次,服务端只存储 SHA-256 哈希值,之后无法再查看明文
  • Python 服务名对应调用方 callerCode = python-agent

配置方式

通过环境变量(推荐)或配置文件注入,不得硬编码到代码中

# 环境变量
export NEARGO_RELAY_SERVICE_KEY="ngr_xxxxxxxxxxxxxxxxxxxxxxxxxxxx"
export NEARGO_RELAY_BASE_URL="http://neargo-token:8080"

使用方式

每次请求均在 HTTP 头中携带:

Authorization: Bearer ngr_xxxxxxxxxxxxxxxxxxxxxxxxxxxx

Key 的吊销与轮换

场景生效时间操作方式
Key 吊销≤ 1 分钟管理端停用 Key,pub/sub 通知各实例清缓存
Key 轮换(平滑)双 Key 并存窗口先创建新 Key,验证成功后吊销旧 Key,零停机
调用方停用≤ 1 分钟管理端停用调用方,所有该方 Key 立即失效

Key 安全纪律

  • 不得写入代码仓库、日志、异常堆栈
  • 在 CI/CD 中用 Secret 管理
  • 如疑似泄露,立即在管理端吊销并轮换

3. RPC v2 请求契约

端点POST /internal/relay/chat Content-Typeapplication/json 下行stream=true(默认/null)时为 text/event-streamstream=false 时为 application/json

请求体字段

字段类型必传说明
braIdstring必传门店标识。RPC 面一律必传,缺失返回 11009。(OpenAI 兼容面 /v1/*X-Bra-Id 头,RPC 面不适用)
idempotencyKeystring必传调用方生成的全局唯一幂等键(推荐 UUID v4)。保证重试场景明细只落库一次。缺失返回 11009
messagesarray<Message>必传全量对话上下文(中转站无状态单轮转发,每轮须携带完整历史)。缺失返回 11009
modelstring可选目录内可用 CHAT 模型名或别名;缺省时路由至默认对话模型(is_default=true
streamboolean可选truenull = SSE 流式;false = 非流式 JSON。默认 true(null 也视为流式)
optionsobject可选透传上游的模型选项,如 {"temperature": 0.7, "maxTokens": 1024}。不参与准入决策
traceIdstring可选调用方链路追踪 ID。缺省由中转站生成并在 meta 首帧和 X-Trace-Id 响应头回传
bizTagstring可选业务标签(如 STORE_SETUPMARKETING),仅作计量维度,不参与路由或准入决策
toolsarray<ToolDef>可选Function Calling 工具描述符列表。中转站只透传给上游,不执行工具逻辑
toolChoicestring可选工具调用策略,如 autonone 或指定工具名;tools 为空时无效

Message 对象

字段类型说明
rolestringsystem / user / assistant / tool
contentstring消息文本;role=assistant 发起工具调用时 content 可为 null
toolCallsarray<ToolCall>role=assistant 且触发工具调用时携带(历史回灌)
toolCallIdstringrole=tool 时使用,关联对应工具调用的 ID

ToolDef 对象

字段类型说明
namestring工具名称
descriptionstring工具功能描述,供模型判断何时调用
parametersobject参数 JSON Schema(符合 JSON Schema Draft-7)

请求示例

{
  "braId": "B0001",
  "idempotencyKey": "550e8400-e29b-41d4-a716-446655440000",
  "model": "qwen-plus",
  "stream": true,
  "messages": [
    {"role": "system", "content": "你是一位专业的零售顾问。"},
    {"role": "user", "content": "我们的库存周转率偏低,有什么建议?"}
  ],
  "options": {"temperature": 0.7},
  "bizTag": "STORE_SETUP",
  "traceId": "trace-abc-123"
}

4. 响应契约

4.1 SSE 流式响应(stream=truenull,默认)

响应头:

Content-Type: text/event-stream
X-Trace-Id: <resolved-traceId>

每帧格式(SSE 标准):

id: <序号>
event: <事件名>
data: <JSON>

注意:帧末为空行(两个换行符),Python 的 SSE 客户端(sseclienthttpx-sse 等)会自动处理。

事件字典

事件名(event: 行)data: JSON 字段顺序说明
meta{"traceId":"...","model":"...","channelCode":"...","requestId":null}首帧实际命中的模型和渠道,model 字段在未指定模型时反映默认模型名称
token{"delta":"<增量文本>"}多帧内容增量,逐 chunk 透传,不缓冲
tool_call{"toolCallId":"...","name":"...","arguments":"<JSON字符串>","finishReason":"tool_calls"}0或多帧模型请求调用工具,本轮流以此结束;执行在调用方
usage{"promptTokens":100,"completionTokens":200,"usageSource":"PROVIDER"}倒数第二真实 token 用量,usageSourcePROVIDERESTIMATED
done{}末帧正常结束
error{"code":"11006","message":"..."}代替上述错误结束,流终止;见第 5 节错误码表
forceClose{"reason":"RISK_SUSPENDED"}代替上述门店风控开关被切断,服务端主动终止在途流

重要说明:实际事件名为 token(代码中 ModelInvokeEvent.Token.eventName() 返回 "token")。契约文档 §2 将其命名为 delta,但实现代码以 token 为准,客户端须订阅 token 事件,而非 delta。详见第 9 节差异说明。

正常流时序

meta → token → token → ... → (tool_call →)* usage → done

错误流时序(不调上游时)

error

风控切断时序

meta → token → ... → forceClose

4.2 非流式响应(stream=false

HTTP 200,JSON 响应体:

{
  "traceId": "trace-abc-123",
  "model": "qwen-plus",
  "channelCode": "channel-deepseek-01",
  "content": "建议您从以下几个维度优化库存周转率...",
  "finishReason": "stop",
  "usage": {
    "inputTokens": 100,
    "outputTokens": 200,
    "source": "PROVIDER",
    "costAmount": null
  }
}

costAmount 目前在非流式响应中由计价服务填充(pricing_missing 时为 null),调用方需做空值容忍。

错误时(非流式),HTTP 状态码对应错误类型,响应体为 R 格式:

{
  "code": "11009",
  "message": "braId is required"
}
错误类型HTTP 状态
11001 凭证401
11002/11003/11009 参数/模型400
10008/11004/11005 限额/限速429
其余(11006 等)502

5. 错误码

错误码含义建议客户端行为
11001调用凭证无效或已吊销(Bearer token 不存在/Key 已停用)检查 Key 是否正确;若已吊销联系运维轮换;不重试
11002模型不可用(不存在/已禁用/不在白名单)检查 model 字段;无指定时检查是否有默认模型;不重试
11003默认对话模型未配置(未传 model 且无默认模型)联系平台运维配置默认模型;不重试
11004限速(RPM 或 TPM 超限)退避重试(建议 1 秒后重试);与 10008 严格区分
11005调用方日金额超限(moneyCapDaily 耗尽)等次日重置;联系运维调整上限;不重试
11006上游不可用(failover 候选全部耗尽)退避重试(建议 5 秒后重试,最多 3 次);持续则告警
11007内容策略阻断(内容安全评估为 BLOCK)检查输入内容;默认策略为 FLAG 不阻断,BLOCK 需管理端开启
11008图片任务不存在或已过期不适用于对话接口,属图片任务接口错误码
11009必传参数缺失(braId/idempotencyKey/messages检查请求体;不重试
10008门店 token 配额超额(半硬限)本轮调用已完成不掐断;下一次才拒。提示用户当日额度耗尽
10009门店暂停(风控主开关触发)/ 在途切断不重试;等待门店风控恢复;通知用户服务暂时不可用

6. Python 代码示例

6.1 非流式调用

import os
import uuid
import httpx

RELAY_BASE_URL = os.environ["NEARGO_RELAY_BASE_URL"]  # e.g. http://neargo-token:8080
SERVICE_KEY    = os.environ["NEARGO_RELAY_SERVICE_KEY"]

def chat_non_streaming(bra_id: str, messages: list[dict], model: str = None) -> dict:
    """非流式对话调用,等待全量结果后返回。"""
    payload = {
        "braId": bra_id,
        "idempotencyKey": str(uuid.uuid4()),
        "messages": messages,
        "stream": False,
    }
    if model:
        payload["model"] = model

    headers = {
        "Authorization": f"Bearer {SERVICE_KEY}",
        "Content-Type": "application/json",
    }

    response = httpx.post(
        f"{RELAY_BASE_URL}/internal/relay/chat",
        json=payload,
        headers=headers,
        timeout=60.0,
    )
    response.raise_for_status()
    data = response.json()

    # 非流式错误也会以 HTTP 4xx/5xx 返回,raise_for_status 处理;
    # 200 时可能仍有业务错误码(理论上非流式错误走 4xx,此处保留防御)
    if "code" in data and data.get("code") not in (None, 0, 200):
        raise RuntimeError(f"Relay error {data['code']}: {data.get('message')}")

    return data  # 含 content / traceId / model / usage 等字段


# 使用示例
if __name__ == "__main__":
    result = chat_non_streaming(
        bra_id="B0001",
        messages=[
            {"role": "system", "content": "你是专业的零售顾问。"},
            {"role": "user", "content": "库存周转率低怎么改善?"},
        ],
    )
    print(result["content"])
    print("trace:", result["traceId"])
    print("tokens:", result["usage"])

6.2 SSE 流式调用

import os
import uuid
import json
import httpx

RELAY_BASE_URL = os.environ["NEARGO_RELAY_BASE_URL"]
SERVICE_KEY    = os.environ["NEARGO_RELAY_SERVICE_KEY"]

def chat_streaming(bra_id: str, messages: list[dict], model: str = None):
    """
    SSE 流式对话调用,逐 token 输出。

    事件类型:
      meta       - 首帧,含实际命中的 model / channelCode / traceId
      token      - 内容增量(data.delta)
      tool_call  - 工具调用请求(data.toolCallId / name / arguments)
      usage      - 用量(data.promptTokens / completionTokens / usageSource)
      done       - 正常结束
      error      - 错误结束(data.code / message)
      forceClose - 风控在途切断(data.reason)
    """
    payload = {
        "braId": bra_id,
        "idempotencyKey": str(uuid.uuid4()),
        "messages": messages,
        "stream": True,
    }
    if model:
        payload["model"] = model

    headers = {
        "Authorization": f"Bearer {SERVICE_KEY}",
        "Content-Type": "application/json",
        "Accept": "text/event-stream",
    }

    with httpx.stream(
        "POST",
        f"{RELAY_BASE_URL}/internal/relay/chat",
        json=payload,
        headers=headers,
        timeout=120.0,
    ) as response:
        response.raise_for_status()

        event_name = None
        content_buffer = []

        for line in response.iter_lines():
            if line.startswith("event:"):
                event_name = line[len("event:"):].strip()
            elif line.startswith("data:"):
                data_raw = line[len("data:"):].strip()
                try:
                    data = json.loads(data_raw)
                except json.JSONDecodeError:
                    continue

                if event_name == "meta":
                    # 首帧:实际命中的模型和渠道
                    print(f"[meta] traceId={data.get('traceId')} model={data.get('model')} "
                          f"channel={data.get('channelCode')}")

                elif event_name == "token":
                    # 内容增量(注意:事件名为 "token",非 "delta")
                    delta = data.get("delta", "")
                    content_buffer.append(delta)
                    print(delta, end="", flush=True)

                elif event_name == "tool_call":
                    # 模型请求调用工具,调用方负责执行
                    print(f"\n[tool_call] id={data.get('toolCallId')} "
                          f"name={data.get('name')} args={data.get('arguments')}")
                    # TODO: 执行工具,下一轮调用时将结果通过 role=tool 消息回灌

                elif event_name == "usage":
                    # 用量统计(流尾,仅出现一次)
                    print(f"\n[usage] prompt={data.get('promptTokens')} "
                          f"completion={data.get('completionTokens')} "
                          f"source={data.get('usageSource')}")

                elif event_name == "done":
                    # 正常结束
                    print("\n[done]")

                elif event_name == "error":
                    # 错误结束
                    raise RuntimeError(
                        f"Relay error {data.get('code')}: {data.get('message')}"
                    )

                elif event_name == "forceClose":
                    # 风控在途切断
                    print(f"\n[forceClose] reason={data.get('reason')}")
                    raise RuntimeError("门店 AI 服务当前不可用(风控切断)")

            elif line == "":
                # 空行表示一帧结束,重置事件名
                event_name = None

        return "".join(content_buffer)


# 使用示例
if __name__ == "__main__":
    full_text = chat_streaming(
        bra_id="B0001",
        messages=[
            {"role": "system", "content": "你是专业的零售顾问。"},
            {"role": "user", "content": "如何提高客单价?"},
        ],
    )

6.3 幂等键生成建议

每次调用均应生成新的 idempotencyKey,推荐使用 UUID v4:

import uuid
idempotency_key = str(uuid.uuid4())  # "550e8400-e29b-41d4-a716-446655440000"

如果在请求发出后网络超时不确定是否到达,重发时复用同一个 idempotencyKey,中转站保证同一 Key 仅落库一条明细(幂等语义)。若确定需要新的调用,则生成新的 Key。


7. 切换步骤

以下步骤覆盖从"Python 直连上游模型"到"经中转站 RPC v2"的完整迁移过程。

步骤一:配置服务 Key

  1. 联系平台运维,为 python-agent 调用方创建服务 Key(管理端 /internal/admin/relay/callers

  2. 接收运维回显的一次性明文 Key(格式:ngr_...

  3. 将 Key 写入 Python 服务的 Secret 管理(不得写入代码仓库):

    # 本地开发(临时)
    export NEARGO_RELAY_SERVICE_KEY="ngr_xxxxxxxxxxxxxxxxxxxxxxxxxxxx"
    
    # 生产(写入 K8s Secret / 云 Secret Manager 等)
    kubectl create secret generic neargo-relay \
      --from-literal=service-key="ngr_xxxxxxxxxxxxxxxxxxxxxxxxxxxx"
    
  4. 写入 Nacos 或应用配置中中转站 base URL:

    # 配置示例
    neargo:
      relay:
        base-url: http://neargo-token:8080
        service-key: ${NEARGO_RELAY_SERVICE_KEY}
    

步骤二:指向 RPC v2 端点

将 Python 代码中直接调用上游模型的部分(OpenAI SDK / httpx 直调 DeepSeek 等)替换为调用中转站:

# 旧(直连,示例)
# client = openai.OpenAI(api_key=os.environ["DEEPSEEK_API_KEY"], base_url="https://api.deepseek.com")
# response = client.chat.completions.create(...)

# 新(经中转站)
result = chat_streaming(bra_id=current_bra_id, messages=messages_history)

关键对应关系

旧方式新方式(RPC v2)
API Key 从环境变量读取服务 Key(ngr_...)从环境变量读取
OpenAI SDK model 参数model 字段(可选,缺省走默认模型)
messages 数组messages 数组(格式相同)
stream=True/Falsestream 字段(相同语义)
无 braId 概念braId 必传(门店标识)
无幂等要求idempotencyKey 必传(每次调用生成 UUID)
直接获得 usageSSE: usage 事件;非流式: usage 字段

步骤三:发送 braId 和 idempotencyKey

RPC v2 面的两个必传字段

  • braId:门店标识,Python 服务需要从调用上下文中传入。若 Python 服务处理的是门店相关请求,需确保每次请求时拿到正确的 braId
  • idempotencyKey:每次业务调用生成一个新的 UUID(重试复用同一个 Key)。

步骤四:解析 SSE 事件

按第 4.1 节的事件字典解析下行事件:

  • meta 首帧:记录实际命中的模型名(meta.model),供对账和调试使用。
  • token 事件(注意:不是 delta):累积 data.delta 字段拼接完整回复。
  • usage 事件:流尾一次,提取 promptTokens/completionTokens/usageSource
  • error 事件:处理各错误码(见第 5 节)。
  • forceClose 事件:门店被切断,向用户展示友好提示。

步骤五:验证

按如下要点验证切换效果:

验证点方法
Key 认证正常发送一次带正确 Key 的请求,观察 meta 首帧
默认模型路由不传 model,确认 meta.model 为默认模型名
幂等去重用同一 idempotencyKey 发两次,确认明细表只有一条
usage 事件可读流尾确认 promptTokens/completionTokens/usageSource 有值
braId 缺失时拒绝不传 braId,预期 error 事件 code=11009

8. 回退步骤

过渡窗口

过渡期间,v1 端点 POST /internal/model/invoke(002 版本,agentType 驱动)与 v2 并存。 在全部调用方(neargo-agent Java + Python)均切换至 v2 并验收通过后,v1 才会返回 410 退役。

回退方法

  1. 配置切换:将 Python 服务的 base URL 和请求逻辑切回 v1 路径:

    # 回退:调用 v1 端点(过渡窗内有效)
    RELAY_URL = f"{NEARGO_TOKEN_BASE_URL}/internal/model/invoke"
    
  2. 无数据残留风险

    • relay_* 系列表为 003 新增,v1 路径不写这些表,回退无数据残留
    • agent_usage_record_* 明细表的 6 个扩展列(caller_code/channel_id/cost_amount/latency_ms/trace_id/safety_flag)为可空列,v1 路径不写这些列,向后兼容
  3. 参考specs/003-model-relay-station/quickstart.md §5 切换与回滚 有完整的过渡窗说明。

注意:回退操作影响计量和限速治理(v1 路径无 RPM/TPM/金额上限治理),请评估影响后决策。


9. 已知差异:实现 vs 契约文档

在阅读 specs/003-model-relay-station/contracts/relay-rpc-chat.md §2 时,需注意一处实现与文档的差异:

项目契约文档(relay-rpc-chat.md §2)实际实现(ModelInvokeEvent.java
增量文本事件名deltatokeneventName() 返回 "token"
增量文本字段{content}{delta}Token record 字段名为 delta

实际客户端须订阅 event: token,数据字段为 data.delta,而非契约文档描述的 event: delta。本文档已按实现(token)描述,优先级高于契约文档。


10. 附录:准入顺序

转发前,中转站按以下顺序逐一检查(任一不过即拒,不调用上游模型):

  1. 凭证校验:Bearer Key 是否有效 → 11001
  2. 必传参数braId / idempotencyKey / messages → 11009
  3. 内容安全:SafetyPipeline 评估(默认 FLAG 不阻断,BLOCK 时 → 11007)
  4. 门店暂停:风控主开关 → 10009
  5. 门店 token 配额:quota_check(半硬限)→ 10008
  6. 调用方金额:money_check → 11005
  7. RPM/TPM 限速:→ 11004
  8. 路由解析:model → 11002 / 11003 / 11006

半硬限语义:已通过准入的调用不中断;超额后下一次拒绝(本轮跑完)。


如有问题,请联系 neargo-token 服务维护团队,并提供 traceId 便于排查。