games-development-ai/docs/plans/2026-07-01-配置控制面-spike与阶段〇生产基建-plan.md
lili a164f4c222 docs(.agents): 配置控制面阶段〇蒸馏 + 累积能力蒸馏批次 + 本轮 plan/设计
- tech-decisions §10:配置控制面阶段〇生产基建落地(Nacos 2.4.3/RocketMQ 5.3.1/Sentinel
  自托管 mini-infra + 三层并发正交承重[Sentinel=入口QPS / consumeThreadMax=dispatch速率 /
  有界worker池=真并发cap follow-up]+ rocketmq-spring consumeThreadMax 无界队列坑:死参数须 min=max)
- AGENTS.md:Nacos/RocketMQ/Sentinel 反转为 MVP 生产 runtime(2026-07-01 build-vs-buy 现货尽调)
- 累积 .agents 蒸馏:agentscope-2.0-facts / build-vs-buy 硬门 / cheap-model-game-generation /
  contract-first-development / game-e2e-cdp-harness + README 索引
- 本轮留痕:配置控制面 spike 与阶段〇 plan + 设计文档

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-01 18:53:39 -07:00

1249 lines
74 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

---
date: 2026-07-01
topic: 配置控制面-spike与阶段〇
status: 已执行完成(SDD 8 任务全绿 + 最终整分支 review With fixes 已修,分支 merge-ready,2026-07-01)· ⚠️ 勘误(C3 review 纠正):本 plan 正文多处「并发权威=consumeThreadMax」framing 有误 —— consumeThreadMax/consumeThreadNumber=15 只 cap dispatch 投递速率;真「在跑生成并发≤15」=有界 worker 池(follow-up);现行口径 + 执行留痕见 `.superpowers/sdd/progress.md` · 原:§6.8 评审已过 → 执行
关联设计: docs/agent-specs/2026-06-30-配置控制面一次性按序实现-设计.md(§3.4 续修 middleware / §3.7 基建选型 / §4 阶段〇)
---
# 配置控制面 · 续修 spike 与阶段〇生产基建 · 实施计划
> **For agentic workers:** REQUIRED SUB-SKILL:用 superpowers:subagent-driven-development(荐)或 superpowers:executing-plans 逐任务执行。步骤用 `- [ ]` 复选框跟踪。**每个改代码的步骤都给了完整代码,不留占位。**
**Goal:** 一次性坐实「续修 middleware(拦 finish + 独立重跑门 + 注入续跑)」这条承重机制,并把生产基建三件套(Nacos / RocketMQ / Sentinel)自托管进 mini-infra、接进 game-cloud 与 AgentScope Service —— 为其后 cheap-worker 归并 Service(阶段一)铺好地基。
**Architecture:** 分四部分,依赖串行 + 部分并行。**Part A 续修 spike**(纯 Python、stub 模型、确定性、零外部依赖,本地跑)独立、先行去风险:若机制不成,回落现有 `control_plane` 外层循环,后续架构其余不受影响。**Part B 基建自托管**(Nacos + RocketMQ 起 mini-infra,lean JVM)。**Part C game-cloud 全量接 Java**(重启用 Nacos 配置/发现、Sentinel 全局并发≤15、RocketMQ 生产/消费),与既有自建准入 `enqueueWithControlPlane` + DB 轮询派发**分层组合**(不推翻 working code)。**Part D wiring**(AgentScope Service 注册进 Nacos discovery、Python worker middleware 经 nacos-sdk-python 热读预算/阈值)。
**Tech Stack:** AgentScope 2.0.2(`MiddlewareBase` 洋葱链)· Python 3.12(cheap-worker/.venv)· nacos-sdk-python **v1(`nacos.NacosClient`,同步 watcher,钉版本)** · Nacos 2.4.x(standalone + MySQL 后端)· Apache RocketMQ(rocketmq-spring 2.3.5,`consumeThreadMax` = 并发≤15 权威)· Sentinel(Spring Cloud Alibaba 2025.0.0.0,`FLOW_GRADE_QPS` 入口保护)· Java 17 + Spring Boot 3.5.14 · Docker(mini-infra 自托管)。
## Global Constraints(全局约束 · 每个任务隐含继承)
- **AgentScope 版本钉死 == 2.0.2**:已装 `cheap-worker/.venv` = 克隆 `.claude/skills/agentscope-skill/agentscope/`(github tag v2.0.2)。读的=跑的;任何 API 疑点以已装 2.0.2 源码为准,不凭 SoT 转述。
- **中文注释 + 可追溯日志**:所有代码带完整简体中文注释;外部交互 / 核心实现 / 错误路径必须有可追溯日志(创始人铁律)。
- **内网代理旁路(栽过的确切坑)**:系统代理是 fake-ip `198.18.x`;内网直连必须绕过它 —— ssh 用真 IP `root@100.64.0.8`(mini-infra)/ `root@100.64.0.7`(mini-desktop),curl 用 `--noproxy '*'`,Python 用 `urllib.request.build_opener(ProxyHandler({}))` 或把 host 塞进 `NO_PROXY/no_proxy`(须在 import openai/httpx 之前)。
- **内网密钥进项目文档、不进 env**:`NEWAPI_KEY` 等在 `docs/内网凭据与端点.md`;Python 侧经 `_bootstrap.ensure_api_key_env()` 解析注入。新增基建端点(Nacos/RocketMQ 地址)写回该文档,不入 env var、不硬编码散落。
- **mini-infra 部署 = 纯增量**:只新增 Nacos/RocketMQ 容器,**绝不动**既有 MySQL / PostgreSQL / Redis / MinIO / RAGflow / Gitea / new-api。mini-infra 实测 15G 总 / 8.6G 可用 / 4 核 / RAGflow 独占 ~9G → **JVM 必须 lean 调**(具体 -Xmx 见 Part B),不得用默认(Nacos 默认 2g / RocketMQ broker 默认 8g 会撑爆)。
- **门放行唯一权威 = middleware 独立重跑**:续修放行与否,唯一权威是 middleware 在 finish 点自己独立跑的门判,**绝不采信 agent 上报的门工具结果**(防便宜档漏调 / 谎报门绿的假绿)。spike 用 fake check 验机制;阶段一换真 `run_gates`。
- **spent_rmb 单 POST 内存累计、无需持久化**:续修移进洋葱 = 整局一个 POST、一个 middleware 实例贯穿,`spent_rmb` 实例内自然累计(`middleware.py:474`);AgentState 无自由 KV 槽,别 fork 它。
- **不动 live 服务无创始人授权**:Part C 改 game-cloud 配置 / Part B 部署,均先本地 / mini-infra 验,真部署到 staging 后端按创始人窗口。
- **测试约定(无 pytest.ini/conftest)**:cheap-worker 侧 `cheap-worker/.venv/bin/python cheap-worker/tests/<t>.py`(`__main__` runner + `sys.exit(1 if _failed else 0)`);tier2 侧 pytest 型 `PYTHONPATH=tier2/gen-worker cheap-worker/.venv/bin/python -m pytest <t> -v`。新单测跟 `cheap-worker/tests/test_budget_gate.py` 同构(mock、直构状态、零网络/LLM/chrome)。
- **提交风格**:`feat(<scope>): <中文> (切片一 阶段〇/PartX)`;结尾 `Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>`。仅创始人要求时才 commit/push。
- **落档散文标准**:本 plan 的散文段(Context/rationale/风险)按资深工程师散文写 —— 无元叙述、无 AI 造词黑话、无套话空强调。
---
## 执行顺序与依赖
```mermaid
graph LR
A[Part A · 续修 spike<br/>本地·独立·先行去风险] -.机制成立才继续续修线.- C2note[阶段一接真门]
B1[B1 · Nacos 自托管] --> B2[B2 · RocketMQ 自托管]
B1 --> C1[C1 · game-cloud 重启用 Nacos]
C1 --> C2[C2 · Sentinel 全局≤15]
B2 --> C3[C3 · RocketMQ 生产/消费]
C1 --> C3
B1 --> D1[D1 · Service 注册 Nacos]
B1 --> D2[D2 · Python middleware 热读 Nacos]
```
- **Part A 独立**,建议**最先做**:它去掉整条续修架构的最大不确定性,失败即回落 `control_plane`、其余不受影响。
- **Part B 是 Part C/D 的前置**:Nacos(B1)先起,Sentinel 规则源 / game-cloud 配置 / Service 发现 / Python 热读全依赖它;RocketMQ(B2)是 C3 前置。
- **Part C/D 在 B 之后**可并行推进(C1→C2 串,C3 依赖 B2+C1,D1/D2 依赖 B1)。
- Part C/D 的真部署验证走 mini-infra / staging 后端窗口;代码 + 单测 + 本地能验的先落。
---
# Part A · 续修 middleware spike(拦 finish + 注入 + 续跑)
**为什么先做:** 设计 §3.4 把「续修 / 门判」从框架外的 `control_plane` 外层循环搬进 AgentScope 洋葱里一个 `on_reasoning` middleware。源码级已核实这条机制稳(落在公开、有类型的 `MiddlewareBase.on_reasoning` 扩展点上 —— 其拦截/透传/改入参有 shipped 单测背书,而「吞 finish Msg 短路续跑」这一步 shipped 未覆盖、正是本 spike 要坐实的):finish = `_reasoning()` yield 一个纯文本 `Msg`(`agent/_agent.py:614-620`),它穿过 middleware 链才到达 reply 循环 → middleware 不 re-yield 即压制;注入靠 `agent.observe(UserMsg)`;续跑靠 reply 循环见本轮无 Msg 自动 `cur_iter++` 重推理(`_agent.py:670`),`max_iters` 封顶。**唯一要显式处理的行为细节**:finish 文本在 `_agent.py:857` 已先存进 `agent.state.context`(早于 yield),压制 yield 抹不掉它 —— 对续修语义是自然的(模型看到自己的声称 + 纠正);要干净血缘需配 `context.pop()`。spike 就是把这条机制在项目里真跑一遍、坐实,并把两个「若在此栽即回落」的点真跑断言掉。
### Task A1:续修 middleware 机制 spike(确定性、stub 模型、fake 门)
**Files:**
- 参照(只读):`.claude/skills/agentscope-skill/agentscope/tests/utils.py`(stub 模型 `MockModel`/`MockCredential` 的 2.0.2 verbatim 范式,约 :38-138)、`.claude/skills/agentscope-skill/agentscope/tests/middleware_test.py`(shipped 单测:`on_reasoning` 透传/见 Msg :117-188、`on_reasoning` 改入参 :563-619、`on_compress_context` 短路不调 next_handler :896-952 —— 「短路」机制的类比证据;on_reasoning **吞 finish Msg 续跑** 无 shipped 单测、正是本 spike 坐实点)
- 参照(只读):`tier2/gen-worker/worker/middleware.py:551`(`Tier2TraceMiddleware.on_reasoning` = 项目内已有的 `on_reasoning` 实现范式)、`cheap-worker/tests/test_budget_gate.py`(单测文件结构 + `__main__` runner 模板)
- Create:`cheap-worker/tests/test_repair_middleware_spike.py`(spike 本体 + 断言;绿了即成机制回归)
**Interfaces:**
- Consumes:AgentScope 2.0.2 公开 API —— `Agent(name, system_prompt, model, toolkit, middlewares, react_config)`、`MiddlewareBase.on_reasoning(self, agent, input_kwargs, next_handler) -> AsyncGenerator`、`agent.observe(UserMsg)`、`agent.reply(UserMsg)`、`agent.state.context: list[Msg]`、`ChatResponse(content=[TextBlock(...)], is_last=True)`、`ReActConfig(max_iters=...)`。注:`UserMsg`/`AssistantMsg` 是**工厂函数**、返回 `Msg` 实例(非类),`isinstance(x, Msg)` 成立。
- Produces:`RepairMiddleware`(spike 内 inline;阶段一提升进 `tier2/gen-worker/worker/middleware.py` 作 `CircuitBreakerMiddleware` 的兄弟)—— `on_reasoning` 里拦 `Msg`→`await check(agent)`→过/耗尽则 `yield item` 放行、否则不 yield + `await agent.observe(UserMsg)` 注入 + `return`。`check` 契约 = `async (agent) -> tuple[bool, str]`(passed, feedback)。
- [ ] **Step 1:核准 2.0.2 的 stub 模型范式(读 shipped 源,不符则以它为准替换)**
先读 `.claude/skills/agentscope-skill/agentscope/tests/utils.py` 与 `middleware_test.py`,核对 `MockModel` / `MockCredential` 的确切构造(`ChatModelBase.__init__` 的参数:`credential` / `model` / `parameters` / `stream`,及是否有 `context_size` 等带默认的参)与 `on_reasoning` 单测用法。已装包(`.venv`)不含 tests/,故从克隆源读。**Step 2 下面给的 stub 是比照 utils.py 的忠实版;若其构造签名与 utils.py 不符,直接用 utils.py 的 `MockModel`/`MockCredential` verbatim 替换 Step 2 那段** —— 目的是 stub 与 shipped 单测逐字一致、避免 API 漂移。
- [ ] **Step 2:写 spike 单测文件(核心断言:拦一次 + 注入 + 续跑放行第二版)**
Create `cheap-worker/tests/test_repair_middleware_spike.py`:
```python
"""
test_repair_middleware_spike.py — 续修 middleware 机制 spike(配置控制面 §3.4 前置去风险)。
坐实 AgentScope 2.0.2 的 MiddlewareBase.on_reasoning 能:
① 拦到 agent 想 finish(_reasoning yield 纯文本 Msg,_agent.py:614)那一刻;
② 独立跑一个外部 check()(此处 fake、零 LLM/chrome,只验机制,不绑真门);
③ 没过 → 压制 finish(不 re-yield)+ 注入「修」(UserMsg 纯文本,agent.observe)→ agent 自动续跑一轮;
④ 过 / 达 max_repairs → 放行 finish,reply 循环收尾。
机制稳:finish Msg 穿过 middleware 链才到 reply 循环(_agent.py:615-620),故 middleware 能在它到达前吞掉;
续跑靠 reply 循环见本轮无 Msg 自动 cur_iter++ 重推理(_agent.py:670),max_iters 封顶。
坑(须处理·非阻断):finish 文本在 _agent.py:857 已先存进 context,压制 yield 抹不掉 —— 要干净血缘配 context.pop()。
绿了即成机制回归;阶段一把 RepairMiddleware 提升进 tier2/gen-worker/worker/middleware.py、fake check 换真 run_gates 独立重跑。
跑:cheap-worker/.venv/bin/python cheap-worker/tests/test_repair_middleware_spike.py
"""
import asyncio
import sys
from pathlib import Path
from typing import Any, AsyncGenerator, Callable, Type
sys.path.insert(0, str(Path(__file__).resolve().parents[1])) # → cheap-worker/
import _bootstrap # noqa: E402,F401 仅加 sys.path(import 时不取 key、不触网)
from agentscope.agent import Agent, ReActConfig # noqa: E402
from agentscope.middleware import MiddlewareBase # noqa: E402
from agentscope.model import ChatModelBase, ChatResponse # noqa: E402
from agentscope.message import Msg, TextBlock, UserMsg # noqa: E402
from agentscope.tool import Toolkit # noqa: E402
from agentscope.credential import CredentialBase # noqa: E402
from pydantic import BaseModel # noqa: E402
# ── stub 模型(确定性,不打真 LLM);按 Step 1 从 tests/utils.py 核准的 2.0.2 范式抄 ──
# 若 ChatModelBase.__init__ 签名与此不符,以 .claude/skills/agentscope-skill/agentscope/tests/utils.py 为准调整。
class _MockCredential(CredentialBase):
@classmethod
def get_chat_model_class(cls) -> Type["ChatModelBase"]:
return MockModel
class MockModel(ChatModelBase):
class Parameters(BaseModel):
...
def __init__(self) -> None:
super().__init__(credential=_MockCredential(), model="mock",
parameters=MockModel.Parameters(), stream=False)
self._resp: list[ChatResponse] = []
self.cnt = 0
def set_responses(self, r: list[ChatResponse]) -> None:
self._resp = r
self.cnt = 0
self.stream = False
async def _call_api(self, *a: Any, **k: Any) -> ChatResponse:
r = self._resp[self.cnt]
self.cnt += 1 # 每次调用取下一条(模拟多轮)
return r
# ── 续修 middleware(spike 核心验证对象;阶段一提升进 worker/middleware.py)──
class RepairMiddleware(MiddlewareBase):
"""on_reasoning 拦 finish → 独立跑 check → 没过则压制 + 注入续跑。check 契约:async (agent)->(passed, feedback)。"""
def __init__(self, check: Callable, max_repairs: int = 3, *, pop_finish_claim: bool = False) -> None:
self._check = check
self._max_repairs = max_repairs
self._pop_finish_claim = pop_finish_claim # True=续修前把「做完了」从 context 弹掉(干净血缘)
self.repairs = 0 # 供断言/预算读
async def on_reasoning(self, agent, input_kwargs, next_handler) -> AsyncGenerator:
async for item in next_handler(**input_kwargs): # 与 Tier2TraceMiddleware.on_reasoning 同款透传入参
if not isinstance(item, Msg):
yield item # 非 finish 的事件流(ModelCallStart/Text* 等)原样透传
continue
# —— 拦到 finish(纯文本 Msg)——
passed, feedback = await self._check(agent)
if passed or self.repairs >= self._max_repairs:
yield item # 放行 finish → reply 循环收尾 return
return
self.repairs += 1
if self._pop_finish_claim and agent.state.context and isinstance(agent.state.context[-1], Msg):
agent.state.context.pop() # 可选:弹掉刚存进的「做完了」声称(_agent.py:857 先存)
# 注入「修」:role=user 纯文本(禁 system/tool/thinking,否则 _handle_incoming_messages 抛 ValueError)
await agent.observe(UserMsg(name="gate", content=f"以下门未通过,请修复后再交付:\n{feedback}"))
return # 吞掉 finish(不 yield)→ 本轮无 Msg → reply 循环 cur_iter++ 续跑
def _text(resp_or_msg) -> str:
"""从 Msg/final 安全取文本(2.0.2 用 get_text_content();若不同以源码为准)。"""
try:
return resp_or_msg.get_text_content() or ""
except Exception: # noqa: BLE001 —— spike 断言辅助,取文本失败按空
return ""
# ───────────────────────── 核心:拦一次 + 注入 + 续跑放行第二版 ─────────────────────────
def test_intercept_inject_and_continue():
"""模型两次都想 finish;check 第一次 False、第二次 True → 门被跑两次(压制一次)、放行第二版、注入消息在 context。"""
calls = {"n": 0}
async def check(agent):
calls["n"] += 1
return (calls["n"] >= 2, "" if calls["n"] >= 2 else "H_progress 门未过(无进展)")
model = MockModel()
model.set_responses([
ChatResponse(content=[TextBlock(text="第一版游戏,做完了")], is_last=True), # iter0 想 finish
ChatResponse(content=[TextBlock(text="已修复进展问题,再交付")], is_last=True), # iter1 想 finish
])
mw = RepairMiddleware(check, max_repairs=3)
agent = Agent(name="cheap_worker", system_prompt="你是便宜档生成 agent",
model=model, toolkit=Toolkit(), middlewares=[mw],
react_config=ReActConfig(max_iters=5)) # 封顶防跑飞
final = asyncio.run(agent.reply(UserMsg(name="user", content="生成一个点击得分游戏")))
assert calls["n"] == 2, f"门应被独立跑两次(压制一次),实际 {calls['n']}"
assert mw.repairs == 1, f"应续修一次,实际 {mw.repairs}"
assert "已修复" in _text(final), f"放行的应是第二版,实际:{_text(final)!r}"
assert any(isinstance(m, Msg) and getattr(m, 'role', None) == "user"
and "H_progress" in _text(m) for m in agent.state.context), "注入的修复消息应在 context"
# ───────────────────────── 边界:达 max_repairs 即放行(预算封顶,不无限续)─────────────────────────
def test_max_repairs_cap_releases():
"""check 恒 False;max_repairs=2 → 续修 2 次后放行(不再压制),reply 收尾、不无限跑。"""
async def check(agent):
return (False, "门恒不过(测封顶)")
model = MockModel()
model.set_responses([ChatResponse(content=[TextBlock(text=f"第{i}版")], is_last=True) for i in range(6)])
mw = RepairMiddleware(check, max_repairs=2)
agent = Agent(name="cw", system_prompt="p", model=model, toolkit=Toolkit(),
middlewares=[mw], react_config=ReActConfig(max_iters=10))
final = asyncio.run(agent.reply(UserMsg(name="user", content="生成")))
assert mw.repairs == 2, f"应恰好续修 max_repairs=2 次,实际 {mw.repairs}"
assert _text(final), "达封顶后应放行一个 finish、非空"
# ───────────────────────── 回落触发点①:空推理(无 Msg 无 tool call)优雅续跑 ─────────────────────────
# Agent1 提示:若 reply 循环对「reasoning 返回空」有额外早退分支(源码未见),spike 真跑才暴露 → 那就回落 control_plane。
# 本用例即 test_intercept_inject_and_continue 的隐含验证(压制那轮 = reasoning 无 Msg 落地);
# 若上面核心用例绿,则 628-670「无 Msg 自动 cur_iter++ 续跑」这条已被真跑坐实,回落触发点①排除。
# ───────────────────────── 回落触发点②:observe 中途 append 不炸 context/formatter ─────────────────────────
def test_observe_midreasoning_no_context_error():
"""续修注入(observe)在半程 context 上 append 后,下一轮推理不因 compress/formatter 报错。"""
calls = {"n": 0}
async def check(agent):
calls["n"] += 1
return (calls["n"] >= 2, "补一个资源环再交付")
model = MockModel()
model.set_responses([
ChatResponse(content=[TextBlock(text="薄循环,做完了")], is_last=True),
ChatResponse(content=[TextBlock(text="加了进货资源环,交付")], is_last=True),
])
mw = RepairMiddleware(check, max_repairs=3)
agent = Agent(name="cw", system_prompt="p", model=model, toolkit=Toolkit(),
middlewares=[mw], react_config=ReActConfig(max_iters=5))
final = asyncio.run(agent.reply(UserMsg(name="user", content="生成经营游戏"))) # 不抛即通过
assert "资源环" in _text(final)
# ───────────────────────── 干净血缘变体:pop 掉 finish 声称 ─────────────────────────
def test_pop_finish_claim_variant():
"""pop_finish_claim=True → 续修前弹掉「做完了」声称;context 里不残留第一版的 finish 文本。"""
calls = {"n": 0}
async def check(agent):
calls["n"] += 1
return (calls["n"] >= 2, "修")
model = MockModel()
model.set_responses([
ChatResponse(content=[TextBlock(text="UNIQUE_CLAIM_一版做完了")], is_last=True),
ChatResponse(content=[TextBlock(text="二版交付")], is_last=True),
])
mw = RepairMiddleware(check, max_repairs=3, pop_finish_claim=True)
agent = Agent(name="cw", system_prompt="p", model=model, toolkit=Toolkit(),
middlewares=[mw], react_config=ReActConfig(max_iters=5))
asyncio.run(agent.reply(UserMsg(name="user", content="生成")))
assert not any("UNIQUE_CLAIM" in _text(m) for m in agent.state.context), "弹掉后 context 不应残留第一版 finish 声称"
if __name__ == "__main__":
_fns = [v for k, v in sorted(globals().items()) if k.startswith("test_") and callable(v)]
_failed = 0
for _fn in _fns:
try:
_fn()
print(f" PASS {_fn.__name__}")
except Exception as e: # noqa: BLE001
_failed += 1
import traceback
print(f" FAIL {_fn.__name__}: {type(e).__name__}: {e}")
traceback.print_exc()
print(f"\n{len(_fns) - _failed}/{len(_fns)} passed")
sys.exit(1 if _failed else 0)
```
- [ ] **Step 3:跑 spike,先看它暴露 API 漂移(spike 的本职)**
Run:`cheap-worker/.venv/bin/python cheap-worker/tests/test_repair_middleware_spike.py`
预期:首跑可能在 2 处报 API 不符 —— ① `MockModel`/`ChatModelBase.__init__` 构造(以 `tests/utils.py` 2.0.2 verbatim 为准调),② `UserMsg(name=,content=)` / `final.get_text_content()` 签名(以 `agentscope/message` 源码为准调)。**这正是 spike 要暴露的**:逐个对齐已装 2.0.2 源码修掉,不是失败、是收敛。
- [ ] **Step 4:全绿 = 机制坐实(4/4 passed)**
Run:`cheap-worker/.venv/bin/python cheap-worker/tests/test_repair_middleware_spike.py`
预期:`4/4 passed`。含义逐条:核心(拦一次+注入+续跑放行第二版)✓、封顶(达 max_repairs 放行)✓、回落点②(observe 半程不炸)✓、干净血缘(pop 变体)✓;回落点①(空推理优雅续跑)由核心用例隐含坐实。→ **续修机制成立,§3.4 的 on_reasoning 洋葱方案可落**,阶段一放心接真门。
- [ ] **Step 5:若 spike 栽(任一核心用例真跑不过且非 API 漂移)→ 记回落,不硬修**
若核心用例暴露的是**机制性**问题(如 reply 循环对空推理有早退分支导致续跑不发生、或 observe 触发 formatter 硬错),**不要**在 spike 里 churn 硬凑:按设计 §6 回落现有 `control_plane` 外层循环(已存在、已跑通),在本 plan 收口记「续修留框架外层、架构其余不受影响」,Part B/C/D 照常。判据:API 签名类 = 修;机制类(压制/注入/续跑三者之一根本不发生)= 回落。
- [ ] **Step 6:Commit**
```bash
git add cheap-worker/tests/test_repair_middleware_spike.py
git commit -m "$(cat <<'EOF'
feat(cheap-worker): 续修 middleware 机制 spike 坐实(拦 finish+注入+续跑)(切片一 阶段一前置)
on_reasoning 洋葱拦 finish Msg + 独立跑 fake check + 没过压制+注入 UserMsg 续跑,
4/4 确定性单测绿(stub 模型、零 LLM/chrome)。机制成立,阶段一接真 run_gates。
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
EOF
)"
```
---
# Part B · 生产基建自托管 mini-infra(Nacos + RocketMQ)
**为什么:** 设计 §3.7 采 Spring Cloud Alibaba 原生三件套 —— 这是 game-cloud 那条栈的原生基建,清一色 Apache-2.0、自托管 mini-infra。**RAM 实测偏紧**(15G 总 / 8.6G 可用 / 4 核 / RAGflow 独占 ~9G),故两件都必须 lean JVM 调、并作纯增量部署,不动既有服务。Nacos 先起(配置版本化 + 服务发现 + Sentinel 规则源,三用一套),RocketMQ 次之(异步 gen 队列)。
### Task B1:Nacos 自托管(standalone + MySQL 后端 + lean JVM)
**Files:**
- Create:`deploy/infra/nacos/docker-compose.yml`(Nacos 2.x standalone,MySQL 后端,lean JVM)
- Create:`deploy/infra/nacos/init-nacos-schema.sh`(在既有 MySQL 建 `nacos` 库 + 灌官方 schema)
- Create:`deploy/infra/nacos/health-check.sh`(curl Nacos 健康,绕代理)
- Modify:`docs/内网凭据与端点.md`(追加 Nacos 端点段)
**Interfaces:**
- Produces:Nacos 控制台 + OpenAPI `http://100.64.0.8:8848/nacos`;配置/发现/Sentinel 规则源都用它。MySQL `nacos` 库复用既有 `100.64.0.8:3306`(root/`ZRH3jwYLOrntBcTAw29MW9BP`,不新增实例)。
- [ ] **Step 1:写健康探针(先失败 —— 服务未起)**
Create `deploy/infra/nacos/health-check.sh`(沿 `deploy/smoke-test.sh` 风格:bash + curl + 退出码):
```bash
#!/usr/bin/env bash
# Nacos 健康探针 —— 绕 fake-ip 系统代理(内网直连铁律)。绿=200,红=非200/超时。
set -uo pipefail
NACOS="${NACOS:-http://100.64.0.8:8848}"
# --noproxy '*' 绕系统代理;/nacos/v1/console/health/readiness 是 Nacos 就绪探针
code="$(curl -s -o /dev/null -w '%{http_code}' --noproxy '*' --max-time 8 \
"${NACOS}/nacos/v1/console/health/readiness" || echo 000)"
if [[ "$code" == "200" ]]; then echo "PASS Nacos readiness 200 @ ${NACOS}"; exit 0; fi
echo "FAIL Nacos readiness=${code} @ ${NACOS}"; exit 1
```
Run:`NACOS=http://100.64.0.8:8848 bash deploy/infra/nacos/health-check.sh`
预期:`FAIL Nacos readiness=000`(还没起 —— 这是对的)。
- [ ] **Step 2:写 MySQL schema 初始化脚本**
Create `deploy/infra/nacos/init-nacos-schema.sh`:
```bash
#!/usr/bin/env bash
# 在既有 mini-infra MySQL 建 nacos 库 + 灌官方 schema(纯增量,不动既有库)。
# 官方 schema 随 Nacos 版本走:容器内 /home/nacos/conf/mysql-schema.sql(2.x)。此处从容器 cp 出再灌,保证版本对齐。
set -euo pipefail
MYSQL_HOST="${MYSQL_HOST:-100.64.0.8}"; MYSQL_PORT="${MYSQL_PORT:-3306}"
MYSQL_ROOT_PW="${MYSQL_ROOT_PW:?需传 MySQL root 密码(见 docs/内网凭据与端点.md)}"
NACOS_IMAGE="${NACOS_IMAGE:-nacos/nacos-server:v2.4.3}"
# 1) 建库(utf8mb4)—— 若已存在则跳过,绝不 drop
docker run --rm mysql:8 mysql -h"$MYSQL_HOST" -P"$MYSQL_PORT" -uroot -p"$MYSQL_ROOT_PW" \
-e "CREATE DATABASE IF NOT EXISTS nacos DEFAULT CHARACTER SET utf8mb4;" 2>/dev/null \
|| { echo "建 nacos 库失败"; exit 1; }
# 2) 从 Nacos 镜像取版本对齐的 schema 并灌入
cid="$(docker create "$NACOS_IMAGE")"; trap 'docker rm -f "$cid" >/dev/null 2>&1 || true' EXIT
docker cp "$cid:/home/nacos/conf/mysql-schema.sql" /tmp/nacos-mysql-schema.sql
docker run --rm -v /tmp/nacos-mysql-schema.sql:/s.sql mysql:8 sh -c \
"mysql -h$MYSQL_HOST -P$MYSQL_PORT -uroot -p'$MYSQL_ROOT_PW' nacos < /s.sql"
echo "PASS nacos schema 已灌入 $MYSQL_HOST:$MYSQL_PORT/nacos"
```
(注:MySQL root 密码经 env 传入、从 `docs/内网凭据与端点.md` 取,不写死进脚本。)
- [ ] **Step 3:写 docker-compose(lean JVM = 关键,防撑爆 8.6G)**
Create `deploy/infra/nacos/docker-compose.yml`:
```yaml
# Nacos 2.x standalone + 外部 MySQL 后端。lean JVM:mini-infra 仅 8.6G 可用,默认 2g 会挤压 RAGflow。
# 纯增量:host 网络下只占 8848/9848(gRPC)/9849;不碰既有容器。
services:
nacos:
image: nacos/nacos-server:v2.4.3
container_name: infra-nacos
restart: unless-stopped
network_mode: host # 与既有 infra-* 一致(见 game-cloud/script/docker 约定),直连 3306
environment:
MODE: standalone
PREFER_HOST_MODE: ip
NACOS_SERVER_PORT: "8848"
# ── lean JVM:standalone 512m 够用(dev/MVP 单操作者),显式压死默认 2g ──
JVM_XMS: "512m"
JVM_XMX: "512m"
JVM_XMN: "256m"
# ── MySQL 后端(复用既有实例的 nacos 库)──
SPRING_DATASOURCE_PLATFORM: mysql
MYSQL_SERVICE_HOST: 100.64.0.8
MYSQL_SERVICE_PORT: "3306"
MYSQL_SERVICE_DB_NAME: nacos
MYSQL_SERVICE_USER: root
MYSQL_SERVICE_PASSWORD: "${MYSQL_ROOT_PW}" # compose 起时经 env 注入,不写死
MYSQL_SERVICE_DB_PARAM: "characterEncoding=utf8&useSSL=false&serverTimezone=UTC&allowPublicKeyRetrieval=true"
# ── 鉴权:开 token(生产基建不裸奔);密钥写回凭据文档 ──
NACOS_AUTH_ENABLE: "true"
NACOS_AUTH_TOKEN: "${NACOS_AUTH_TOKEN}" # base64,≥32 字节;见凭据文档
NACOS_AUTH_IDENTITY_KEY: "${NACOS_AUTH_IDENTITY_KEY}"
NACOS_AUTH_IDENTITY_VALUE: "${NACOS_AUTH_IDENTITY_VALUE}"
# ── Nacos 2.4.x 移除了内置 nacos/nacos 默认账号,首启需初始化 admin 密码(见 Step 4)──
# 部分 2.4.x 镜像认此 env 首启建 admin;若不认则走 Step 4 的 API 初始化。执行时以镜像实际为准。
NACOS_AUTH_ADMIN_PASSWORD: "${NACOS_ADMIN_PASSWORD}"
volumes:
- ./data/logs:/home/nacos/logs
```
- [ ] **Step 4:部署 + 灌 schema + 探活(真跑,mini-infra)**
Run(ssh 真 IP 绕代理;env 从凭据文档取):
```bash
# 1) 灌 schema
ssh root@100.64.0.8 'MYSQL_ROOT_PW=<见凭据文档> bash -s' < deploy/infra/nacos/init-nacos-schema.sh
# 2) 起 Nacos(env 注入密钥 + admin 密码)
scp -r deploy/infra/nacos root@100.64.0.8:/opt/infra-nacos
ssh root@100.64.0.8 'cd /opt/infra-nacos && MYSQL_ROOT_PW=<..> NACOS_AUTH_TOKEN=<..> \
NACOS_AUTH_IDENTITY_KEY=<..> NACOS_AUTH_IDENTITY_VALUE=<..> NACOS_ADMIN_PASSWORD=<..> docker compose up -d'
# 3) 探活(等 ~30s 启动)
NACOS=http://100.64.0.8:8848 bash deploy/infra/nacos/health-check.sh
# 4) 初始化 admin 账号(Nacos 2.4.x 无默认 nacos/nacos;若镜像未经 env 建 admin,则首启用 API 建一次)
ssh root@100.64.0.8 'curl -s --noproxy "*" -X POST \
"http://127.0.0.1:8848/nacos/v1/auth/users/admin?password=<NACOS_ADMIN_PASSWORD>"' || true
# 5) 验鉴权:admin 登录取 token 应成功(200 + accessToken)
ssh root@100.64.0.8 'curl -s --noproxy "*" -X POST "http://127.0.0.1:8848/nacos/v1/auth/login" \
-d "username=nacos&password=<NACOS_ADMIN_PASSWORD>"' | grep -q accessToken \
&& echo "PASS admin 登录取 token" || echo "FAIL admin 未初始化/口令错"
```
预期:`PASS Nacos readiness 200` + `PASS admin 登录取 token`。同时 `ssh root@100.64.0.8 'docker stats --no-stream'` 确认 infra-nacos 内存 <700M、既有容器未被挤崩。
> 注:Nacos 2.4.x 移除内置 nacos/nacos,admin 初始化方式(env vs 首启 API)在 2.4.x 各小版本有差异 —— 真部署前先核 `nacos/nacos-server:v2.4.3` 镜像的确切机制,admin 口令写回凭据文档,**C1/C2/D 的 `${NACOS_PASSWORD}` 均指它**。
- [ ] **Step 5:把端点写回凭据文档 + commit**
Modify `docs/内网凭据与端点.md`:追加「Nacos:控制台 `http://100.64.0.8:8848/nacos`,**admin 账号 `nacos` / 口令 `<NACOS_ADMIN_PASSWORD>`**(C1/C2/D 的 `${NACOS_PASSWORD}` 均指它),鉴权 token/identity 见下,MySQL 后端 `nacos` 库」段(口令/token 入库 —— 内网阶段密钥进文档铁律)。
```bash
git add deploy/infra/nacos/ docs/内网凭据与端点.md
git commit -m "$(cat <<'EOF'
feat(infra): Nacos 2.4.3 自托管 mini-infra(standalone+MySQL 后端+lean 512m JVM)(切片一 阶段〇 B1)
纯增量部署、host 网络、复用既有 MySQL nacos 库;鉴权开 token。探针 200 绿,内存 <700M 不挤 RAGflow。
端点/密钥写回 docs/内网凭据与端点.md。
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
EOF
)"
```
### Task B2:RocketMQ 自托管(NameServer + Broker + lean JVM)
**Files:**
- Create:`deploy/infra/rocketmq/docker-compose.yml`(NameServer + 单 Broker,lean JVM)
- Create:`deploy/infra/rocketmq/broker.conf`(broker 配置:brokerIP、autoCreateTopic)
- Create:`deploy/infra/rocketmq/health-check.sh`(探 NameServer + Broker)
- Modify:`docs/内网凭据与端点.md`(追加 RocketMQ 端点段)
**Interfaces:**
- Produces:NameServer `100.64.0.8:9876`、Broker `100.64.0.8:10911`。game-cloud(C3)`rocketmq.name-server: 100.64.0.8:9876`。gen 队列 topic = `AIGC_GEN_TOPIC`(C3 生产/消费)。
- [ ] **Step 1:写健康探针(先失败)**
Create `deploy/infra/rocketmq/health-check.sh`:
```bash
#!/usr/bin/env bash
# RocketMQ 探针:NameServer 端口通 + Broker 在 NameServer 注册。绕代理直连。
set -uo pipefail
NS="${NS:-100.64.0.8:9876}"
# 端口连通(nc);Broker 注册用 mqadmin clusterList 查(容器内)
if ! nc -z -w5 ${NS/:/ } 2>/dev/null; then echo "FAIL NameServer ${NS} 端口不通"; exit 1; fi
brokers="$(ssh root@100.64.0.8 'docker exec infra-rmq-broker sh mqadmin clusterList -n 127.0.0.1:9876 2>/dev/null | grep -c broker' || echo 0)"
if [[ "${brokers:-0}" -ge 1 ]]; then echo "PASS RocketMQ NameServer 通 + Broker 已注册"; exit 0; fi
echo "FAIL Broker 未在 NameServer 注册"; exit 1
```
Run:`bash deploy/infra/rocketmq/health-check.sh` → 预期 `FAIL`(未起)。
- [ ] **Step 2:写 broker.conf**
Create `deploy/infra/rocketmq/broker.conf`:
```properties
# RocketMQ Broker —— 单机、自动建 topic(MVP)、brokerIP 指 mini-infra Tailscale 地址(否则跨机拿不到 broker 地址)。
brokerClusterName = DefaultCluster
brokerName = broker-a
brokerId = 0
deleteWhen = 04
fileReservedTime = 48
brokerRole = ASYNC_MASTER
flushDiskType = ASYNC_FLUSH
# 关键:brokerIP1 必须是消费方(mini-desktop / game-cloud)能连到的地址,否则 NameServer 返回的 broker 地址不可达
brokerIP1 = 100.64.0.8
autoCreateTopicEnable = true
autoCreateSubscriptionGroup = true
```
- [ ] **Step 3:写 docker-compose(lean JVM:broker 默认 8g → 压到 1g)**
Create `deploy/infra/rocketmq/docker-compose.yml`:
```yaml
# RocketMQ NameServer + 单 Broker。lean JVM:broker 默认 -Xmx8g 会瞬间撑爆 8.6G,压到 1g;nameserver 512m。
services:
namesrv:
image: apache/rocketmq:5.3.1
container_name: infra-rmq-namesrv
restart: unless-stopped
network_mode: host
environment:
JAVA_OPT_EXT: "-Xms512m -Xmx512m -Xmn256m" # 压死默认
command: sh mqnamesrv
broker:
image: apache/rocketmq:5.3.1
container_name: infra-rmq-broker
restart: unless-stopped
network_mode: host
depends_on: [namesrv]
environment:
JAVA_OPT_EXT: "-Xms1g -Xmx1g -Xmn512m" # 压死默认 8g
NAMESRV_ADDR: "127.0.0.1:9876"
volumes:
- ./broker.conf:/home/rocketmq/rocketmq-5.3.1/conf/broker.conf
- ./data/store:/home/rocketmq/store
- ./data/logs:/home/rocketmq/logs
command: sh mqbroker -c /home/rocketmq/rocketmq-5.3.1/conf/broker.conf
```
- [ ] **Step 4:部署 + 探活(真跑)**
Run:
```bash
scp -r deploy/infra/rocketmq root@100.64.0.8:/opt/infra-rocketmq
ssh root@100.64.0.8 'cd /opt/infra-rocketmq && docker compose up -d'
sleep 20 && bash deploy/infra/rocketmq/health-check.sh
ssh root@100.64.0.8 'docker stats --no-stream | grep rmq' # 确认 broker <1.3G、namesrv <700M
```
预期:`PASS RocketMQ NameServer 通 + Broker 已注册`;两容器内存在 lean 预算内、mini-infra 总可用未见崩(`free -h` 复查 available 仍 >2G)。
- [ ] **Step 5:端点写回凭据文档 + commit**
Modify `docs/内网凭据与端点.md`:追加「RocketMQ:NameServer `100.64.0.8:9876`、Broker `100.64.0.8:10911`、gen 队列 topic `AIGC_GEN_TOPIC`」。
```bash
git add deploy/infra/rocketmq/ docs/内网凭据与端点.md
git commit -m "$(cat <<'EOF'
feat(infra): RocketMQ 5.3.1 自托管 mini-infra(NameServer+单 Broker+lean 1g JVM)(切片一 阶段〇 B2)
纯增量、host 网络、brokerIP1 指 Tailscale 地址;探针绿、Broker 注册成功、内存在 lean 预算内。
端点写回 docs/内网凭据与端点.md。
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
EOF
)"
```
---
# Part C · game-cloud 全量接 Java(Nacos 重启用 + Sentinel≤15 + RocketMQ)
**为什么 + 组合决策(经 §6.8 Opus 评审修正):** game-cloud 的 Java 侧不是白地 —— 已有自建准入 `enqueueWithControlPlane`(`AigcTaskServiceImpl.java:310`,配额/并发/背压三门·纯 DB count)+ DB 轮询派发(`AigcExecutorConfiguration.java` 的 `@Scheduled` tick + `game_aigc_task` 表状态机),Nacos/RocketMQ 依赖已声明但被关停(`application.yaml:50-58`/`:157-161`),Sentinel 完全缺失。故本 Part 按**分层组合**收敛(非推翻既有 working code),三层各管一件、职责不重叠:
- **在跑生成并发≤15 的权威 = RocketMQ 消费者 `consumeThreadMax=15`**(C3)—— 有界消费线程池 = 至多 15 个生成同时在跑,RocketMQ 原生、结构性保证。**MVP 单实例部署**下这就是并发硬上限;集群级全局限流(多实例)留放量后上 Sentinel cluster-flow Token Server 或 Redis 全局配额(那时才有多实例,YAGNI 前置无意义)。
- **Sentinel = 入口 QPS 突发保护**(C2)—— 挂在入队入口,`FLOW_GRADE_QPS` 限「每秒入队请求数」防洪峰,规则热源 Nacos、可不重启热调。**它不是并发权威**:空体入口的 entry/exit 是纳秒级,thread-grade 统计不到真正在跑的生成(那发生在 `admit()` 返回之后的消费线程里)—— Opus 评审纠此承重错。Sentinel 在此的价值 = 热可调的入口速率闸 + 未来多实例时升级成集群限流的落点。
- **既有 `enqueueWithControlPlane` DB 三门保留作业务级 per-creator 配额**(每日额度 + 每创作者并发)+ 全局背压 `queue-depth-limit` 作纵深防御(设 ≥15 兜底)—— Sentinel/consumeThreadMax 做不了业务配额,正交。
- **RocketMQ = 异步投递机制**,`game_aigc_task` 表**保留作任务记录 + CAS 幂等**(`claimQueuedTask` queued→running 仍用);生产者入队发消息 → Java 有界消费者消费 → 调 Service。**替换**的是 `@Scheduled` DB 轮询这个**触发方式**,不是任务表本身。
> 「并发权威=consumeThreadMax、Sentinel=入口 QPS、单实例 MVP」是本 Part 的关键设计决策(Opus 评审把原稿「Sentinel=全局并发≤15 权威」的承重错纠正到此)。它牵动 spec §3.7 的「Sentinel 全局≤15」表述 → 见文末收口 TODO。若创始人改判(如本轮就上 Token Server 做集群全局限流),按判调 C2/C3。
### Task C1:game-cloud 重启用 Nacos 配置中心 + 服务发现
**Files:**
- Modify:`game-cloud/huijing-server/src/main/resources/application.yaml:50-58`(`nacos.discovery.enabled`/`config.enabled` false→true + server-addr)
- Modify:`game-cloud/huijing-server/src/main/resources/application-staging.yaml`(staging profile 的 nacos server-addr → mini-infra)
- (依赖已在)`game-cloud/huijing-server/pom.xml:224,229`(nacos-discovery/config starter 已声明,无需新增)
**Interfaces:**
- Produces:game-cloud 各 server 注册进 Nacos discovery + 从 Nacos 读配置。Service 发现名对齐 `spring.application.name`。
- [ ] **Step 1:改 application.yaml 开 Nacos + 指 mini-infra**
Modify `game-cloud/huijing-server/src/main/resources/application.yaml`(:50-58 段):把 `spring.cloud.nacos.discovery.enabled` 与 `config.enabled` 由 `false` 改 `true`,补 `server-addr` 与鉴权:
```yaml
spring:
cloud:
nacos:
server-addr: 100.64.0.8:8848 # mini-infra Nacos(B1)
username: nacos # 鉴权(B1 开了 token 模式,用户名口令见凭据文档)
password: ${NACOS_PASSWORD}
discovery:
enabled: true # was false —— 重启用服务发现
namespace: public
config:
enabled: true # was false —— 重启用配置中心
namespace: public
file-extension: yaml
```
(staging profile 同改 `application-staging.yaml` 的 nacos 段;`application-staging.yaml:2` 那句「nacos 关」注释同步更新为「指 mini-infra」。)
- [ ] **Step 2:本地编译验证(不需真连 Nacos)**
Run:`cd game-cloud && ./mvnw -q -pl huijing-server -am compile 2>&1 | tail -20`
预期:`BUILD SUCCESS`(配置改动不破坏编译;真连 Nacos 的验证在 Step 3 随后端窗口)。
- [ ] **Step 3:真连验证(随后端窗口)—— 服务注册 + 配置读**
Run(staging 后端窗口):启动一个 game-cloud server(如 `huijing-server` 的 infra 模块),然后:
```bash
# 该服务应出现在 Nacos 服务列表
curl -s --noproxy '*' "http://100.64.0.8:8848/nacos/v1/ns/catalog/services?pageNo=1&pageSize=20&namespaceId=public" \
-H "..." | grep -q "$(应用名)" && echo "PASS 已注册" || echo "FAIL 未注册"
```
预期:`PASS 已注册`。(此步依赖后端窗口,代码 + 编译先落。)
- [ ] **Step 4:Commit**
```bash
git add game-cloud/huijing-server/src/main/resources/application.yaml game-cloud/huijing-server/src/main/resources/application-staging.yaml
git commit -m "$(cat <<'EOF'
feat(game-cloud): 重启用 Nacos 配置中心+服务发现指向 mini-infra(切片一 阶段〇 C1)
application.yaml nacos discovery/config enabled false→true + server-addr 100.64.0.8:8848 + 鉴权。
依赖零新增(starter 已声明)。本地编译绿;真注册验证随后端窗口。
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
EOF
)"
```
### Task C2:Sentinel 入口 QPS 突发保护(FLOW_GRADE_QPS,规则源 Nacos)
**Files:**
- Modify:`game-cloud/huijing-dependencies/pom.xml`(BOM 已托管 SCA 2025.0.0.0;此处在 aigc 模块 pom 引 sentinel starter + sentinel-datasource-nacos)
- Modify:`game-cloud/game-module-aigc/.../pom.xml`(引 `spring-cloud-starter-alibaba-sentinel` + `sentinel-datasource-nacos`)
- Create:`game-cloud/game-module-aigc/.../admission/GenAdmissionResource.java`(生成入口的 `@SentinelResource` + blockHandler)
- Modify:`game-cloud/game-module-aigc/.../service/AigcTaskServiceImpl.java:310`(`enqueueWithControlPlane` 入口套 Sentinel 资源,DB 三门保留在后)
- Modify:`game-cloud/huijing-server/src/main/resources/application.yaml`(sentinel datasource from Nacos + transport)
- Nacos 配置:新增 dataId `sentinel-gen-flow-rules`(JSON,`FLOW_GRADE_QPS` count=10 入口速率闸)
**Interfaces:**
- Consumes:C1 的 Nacos(规则源 + 鉴权)。
- Produces:生成入队入口的 QPS 被 Sentinel 限(防突发洪峰,`FLOW_GRADE_QPS`);越限走 `blockHandler` 优雅拒(返「排队中/稍后再试」,非 500)。规则改 Nacos dataId `sentinel-gen-flow-rules` 即热生效(不重启)。**并发≤15 不在本 Task**,由 C3 `consumeThreadMax` 保证。
- [ ] **Step 1:引 Sentinel 依赖**
Modify aigc 模块 pom,新增(版本由 SCA 2025.0.0.0 BOM 托管、不写死):
```xml
<dependency>
<groupId>com.alibaba.cloud</groupId>
<artifactId>spring-cloud-starter-alibaba-sentinel</artifactId>
</dependency>
<dependency>
<groupId>com.alibaba.csp</groupId>
<artifactId>sentinel-datasource-nacos</artifactId>
</dependency>
```
- [ ] **Step 2:写 Sentinel 资源 + blockHandler(入队入口的 QPS 闸)**
Create `game-cloud/game-module-aigc/.../admission/GenAdmissionResource.java`:
```java
package com.wanxiang.huijing.game.module.aigc.admission;
import com.alibaba.csp.sentinel.annotation.SentinelResource;
import com.alibaba.csp.sentinel.slots.block.BlockException;
import com.wanxiang.huijing.framework.common.exception.ServiceException;
import org.springframework.stereotype.Component;
/**
* 生成入队的 Sentinel 资源门 —— 入口 QPS 突发保护(FLOW_GRADE_QPS,防洪峰)。
* <p>职责边界(经 Opus 评审纠正):本门统计「每秒入队请求数」,不是「在跑生成并发」——
* 后者由 RocketMQ 消费者 consumeThreadMax=15 保证(空体入口 entry/exit 纳秒级,thread-grade 统计不到
* 真正在跑的生成,那发生在 admit() 返回之后的消费线程里)。与 DB 三门分层:本门=入口速率、
* consumeThreadMax=并发、DB 门=per-creator 业务配额,三层正交。
* <p>越限行为:Sentinel 抛 BlockException → blockHandler 优雅转成业务异常「生成繁忙,请稍后重试」,
* 绝不 500、不静默丢。
*/
@Component
public class GenAdmissionResource {
/** 生成准入资源名(与 Nacos dataId sentinel-gen-flow-rules 里的 resource 对齐)。 */
public static final String RESOURCE = "aigc:gen:admission";
/**
* 入口速率闸:受 Sentinel FLOW_GRADE_QPS 规则约束(入队 QPS 超限即 block)。
* 通过则返回(交由调用方继续走 DB 三门 + 入队);被 block 则走 blockHandler。
*/
@SentinelResource(value = RESOURCE, blockHandler = "onAdmissionBlocked")
public void admit() {
// 空体:让 Sentinel 统计「入队 QPS」(FLOW_GRADE_QPS 计每秒通过数,与方法耗时无关,空体正合适)。
// 真正的业务准入(配额)在 DB 三门;在跑生成并发≤15 在 C3 consumeThreadMax,不在这里。
}
/** 越 QPS 闸的优雅降级:转成可读业务异常(错误码段 aigc,前端提示排队),不抛 500。 */
public void onAdmissionBlocked(BlockException ex) {
// 可追溯日志(创始人铁律:错误路径必须留痕)
org.slf4j.LoggerFactory.getLogger(GenAdmissionResource.class)
.warn("[aigc-admission] Sentinel 入口 QPS 闸拦截:入队速率超限,拒绝新请求 rule={}", ex.getRule());
throw new ServiceException(1_002_100_000, "生成服务繁忙(入队速率超限),请稍后重试");
}
}
```
(错误码 `1_002_100_000` 按 aigc 模块既有错误码段对齐 —— 执行时查 `game-module-aigc` 的 ErrorCodeConstants,用其规范号段。)
- [ ] **Step 3:在生成入口调准入门(套在 DB 三门之前)**
Modify `AigcTaskServiceImpl.java`(`enqueueWithControlPlane` :310 处):在方法体最前注入 `genAdmissionResource.admit();`,让 Sentinel 入口 QPS 闸**先于** DB 三门执行:
```java
// 注入(类字段):private final GenAdmissionResource genAdmissionResource;
// enqueueWithControlPlane(AigcTaskDO task) 方法体开头:
genAdmissionResource.admit(); // ① 入口 QPS 突发闸(Sentinel FLOW_GRADE_QPS,规则热源 Nacos);越限抛 ServiceException 排队提示
// ② 既有 DB 三门保留:配额(daily)/ per-creator 并发 / 背压 queue-depth-limit(业务配额 + 纵深防御)
// ③ 在跑生成并发≤15 = C3 RocketMQ consumeThreadMax(结构性上限,MVP 单实例);不在此处
// ...原有 enqueueWithControlPlane 逻辑不动...
```
- [ ] **Step 4:配 Sentinel datasource from Nacos + transport**
Modify `application.yaml`,加 sentinel 段:
```yaml
spring:
cloud:
sentinel:
transport:
dashboard: ${SENTINEL_DASHBOARD:} # 留空 = 不部面板(RAM 紧,规则直存 Nacos)
datasource:
gen-flow:
nacos:
server-addr: 100.64.0.8:8848
username: nacos
password: ${NACOS_PASSWORD}
dataId: sentinel-gen-flow-rules
groupId: DEFAULT_GROUP
data-type: json
rule-type: flow # 流控规则
```
在 Nacos 建 dataId `sentinel-gen-flow-rules`(JSON,QPS grade = 每秒入队数;`grade:1` = FLOW_GRADE_QPS):
```json
[
{
"resource": "aigc:gen:admission",
"grade": 1,
"count": 10,
"limitApp": "default",
"strategy": 0,
"controlBehavior": 0
}
]
```
(`grade:1` = `RuleConstant.FLOW_GRADE_QPS`,`count:10` = 入队 ≤10 QPS(突发保护,方向性值,按真实入队速率调)。改这个 dataId 即热调、不重启。**并发≤15 不在这里** —— 由 C3 `consumeThreadMax=15` 保证。)
- [ ] **Step 5:单测(mock Sentinel,验 blockHandler 语义)**
Create `game-cloud/game-module-aigc/src/test/java/.../admission/GenAdmissionResourceTest.java`:直接调 `onAdmissionBlocked(mockBlockException)` 断言抛 `ServiceException` 且 code/msg 正确(资源门 admit() 的真 QPS 限流由集成验证覆盖,单测只锁降级语义):
```java
@Test
void onAdmissionBlocked_throwsServiceException_notServerError() {
GenAdmissionResource r = new GenAdmissionResource();
BlockException be = new FlowException("aigc:gen:admission");
ServiceException ex = assertThrows(ServiceException.class, () -> r.onAdmissionBlocked(be));
assertTrue(ex.getMessage().contains("繁忙")); // 优雅提示、非 500
}
```
Run:`cd game-cloud && ./mvnw -q -pl game-module-aigc test -Dtest=GenAdmissionResourceTest 2>&1 | tail -15`
预期:`BUILD SUCCESS` / 测试通过。
- [ ] **Step 6:集成验证(随后端窗口)——入队 QPS 超限被拦(并发由 C3 consumeThreadMax 兜)**
Run(后端窗口):以 >10 QPS 猛发入队请求,预期超速率的请求走 blockHandler 返「繁忙」(非 500);改 Nacos dataId count→20 后热生效。**并发≤15 的验证在 C3**(consumeThreadMax:并发提交 20,消费侧同时在跑 ≤15)。
- [ ] **Step 7:Commit**
```bash
git add game-cloud/game-module-aigc/ game-cloud/huijing-server/src/main/resources/application.yaml
git commit -m "$(cat <<'EOF'
feat(game-cloud): Sentinel 入口 QPS 突发保护(FLOW_GRADE_QPS·规则热源 Nacos)(切片一 阶段〇 C2)
入队入口套 @SentinelResource(aigc:gen:admission)做入口速率闸,与 DB 三门(per-creator 配额)分层;
越限 blockHandler 优雅拒非 500。并发≤15 权威=C3 consumeThreadMax(非本 Task)。规则改 Nacos dataId 热生效。
单测锁降级语义;真限流集成验证随后端窗口。
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
EOF
)"
```
### Task C3:RocketMQ 异步 gen 队列(生产者 + 有界消费者≤15)
**Files:**
- Modify:`game-cloud/huijing-server/src/main/resources/application.yaml:157-161`(`rocketmq.name-server` → mini-infra)
- (依赖已在)`huijing-spring-boot-starter-mq`(rocketmq-spring 2.3.5,`RocketMQTemplate` 已可用)
- Create:`game-cloud/game-module-aigc/.../mq/GenTaskProducer.java`(入队时发消息)
- Create:`game-cloud/game-module-aigc/.../mq/GenTaskConsumer.java`(`@RocketMQMessageListener`,`consumeThreadMax` 有界)
- Modify:`AigcTaskServiceImpl.java`(入队成功后发 MQ,替 `@Scheduled` 轮询触发)
- Modify:`AigcExecutorConfiguration.java`(`@Scheduled` tick 降级为 MQ 死信/兜底补偿,不再作主触发)
**Interfaces:**
- Consumes:B2 的 RocketMQ(`100.64.0.8:9876`,topic `AIGC_GEN_TOPIC`)+ 既有 `game_aigc_task` 表 + CAS `claimQueuedTask`。
- Produces:入队即发 MQ → Java 有界消费者(`consumeThreadMax` ≤15)消费 → CAS 认领(幂等)→ 调派发。任务记录/状态仍以 `game_aigc_task` 表为准(MQ 只做投递触发)。
- [ ] **Step 1:指 RocketMQ 到 mini-infra**
Modify `application.yaml`(:157-161 段)+ 各 profile:`rocketmq.name-server: 100.64.0.8:9876`(替 `127.0.0.1:9876`);producer.group 保留既有 `${spring.application.name}_PRODUCER`。
- [ ] **Step 2:写生产者(入队发消息,消息体 = taskId,不搬业务数据)**
Create `game-cloud/game-module-aigc/.../mq/GenTaskProducer.java`:
```java
package com.wanxiang.huijing.game.module.aigc.mq;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.stereotype.Component;
/**
* 生成任务 MQ 生产者 —— 入队成功后发一条「有新任务」信号,把「触发」从 @Scheduled DB 轮询改为事件驱动。
* <p>消息体只放 taskId(轻信号),任务真身/状态仍以 game_aigc_task 表为准 —— MQ 是投递触发,不是数据搬运。
* best-effort:发 MQ 失败不回滚入队(任务已落库),由 @Scheduled 兜底补偿轮询兜住(见 AigcExecutorConfiguration)。
*/
@Component
public class GenTaskProducer {
/** gen 队列 topic(与 B2 broker 自动建的 topic 对齐)。 */
public static final String TOPIC = "AIGC_GEN_TOPIC";
private final RocketMQTemplate rocketMQTemplate;
public GenTaskProducer(RocketMQTemplate rocketMQTemplate) {
this.rocketMQTemplate = rocketMQTemplate;
}
/** 发「新任务」信号(同步发,拿发送结果落日志;失败只告警、不抛,由 DB 轮询兜底补偿)。 */
public void publishNewTask(Long taskId) {
try {
rocketMQTemplate.syncSend(TOPIC, String.valueOf(taskId));
org.slf4j.LoggerFactory.getLogger(GenTaskProducer.class)
.info("[aigc-mq] 发送 gen 任务信号 taskId={}", taskId);
} catch (Exception e) { // best-effort:MQ 抖动不连累入队(已落库),兜底轮询会捡起
org.slf4j.LoggerFactory.getLogger(GenTaskProducer.class)
.warn("[aigc-mq] 发送 gen 信号失败(不回滚入队,由兜底轮询补偿)taskId={}", taskId, e);
}
}
}
```
- [ ] **Step 3:写有界消费者(consumeThreadMax≤15 + CAS 认领幂等)**
Create `game-cloud/game-module-aigc/.../mq/GenTaskConsumer.java`:
```java
package com.wanxiang.huijing.game.module.aigc.mq;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
/**
* 生成任务 MQ 消费者 —— 有界并发消费(consumeThreadMax=15 = 在跑生成并发的权威硬上限,MVP 单实例)。
* <p>consumeThreadMax 是「同时在跑的生成」的结构性上限(线程池占满即 15 个生成并行,新消息排队等空位)——
* 这才是真正的并发≤15(区别于 C2 Sentinel 的入口 QPS 闸)。多实例扩容后的集群级全局限流留放量后(Token Server)。
* <p>收到 taskId 信号 → 用既有 CAS claimQueuedTask(queued→running)认领,认领成功才派发,认领失败(已被别人拿走/
* 非 queued)即幂等丢弃 —— 天然去重,同一 taskId 重复投递不会双跑。
*/
@Component
@RocketMQMessageListener(
topic = GenTaskProducer.TOPIC,
consumerGroup = "aigc_gen_consumer",
consumeThreadMax = 15, // 在跑生成并发权威上限(结构性 ≤15;C2 Sentinel 是入口 QPS、不是这个)
consumeThreadNumber = 4 // 起始消费线程
)
public class GenTaskConsumer implements RocketMQListener<String> {
private final AigcGenDispatchService dispatchService; // 封装既有 claim + 派发(见 Step 4)
public GenTaskConsumer(AigcGenDispatchService dispatchService) {
this.dispatchService = dispatchService;
}
@Override
public void onMessage(String taskIdStr) {
Long taskId = Long.valueOf(taskIdStr);
// CAS 认领(既有 claimQueuedTask:仅 queued→running 才成功)→ 认领成功才跑,失败即幂等丢弃
boolean claimed = dispatchService.claimAndDispatch(taskId);
org.slf4j.LoggerFactory.getLogger(GenTaskConsumer.class)
.info("[aigc-mq] 消费 gen 信号 taskId={} claimed={}", taskId, claimed);
// 认领失败不抛(幂等丢弃);派发内部异常由 dispatchService 落 task 失败态,不让 MQ 无限重投
}
}
```
- [ ] **Step 4:抽既有认领+派发为 `AigcGenDispatchService.claimAndDispatch`,入队处发 MQ,兜底补偿覆盖 stale queued + 卡死 running**
Modify `AigcTaskServiceImpl.java`:入队成功(`game_aigc_task` 落库 queued)后调 `genTaskProducer.publishNewTask(task.getId())`。把既有 `claimQueuedTask` + `SaaGraphDispatcher` 派发逻辑收进 `AigcGenDispatchService.claimAndDispatch(taskId)`(复用现有 CAS + 派发,不重写):**认领成功(queued→running)但派发抛异常时,必须在 catch 里把状态从 `running` 落回 `failed`(或按重试策略回 `queued`),绝不留悬空 `running`**(Opus 评审:CAS 成功 + 派发失败会永久卡 running)。Modify `AigcExecutorConfiguration.java`:`@Scheduled` tick **降级为兜底补偿**,周期拉大,捞**两类**卡住任务:① stale `queued`(入队了但 MQ 信号丢)重发信号;② **`running` 超 `updateTime` 阈值(如 >30min 无进展)的僵死任务**(认领成功后进程崩溃 / 派发失败未落状态的漏网)→ 按超时重置为 `failed` 或 `queued` 重投,不再作主触发路径。
- [ ] **Step 5:单测(mock RocketMQTemplate + dispatch,验幂等丢弃 + best-effort)**
Create `GenTaskProducerTest.java` / `GenTaskConsumerTest.java`:
- Producer:mock `RocketMQTemplate.syncSend` 抛异常 → 断言 `publishNewTask` **不抛**(best-effort,不回滚入队)。
- Consumer:mock `dispatchService.claimAndDispatch` 返 false(已被认领)→ 断言 `onMessage` 正常返回、不抛(幂等丢弃)。
Run:`cd game-cloud && ./mvnw -q -pl game-module-aigc test -Dtest=GenTask*Test 2>&1 | tail -15`
预期:通过。
- [ ] **Step 6:集成验证(随后端窗口)——入队→MQ→有界消费→派发**
Run(后端窗口):提交一个生成任务 → 观察 `game_aigc_task` 状态 queued→running(经 MQ 消费认领,非等 `@Scheduled` tick);并发提交 20 个,消费侧并发不超 15;杀 MQ 信号(直接 DB 插 queued 不发 MQ)→ 兜底轮询补偿捡起。
- [ ] **Step 7:Commit**
```bash
git add game-cloud/game-module-aigc/ game-cloud/huijing-server/src/main/resources/application.yaml
git commit -m "$(cat <<'EOF'
feat(game-cloud): RocketMQ 异步 gen 队列(生产者+有界消费≤15+CAS 幂等)(切片一 阶段〇 C3)
入队发 AIGC_GEN_TOPIC 信号替 @Scheduled 轮询触发;消费者 consumeThreadMax=15(在跑并发权威)+ 既有
claimQueuedTask CAS 认领幂等;派发失败落 failed 不留悬空 running;@Scheduled 降级兜底(捞 stale queued +
超时 running)。发 MQ best-effort 不回滚入队。单测锁幂等/best-effort;端到端随后端窗口。
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
EOF
)"
```
---
# Part D · Wiring(Service 注册 Nacos + Python 热读 Nacos)
**为什么:** 生产形态是「Java 管准入+队列,Python 只服务 `/chat`」(设计 §3.7)。为此:AgentScope Service 要注册进 Nacos discovery(让 Java 消费者按服务名发现它、替硬编码 `@8200`);Python worker 的 middleware 要经 nacos-sdk-python 热读预算/阈值(消灭原自建薄读接口,§3.3)。
### Task D1:AgentScope Service 注册进 Nacos discovery
**Files:**
- Modify:`tier2/gen-worker/service/app.py:264-285`(`main()` uvicorn 起服务前后加 Nacos 实例注册 + 优雅注销)
- Create:`tier2/gen-worker/service/nacos_registry.py`(注册/心跳/注销封装 + 代理旁路)
- Modify:`cheap-worker/requirements.txt` 或 tier2 依赖清单(加 `nacos-sdk-python<2`,钉 v1)
**Interfaces:**
- Consumes:B1 Nacos(`100.64.0.8:8848`)。tier2 Service 现绑 `TIER2_SERVICE_HOST/PORT`(默认 `0.0.0.0:8200`,`app.py:272-273`,非死硬编码)。
- Produces:Service 以 `ip=<mini-desktop 100.64.0.7>, port=8200, service-name=agentscope-gen-service` 注册进 Nacos;Java 侧按服务名发现。
- [ ] **Step 1:装 nacos-sdk-python**
Run:`cheap-worker/.venv/bin/pip install "nacos-sdk-python<2"` → **钉 v1**(`import nacos` 的 `NacosClient`,含 `add_naming_instance`(注册)+ `add_config_watcher`(同步 watcher,内部起后台轮询线程,免 asyncio event-loop))。装后 `cheap-worker/.venv/bin/python -c "import nacos; print(nacos.__version__)"` 确认 1.x,记进依赖清单钉版本(与 AgentScope 钉 2.0.2 同纪律)。
- [ ] **Step 2:写注册封装(含代理旁路 + 优雅注销)**
Create `tier2/gen-worker/service/nacos_registry.py`(注册 naming 实例 + 心跳;nacos client 内部走 HTTP,须绕 fake-ip 代理):
```python
"""nacos_registry.py —— AgentScope Service 注册进 Nacos discovery(替硬编码 @8200,让 Java 按服务名发现)。
内网铁律:nacos client 走 HTTP,必须绕系统 fake-ip 代理(把 Nacos host 塞进 NO_PROXY,须在建 client 前)。
best-effort:注册失败只告警不中断 Service 起(生成主链不依赖被发现;发现只为 Java 消费方定位)。
"""
import os
import atexit
_SERVICE_NAME = "agentscope-gen-service"
def _bypass_proxy_for(host: str) -> None:
"""把 host 并进 NO_PROXY/no_proxy(须在建 nacos client 前;否则内网直连被 fake-ip 代理拦)。"""
for key in ("NO_PROXY", "no_proxy"):
cur = os.environ.get(key, "")
if host not in cur:
os.environ[key] = (cur + "," + host).strip(",") if cur else host
def register(service_host: str, service_port: int) -> None:
"""把本 Service 实例注册进 Nacos naming(best-effort;失败告警不抛)。"""
nacos_addr = os.environ.get("NACOS_SERVER_ADDR", "100.64.0.8:8848")
_bypass_proxy_for(nacos_addr.split(":")[0])
try:
import nacos # v1 NacosClient(同步,统一 naming + config;免 asyncio event-loop)
client = nacos.NacosClient(
nacos_addr,
namespace=os.environ.get("NACOS_NAMESPACE", "public"),
username=os.environ.get("NACOS_USERNAME", "nacos"),
password=os.environ.get("NACOS_PASSWORD", ""),
)
# ip 用「Java 侧能连到的」地址:mini-desktop Tailscale 100.64.0.7(不是 0.0.0.0)
advertise_ip = os.environ.get("SERVICE_ADVERTISE_IP", "100.64.0.7")
client.add_naming_instance(_SERVICE_NAME, advertise_ip, service_port, healthy=True)
print(f"[nacos-registry] 已注册 {_SERVICE_NAME} @ {advertise_ip}:{service_port}", flush=True)
def _deregister(): # 优雅注销(进程退出时摘实例,避免 Nacos 留死实例)
try:
client.remove_naming_instance(_SERVICE_NAME, advertise_ip, service_port)
print(f"[nacos-registry] 已注销 {_SERVICE_NAME}", flush=True)
except Exception as e: # noqa: BLE001
print(f"[nacos-registry] 注销失败(忽略):{e}", flush=True)
atexit.register(_deregister)
except Exception as e: # noqa: BLE001 —— 注册失败不连累 Service 起
print(f"[nacos-registry] 注册失败(best-effort,不中断 Service):{type(e).__name__}: {e}", flush=True)
```
- [ ] **Step 3:在 main() 起服务前调注册**
Modify `tier2/gen-worker/service/app.py`(`main()` :264 附近,`uvicorn.run` 之前):
```python
# Nacos discovery 注册(best-effort;NACOS_REGISTER=1 才注册,默认关,不影响本地/CLI 跑)
if os.environ.get("NACOS_REGISTER") == "1":
from service import nacos_registry
nacos_registry.register(host, port) # host/port 为 :272-273 读出的 TIER2_SERVICE_HOST/PORT
```
- [ ] **Step 4:单测(mock nacos client,验注册调用 + 代理旁路 + best-effort)**
Create `tier2/gen-worker/tests/test_nacos_registry.py`(pytest 型;mock `nacos.NacosClient`):
- 断言 `register` 调 `add_naming_instance(service_name, advertise_ip, port, ...)`。
- 断言调用后 `os.environ["NO_PROXY"]` 含 Nacos host。
- 断言 `NacosClient` 构造抛异常时 `register` **不抛**(best-effort)。
Run:`PYTHONPATH=tier2/gen-worker cheap-worker/.venv/bin/python -m pytest tier2/gen-worker/tests/test_nacos_registry.py -v`
预期:通过。
- [ ] **Step 5:Commit**
```bash
git add tier2/gen-worker/service/nacos_registry.py tier2/gen-worker/service/app.py tier2/gen-worker/tests/test_nacos_registry.py
git commit -m "$(cat <<'EOF'
feat(gen-worker): AgentScope Service 注册进 Nacos discovery(替硬编码@8200)(切片一 阶段〇 D1)
nacos_registry 注册实例+心跳+优雅注销,advertise_ip 指 mini-desktop Tailscale;代理旁路;
NACOS_REGISTER=1 才开、best-effort 失败不中断 Service。单测 mock client 锁注册+旁路+best-effort。
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
EOF
)"
```
### Task D2:Python worker middleware 经 nacos-sdk-python 热读预算/阈值
**为什么:** 预算目标、门阈值这类不在 `/agent`+`/session` 的 middleware 参数,写进 Nacos 独立 dataId,Python worker 的 middleware 用 **v1 `nacos.NacosClient.add_config_watcher`** 订阅长轮询即热更(v1 watcher **同步**、内部自起后台轮询线程,回调更新线程安全 dict 给同步 middleware 读 —— 不需要 asyncio event loop,正合线程模型的 worker;这也是 D1/D2 统一钉 v1 而非 v2 async `add_listener` 的原因)。fail-safe 由 `get_config` 内建三级兜底(本地 failover → Nacos 服务器 → 快照目录),Nacos 挂了不连累生成,消灭原自建薄读接口。**执行前先核真契约(比照 Part A 对 AgentScope 的严谨):** 装好 v1 后读 `nacos-sdk-python` 源确认 `add_config_watcher(data_id, group, cb)` 的**回调参数真实结构**(v1 回调收一个 dict,键含 `data_id`/`group`/`content`/`md5` 等)与 `get_config` 的三级兜底行为,把下面 fake 的契约对齐到真实签名再落。
**Files:**
- Create:`tier2/gen-worker/worker/nacos_hotconfig.py`(v1 同步 watcher + 线程安全 config store + 快照兜底)
- Create:`tier2/gen-worker/tests/test_nacos_hotconfig.py`(fake nacos client,验热更 + 快照兜底 + 线程安全读)
**Interfaces:**
- Consumes:B1 Nacos。dataId `gen-hot-params`(JSON:`{rmb_hard_limit, gate_thresholds, ...}`)。
- Produces:`NacosHotConfig.get(key, default)` —— 同步、线程安全、任何时候可读(冷启动读快照/默认,listener 推送后读新值)。middleware(如 `CircuitBreakerMiddleware`)构造时可从它取 `rmb_hard_limit` 等热参,替 `cheap_studio.py:156` 的硬编码 `10.0`。
- [ ] **Step 1:写失败测试(热更 + 快照兜底 + 线程安全)**
Create `tier2/gen-worker/tests/test_nacos_hotconfig.py`(用 fake nacos client,不触真网络):
```python
"""test_nacos_hotconfig.py —— Nacos 热配读(listener 推送热更 + get_config 快照兜底 + 线程安全)。
用 fake nacos client(不触网):
· 初始 get_config 返 {"rmb_hard_limit": 10.0} → get() 读到 10.0
· 触发 listener 回调(模拟推送){"rmb_hard_limit": 20.0} → get() 热读到 20.0(不重启)
· get_config 抛(Nacos 不可达)→ 回落快照/默认,get() 仍返默认、不抛
跑:PYTHONPATH=tier2/gen-worker cheap-worker/.venv/bin/python -m pytest tier2/gen-worker/tests/test_nacos_hotconfig.py -v
"""
import json
import time
import pytest
class _FakeClient:
"""fake nacos client:get_config 返预置内容;add_listener 存回调供测试手动触发。"""
def __init__(self, content: str | None):
self._content = content
self._cb = None
self.unreachable = False
def get_config(self, data_id, group, timeout=None):
if self.unreachable:
raise RuntimeError("nacos unreachable")
return self._content
def add_config_watcher(self, data_id, group, cb):
self._cb = cb # v1 真回调收一个 dict(键含 data_id/group/content/md5);此处比照
def push(self, new_content: str): # 测试用:模拟 Nacos 推送(按 v1 回调 dict 形状)
self._content = new_content
if self._cb:
self._cb({"data_id": "gen-hot-params", "group": "DEFAULT_GROUP",
"content": new_content, "md5": ""})
def test_initial_read_from_get_config():
from worker.nacos_hotconfig import NacosHotConfig
fake = _FakeClient(json.dumps({"rmb_hard_limit": 10.0}))
hc = NacosHotConfig(client=fake, data_id="gen-hot-params", defaults={"rmb_hard_limit": 3.0})
hc.start()
assert hc.get("rmb_hard_limit", 3.0) == 10.0
def test_hot_update_via_listener():
from worker.nacos_hotconfig import NacosHotConfig
fake = _FakeClient(json.dumps({"rmb_hard_limit": 10.0}))
hc = NacosHotConfig(client=fake, data_id="gen-hot-params", defaults={"rmb_hard_limit": 3.0})
hc.start()
fake.push(json.dumps({"rmb_hard_limit": 20.0}))
time.sleep(0.2) # 让回调线程处理
assert hc.get("rmb_hard_limit", 3.0) == 20.0
def test_unreachable_falls_back_to_default_no_throw():
from worker.nacos_hotconfig import NacosHotConfig
fake = _FakeClient(None)
fake.unreachable = True
hc = NacosHotConfig(client=fake, data_id="gen-hot-params", defaults={"rmb_hard_limit": 3.0})
hc.start() # 不应抛
assert hc.get("rmb_hard_limit", 3.0) == 3.0 # 回落默认
```
Run:`PYTHONPATH=tier2/gen-worker cheap-worker/.venv/bin/python -m pytest tier2/gen-worker/tests/test_nacos_hotconfig.py -v`
预期:FAIL(`ModuleNotFoundError: worker.nacos_hotconfig`)。
- [ ] **Step 2:写 NacosHotConfig(线程安全 store + 快照兜底 + listener)**
Create `tier2/gen-worker/worker/nacos_hotconfig.py`:
```python
"""nacos_hotconfig.py —— middleware 热参(预算/阈值)经 Nacos 热读,消灭原自建薄读接口(设计 §3.3)。
机制:构造时 get_config 拿初值(get_config 内建三级兜底:本地 failover→Nacos→快照,Nacos 挂了不连累生成);
add_config_watcher 订阅长轮询推送 → 回调把新内容解析进线程安全 store → 同步 middleware 随时 get() 读当前值。
用 v1 `nacos.NacosClient.add_config_watcher`:同步 watcher、内部自起后台轮询线程,回调更新线程安全 dict,
不需 asyncio event loop(正合线程模型 worker)。best-effort:解析/取值任何异常都回落默认,绝不抛。
"""
import json
import threading
class NacosHotConfig:
def __init__(self, *, client, data_id: str, group: str = "DEFAULT_GROUP", defaults: dict | None = None):
self._client = client
self._data_id = data_id
self._group = group
self._defaults = dict(defaults or {})
self._lock = threading.RLock()
self._store: dict = dict(self._defaults) # 冷启动先摆默认,start() 再拉真值
def start(self) -> None:
"""拉初值 + 订阅推送。best-effort:Nacos 不可达则留默认、不抛。"""
self._load_once()
try:
self._client.add_config_watcher(self._data_id, self._group, self._on_change)
except Exception as e: # noqa: BLE001 —— 订阅失败:留一次性初值,不热更,但不炸
print(f"[nacos-hotconfig] 订阅失败(留静态初值,不热更):{e}", flush=True)
def _load_once(self) -> None:
try:
content = self._client.get_config(self._data_id, self._group) # 内建快照兜底
self._apply(content)
except Exception as e: # noqa: BLE001 —— Nacos 不可达:留默认
print(f"[nacos-hotconfig] 初读失败(回落默认):{e}", flush=True)
def _on_change(self, event) -> None:
"""v1 watcher 回调:event 是 dict(v1 传 {data_id,group,content,md5,...}),取 content 解析进 store(best-effort)。"""
content = event.get("content") if isinstance(event, dict) else event
self._apply(content)
def _apply(self, content) -> None:
if not content:
return
try:
parsed = json.loads(content)
if isinstance(parsed, dict):
with self._lock:
self._store.update(parsed)
print(f"[nacos-hotconfig] 热参更新:{list(parsed.keys())}", flush=True)
except Exception as e: # noqa: BLE001 —— 脏内容不覆盖当前值,不抛
print(f"[nacos-hotconfig] 解析失败(保留旧值):{e}", flush=True)
def get(self, key: str, default=None):
"""同步、线程安全读当前热参;缺则返 default。"""
with self._lock:
return self._store.get(key, default)
```
(注:接口按 **v1 `nacos.NacosClient`** 的 `get_config(data_id, group)` / `add_config_watcher(data_id, group, cb)` 定,二者都同步、无需 event loop;真 client 直接传进 `NacosHotConfig(client=...)`。执行 Step 3 前按 D2「为什么」段先核 v1 回调 dict 的真实键名,与上面 fake 对齐即可,store/get 语义不变。)
- [ ] **Step 3:跑测试到绿**
Run:`PYTHONPATH=tier2/gen-worker cheap-worker/.venv/bin/python -m pytest tier2/gen-worker/tests/test_nacos_hotconfig.py -v`
预期:3 passed(初读 / 热更 / 不可达回落)。
- [ ] **Step 4:真连烟测(随后端窗口)—— 改 Nacos dataId,worker 热读到新值**
Run(后端窗口):Nacos 建 dataId `gen-hot-params` = `{"rmb_hard_limit": 10.0}`;起一个接了 `NacosHotConfig` 的最小脚本 → `get("rmb_hard_limit")` 读到 10.0;Nacos 改成 20.0 → 数秒内 `get` 读到 20.0(不重启)。记为随窗口的真连烟测。
- [ ] **Step 5:Commit**
```bash
git add tier2/gen-worker/worker/nacos_hotconfig.py tier2/gen-worker/tests/test_nacos_hotconfig.py
git commit -m "$(cat <<'EOF'
feat(gen-worker): middleware 热参经 Nacos v1 同步 watcher 热读(快照兜底)(切片一 阶段〇 D2)
NacosHotConfig 线程安全 store + get_config 三级兜底 + v1 add_config_watcher 同步热更(免 event-loop),
消灭自建薄读接口。3 单测绿(初读/热更/不可达回落默认不抛)。真连烟测随后端窗口。
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
EOF
)"
```
---
## 端到端验证(整体收口判据)
1. **续修机制坐实(Part A)**:`test_repair_middleware_spike.py` 4/4 绿 —— 拦 finish + 注入 UserMsg + 续跑放行第二版、max_repairs 封顶、observe 半程不炸、干净血缘 pop。→ §3.4 on_reasoning 洋葱方案可落,阶段一接真门。(若栽在机制类问题 → 回落 `control_plane`,记录在案。)
2. **基建自托管(Part B)**:mini-infra 上 Nacos readiness 200 + RocketMQ Broker 注册;两者 lean JVM 内存在预算内(Nacos <700M / RocketMQ <2G),`free -h` 复查 available 仍 >2G、既有 RAGflow/MySQL 等未被挤崩。
3. **game-cloud 接 Java(Part C)**:本地 `./mvnw compile` + 相关单测绿(GenAdmissionResource 降级语义、GenTask 幂等/best-effort);真限流(击穿 ≤15 被拦)、真队列(入队→MQ→有界消费→派发)、真 Nacos 注册/配置读 = 随后端窗口的集成验证,代码 + 编译 + 单测先落。
4. **Wiring(Part D)**:Service Nacos 注册单测绿(mock client);Python 热配 3 单测绿(热更/快照兜底);真连烟测(改 Nacos dataId worker 热读到新值)随后端窗口。
5. **回归**:cheap-worker 既有单测全绿(`test_budget_gate` 等);tier2 既有单测全绿;品牌门 / 入口卫生门 / canonical 门干净。
## 风险与回滚
- **续修 spike 栽(最需先坐实)**:API 签名类当场对齐已装 2.0.2 修;机制类(压制/注入/续跑之一根本不发生)= 回落现有 `control_plane` 外层循环(已存在、已跑通),架构其余(Part B/C/D)不受影响。判据见 Task A1 Step 5。
- **mini-infra RAM 顶到边(8.6G 可用 / 4 核)**:Nacos/RocketMQ 已 lean 调(512m/512m/1g);仍紧则进一步压 broker `-Xmx` 或临时暂停非必需容器(不动数据面)。4 核 CPU 争用在 MVP 单操作者下可接受;放量前若压不住,再议 RocketMQ 迁独立机。
- **Part C 并发权威定位(Opus 评审已纠)**:在跑并发≤15 权威 = C3 `consumeThreadMax`(MVP 单实例),Sentinel = 入口 QPS 保护(非并发权威),既有 DB 门 = per-creator 配额,三层组合。多实例集群全局限流留放量后(Token Server / Redis)。若创始人要本轮就上 Token Server 做集群全局限流,按判调 C2/C3;不影响 Part A/B/D。
- **Java 集成真验依赖后端窗口**:C1/C2/C3、D1/D2 的真部署/端到端验证排后端窗口;本轮交付 = 代码 + 本地编译/单测 + 探针脚本,窗口到即可跑集成。不动 live 服务无创始人授权。
- **nacos-sdk-python v1 回调契约**:D1/D2 钉 v1 `NacosClient`(同步 watcher,免 event-loop)。真连烟测须验 `add_config_watcher` 回调真触发 + worker 线程读到热值;执行前先核 v1 回调 dict 键名(比照 Part A 对 AgentScope 的严谨),与 fake 对齐。
## 不在本轮范围(明确划界)
- **阶段一主体**:cheap-worker 从 `cheap_studio` 归并到 Service `/chat`、续修 middleware 换真 `run_gates` 独立重跑、九门改 workspace MCP、`_extract_function_tools` 补 mcps —— 待 Part A spike 结论后规划(设计 §3.1/§3.4)。
- **两条门线归一**:tier2 `layerResults.L1` vs cheap `guards` 的 verdict 形状差异,续修 middleware 接真门时的归一 = 阶段一的活。
- **阶段二/三/四**:配置中心 yudao⊕Nacos 版本化 UI、观测 Studio+genMonitor —— 后续各出 plan。
- **SoT 收口(doc,非本 plan 执行项,交横切一致性主人)**:① spec §关联 的 `types/_hook.py` 是 2.0.2 废弃死代码,改标 `MiddlewareBase`;② spec §3.7「Sentinel 全局≤15 权威」软化为「Sentinel=入口 QPS + consumeThreadMax=并发≤15(单实例)+ 集群全局限流放量后」(Opus 评审纠承重错);③ Nacos 反转同步 5 架构档(设计 §6 pending);④ resume=外层循环旧叙述降 fallback。