Skip to content

Commit e4f34d6

Browse files
committed
fix: harden multi-tenant gateway edge cases
1 parent 7ab06cf commit e4f34d6

15 files changed

Lines changed: 629 additions & 141 deletions

File tree

‎examples/multi_tenant_im_agent/ARCHITECTURE.zh_CN.md‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -119,6 +119,8 @@ Session event → state → summary 的更新规则:原始事件先获得不
119119

120120
推荐生产组合:SQL 保存控制面和不可丢事件,Redis 保存热 Session,向量库保存 Knowledge/Memory,对象存储保存 Artifact。Redis 更新使用 Lua 或 CAS 版本,SQL 使用 `SELECT FOR UPDATE`/乐观版本,禁止“读整个 Session 后无条件覆盖”的丢更新模式。
121121

122+
本地 SQLite 仅用于单进程验收;仓库在进程内用写锁补足 SQLite 忽略 `FOR UPDATE` 的线程语义,但不宣称支持多个进程共享同一 SQLite 文件。多副本部署必须使用 MySQL/PostgreSQL,迁移中的 MySQL 时间列保留微秒精度。
123+
122124
跨节点可见性由共享 Redis/SQL Session 后端直接保证。若生产部署另加 Worker 本地只读缓存,建议在 Memory 写入成功后发布 `tenant/session/memory_version` 失效通知,并用短 TTL 与读取版本号兜底;该本地缓存层不属于本示例的已实现范围。
123125

124126
## 7. 数据迁移
@@ -170,6 +172,7 @@ Redis → SQL:按 Session 扫描,不使用生产 `KEYS *`;以版本 CAS
170172
- 日志和 Trace 禁止记录正文、Authorization、IM token、模型 Key 和数据库密码。
171173
- HMAC namespace secret 定期轮换时需支持 current/previous 两个版本的读取窗口。
172174
- Admin API 生产中接入 mTLS/OIDC、RBAC、来源网段和操作审计;示例 token 仅用于最小演示。
175+
- `/metrics` 同样要求 `X-Metrics-Token`,避免公开按租户聚合的业务量、Token 和成本标签。
173176
- 容器只读根文件系统、非 root、丢弃 Linux capabilities;工具执行放独立沙箱池,禁止与 Gateway 共进程。
174177

175178
## 10. 可观测性与审计
@@ -208,7 +211,7 @@ im.callback → tenant.resolve → signature.verify → session.lease
208211
| IM 回复失败 | Agent 结果和 Outbox 已提交,后台重试,不再次运行 Agent |
209212
| 配置错误 | 原子热更新拒绝整批错误配置,继续使用上一版本 |
210213

211-
Outbox 第 8 次投递仍失败后转为 `dead_letter` 状态,指标可直接告警;人工重放必须保留原 outbox_id。
214+
Outbox 第 8 次投递仍失败后转为 `dead_letter` 状态,指标可直接告警;人工重放必须保留原 outbox_id。外发是 at-least-once:若 IM 平台已接收、但确认响应丢失,后台恢复可能再次发送;能提供幂等键的平台应绑定 `outbox_id`,否则需结合发送回执对账。
212215

213216
## 12. 部署、灰度和回滚
214217

‎examples/multi_tenant_im_agent/EVALUATION.zh_CN.md‎

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@
88
python examples/multi_tenant_im_agent/scripts/judge_demo.py
99
```
1010

11-
脚本启动独立 FastAPI Gateway 和临时数据库,从 HTTP 边界验证健康检查、Telegram/企业微信验签、账号路由、重复投递抑制、同 ID 异载荷冲突、Admin 鉴权和 Prometheus 指标,然后自动停止服务并清理数据库。离线确定性模型仍经过真实 tRPC-Agent `LlmAgent + Runner`,不需要模型或 IM 凭据。
11+
脚本启动独立 FastAPI Gateway 和临时数据库,从 HTTP 边界验证健康检查、Telegram/企业微信验签、账号路由、重复投递抑制、同 ID 异载荷冲突、Admin/Prometheus 鉴权,然后自动停止服务并清理数据库。离线确定性模型仍经过真实 tRPC-Agent `LlmAgent + Runner`,不需要模型或 IM 凭据。
1212

1313
完整单元与契约测试:
1414

@@ -25,21 +25,23 @@ pytest examples/multi_tenant_im_agent/tests -q
2525
| 水平扩展 | 无 sticky session;确定性 Session ID;SQL Session 租约串行化;Redis/SQL 共享 Session | `domain.py`、`repository.py`、`runtime.py` | 租约互斥、用户/群聊/租户隔离测试 |
2626
| 租户隔离 | 配置、复合外键、IM 账号、Session HMAC、工具 allowlist、日志脱敏、密钥仅引用环境变量 | `repository.py`、`runtime.py`、`governance.py` | 跨租户相同消息号、相同应用 ID、未知工具 fail-closed 测试 |
2727
| 多后端 | tRPC-Agent InMemory、Redis、SQL Session Service;Memory、Summary、Artifact、Knowledge、Audit 数据模型 | `runtime.py`、`repository.py` | 后端配置校验及 Alembic/ORM 契约测试 |
28-
| 数据一致性 | Session 行锁递增序列、幂等声明、事务内消息完成与 Outbox 提交、租约恢复 | `repository.py` | 重复消息、payload conflict、失败重试、Outbox 恢复测试 |
28+
| 数据一致性 | Session 行锁递增序列、幂等声明、消息完成/Outbox/预算结算同事务、处理中断与租约恢复 | `repository.py` | 并发重复、payload conflict、处理中断、事务回滚、Outbox 恢复测试 |
2929
| IM 软件接入 | Telegram 与企业微信 Adapter;验签、解析、群/单聊 Session、长度限制和外发转换 | `adapters.py` | 两通道签名与归一化测试、真实 HTTP 黑盒验收 |
3030
| 幂等范围 | `(tenant_id, channel, account_id, external_message_id)`,同租户多机器人不会误冲突 | `repository.py` | 跨账号相同外部消息 ID 测试 |
3131
| 治理安全 | 用户白名单、请求大小、输入和 Token 预算、模型超时、工具默认禁用、Runner Filter 二次校验 | `governance.py`、`runtime.py`、`app.py` | 策略拒绝、预算拒绝、未知工具拒绝测试 |
32-
| 预算与成本 | 数据库行锁原子预留月预算,模型完成后按 usage metadata 结算;审计和 Prometheus 汇总 Token/成本 | `repository.py`、`service.py`、`telemetry.py` | 实际用量结算和 provider usage metadata 测试 |
32+
| 预算与成本 | 数据库行锁原子预留月预算,模型完成后按 usage metadata 结算(缺失时含输出长度估算);结算与结果事务原子提交 | `repository.py`、`service.py`、`telemetry.py` | 实际/未知用量、提交失败回滚和 provider usage metadata 测试 |
3333
| 可观测性 | OpenTelemetry 根 Span 并继承 Runner/模型/工具 Span;请求、阶段延迟、投递、审计、Token、成本指标 | `telemetry.py`、`service.py` | Metrics HTTP 验收与安全 Trace 测试 |
3434
| 审计 | 题目要求字段齐全;用户 ID HMAC;不保存正文、Key、Token;审计故障不触发模型重放 | `repository.py`、`service.py` | 审计字段与审计库故障测试 |
35-
| 故障恢复 | 模型失败可重投;回复结果先入 Outbox;指数退避、稳定抖动、8 次后死信;崩溃租约恢复 | `repository.py`、`service.py` | 模型失败恢复、投递恢复、死信测试 |
35+
| 故障恢复 | 模型失败可重投;处理中事件凭新租约恢复;回复先入 Outbox;指数退避、稳定抖动、8 次后死信 | `repository.py`、`service.py` | 模型/处理中断恢复、投递恢复、死信测试 |
3636
| 灰度和回滚 | 配置版本、按 tenant/account 路由 canary;迁移先行;RollingUpdate、HPA、PDB、探针 | `ARCHITECTURE.zh_CN.md`、`deploy/kubernetes.yaml` | 部署清单可静态检查 |
3737
| 最小/生产部署 | Docker Compose 最小栈;Kubernetes 迁移 Job、Gateway、Redis、MySQL、Collector 拓扑 | `docker-compose.yml`、`deploy/kubernetes.yaml` | `/healthz`、`/readyz` 与迁移契约测试 |
3838

3939
## 3. 最容易忽略的工程细节
4040

4141
- IM 平台重投已完成消息时只返回幂等成功,不会再次调用模型、工具或发送回复。
42-
- Agent 结果和投递 Outbox 在同一数据库事务完成;IM 故障不会导致昂贵的模型调用重放。
42+
- Agent 结果、投递 Outbox 和预算结算在同一数据库事务完成;IM 故障不会导致昂贵的模型调用重放。
43+
- 流式模型优先使用终态完整文本,超时显式关闭异步流,空文本有安全回退。
44+
- `/metrics` 与 Admin API 都要求常量时间比较的 Token,不公开租户用量标签。
4345
- 审计写入故障被独立计数并脱敏记录,不会把已经完成的消息错误标成 failed。
4446
- 月度预算不是“先查询再判断”,而是在数据库事务中预留,避免多 Worker 并发超额。
4547
- 离线验收使用确定性模型,但模型仍由真正的 Runner 执行;在线模式只替换模型实例和 Session 后端。
@@ -51,5 +53,6 @@ pytest examples/multi_tenant_im_agent/tests -q
5153
- Memory/Knowledge/Artifact 的跨介质迁移给出可执行阶段、数据模型和一致性规则,但不附带特定云厂商账号。
5254
- 示例工具注册表默认为空并 fail-closed;真实危险工具需要业务审批系统,不能用演示确认按钮代替。
5355
- 多区域需要 home-region 或全局序列服务;本示例不声称单数据库能够解决跨区域一致性。
56+
- IM 外发采用业界常见的 at-least-once 语义;若平台已接收但响应在网络中丢失,仍可能重复投递,生产应优先传递平台支持的幂等键或做回执对账。
5457

5558
以上边界不会影响本地自动验收,并避免为了演示而提交平台账号、生产密钥或不可验证的云资源。

‎examples/multi_tenant_im_agent/README.zh_CN.md‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@
1212
- SQL 唯一幂等键,防止 IM 重复投递造成模型和工具重复执行。
1313
- 数据库 Session 租约串行化同一会话,Worker 无状态且不依赖 sticky session。
1414
- tRPC-Agent Session 后端可按租户选择 InMemory、Redis 或 SQL。
15-
- 回复与 Outbox 同事务提交;投递失败后台重试,Worker 崩溃后可恢复过期任务。
15+
- 回复、Outbox 与 Token 预算结算同事务提交;投递失败后台重试,Worker 崩溃后可恢复过期消息和投递任务。
1616
- 租户级用户白名单、输入长度、单请求/月度 Token 预算,并在 Runner 内增加真实的 tRPC-Agent Filter 二次 fail-closed 校验。
1717
- 全字段审计表、Prometheus 指标、OpenTelemetry OTLP Trace。
1818
- Docker Compose 最小部署与 Kubernetes 生产部署样例。
@@ -43,7 +43,7 @@ examples/multi_tenant_im_agent/scripts/run_offline_demo.ps1
4343

4444
- `GET /healthz`:进程存活探针。
4545
- `GET /readyz`:数据库就绪探针。
46-
- `GET /metrics`:Prometheus 文本指标。
46+
- `GET /metrics`:需 `X-Metrics-Token`(与 Admin Token 同值)的 Prometheus 文本指标。
4747
- `GET /admin/tenants`:需 `X-Admin-Token`,仅返回不含密钥的租户摘要。
4848
- `POST /webhooks/telegram/acme-support-bot`:Telegram 回调。
4949
- `POST /webhooks/wecom/acme-wecom-app`:企业微信回调。
@@ -97,7 +97,7 @@ uv run --no-project --with-editable . python examples/multi_tenant_im_agent/scri
9797
pytest examples/multi_tenant_im_agent/tests -q
9898
```
9999

100-
测试覆盖租户路由冲突、会话隔离、Telegram/企业微信验签、用户策略、幂等重投、载荷冲突、失败恢复、Session 租约、Outbox 重试和 HTTP 健康检查。
100+
当前 39 项测试覆盖租户路由冲突、会话隔离、Telegram/企业微信验签、并发重复、处理中断恢复、用户策略、幂等重投、载荷冲突、预算事务回滚、流式终态回复、超时资源关闭、Session 租约、Outbox 重试、MySQL 迁移契约、请求体上限和 HTTP 鉴权。
101101

102102
## 目录
103103

‎examples/multi_tenant_im_agent/adapters.py‎

Lines changed: 12 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -76,7 +76,9 @@ def parse(
7676
) -> InboundMessage:
7777
try:
7878
update = json.loads(raw_body)
79-
message = update.get("message") or update.get("edited_message")
79+
# Edits carry a new update_id and would otherwise rerun the model
80+
# for the same logical message. Treat them as unsupported events.
81+
message = update.get("message")
8082
if not isinstance(message, dict):
8183
raise UnsupportedMessageError(
8284
"Telegram update has no supported message"
@@ -253,11 +255,14 @@ class HttpChannelSender:
253255
"""Production delivery client with bounded timeouts and no secret logging."""
254256

255257
def __init__(self, timeout_seconds: float = 10.0):
256-
self.timeout_seconds = timeout_seconds
257-
258-
async def send(self, request: DeliveryRequest) -> Mapping[str, Any]:
259258
import httpx
260259

260+
self._client = httpx.AsyncClient(timeout=timeout_seconds)
261+
262+
async def close(self) -> None:
263+
await self._client.aclose()
264+
265+
async def send(self, request: DeliveryRequest) -> Mapping[str, Any]:
261266
if request.channel == "telegram":
262267
token = require_secret(request.credentials_env.get("bot_token", ""))
263268
url = f"https://api.telegram.org/bot{token}/sendMessage"
@@ -275,10 +280,9 @@ async def send(self, request: DeliveryRequest) -> Mapping[str, Any]:
275280
f"no HTTP sender configured for channel: {request.channel}"
276281
)
277282

278-
async with httpx.AsyncClient(timeout=self.timeout_seconds) as client:
279-
response = await client.post(url, json=payload)
280-
response.raise_for_status()
281-
result = response.json()
283+
response = await self._client.post(url, json=payload)
284+
response.raise_for_status()
285+
result = response.json()
282286
return result if isinstance(result, dict) else {"accepted": True}
283287

284288

‎examples/multi_tenant_im_agent/app.py‎

Lines changed: 24 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,9 @@ async def lifespan(app: FastAPI):
5252
close = getattr(service.runtime, "close", None)
5353
if close is not None:
5454
await close()
55+
close_sender = getattr(service.sender, "close", None)
56+
if close_sender is not None:
57+
await close_sender()
5558
await asyncio.to_thread(service.repository.close)
5659

5760
app = FastAPI(
@@ -73,7 +76,10 @@ async def readyz():
7376
return {"status": "ready"}
7477

7578
@app.get("/metrics", response_class=PlainTextResponse, tags=["operations"])
76-
async def metrics():
79+
async def metrics(x_metrics_token: str = Header(default="")):
80+
expected = require_secret(admin_token_env)
81+
if not hmac.compare_digest(x_metrics_token, expected):
82+
raise HTTPException(status_code=403, detail="forbidden")
7783
return service.metrics.render_prometheus()
7884

7985
@app.get("/admin/tenants", tags=["admin"])
@@ -85,9 +91,18 @@ async def tenants(x_admin_token: str = Header(default="")):
8591

8692
@app.post("/webhooks/{channel}/{account_id}", tags=["webhooks"])
8793
async def webhook(channel: str, account_id: str, request: Request):
88-
raw_body = await request.body()
89-
if len(raw_body) > 1_048_576:
94+
maximum = 1_048_576
95+
content_length = request.headers.get("content-length", "")
96+
if content_length.isdigit() and int(content_length) > maximum:
9097
raise HTTPException(status_code=413, detail="callback body too large")
98+
chunks: list[bytes] = []
99+
received = 0
100+
async for chunk in request.stream():
101+
received += len(chunk)
102+
if received > maximum:
103+
raise HTTPException(status_code=413, detail="callback body too large")
104+
chunks.append(chunk)
105+
raw_body = b"".join(chunks)
91106
response = await service.handle_webhook(
92107
channel=channel,
93108
account_id=account_id,
@@ -115,14 +130,18 @@ def build_app_from_env() -> FastAPI:
115130
)
116131
registry = load_tenant_registry(config_path)
117132
offline = os.environ.get("OFFLINE_ECHO_MODE", "false").lower() == "true"
133+
require_secret("ADMIN_API_TOKEN")
134+
for tenant in registry.all():
135+
for binding in tenant.bindings:
136+
# Callback authentication is mandatory in both offline and online
137+
# modes; only outbound/provider credentials may be skipped offline.
138+
require_secret(binding.webhook_secret_env)
118139
if not offline:
119-
require_secret("ADMIN_API_TOKEN")
120140
for tenant in registry.all():
121141
require_secret(tenant.model_api_key_env)
122142
if tenant.session_backend is not StorageBackend.MEMORY:
123143
require_secret(tenant.session_dsn_env)
124144
for binding in tenant.bindings:
125-
require_secret(binding.webhook_secret_env)
126145
if binding.channel.lower() == "telegram":
127146
require_secret(binding.bot_token_env)
128147
elif binding.channel.lower() == "wecom":

‎examples/multi_tenant_im_agent/domain.py‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -102,7 +102,9 @@ def payload_hash(self) -> str:
102102
@dataclass(frozen=True)
103103
class AgentReply:
104104
text: str
105-
token_count: int = 0
105+
# ``None`` means that the provider omitted usage metadata. A real zero is
106+
# kept distinct so callers never mistake it for an unknown value.
107+
token_count: int | None = None
106108
cost: float = 0.0
107109
tool_names: tuple[str, ...] = ()
108110

0 commit comments

Comments
 (0)