新建 plan/设计档 frontmatter 必带一行 上级:(仓根相对路径),沿链 6 跳内达 canonical SoT; 全链/树/当前位置都是查询结果(plan-tree.py)而非维护对象,_index 在飞板保持线级、不加逐叶登记。 G7 拦三种腐坏:缺字段 / 上级死链 / 链断(中途档既非 canonical 又无上级),成环与超跳同挡(均已植坏档实测命中)。 存量修复:统一执行计划基建线补认领配置控制面设计(治真断链、含 07-01 SCA 选型反转口径), 配置控制面设计与阶段〇/一① plan 补 上级:;AGENTS.md/.agents README/feature-design-doc 模板同步七检口径。 Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
74 KiB
date, topic, status, 关联设计, 上级
| date | topic | status | 关联设计 | 上级 |
|---|---|---|---|---|
| 2026-07-01 | 配置控制面-spike与阶段〇 | 已执行完成(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 阶段〇) | docs/agent-specs/2026-06-30-配置控制面一次性按序实现-设计.md |
配置控制面 · 续修 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 用真 IProot@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 造词黑话、无套话空强调。
执行顺序与依赖
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:
"""
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
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 规则源都用它。MySQLnacos库复用既有100.64.0.8:3306(root/ZRH3jwYLOrntBcTAw29MW9BP,不新增实例)。 -
Step 1:写健康探针(先失败 —— 服务未起)
Create deploy/infra/nacos/health-check.sh(沿 deploy/smoke-test.sh 风格:bash + curl + 退出码):
#!/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:
#!/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:
# 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 从凭据文档取):
# 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 入库 —— 内网阶段密钥进文档铁律)。
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、Broker100.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:
#!/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:
# 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:
# 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:
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」。
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 在此的价值 = 热可调的入口速率闸 + 未来多实例时升级成集群限流的落点。 - 既有
enqueueWithControlPlaneDB 三门保留作业务级 per-creator 配额(每日额度 + 每创作者并发)+ 全局背压queue-depth-limit作纵深防御(设 ≥15 兜底)—— Sentinel/consumeThreadMax 做不了业务配额,正交。 - RocketMQ = 异步投递机制,
game_aigc_task表保留作任务记录 + CAS 幂等(claimQueuedTaskqueued→running 仍用);生产者入队发消息 → Java 有界消费者消费 → 调 Service。替换的是@ScheduledDB 轮询这个触发方式,不是任务表本身。
「并发权威=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.enabledfalse→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 与鉴权:
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 模块),然后:
# 该服务应出现在 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
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_QPScount=10 入口速率闸)
Interfaces:
-
Consumes:C1 的 Nacos(规则源 + 鉴权)。
-
Produces:生成入队入口的 QPS 被 Sentinel 限(防突发洪峰,
FLOW_GRADE_QPS);越限走blockHandler优雅拒(返「排队中/稍后再试」,非 500)。规则改 Nacos dataIdsentinel-gen-flow-rules即热生效(不重启)。并发≤15 不在本 Task,由 C3consumeThreadMax保证。 -
Step 1:引 Sentinel 依赖
Modify aigc 模块 pom,新增(版本由 SCA 2025.0.0.0 BOM 托管、不写死):
<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:
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 三门执行:
// 注入(类字段):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 段:
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):
[
{
"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 限流由集成验证覆盖,单测只锁降级语义):
@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
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(@Scheduledtick 降级为 MQ 死信/兜底补偿,不再作主触发)
Interfaces:
-
Consumes:B2 的 RocketMQ(
100.64.0.8:9876,topicAIGC_GEN_TOPIC)+ 既有game_aigc_task表 + CASclaimQueuedTask。 -
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:
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:
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
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 代理):
"""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 之前):
# 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
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,不触真网络):
"""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:
"""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
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
)"
端到端验证(整体收口判据)
- 续修机制坐实(Part A):
test_repair_middleware_spike.py4/4 绿 —— 拦 finish + 注入 UserMsg + 续跑放行第二版、max_repairs 封顶、observe 半程不炸、干净血缘 pop。→ §3.4 on_reasoning 洋葱方案可落,阶段一接真门。(若栽在机制类问题 → 回落control_plane,记录在案。) - 基建自托管(Part B):mini-infra 上 Nacos readiness 200 + RocketMQ Broker 注册;两者 lean JVM 内存在预算内(Nacos <700M / RocketMQ <2G),
free -h复查 available 仍 >2G、既有 RAGflow/MySQL 等未被挤崩。 - game-cloud 接 Java(Part C):本地
./mvnw compile+ 相关单测绿(GenAdmissionResource 降级语义、GenTask 幂等/best-effort);真限流(击穿 ≤15 被拦)、真队列(入队→MQ→有界消费→派发)、真 Nacos 注册/配置读 = 随后端窗口的集成验证,代码 + 编译 + 单测先落。 - Wiring(Part D):Service Nacos 注册单测绿(mock client);Python 热配 3 单测绿(热更/快照兜底);真连烟测(改 Nacos dataId worker 热读到新值)随后端窗口。
- 回归: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.L1vs cheapguards的 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。