diff --git a/tier2/config/infra.yaml b/tier2/config/infra.yaml new file mode 100644 index 00000000..ec9f2ac0 --- /dev/null +++ b/tier2/config/infra.yaml @@ -0,0 +1,58 @@ +# tier2/config/infra.yaml —— tier2 基建数据面端点 + 凭据(运行时读;内网入仓) +# ════════════════════════════════════════════════════════════════════════════ +# 【这份解决什么问题】 +# tier2 P4 服务化(Agent Service / create_app)与落库要连三套内网基建:Redis(session 状态 + +# MessageBus)、MySQL(源工程版本 manifest 表)、MinIO(源文件全文 S3 兼容 OSS)。这些端点 + 口令 +# 按创始人铁律(2026-06-18:内网阶段所有决策 + 密钥进项目文档、内网 Tailscale 凭据允许入仓) +# 集中落在一份配置文件里、运行时读;代码从配置取,绝不把端点/口令硬编码进 .py。 +# +# 【凭据出处】 +# 值誊自 docs/内网凭据与端点.md「基建数据面(mini-infra 100.64.0.8)」段(2026-06-24 从容器 env 抠出 +# 落档)。该档是 single source of truth;此处是给 tier2 运行时读的副本。若凭据轮换,先改那份档, +# 再同步本文件(横切一致性主人 = 创始人 + 6c6g 文档线负责对账)。 +# +# 【运行时读 + 三级回落】(口径对齐 worker/genconfig.py) +# service/infra_config.py 的 get(section, key, default) 每次调用按以下顺序取值,任一级取到即用: +# ① env 单点覆盖:TIER2_INFRA__
__(全大写、双下划线分隔)——临时压一个值不必改文件 +# (例:换 Redis 端口 export TIER2_INFRA__REDIS__PORT=6380)。 +# ② 本 YAML 的 section.key。 +# ③ 调用方传入的 default(连本文件都缺该 key 时的最后防线)。 +# 读文件 / 解析 YAML / env 转型的任何异常都不抛、不中断(best-effort):读不到或脏 → 落调用方 default。 +# +# 【安全口径】内网 Tailscale 段(100.64.0.x),口令明文入仓是创始人明确授权的内网阶段做法(见上)。 +# 一旦 tier2 走出内网 / 上公网,口令必须迁出仓库改走密钥管理 —— 这条作为 followup 长期挂在 P4 收口单。 +# ════════════════════════════════════════════════════════════════════════════ + +# ── Redis:Agent Service(create_app)的 session 状态持久化 + MessageBus 跨会话事件总线 ── +# AgentScope 2.0.2 的 RedisStorage(host/port/db/password) 与 RedisMessageBus(host/port/db/password) +# 都只有 Redis 实现(无 in-memory/SQLite 版本),生产必须配 Redis。源码核验: +# storage/_redis_storage.py:77 RedisStorage.__init__(host, port, db, password, ...) +# message_bus/_redis_message_bus.py:44 RedisMessageBus.__init__(host, port, db, password, ...) +# Redis 需 AUTH(NOAUTH 已实证),故必须带 password。storage 与 message_bus 故意解耦,可走不同 db +# 隔离键空间(本配置默认同实例:storage→db 0,message_bus→db 1;按需调)。 +redis: + host: "100.64.0.8" # mini-infra Tailscale IP + port: 6379 + password: "9ea28f5d28d68b09bfd7ccfc31216a52" # requirepass(NOAUTH 已实证,必带) + storage_db: 0 # RedisStorage 用的逻辑库(session/agent/team 记录) + message_bus_db: 1 # RedisMessageBus 用的逻辑库(事件日志/inbox/run-lock;与 storage 隔离) + +# ── MySQL:tier2 落库源工程版本表(manifest)——P4 第二装载(U2)用 ── +# 现状:Agent Service(U1)本身不直接连 MySQL(它的 session 状态走 Redis);MySQL 是 U2 落库面用的 +# (源工程版本 manifest + 寻址)。先在此登记端点/口令,U2 的 store 实现从这里读、不重复抠凭据。 +mysql: + host: "100.64.0.8" + port: 3306 + user: "root" + password: "ZRH3jwYLOrntBcTAw29MW9BP" + database: "tier2" # tier2 专用库(与 game-cloud 业务库隔离;U2 建表前先 CREATE DATABASE) + +# ── MinIO(S3 兼容 OSS):tier2 落库源文件全文 —— P4 第二装载(U2)用 ── +# key 约定:tier2-src///<工程内相对路径>(见凭据档)。与 ragflow 共用实例, +# tier2 另开 bucket(tier2-src)隔离。现状同 MySQL:U1 不直接用,U2 落库面读它。 +minio: + endpoint: "100.64.0.8:9000" # S3 API 端点(控制台在 9001) + access_key: "ragflow" # root 即 ragflow(与 ragflow 共用实例) + secret_key: "6c4b77b2f055c8a66e00e3c38ef3d818c166951d478cc7cc" + bucket: "tier2-src" # tier2 源文件全文桶(U2 落库前先 mb 建桶) + secure: false # 内网 http(非 https) diff --git a/tier2/config/schema/tier2_source_project_version.sql b/tier2/config/schema/tier2_source_project_version.sql new file mode 100644 index 00000000..fb0690bc --- /dev/null +++ b/tier2/config/schema/tier2_source_project_version.sql @@ -0,0 +1,33 @@ +-- tier2/config/schema/tier2_source_project_version.sql +-- ════════════════════════════════════════════════════════════════════════════ +-- tier2 富游戏自治生成线 · 第二装载落库(U2)源工程版本表。 +-- +-- 【这张表解决什么】 +-- F 族要素⑦ addressing:把 finish 交付的源工程 manifest(七要素去掉 fileTree 各文件 content) +-- 落 MySQL 一行,按 (game_id, source_hash) 幂等、按 (game_id, version_id) 取回寻址。源文件全文 +-- (fileTree 各文件 content)不进本表,落 MinIO(S3 兼容 OSS),manifest.addressing.sourceUrl 指向 +-- MinIO 版本前缀。改源重建 → 新 version_id,旧版本保留(F1 版本寻址 / 游戏=长生命周期项目)。 +-- +-- 【与代码的关系】 +-- worker/store.py 的 BackendStore 首次 save 时会执行等价的 CREATE TABLE IF NOT EXISTS(代码内自建, +-- 见 _DDL_SOURCE_PROJECT_VERSION),故部署不强制先手动跑本文件。本文件是同源的【可审阅 DDL 留档】: +-- DBA / 后端接线 / 评审按它对账表结构;若手工预建,在 tier2 库执行本文件即可(代码再跑 IF NOT EXISTS 不会重建)。 +-- 表所在库 = infra.yaml mysql.database(默认 tier2;与 game-cloud 业务库隔离)。 +-- +-- 【MinIO key 约定(与本表配套,不在本 DDL 内)】 +-- object key = //<工程内相对路径>,bucket = infra.yaml minio.bucket(默认 tier2-src)。 +-- ════════════════════════════════════════════════════════════════════════════ + +CREATE TABLE IF NOT EXISTS tier2_source_project_version ( + id BIGINT NOT NULL AUTO_INCREMENT COMMENT '自增主键', + game_id VARCHAR(128) NOT NULL COMMENT '游戏稳定标识(同款多版本共享;落库 id)', + version_id VARCHAR(96) NOT NULL COMMENT '版本号 vXXXX-hash12(改源重建生成新版本)', + source_hash CHAR(64) NOT NULL COMMENT '源工程内容指纹 sha256(= source_project.contentHash;幂等键)', + manifest_json LONGTEXT NOT NULL COMMENT '源工程 manifest 七要素形状(去 fileTree content;含 addressing)', + addressing VARCHAR(512) NULL COMMENT '落库寻址摘要(store/sourceUrl;冗余出 manifest 便于检索)', + created_at DATETIME NOT NULL COMMENT '落库时刻(now_ts 转 UTC datetime)', + PRIMARY KEY (id), + UNIQUE KEY uk_game_source (game_id, source_hash) COMMENT '幂等键:同款同源不重复落', + UNIQUE KEY uk_game_version (game_id, version_id) COMMENT '版本寻址键:同款版本号唯一', + KEY idx_game_created (game_id, created_at) COMMENT 'fetch 取最新版本走它' +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='tier2 富游戏源工程版本表(U2 第二装载落库)'; diff --git a/tier2/gen-worker/requirements.txt b/tier2/gen-worker/requirements.txt index b9eaf2a3..9beba7ec 100644 --- a/tier2/gen-worker/requirements.txt +++ b/tier2/gen-worker/requirements.txt @@ -5,7 +5,12 @@ agentscope==2.0.2 openai>=2.41.1 # 便宜档(deepseek-v4 等)走 OpenAI 兼容端点(new-api) json-repair # 脏 JSON 兜底(承袭裸路) -pyyaml # 模型路由 models.yaml +pyyaml # 模型路由 models.yaml + infra.yaml 基建配置 + +# U2 第二装载落库(BackendStore,仅 TIER2_STORE=backend 启用;worker/store.py 内惰性 import, +# 6c6g 不装这俩也能 py_compile / import store.py;真跑/部署在 mini-desktop 时装)。 +pymysql # MySQL 源工程版本 manifest 表(tier2_source_project_version) +minio # MinIO(S3 兼容 OSS)源文件全文落库 # 注:M3 走 AgentScope 内置 AnthropicChatModel(已含在 agentscope)。 # 注:Phaser / esbuild 是 JS 侧工具链(game-runtime / tier2/harness),不在此 Python 依赖内。 diff --git a/tier2/gen-worker/service/__init__.py b/tier2/gen-worker/service/__init__.py new file mode 100644 index 00000000..801f9811 --- /dev/null +++ b/tier2/gen-worker/service/__init__.py @@ -0,0 +1,15 @@ +"""tier2/gen-worker/service —— tier2 富游戏自治线 · Agent Service 服务壳(U1)。 + +把现在的「本地 CLI runner 直驱裸 Agent」(run_engine.py / agent_loop.studio.run_studio)升级成经 +AgentScope 2.0.2 官方 Agent Service(`agentscope.app.create_app`)接入的服务:多租户 / 多会话 / +REST 触发 + SSE 事件流,session 状态 + 跨会话 MessageBus 走 Redis(mini-infra)。 + +模块: + - infra_config.py:基建数据面配置加载器(运行时读 tier2/config/infra.yaml,不 import 重依赖)。 + - app.py:create_app 服务壳 —— 复用 studio 的 agent 装配件(roles 系统提示词 + 九工具 Toolkit + + 四熔断 / trace 中间件 + BYPASS 权限 + 历史压缩),Redis 连接从 infra_config 取。重依赖惰性 import。 + - bootstrap.py:session 三路初始化(新建 / 续接 / 加载已有工程迭代)经 REST API 的客户端编排 + 口径说明。 + +重依赖(agentscope[full] 的 app、redis 客户端)只在函数内惰性 import,保证 6c6g(未装重依赖)能 +py_compile + import 本包;真跑 / 部署在 mini-desktop(见 app.py 文件头 deploySteps)。 +""" diff --git a/tier2/gen-worker/service/app.py b/tier2/gen-worker/service/app.py new file mode 100644 index 00000000..3db48e92 --- /dev/null +++ b/tier2/gen-worker/service/app.py @@ -0,0 +1,278 @@ +"""service/app.py —— tier2 富游戏自治线 · Agent Service 服务壳(U1;on AgentScope v2.0.2)。 + +把「本地 CLI 直驱裸 Agent」升级成经官方 `agentscope.app.create_app` 的 Agent Service:多租户 / 多会话 / +REST 触发(POST /chat,fire-and-forget)+ SSE 事件流(GET /sessions/{sid}/stream),session 状态与 +跨会话 MessageBus 走 Redis(mini-infra)。 + +════════════════════════════════════════════════════════════════════════════ +【create_app 的真实形态(源码核验 app/_app.py:33 + examples/agent_service/main.py)】 + create_app(storage, message_bus, workspace_manager, *, extra_credentials, extra_middlewares, + extra_agent_middlewares, extra_agent_tools, custom_subagent_templates, + custom_agent_cls, title, version) -> fastapi.FastAPI + 关键事实(决定本壳怎么接,逐条源码核验): + 1. **create_app 不收 agent 实例,也不收 agent 工厂**。agent 由框架在【每个 chat 回合】内部装配 + (_service/_chat.py:_run_impl 第 317 行 `self._agent_cls(...)`)。我们能注入的只有装配的「零件」: + · custom_agent_cls —— Agent 子类(默认就用内置 Agent,本壳传 None); + · extra_agent_tools —— 异步工厂 (user_id, agent_id, session_id) -> list[ToolBase],每回合调一次, + 产出的工具加进 toolkit 的 "basic" 组(_app.py:104-111,_chat.py:240-251 的 extra_factory); + · extra_agent_middlewares —— 异步工厂同签名,产出的中间件接在框架中间件(InboxMiddleware / + StateChangeMiddleware / ToolOffloadMiddleware)之后(_app.py:95-103,_chat.py:273-280); + · custom_subagent_templates —— 工作室 Team 的 worker 蓝图(本壳暂留空,见下「工作室 Team 现状」)。 + 2. **模型不在 create_app 配,在 session 的 chat_model_config 配**,且必须先注册成 AgentScope + Credential。_chat.py:297-303 从 `session_record.config.chat_model_config` 经 get_model 解析模型, + 凭据来自存储里的 CredentialRecord。**故 M3(MiniMax-M3 经 new-api 走 Anthropic 原生)要先 + POST /credential 注册成 anthropic_credential,再在建 session 时引用**(端点/口令见 bootstrap.py 与 + docs/内网凭据与端点.md;凭据档由 new-api 网关托管,base_url=host 根、key=NEWAPI_KEY)。 + 3. **系统提示词 / context_config / react_config 在 AgentRecord 配**(POST /agent),_chat.py:319/323/324 + 从 agent_record.data 取。故 tier2 单写 agent 的 system_prompt(roles.writer_system)、历史压缩 + (config.build_context_config)、ReAct 轮数(ReActConfig(max_iters=...))经建 agent 时写进 AgentRecord。 + 4. **session 状态(AgentState:context/summary/cur_iter/permission/tool/tasks)持久化在 Redis**,每回合 + reload(_chat.py:315)→ 跑 → 回写(_chat.py:441)。这就是 session「续接」的原生机制,无需自写。 + 5. **storage 与 message_bus 都只有 Redis 实现,无 in-memory/SQLite 版本**,生产必须配 Redis 且需 AUTH。 + 6. **无 CLI / serve 命令**:启动 = 在本文件里调 uvicorn.run(见文末 main + deploySteps)。 + +【与 CLI 单写主链(agent_loop.studio.run_studio)的关系 —— 诚实差距,见 followups】 + run_studio 是「外层 Python 编排」:它在 Agent 的 ReAct 循环【之外】套了一圈有界 resume(agent 过早 + 停下就带 verdict 反馈踹回去续修)、阶段 1 工作室多 agent 设计、收口后的 L2/L3 软检、成本台账、落库寻址。 + 而 create_app 的 chat 是【每回合一次 agent.reply_stream】、由前端经 REST 决定要不要再发一轮—— + 外层编排的位置变成了「服务消费方 / 前端 / 控制面」。**所以本壳不可能、也不应该把 run_studio 整个塞进 + create_app**(那会与框架的 per-turn 模型打架)。本壳做的是:**复用 run_studio 用的同一批装配零件** + (roles 系统提示词 / build_toolkit 九工具 / 四熔断 + trace 中间件 / BYPASS 权限 / 历史压缩),把它们 + 接到 create_app 的扩展点上,**绝不重写任何生成逻辑**;有界 resume / 设计团队 / L2/L3 / 成本 / 落库 + 这些「循环外」能力如何在服务态落位,作为 followups 明确标差距(见文末)。 + +【session 三路初始化 —— 汇到同一 agent 装配点】(实现见 bootstrap.py) + 三路全部经 create_app 现成的 REST 端点达成(不自造端点、不旁路框架),最终都走 _chat.py 的同一 + 装配点(custom_agent_cls + extra_agent_tools + extra_agent_middlewares + session 的模型/系统提示词): + ① 新建(模板起手):POST /agent → POST /sessions(新 workspace_id)→ POST /chat(kick 文本引导 agent + scaffold_init 起手)。 + ② 续接会话:对已存在的 (agent_id, session_id) 直接再 POST /chat;框架自动 reload 该 session 的 + AgentState(含已写历史 / cur_iter / tasks)续跑——即「从 SessionRecord.state resume」。 + ③ 加载已有工程迭代:绑到「工程所在 workdir 对应的 agent_id」再开 session。AgentScope 的 + LocalWorkspaceManager workdir = basedir/agent_id(_local_workspace_manager.py:116,**按 agent_id 不按 + workspace_id**),故一款游戏工程 = 一个 agent_id;迭代 = 在该 agent_id 上新开 / 续用 session。 + (注:tier2 九工具自己的工程目录另在 game-runtime/games/_tier2-gen/,与 AgentScope + workspace 解耦——见下「九工具 workdir 与 AgentScope workspace 的关系」。) + +【九工具 workdir 与 AgentScope workspace 的关系(重要,避免误解)】 + tier2 九工具(scaffold_init/write_source/build/run_gates/finish 等)是 worker.toolkit.build_toolkit 产出的 + 纯 Python 闭包,它们读写自己的工程目录 game-runtime/games/_tier2-gen/(run._workdir), + **不依赖 AgentScope 的 workspace.workdir**。本壳把 game_id 绑成 AgentScope 的 session_id,九工具据此 + 各自管文件。AgentScope 的 LocalWorkspaceManager 是 create_app 的硬性必填项(每个 chat run 都要它), + 但本线生成产物落在九工具自己的目录,AgentScope workspace 基本只承载框架内置的 filesystem 工具 + (本壳不靠它们生成)。这是有意为之:九工具是 spike 已验证的生成主链,零改接进来。 + +【工作室 Team(多 agent 设计)现状】 + run_studio 阶段 1 的工作室星形多 agent 设计在 CLI 线是「纯库 import Agent」实现(design_team.py:13 明示 + studio.py 无 create_app、拿不到部署态 Team 原语)。create_app 自带部署态 Team(AgentCreate/TeamCreate/ + TeamSay + SubAgentTemplate),但语义与 CLI 线的 design_team 不同构。本壳先把 custom_subagent_templates + 留空(单写 agent 主链先服务化跑通),把「设计阶段在服务态如何落位(用部署态 Team 重写 or 设计阶段作为 + 前置 REST 调用产出 design_text 再塞进单写 agent 的 system_prompt)」列为 followup。 + +【惰性 import 红线】agentscope[full] 的 app、redis 客户端在 6c6g 未装。本模块顶层【绝不】import + agentscope.app / redis;所有重依赖只在 build_app / 工厂函数体内 import,保证 6c6g py_compile + import 过。 + worker.config 在顶层 import agentscope(模型客户端),故对它也惰性 import(在工厂体内),避免顶层连带炸。 +════════════════════════════════════════════════════════════════════════════ +""" + +from __future__ import annotations + +import os +import sys +from pathlib import Path +from typing import TYPE_CHECKING, Any + +# infra_config 不 import 重依赖,顶层 import 安全(6c6g 可用)。 +from . import infra_config + +# 仅类型检查期引用 FastAPI / agentscope 类型,运行时不 import(6c6g 无这些包也能 py_compile)。 +if TYPE_CHECKING: # pragma: no cover + from fastapi import FastAPI + + +# ── 包内/直跑兼容:让顶层包 `worker` / `observability` 可解析(与 studio.py 同款兜底)── +# 本文件在 tier2/gen-worker/service/app.py;把 gen-worker/ 加进 sys.path,使 `worker.*` 可 import +# (目录名 gen-worker 含连字符不可直接 import,但子目录 worker/ 是合法包名)。 +_GEN_WORKER_DIR = Path(__file__).resolve().parents[1] +if str(_GEN_WORKER_DIR) not in sys.path: + sys.path.insert(0, str(_GEN_WORKER_DIR)) + + +# 服务级常量:tier2 Agent Service 的标题(OpenAPI docs 显示)。 +SERVICE_TITLE = "tier2-rich-game-agent-service" + + +def _extract_function_tools(toolkit: Any) -> list: + """从 worker.toolkit.build_toolkit 返回的 Toolkit 里抠出 FunctionTool 列表(供 extra_agent_tools 工厂用)。 + + 为什么要抠:create_app 的 extra_agent_tools 工厂约定返回 `list[ToolBase]`(_types.py:21),框架会把 + 它们加进【框架自建的】toolkit 的 basic 组;而 tier2 的 build_toolkit 直接返回一个【已组装好的】 + Toolkit 实例。本壳不想绕过框架自建 toolkit(那会丢掉 workspace 工具 / planning / schedule 等),所以 + 取折中:从 tier2 Toolkit 里取出九个 FunctionTool 实例,交给框架合并。 + + 2.0.2 Toolkit 内部把工具存在 `tools`(dict[str, ToolBase])或等价结构;不同小版本字段名可能微调, + 故按几个候选名探测,全失败则诚实抛(让启动期就暴露,而非运行期静默少工具)。 + """ + # Toolkit 把注册的工具放在内部映射里(2.0.2:Toolkit.tools 为 dict[name->ToolBase])。按候选探测,容错。 + for attr in ("tools", "_tools", "function_tools"): + container = getattr(toolkit, attr, None) + if isinstance(container, dict) and container: + return list(container.values()) + if isinstance(container, (list, tuple)) and container: + return list(container) + raise RuntimeError( + "无法从 tier2 Toolkit 抠出 FunctionTool 列表(探测 tools/_tools/function_tools 均空):" + "AgentScope 2.0.2 Toolkit 内部结构可能与预期不符,请核对 build_toolkit 产物结构后调整 _extract_function_tools。" + ) + + +async def _tier2_tools_factory(user_id: str, agent_id: str, session_id: str) -> list: + """extra_agent_tools 工厂:每个 chat 回合产出 tier2 九工具(绑定 game_id = session_id)。 + + 签名严格对齐 AgentToolFactory((user_id, agent_id, session_id) -> Awaitable[list[ToolBase]],_types.py:21)。 + 把 AgentScope 的 session_id 当成 tier2 的 game_id —— 九工具据此各自管自己的工程目录 + (game-runtime/games/_tier2-gen/),与 AgentScope workspace 解耦(见文件头说明)。 + + 重依赖(worker.toolkit → 经 worker.config 连带 import agentscope)在此惰性 import,保证 6c6g import app.py 不炸。 + + Args: + user_id: 框架注入的用户 id(本线暂不按用户隔离工程目录,留作 followup;先用 session_id 作 game_id)。 + agent_id: 框架注入的 agent id。 + session_id: 框架注入的 session id —— 即本款游戏工程的 game_id。 + + Returns: + list[ToolBase]:九个 FunctionTool 实例(scaffold_init/write_source/validate_datatable/build/ + run_gates/read_verdict/finish 等;具体面见 worker.toolkit.build_toolkit)。 + """ + # 惰性 import:worker.toolkit 顶层经 worker.config 会 import agentscope(6c6g 无);故只在回合内 import。 + from worker.toolkit import Tier2Session, build_toolkit # noqa: PLC0415 —— 惰性 import 红线 + + # game_id 绑 session_id:一个 AgentScope session = 一款游戏工程(九工具据 game_id 管目录)。 + # play_spec 在服务态由控制面/前端经后续接口下发(spike 期先 None;真玩门 run_gates 仍可跑, + # 只是无定制驱动规格——按品类注册表默认)。这是与 CLI 的差异点之一,列入 followup。 + session = Tier2Session(session_id) + toolkit = build_toolkit(session) + return _extract_function_tools(toolkit) + + +async def _tier2_middlewares_factory(user_id: str, agent_id: str, session_id: str) -> list: + """extra_agent_middlewares 工厂:每个 chat 回合产出 tier2 的四熔断 + trace 中间件。 + + 签名严格对齐 AgentMiddlewareFactory((user_id, agent_id, session_id) -> Awaitable[list[MiddlewareBase]])。 + 产出的中间件接在框架中间件之后(_chat.py:273-280);trace 列在 breaker 前 → 它是更外层洋葱 + (先 ingest 事件再进熔断巡检,与 studio.py 同序)。trace_id 用 session_id 贯穿本款生成。 + + 重依赖(worker.middleware → 经包顶层连带 agentscope)在此惰性 import。 + + Returns: + list[MiddlewareBase]:[Tier2TraceMiddleware(trace_id=session_id), CircuitBreakerMiddleware()]。 + """ + # 惰性 import:worker.middleware 经 worker 包顶层会牵出 agentscope;故只在回合内 import。 + from worker.middleware import ( # noqa: PLC0415 —— 惰性 import 红线 + CircuitBreakerMiddleware, + Tier2TraceMiddleware, + ) + + # trace 在外、breaker 在内(与 studio.py middlewares=[tracer, breaker] 同序);sink=None → 步留内存, + # 真落库 sink 随控制面 phase-1 接(同 studio 现状)。每回合一组新实例(熔断计数按回合,非跨回合累计—— + # 这是与 CLI 单局累计的差异点,列入 followup:跨回合累计 ¥ 硬闸需把计数挂到 session 维度)。 + tracer = Tier2TraceMiddleware(trace_id=session_id) + breaker = CircuitBreakerMiddleware() + return [tracer, breaker] + + +def build_app(*, title: str = SERVICE_TITLE) -> "FastAPI": + """组装 tier2 Agent Service 的 FastAPI app(create_app + Redis storage/message_bus + 本地 workspace)。 + + 这是部署入口要调的工厂:`from service.app import build_app; app = build_app()`,再 uvicorn 起(见文末 main)。 + 所有重依赖(agentscope.app / redis 客户端,经 LocalWorkspaceManager 还要 fastapi)在此惰性 import, + 保证 6c6g 仅 import 本模块(不调 build_app)时不炸——真正连基建在 mini-desktop 调 build_app 时才发生。 + + 装配口径(逐条对应 create_app 扩展点,见文件头): + · storage = RedisStorage(**infra.redis storage 库) —— session/agent/team 记录持久化 + · message_bus = RedisMessageBus(**infra.redis message_bus 库) —— 跨会话事件 / inbox / run-lock + · workspace_manager = LocalWorkspaceManager(basedir=...) —— 框架必填(本线产物另落九工具目录) + · extra_agent_tools = _tier2_tools_factory —— 注入九工具(per-turn) + · extra_agent_middlewares = _tier2_middlewares_factory —— 注入四熔断 + trace(per-turn) + · custom_subagent_templates = [] (工作室 Team 暂留空,见文件头 followup) + · custom_agent_cls = None (用内置 Agent;tier2 装配靠零件注入,不必子类化) + + Returns: + fastapi.FastAPI:create_app 组装好的 app,直接交 uvicorn。 + """ + # ── 惰性 import 重依赖(6c6g 无;只在真正建 app 时 import)── + from agentscope.app import create_app # noqa: PLC0415 + from agentscope.app.message_bus import RedisMessageBus # noqa: PLC0415 + from agentscope.app.storage import RedisStorage # noqa: PLC0415 + from agentscope.app.workspace_manager import LocalWorkspaceManager # noqa: PLC0415 + + # Redis 连接参数从 infra.yaml 取(storage 与 message_bus 走不同逻辑库隔离键空间;均带 password)。 + storage_params = infra_config.redis_params(for_message_bus=False) + bus_params = infra_config.redis_params(for_message_bus=True) + + # 本地 workspace 根目录:放在 gen-worker 下的 _service-workspaces(与九工具的 _tier2-gen 目录分开; + # 这里只承载框架内置工具的工作区,本线生成产物不落这里——见文件头)。可被 env 覆盖换盘。 + ws_basedir = os.environ.get( + "TIER2_SERVICE_WORKSPACES", + str(_GEN_WORKER_DIR / "_service-workspaces"), + ) + + # 可追溯启动日志(口令脱敏):连的哪个 Redis、workspace 落哪。 + print( + f"[tier2-service] build_app: redis storage={storage_params['host']}:{storage_params['port']}" + f"/db{storage_params['db']} bus=db{bus_params['db']} " + f"auth={'on' if storage_params['password'] else 'off'} ws_basedir={ws_basedir}", + flush=True, + ) + + return create_app( + storage=RedisStorage(**storage_params), + message_bus=RedisMessageBus(**bus_params), + workspace_manager=LocalWorkspaceManager(basedir=ws_basedir), + # tier2 装配零件经扩展点注入(每个 chat 回合调一次)。 + extra_agent_tools=_tier2_tools_factory, + extra_agent_middlewares=_tier2_middlewares_factory, + # 工作室 Team worker 蓝图暂留空(单写主链先服务化;设计阶段服务态落位见文件头 followup)。 + custom_subagent_templates=[], + # 用内置 Agent(tier2 单写 agent 靠零件注入装配,不必子类化)。 + custom_agent_cls=None, + title=title, + ) + + +def main() -> None: + """部署入口:起 uvicorn 跑 tier2 Agent Service(只在 mini-desktop 装好 agentscope[full]+redis 后跑)。 + + 端口 / host 经 env 调(默认 0.0.0.0:8200,避开 wg1/其它服务常用口)。reload=False(生产形态; + 调试可 export TIER2_SERVICE_RELOAD=1)。注意 reload 模式需用 import 字符串而非 app 对象,故两分支。 + """ + import uvicorn # noqa: PLC0415 —— 惰性 import(6c6g 无 uvicorn) + + host = os.environ.get("TIER2_SERVICE_HOST", "0.0.0.0") + port = int(os.environ.get("TIER2_SERVICE_PORT", "8200")) + reload = os.environ.get("TIER2_SERVICE_RELOAD", "0").strip().lower() in ("1", "true", "yes", "on") + + print(f"[tier2-service] 启动 Agent Service:http://{host}:{port} reload={reload}", flush=True) + if reload: + # reload 模式必须传 import 字符串(uvicorn 要能在 worker 子进程重导入);指向本模块的 module-level `app`。 + # 子进程重导入 service.app 时,module-level `app` 只在 TIER2_SERVICE_EAGER_APP=1 时才建——故在此 + # 先显式置位,保证 reload worker 进程导入即建出真 app(否则会服务到 app=None)。非 reload 分支不需要。 + os.environ["TIER2_SERVICE_EAGER_APP"] = "1" + uvicorn.run("service.app:app", host=host, port=port, reload=True) + else: + # 生产形态:直接传 app 对象(此处才真正连 Redis 建 app)。 + uvicorn.run(build_app(), host=host, port=port) + + +# reload 模式 / `uvicorn service.app:app` 直起时用的 module-level app。 +# ⚠️ 仅当显式设置 TIER2_SERVICE_EAGER_APP=1 时才在 import 期建 app(会连 Redis)——默认不建, +# 保证 6c6g `import service.app` 不触发重依赖 import / 不连基建(惰性红线)。mini-desktop 用 +# `uvicorn service.app:app` 起时,设这个 env 让 module-level app 就绪;或直接 `python -m service.app` +# 走 main()(推荐,无需 eager)。 +app = None +if os.environ.get("TIER2_SERVICE_EAGER_APP", "0").strip().lower() in ("1", "true", "yes", "on"): + app = build_app() # pragma: no cover —— 仅 mini-desktop 显式开启时执行 + + +if __name__ == "__main__": + main() diff --git a/tier2/gen-worker/service/bootstrap.py b/tier2/gen-worker/service/bootstrap.py new file mode 100644 index 00000000..7c001b1a --- /dev/null +++ b/tier2/gen-worker/service/bootstrap.py @@ -0,0 +1,262 @@ +"""service/bootstrap.py —— tier2 Agent Service 的 session 三路初始化客户端(经 REST 编排,不旁路框架)。 + +【为什么是 REST 客户端而不是直接调内部函数】 + create_app 的 chat 是 fire-and-forget 的 REST 服务(_router/_chat.py):前端/控制面经 HTTP 触发一个 + chat run,事件从 SSE 流收。session 三路初始化属于「服务消费方」的编排——先建凭据/agent/session,再触发 + chat。本模块就是这层编排的参考实现 + 口径说明:用 httpx 打 create_app 暴露的现成端点,**不自造端点、 + 不旁路框架的装配点**。这样三路最终都汇到 _service/_chat.py 的同一 agent 装配 + (extra_agent_tools 注入九工具 + extra_agent_middlewares 注入熔断/trace + session 的模型/系统提示词)。 + +【三路对照(全部经现成 REST 端点)】 + ① 新建(模板起手): + POST /credential 注册 M3 凭据(若尚未注册;anthropic_credential,base_url=new-api host 根) + → POST /agent 建单写 agent(system_prompt=roles.writer_system,react_config 放开轮数, + context_config=tier2 历史压缩) + → POST /sessions 建新 session(新 workspace_id;引用上面的模型凭据) + → POST /chat 发 kick 文本(引导 agent scaffold_init 起手 → write_source→build→run_gates→finish) + ② 续接会话(从 SessionRecord.state resume): + 对已存在的 (agent_id, session_id) 直接 POST /chat;框架自动 reload 该 session 的 AgentState + (含已写历史 / cur_iter / tasks)接着跑。无需自写 resume 逻辑——这是 _chat.py:315 的原生行为。 + ③ 加载已有工程迭代: + 绑到「该工程 workdir 对应的 agent_id」(AgentScope workdir=basedir/agent_id),POST /sessions 新开 + (或复用)该 agent_id 的 session → POST /chat 发「迭代指令」。一款游戏工程 = 一个 agent_id; + 迭代 = 在该 agent_id 上继续。 + +【诚实差距(见 app.py 文件头 + 本文件 followups)】 + - 模型凭据:_chat.py 要求 session.chat_model_config 引用一个已注册的 Credential。本模块的 ensure_* + 把 M3 凭据(MiniMax-M3 经 new-api 走 Anthropic 原生)注册成 anthropic_credential。凭据参数(base_url / + key)从 docs/内网凭据与端点.md 取;此处用 worker.client 解析(与 CLI 线同口径),避免硬编码。 + - kick 文本沿用 studio.py 的同一段(引导 scaffold_init→...→finish),保证服务态与 CLI 态对 agent 的 + 指令一致。但「agent 过早停下再踹回去续修」的有界 resume 在服务态由消费方决定是否再 POST /chat + (前端/控制面轮询 verdict 决定续不续)——这是 followup:控制面需实现「读 SSE 末态 verdict→未 finish + 且门未绿且有预算→再触发一轮」的循环,等价 CLI 的 max_resumes。 + +【惰性 import】httpx / worker.client 在函数体内 import(6c6g 未必装 httpx;worker.client 牵出 agentscope)。 + 本模块顶层零重依赖,6c6g 可 py_compile + import。 +""" + +from __future__ import annotations + +from typing import Any + + +# 单写 agent 的 kick 文本(与 worker/agent_loop/studio.py 同口径,保证服务态与 CLI 态指令一致)。 +WRITER_KICK_TEXT = ( + "开始实现这款富游戏。先 scaffold_init 起手,然后在循环里 write_source→(validate_datatable)→" + "build→run_gates→read_verdict→针对失败门 write_source 修→再 run_gates……门全绿后调 finish 交付。" + "切记:run_gates 的 verdict 是机器判的,看到 decision=fix 不要停下来收尾,要按失败门继续修。" +) + + +def _resolve_m3_credential_payload(model_name: str = "MiniMax-M3") -> dict[str, Any]: + """组装注册 M3 凭据(POST /credential)的请求体——base_url/key 经 worker.client 解析,绝不硬编码。 + + M3 走 Anthropic 原生 /v1/messages(thinking 分离),故凭据 type = anthropic_credential,base_url = new-api + host 根(不带 /v1,SDK 自拼 /v1/messages;与 worker.config.build_model 同口径)。key = NEWAPI_KEY + (worker.client.get_api_key 从 env/.env 读;凭据档授权内网入仓)。 + + 返回的形状对齐 create_app 的 credential 注册端点期望的字段(type / 凭据明文)。具体字段名以 + /credential 端点 schema(_router/_schema/_credential.py)为准——本函数给出 anthropic 凭据的最小集, + 部署 smoke 时若字段名有出入,按该 schema 微调(列入 followup:对齐 credential 端点 schema)。 + """ + # 惰性 import:worker.client 解析代理旁路时会牵出环境(且 worker 包顶层会 import agentscope)。 + from worker import client # noqa: PLC0415 + + base_url = client.resolve_base_url() # new-api host 根(已装代理旁路) + api_key = client.get_api_key() # NEWAPI_KEY(凭据档 / env) + return { + # 字段名以 /credential 端点 schema 为准;此处给 anthropic 凭据的最小集(type + base_url + api_key)。 + "type": "anthropic_credential", + "base_url": base_url, + "api_key": api_key, + "model": model_name, + } + + +async def _post(client: Any, base: str, path: str, json_body: dict, *, user_header: str) -> dict: + """打一个 POST 到运行中的 Agent Service,返回 JSON;非 2xx 抛(带可追溯日志)。 + + user_header:Agent Service 经依赖注入解析 user_id(get_current_user_id);多租户下需带用户标识头。 + 具体头名以 app/deps.py 的 get_current_user_id 实现为准(可能是某个 header / 默认匿名)——本函数留 + 一个口子,部署 smoke 时按真实 deps 调整(列入 followup:对齐鉴权头)。 + """ + headers = {"X-User-Id": user_header} if user_header else {} + resp = await client.post(f"{base}{path}", json=json_body, headers=headers, timeout=30.0) + if resp.status_code >= 300: + # 可追溯日志:哪个端点、什么码、返回体前 300 字(排障用)。 + print(f"[tier2-bootstrap] POST {path} -> {resp.status_code}: {resp.text[:300]}", flush=True) + resp.raise_for_status() + return resp.json() + + +async def start_new_game( + base_url: str, + brief: str, + *, + user_id: str = "tier2", + model_name: str = "MiniMax-M3", + writer_max_iters: int = 40, + credential_id: str | None = None, +) -> dict: + """【路①】新建:注册凭据(若需)→ 建 agent → 建 session → 发 kick,启动一款新富游戏的生成。 + + 经 create_app 现成 REST 端点完成,最终汇到 _chat.py 的同一装配点(九工具 + 熔断/trace 经工厂注入)。 + + Args: + base_url: 运行中的 Agent Service 根地址(如 http://100.64.0.7:8200)。 + brief: 一句话题面。 + user_id: 多租户用户标识(经 X-User-Id 头传;默认 tier2)。 + model_name: M3 模型名(默认 MiniMax-M3;经 new-api 走 Anthropic 原生)。 + writer_max_iters: 单写 ReAct 放开的最大轮数(写进 AgentRecord.react_config)。 + credential_id: 已注册的模型凭据 id;None → 本函数先 POST /credential 注册一个 M3 凭据。 + + Returns: + {credential_id, agent_id, session_id, chat: }。后续可凭 session_id 订阅 + GET /sessions/{session_id}/stream 收事件,或凭 (agent_id, session_id) 走 resume_session 续跑。 + """ + # 惰性 import:httpx(6c6g 未必装)+ worker.config(组装 system_prompt/context_config,牵出 agentscope)。 + import httpx # noqa: PLC0415 + from worker import config, roles # noqa: PLC0415 + + async with httpx.AsyncClient() as http: + # ① 注册 M3 凭据(若未提供 credential_id)。 + if credential_id is None: + cred = await _post(http, base_url, "/credential/", _resolve_m3_credential_payload(model_name), + user_header=user_id) + credential_id = cred.get("credential_id") or cred.get("id") + + # ② 建单写 agent:system_prompt=roles.writer_system(题面 + 空 design_text;设计阶段服务态落位见 followup), + # react_config 放开轮数,context_config=tier2 富游戏历史压缩(防长程多文件丢约定)。 + # context_config / react_config 经 model_dump 转 JSON 传(端点 schema 收 ContextConfig/ReActConfig)。 + ctx_cfg = config.build_context_config() + agent_body = { + "name": "tier2-writer", + "system_prompt": roles.writer_system(brief, "", fixture_hint=""), + "context_config": ctx_cfg.model_dump(mode="json"), + "react_config": {"max_iters": writer_max_iters}, + } + agent = await _post(http, base_url, "/agent/", agent_body, user_header=user_id) + agent_id = agent.get("agent_id") or agent.get("id") + + # ③ 建新 session(新 workspace_id 由框架分配;引用模型凭据)。 + # chat_model_config 的字段对齐 ChatModelConfig(type/credential_id/model/parameters)。 + session_body = { + "agent_id": agent_id, + "chat_model_config": { + "type": "anthropic_credential", + "credential_id": credential_id, + "model": model_name, + # parameters:M3 thinking 分离 + max_tokens(口径同 worker.config.build_model; + # max_tokens 必须 > thinking_budget)。 + "parameters": {"max_tokens": 16000, "thinking_enable": True, "thinking_budget": 8000}, + }, + } + session = await _post(http, base_url, "/sessions/", session_body, user_header=user_id) + session_id = session.get("session_id") or session.get("id") + + # ④ 发 kick 触发首轮 chat run(fire-and-forget;事件从 SSE 收)。input 为 UserMsg 形状。 + chat_body = { + "agent_id": agent_id, + "session_id": session_id, + "input": {"name": "user", "role": "user", "content": WRITER_KICK_TEXT}, + } + chat = await _post(http, base_url, "/chat/", chat_body, user_header=user_id) + + return {"credential_id": credential_id, "agent_id": agent_id, + "session_id": session_id, "chat": chat} + + +async def resume_session( + base_url: str, + agent_id: str, + session_id: str, + *, + user_id: str = "tier2", + feedback: str | None = None, +) -> dict: + """【路②】续接会话:对已存在的 (agent_id, session_id) 再发一轮 chat,从 SessionRecord.state resume。 + + 框架自动 reload 该 session 的 AgentState(_chat.py:315)续跑——无需自写 resume。feedback 为 None 时 + 发一条「继续」指令;给了 feedback(如上一轮 verdict 失败门摘要)则把它作为本轮输入(等价 CLI 的 + 「带 verdict 反馈踹回去续修」,只是触发权在消费方/控制面)。 + + Args: + feedback: 续跑指令文本;None → 用默认「继续按失败门修,门绿再 finish」。 + """ + import httpx # noqa: PLC0415 + + text = feedback or ( + "继续。如果上一轮 run_gates 的门还没全绿,按失败门 write_source 针对性修,build→run_gates," + "直到门绿再 finish;若已 finish 则无需再改。" + ) + chat_body = { + "agent_id": agent_id, + "session_id": session_id, + "input": {"name": "user", "role": "user", "content": text}, + } + async with httpx.AsyncClient() as http: + chat = await _post(http, base_url, "/chat/", chat_body, user_header=user_id) + return {"agent_id": agent_id, "session_id": session_id, "chat": chat} + + +async def iterate_existing_project( + base_url: str, + agent_id: str, + iterate_instruction: str, + *, + user_id: str = "tier2", + model_name: str = "MiniMax-M3", + credential_id: str | None = None, +) -> dict: + """【路③】加载已有工程迭代:在「工程所在 workdir 对应的 agent_id」上新开 session 并发迭代指令。 + + AgentScope LocalWorkspaceManager 的 workdir = basedir/agent_id(按 agent_id,不按 workspace_id),故 + 一款游戏工程 = 一个 agent_id;迭代 = 在该 agent_id 上开新 session(或复用其既有 session)继续生成。 + 本函数针对「在该 agent_id 上开一个新 session 做新一轮迭代」——若想在同一 session 续接,用 resume_session。 + + 注:tier2 九工具的工程目录是 game-runtime/games/_tier2-gen/(按 session_id=game_id),与 + AgentScope workspace(按 agent_id)是两套目录。要让「迭代」真正读到上一版源文件,九工具侧需按 + 「同一工程 → 同一 game_id」复用目录——这要求迭代时复用同一 session_id(走 resume_session)而非新开 + session。**故路③在「真改上一版文件」语义下,优先用 resume_session(同 session_id=同 game_id 目录); + 本函数的「新 session 迭代」适用于「同一 agent 的工作区里另起一款变体」。** 这个 game_id↔工程目录的 + 绑定差异列入 followup(见文末),是服务态与 CLI 态最需要校准的一处。 + + Args: + agent_id: 既有工程对应的 agent_id。 + iterate_instruction: 迭代需求(如「把订单耐心调长、加一个连击系统」)。 + """ + import httpx # noqa: PLC0415 + + async with httpx.AsyncClient() as http: + if credential_id is None: + cred = await _post(http, base_url, "/credential/", _resolve_m3_credential_payload(model_name), + user_header=user_id) + credential_id = cred.get("credential_id") or cred.get("id") + + session_body = { + "agent_id": agent_id, + "chat_model_config": { + "type": "anthropic_credential", + "credential_id": credential_id, + "model": model_name, + "parameters": {"max_tokens": 16000, "thinking_enable": True, "thinking_budget": 8000}, + }, + } + session = await _post(http, base_url, "/sessions/", session_body, user_header=user_id) + session_id = session.get("session_id") or session.get("id") + + # 迭代指令:先 scaffold_init 起手(已存在文件不覆盖),再按 iterate_instruction 改/加,跑门到绿 finish。 + text = ( + "这是对一款已有富游戏工程的迭代。先 scaffold_init 起手(已存在的文件不会被覆盖)," + f"然后按以下迭代需求改/加源文件:\n{iterate_instruction}\n" + "改完 build→run_gates→read_verdict,门绿后 finish 交付新版本。" + ) + chat_body = { + "agent_id": agent_id, + "session_id": session_id, + "input": {"name": "user", "role": "user", "content": text}, + } + chat = await _post(http, base_url, "/chat/", chat_body, user_header=user_id) + + return {"credential_id": credential_id, "agent_id": agent_id, + "session_id": session_id, "chat": chat} diff --git a/tier2/gen-worker/service/infra_config.py b/tier2/gen-worker/service/infra_config.py new file mode 100644 index 00000000..aae01425 --- /dev/null +++ b/tier2/gen-worker/service/infra_config.py @@ -0,0 +1,251 @@ +"""service/infra_config.py —— tier2 基建数据面配置加载器(运行时读 tier2/config/infra.yaml)。 + +【这份解决什么问题】 + Agent Service(U1)要连 Redis(session + MessageBus),U2 落库要连 MySQL / MinIO。这些端点 + 口令 + 按创始人铁律集中在 tier2/config/infra.yaml(内网入仓),**运行时读**;本模块是它的读取层, + 让代码从配置取、绝不把端点/口令硬编码进 .py。口径刻意对齐 worker/genconfig.py(同款三级回落 + + best-effort + mtime 缓存),便于维护者一眼复用既有心智。 + +【取值三级回落(运行时读,失败绝不中断)】 + get(section, key, default) 每次调用按以下顺序取值,任一级取到即用: + ① env 单点覆盖:TIER2_INFRA__
__(全大写、双下划线分隔)——临时压一个值不必改文件; + ② 外部 YAML:tier2/config/infra.yaml(路径可被 env TIER2_INFRACONFIG 指向别处)的 section.key; + ③ 调用方传入的 default(连本文件都缺该 key 时的最后防线)。 + 与 genconfig 的区别:本模块不登记「内置默认副本」——基建端点/口令没有「字节不变的现行硬编码值」可对齐 + (它们本就该只存在于配置文件 / env,不该在代码里有副本),故缺失时直接落调用方 default(通常 None)。 + +【best-effort 铁律(对齐 genconfig.py / client.py 同款纪律)】 + 读文件 / 解析 YAML / env 转型的任何异常都不抛、不中断:读不到或脏 → 落 default + 一次性告警。 + YAML 只在「首次访问 + 文件 mtime 变化」时重读(改了文件下一次 get 自动生效)。 + +【依赖】pyyaml(声明在 requirements.txt);import 失败则整层降级为「只用 env + default」,绝不让缺解析库 + 阻断服务启动(读不到 = 用 default,与脏文件同处置)。本模块**不 import agentscope / redis**,故 6c6g + 可直接 import + py_compile(重依赖在 app.py 里惰性 import)。 +""" + +from __future__ import annotations + +import os +import threading +from pathlib import Path +from typing import Any + +# YAML 解析 best-effort 兜底:import 失败 → 整层降级(只用 env + default),不抛。 +try: + import yaml as _yaml # type: ignore +except Exception: # pragma: no cover —— 极端环境缺 pyyaml:降级为 env + default(不抛) + _yaml = None + + +# ── 配置文件默认路径(可被 env TIER2_INFRACONFIG 指向别处)────────────────────────── +# 本文件在 tier2/gen-worker/service/infra_config.py → 上溯 2 级到 gen-worker/,再上 1 级到 tier2/, +# 配置在 tier2/config/infra.yaml(与 genconfig 推 tier2/ 根同口径,只是层级少一层 worker/)。 +_TIER2_DIR = Path(__file__).resolve().parents[2] +_DEFAULT_CONFIG_PATH = _TIER2_DIR / "config" / "infra.yaml" + +# env 覆盖前缀:TIER2_INFRA__
__(双下划线分隔 section / key;全大写)。 +_ENV_PREFIX = "TIER2_INFRA__" + + +# ── 运行时读状态(缓存 + mtime 守门;线程安全;对齐 genconfig)───────────────────────── +_lock = threading.Lock() +_cache: dict[str, Any] | None = None +_cache_mtime: float | None = None +_cache_path: str | None = None +_warned: set[str] = set() + + +def _warn_once(tag: str, msg: str) -> None: + """同一 tag 只告警一次(防 best-effort 路径刷屏);可追溯日志,带 [tier2-infra] 前缀。""" + if tag in _warned: + return + _warned.add(tag) + print(f"[tier2-infra] {msg}", flush=True) + + +def config_path() -> Path: + """当前生效的配置文件路径:env TIER2_INFRACONFIG 指定 > 默认 tier2/config/infra.yaml。""" + p = os.environ.get("TIER2_INFRACONFIG") + return Path(p) if p else _DEFAULT_CONFIG_PATH + + +def _load_yaml() -> dict[str, Any]: + """运行时读 infra.yaml(带 mtime 缓存 + best-effort)。返回解析后的 dict(读不到/脏 → 空 dict)。 + + 任何异常(文件缺失 / YAML 语法错 / 顶层非 dict / 缺 pyyaml)都不抛——返回空 dict(由上层落 default)。 + """ + global _cache, _cache_mtime, _cache_path + + path = config_path() + path_str = str(path) + + if _yaml is None: + _warn_once("no-yaml", "pyyaml 不可用 → 配置层降级为「env 覆盖 + 调用方 default」。") + return {} + + with _lock: + try: + mtime = path.stat().st_mtime if path.exists() else None + except OSError: + mtime = None + + # 命中缓存:路径与 mtime 都没变 → 直接返回上次解析结果(运行时读但不重复开文件)。 + if _cache is not None and _cache_path == path_str and _cache_mtime == mtime: + return _cache + + # 文件不存在:缓存空 dict(按 default 跑);一次性告警(这是严重的——基建端点全得靠 default/env)。 + if mtime is None: + _warn_once( + f"missing:{path_str}", + f"基建配置文件不存在({path_str})→ 所有端点/口令只能靠 env 覆盖或调用方 default;" + "服务大概率连不上基建,请确认 tier2/config/infra.yaml 已就位。", + ) + _cache, _cache_mtime, _cache_path = {}, None, path_str + return _cache + + # 读 + 解析(best-effort:任何异常都落空 dict + 告警,绝不抛)。 + try: + raw = path.read_text(encoding="utf-8") + data = _yaml.safe_load(raw) + if data is None: + data = {} + if not isinstance(data, dict): + _warn_once( + f"nonmap:{path_str}", + f"基建配置顶层不是 YAML 映射({path_str})→ 视为脏,降级为 env + default。", + ) + data = {} + except Exception as exc: # noqa: BLE001 —— 脏 YAML / IO 异常都降级,绝不中断服务启动 + _warn_once( + f"parse:{path_str}", + f"基建配置解析失败({path_str}: {type(exc).__name__}: {exc})→ 降级为 env + default。", + ) + data = {} + + _cache, _cache_mtime, _cache_path = data, mtime, path_str + return _cache + + +def _coerce(value: Any, like: Any) -> Any: + """把 env / YAML 原始值按「调用方 default(like)的类型」做最小转型(best-effort,失败回 None)。 + + env 取出来全是字符串;按 default 的类型对齐,避免把 "6379" 当字符串塞进期望 int 的端口。 + bool 特殊处理(避免 bool('false') == True 的经典坑)。转型失败 → None,由上层视为「没取到」回落下一级。 + """ + if value is None: + return None + try: + if isinstance(like, bool): + if isinstance(value, bool): + return value + s = str(value).strip().lower() + if s in ("1", "true", "yes", "on"): + return True + if s in ("0", "false", "no", "off"): + return False + return None + if isinstance(like, int) and not isinstance(like, bool): + return int(value) + if isinstance(like, float): + return float(value) + if isinstance(like, str): + return str(value) + except (TypeError, ValueError): + return None + # like 是 None / 其它类型:原样返回(不强转)。 + return value + + +def get(section: str, key: str, default: Any = None) -> Any: + """取一个基建配置项:env 覆盖 > 外部 YAML > 调用方 default。绝不抛(best-effort)。 + + Args: + section: 分区名(redis / mysql / minio)。 + key: 分区内的项名(host / port / password / ...)。 + default: 调用方传入的兜底值(连本文件都缺该 key 时的最后防线;转型基准也用它的类型)。 + + Returns: + 配置值(类型按 default 对齐转型;全级都取不到 → 返回 default)。 + """ + like = default # 本模块无内置默认副本,转型基准 = 调用方 default 的类型。 + + # ① env 单点覆盖(最高优先级)。 + env_name = f"{_ENV_PREFIX}{section.upper()}__{key.upper()}" + env_raw = os.environ.get(env_name) + if env_raw is not None: + coerced = _coerce(env_raw, like) + if coerced is not None: + return coerced + _warn_once( + f"env-bad:{env_name}", + f"env {env_name}={env_raw!r} 无法转成 {type(like).__name__} → 忽略,回落 YAML/default。", + ) + + # ② 外部 YAML 的 section.key。 + data = _load_yaml() + sec = data.get(section) if isinstance(data, dict) else None + if isinstance(sec, dict) and key in sec: + coerced = _coerce(sec.get(key), like) + if coerced is not None: + return coerced + _warn_once( + f"yaml-bad:{section}.{key}", + f"配置 {section}.{key}={sec.get(key)!r} 无法转成 {type(like).__name__} → 忽略,回落 default。", + ) + + # ③ 调用方 default。 + return default + + +def redis_params(*, for_message_bus: bool = False) -> dict[str, Any]: + """汇出一组可直接展开给 AgentScope RedisStorage / RedisMessageBus 构造器的连接参数。 + + 源码核验的构造器形参(均支持 password): + RedisStorage(host, port, db, password, ...) storage/_redis_storage.py:77 + RedisMessageBus(host, port, db, password, ...) message_bus/_redis_message_bus.py:44 + + Args: + for_message_bus: False → storage 用(db=storage_db,默认 0);True → message_bus 用 + (db=message_bus_db,默认 1)。storage 与 message_bus 故意走不同逻辑库隔离键空间。 + + Returns: + {host, port, db, password};调用方 `RedisStorage(**redis_params())` 即可。 + """ + db_key = "message_bus_db" if for_message_bus else "storage_db" + db_default = 1 if for_message_bus else 0 + return { + "host": get("redis", "host", "127.0.0.1"), + "port": get("redis", "port", 6379), + "db": get("redis", db_key, db_default), + "password": get("redis", "password", None), + } + + +def reload() -> None: + """强制丢弃缓存,下次 get 重新读文件(测试 / 手动热更;运行时 mtime 已自动重读,一般不必显式调)。""" + global _cache, _cache_mtime, _cache_path + with _lock: + _cache, _cache_mtime, _cache_path = None, None, None + + +# ── __main__ 自测块:不依赖 agentscope / redis / 网络,只校验三级回落与转型(6c6g 可直接跑)── +if __name__ == "__main__": # pragma: no cover —— 本地自测 + # ① YAML 流通:读真实 infra.yaml 的 redis.host / redis.port(若文件就位)。 + reload() + print("[infra-selftest] redis.host =", get("redis", "host", "MISSING")) + print("[infra-selftest] redis.port =", get("redis", "port", -1), type(get("redis", "port", -1))) + print("[infra-selftest] redis_params(storage) =", + {**redis_params(), "password": "***" if redis_params()["password"] else None}) + print("[infra-selftest] redis_params(bus) =", + {**redis_params(for_message_bus=True), + "password": "***" if redis_params(for_message_bus=True)["password"] else None}) + + # ② env 覆盖 + 转型(字符串 → int)。 + os.environ["TIER2_INFRA__REDIS__PORT"] = "6380" + assert get("redis", "port", 6379) == 6380, "env 覆盖 int 失败" + os.environ.pop("TIER2_INFRA__REDIS__PORT") + + # ③ 缺失 key 回落 default。 + assert get("redis", "不存在的key", "fallback") == "fallback", "缺失 key 回落 default 失败" + + print("[infra-selftest] 三级回落 + 转型自测通过。") diff --git a/tier2/gen-worker/worker/store.py b/tier2/gen-worker/worker/store.py index 49dc44b8..c21ffb15 100644 --- a/tier2/gen-worker/worker/store.py +++ b/tier2/gen-worker/worker/store.py @@ -35,8 +35,8 @@ from __future__ import annotations import hashlib import json -import shutil -import time +import os +import sys from abc import ABC, abstractmethod from pathlib import Path from typing import Any, Optional @@ -83,6 +83,22 @@ def derive_version_id(source_hash: str, now_ts: float) -> str: return f"v{sec}-{short}" +def _is_safe_rel_path(rel: str) -> bool: + """落库源文件相对路径安全校验(禁 '..' / 绝对路径越界)。 + + 与 LocalFsStore.save / run.scaffold / A3 path pattern 同口径:工程内相对路径(POSIX 正斜杠), + 不得含 '..' 段、不得以 '/' 开头。BackendStore 拼 MinIO object key 时复用此校验, + 避免恶意/异常路径把对象写到 bucket 外的非预期前缀。 + """ + if not rel: + return False + if rel.startswith("/"): + return False + if ".." in rel.split("/"): + return False + return True + + class SourceProjectStore(ABC): """源项目落库取回接口(F1 要素⑦ addressing 的代码契约 · 落库后端无关)。 @@ -303,55 +319,336 @@ class LocalFsStore(SourceProjectStore): return sp +# ── BackendStore 落库后端常量(MySQL 表名 / MinIO key 前缀;集中此处便于对账)───────────── +# 表名与 key 前缀是 doc↔code 契约点:apiNotes / 凭据档 / Java 后端取回都按这两个常量对齐。 +_MYSQL_TABLE = "tier2_source_project_version" +# MinIO object key 形态 = //<工程内相对路径>(bucket 由 infra.yaml minio.bucket 决定, +# 默认 tier2-src)。注:bucket 名不进 key(MinIO/S3 的 object key 不含 bucket 前缀),与凭据档 +# "key=tier2-src//..." 的写法差异 = 那处把 bucket 也写进了人读路径示意,实际 SDK put/get 只传 +# bucket + 不含 bucket 的 object key,故此处 key 从 / 起。 + + +# ── MySQL 建表 DDL(代码内执行 CREATE TABLE IF NOT EXISTS;另有一份同源 .sql 见 config/schema/)── +# 幂等键 = (game_id, source_hash):同款同源(改了源 source_hash 才变)只落一行,命中不重写。 +# version_id 另加唯一索引(同一 id 下版本号唯一,且取回按 (game_id, version_id) 寻址)。 +_DDL_SOURCE_PROJECT_VERSION = f""" +CREATE TABLE IF NOT EXISTS {_MYSQL_TABLE} ( + id BIGINT NOT NULL AUTO_INCREMENT COMMENT '自增主键', + game_id VARCHAR(128) NOT NULL COMMENT '游戏稳定标识(同款多版本共享;落库 id)', + version_id VARCHAR(96) NOT NULL COMMENT '版本号 vXXXX-hash12(改源重建生成新版本)', + source_hash CHAR(64) NOT NULL COMMENT '源工程内容指纹 sha256(= source_project.contentHash;幂等键)', + manifest_json LONGTEXT NOT NULL COMMENT '源工程 manifest 七要素形状(去 fileTree content;含 addressing)', + addressing VARCHAR(512) NULL COMMENT '落库寻址摘要(store/sourceUrl;冗余出 manifest 便于检索)', + created_at DATETIME NOT NULL COMMENT '落库时刻(now_ts 转 UTC datetime)', + PRIMARY KEY (id), + UNIQUE KEY uk_game_source (game_id, source_hash) COMMENT '幂等键:同款同源不重复落', + UNIQUE KEY uk_game_version (game_id, version_id) COMMENT '版本寻址键:同款版本号唯一', + KEY idx_game_created (game_id, created_at) COMMENT 'fetch 取最新版本走它' +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='tier2 富游戏源工程版本表(U2 第二装载落库)'; +""" + + class BackendStore(SourceProjectStore): - """后端落库占位(seam)—— MySQL manifest + OSS 源文件,由 game-cloud Java 后端实现。 + """后端落库真实现 —— manifest 落 MySQL、源文件全文落 MinIO(S3 兼容 OSS)。 ════════════════════════════════════════════════════════════════════════ - 【这是 seam,不是本分支要实连的东西】 - tier2 是独立 Python 生成 service,落库后端(MySQL + 对象存储 OSS)是 game-cloud Java 后端的职责, - forbidden-import CI(tier2/ci/check-forbidden-import.sh)守着 tier2 不实连后端 / 不 import 廉价线运行态。 - 本类只声明【后端要兑现的接口形状 + 填充点】,让接线契约清晰、不留孤儿设计;真实现落在 Java 后端。 + 【这个类在 F 族里兑现什么】 + F 族要素⑦ addressing:源工程怎么存进 MySQL+对象存储、怎么按 (game_id, versionId) 取回来重建。 + P3 时这是个占位 seam(raise NotImplementedError);U2 把它落成真实现,供部署时经 default_store() + 的 TIER2_STORE=backend 切入。LocalFsStore 不动(spike 期本地往返仍用它)。 - 【Java 后端填充点(对接 contracts/agent-loop/tier2-source-project.schema.json 的 addressing)】 - save(source_project, now_ts): - 1) manifest(七要素去掉源文件全文)落 MySQL 一张源工程版本表(game_id, version_id, source_hash, - build_profile, dep_lock, content_hash, addressing, created_at);幂等键 = (game_id, source_hash)。 - 2) 源文件全文(file_list)落对象存储 OSS,key 形如 tier2-src///<相对路径>; - manifest 记 sourceUrl(addressing.sourceUrl,OSS 只读镜像,同 GamePackage.packageUrl 范式)。 - 3) versionId 由后端按 derive_version_id 同口径(sourceHash + 时间序)或 DB 序列生成, - 改源重建 → 新 versionId,旧版本保留(F1 版本寻址 / 长生命周期项目)。 - 4) 返回 {id, versionId, sourceHash}。 - fetch(id, versionId): - 1) versionId=None → 查该 game_id 最新版本;给定 → 查该版本。 - 2) 读 MySQL manifest + 据 sourceUrl 从 OSS 拉源文件全文,回填进 fileTree[].content,返回可重建的 source_project。 + 【落库后端约定(对接 contracts/agent-loop/tier2-source-project.schema.json 的 addressing)】 + save(source_project, now_ts, file_list, id=game_id): + 1) manifest(source_project 七要素,fileTree 只含 path/role 不含 content)落 MySQL 一行 + tier2_source_project_version(game_id, version_id, source_hash, manifest_json, addressing, + created_at);幂等键 (game_id, source_hash) 命中 → 不重写,返回已有版本。 + 2) 源文件全文(file_list 每个 {path, content})落 MinIO bucket(默认 tier2-src), + object key = //<工程内相对路径>。 + 3) manifest.addressing.sourceUrl 记 MinIO 版本前缀 URL(只读镜像,同 GamePackage.packageUrl 范式)。 + 4) versionId 按 derive_version_id 同口径(sourceHash + 时间序)生成,改源重建 → 新 versionId, + 旧版本保留(F1 版本寻址 / 长生命周期项目)。返回 {id, versionId, sourceHash}。 + fetch(id, versionId=None): + 1) versionId=None → 查该 game_id 最新版本(created_at 最大);给定 → 查该版本。 + 2) 读 MySQL manifest + 据 object key 从 MinIO 拉源文件全文,回填进 fileTree[].content, + 返回可重建的 source_project(形状与 save 入参一致)。 - 【为什么留这个占位类而非只写注释】 - AGENTS §6 条款 7「不留孤儿设计」:落库接口要连同其【后端入口】一起交付,接缝才闭得上。 - 本类是那个明确的后端入口锚点——主链需要时可注入 BackendStore(经依赖注入切换实现), - 接线时一眼看清后端该兑现什么;现在调它会 raise NotImplementedError(诚实未实现,不静默假成功)。 + 【惰性 import + best-effort 纪律(本分支硬约束)】 + pymysql / minio 在【方法体内】惰性 import(模块顶层不 import),保证 6c6g(未装这俩)能 + py_compile + import store.py;真跑 / 部署在 mini-desktop。连接 / 读写失败【不静默吞】: + 响亮记日志(审计点)+ 抛给调用方。上游 run.persist_source_project 已 try/except 兜住 save 异常 + (落库失败 → 告警,产物仍在 GEN_DIR workdir,不中断生成主链),与 LocalFsStore 同口径 best-effort。 ════════════════════════════════════════════════════════════════════════ """ + def __init__(self) -> None: + """构造不连基建(惰性):连接在每次 save/fetch 时按 infra.yaml 现取现连,免持有长连接跨会话失效。 + + 缺 minio.bucket 配置时回落默认 'tier2-src'(与 infra.yaml / 凭据档默认一致)。 + """ + self._table = _MYSQL_TABLE + + # ── 配置读取(经 service.infra_config 三级回落:env > infra.yaml > default;best-effort)── + def _infra(self): + """惰性取 infra_config 模块(运行时读 tier2/config/infra.yaml)。 + + store.py 在 worker 包内,infra_config 在 service 包内;两者都是挂在 gen-worker/(已加进 sys.path) + 下的顶层包(目录名 gen-worker 含连字符不可直接 import,故走 sys.path + 顶层包名,与 app.py 同款)。 + 若调用方入口(如 CLI run_studio)未把 gen-worker/ 加进 sys.path,这里兜底补一次再 import, + 让 BackendStore 不依赖调用方的 sys.path 布置也能取到配置。 + """ + try: + from service import infra_config # noqa: PLC0415 —— 惰性 import(顶层包,需 gen-worker 在 sys.path) + return infra_config + except Exception: # noqa: BLE001 —— 大概率 gen-worker 不在 sys.path;补一次再 import + gen_worker_dir = Path(__file__).resolve().parents[2] # worker/store.py → gen-worker/ + if str(gen_worker_dir) not in sys.path: + sys.path.insert(0, str(gen_worker_dir)) + from service import infra_config # noqa: PLC0415 + return infra_config + + def _mysql_conf(self) -> dict: + """从 infra.yaml 取 MySQL 连接参数(host/port/user/password/database)。""" + ic = self._infra() + return { + "host": ic.get("mysql", "host", "127.0.0.1"), + "port": ic.get("mysql", "port", 3306), + "user": ic.get("mysql", "user", "root"), + "password": ic.get("mysql", "password", ""), + "database": ic.get("mysql", "database", "tier2"), + } + + def _minio_conf(self) -> dict: + """从 infra.yaml 取 MinIO 连接参数(endpoint/access_key/secret_key/bucket/secure)。""" + ic = self._infra() + return { + "endpoint": ic.get("minio", "endpoint", "127.0.0.1:9000"), + "access_key": ic.get("minio", "access_key", ""), + "secret_key": ic.get("minio", "secret_key", ""), + "bucket": ic.get("minio", "bucket", "tier2-src"), + "secure": ic.get("minio", "secure", False), + } + + def _connect_mysql(self): + """惰性连 MySQL(pymysql)。失败响亮抛(连接是落库审计点,不静默)。""" + try: + import pymysql # noqa: PLC0415 —— 惰性 import(6c6g 无 pymysql 也能 py_compile / import 本模块) + except Exception as e: # noqa: BLE001 + print(f"[tier2-store] ❌ BackendStore 缺 pymysql(惰性 import 失败): " + f"{type(e).__name__}: {e};请在运行机 pip install pymysql。", flush=True) + raise + conf = self._mysql_conf() + # autocommit=True:落库是单行 upsert + select,不需要显式事务边界;cursorclass 用 DictCursor 便于按列名取。 + return pymysql.connect( + host=conf["host"], port=int(conf["port"]), user=conf["user"], + password=conf["password"], database=conf["database"], + charset="utf8mb4", autocommit=True, + cursorclass=pymysql.cursors.DictCursor, + ) + + def _connect_minio(self): + """惰性连 MinIO,返回 (client, bucket)。失败响亮抛。best-effort 建 bucket(不存在则建)。""" + try: + from minio import Minio # noqa: PLC0415 —— 惰性 import(6c6g 无 minio 也能 py_compile / import) + except Exception as e: # noqa: BLE001 + print(f"[tier2-store] ❌ BackendStore 缺 minio(惰性 import 失败): " + f"{type(e).__name__}: {e};请在运行机 pip install minio。", flush=True) + raise + conf = self._minio_conf() + client = Minio( + conf["endpoint"], access_key=conf["access_key"], secret_key=conf["secret_key"], + secure=bool(conf["secure"]), + ) + bucket = conf["bucket"] + # best-effort 建桶:不存在则建(建桶失败响亮抛——后续 put 必然失败,早暴露好过晚静默)。 + try: + if not client.bucket_exists(bucket): + client.make_bucket(bucket) + print(f"[tier2-store] BackendStore 建 MinIO bucket: {bucket}", flush=True) + except Exception as e: # noqa: BLE001 —— 建桶/探测失败响亮抛(连不上对象存储,落库无意义) + print(f"[tier2-store] ❌ BackendStore MinIO bucket 探测/创建失败: bucket={bucket} " + f"{type(e).__name__}: {e}", flush=True) + raise + return client, bucket + + def _source_url(self, endpoint: str, secure: bool, bucket: str, prefix: str) -> str: + """拼 manifest.addressing.sourceUrl(MinIO 版本前缀只读镜像 URL;同 GamePackage.packageUrl 范式)。 + + 形态:http(s)://////(指向该版本所有源文件的前缀)。 + 仅作 manifest 记账用(人/后端按它定位 OSS 前缀),fetch 不靠它寻址(fetch 按 game_id+version_id 拼 key)。 + """ + scheme = "https" if secure else "http" + return f"{scheme}://{endpoint}/{bucket}/{prefix}" + def save(self, source_project: dict, *, now_ts: float, file_list: Optional[list[dict]] = None, id: Optional[str] = None) -> dict: - # Java 后端填(MySQL manifest + OSS 源文件);tier2 内不实连后端(forbidden-import 守着)。 - # 入参与抽象契约一致:source_project(manifest)+ file_list(落 OSS 的源文件全文)+ id(game_id)+ now_ts。 - raise NotImplementedError( - "BackendStore.save 待 Java 后端实现(MySQL manifest + OSS 源文件落库,见类注释填充点 1-4);" - "tier2 spike 期用 LocalFsStore。") + """落库一份源工程(后端实现:manifest 落 MySQL、源文件全文落 MinIO,幂等 + 版本寻址)。 + + Args / Returns 见 SourceProjectStore.save 抽象契约。连接/读写失败响亮记日志 + 抛(不静默); + 上游 run.persist_source_project 已 try/except 兜住(落库失败 → 告警,不中断主链)。 + """ + import datetime # noqa: PLC0415 —— 仅落库时把 now_ts 转 UTC datetime,延迟 import + + source_hash = (source_project or {}).get("contentHash") or _canonical_content_hash(file_list or []) + # id 缺省回落:无外部 game_id 时用 contentHash 前缀(可落库,但同款改源会换 id;接线方应显式传 game_id)。 + gid = id or f"sp-{source_hash[:16]}" + + conn = self._connect_mysql() + try: + # ── 建表(幂等;CREATE TABLE IF NOT EXISTS,首次落库自建)── + with conn.cursor() as cur: + cur.execute(_DDL_SOURCE_PROJECT_VERSION) + + # ── 幂等:(game_id, source_hash) 命中 → 不重写,返回已有版本 ── + with conn.cursor() as cur: + cur.execute( + f"SELECT version_id FROM {self._table} WHERE game_id=%s AND source_hash=%s LIMIT 1", + (gid, source_hash), + ) + row = cur.fetchone() + if row: + existing_vid = row.get("version_id") + print(f"[tier2-store] BackendStore.save 幂等命中(源未变,不重复落库): id={gid} " + f"versionId={existing_vid} sourceHash={source_hash[:12]}", flush=True) + return {"id": gid, "versionId": existing_vid, "sourceHash": source_hash} + + # ── 新版本:派生 versionId,先落源文件全文到 MinIO,再写 MySQL manifest 行 ── + version_id = derive_version_id(source_hash, now_ts) + prefix = f"{gid}/{version_id}/" # MinIO object key 前缀(bucket 内) + + client, bucket = self._connect_minio() + mc = self._minio_conf() + written = self._put_files(client, bucket, prefix, file_list or []) + + # manifest 落 MySQL:source_project 的形状原样存(fileTree 只含 path/role,不含 content—— + # content 在 MinIO);并把 addressing.sourceUrl 指向 MinIO 版本前缀。 + sp_for_db = dict(source_project or {}) + addressing = dict(sp_for_db.get("addressing") or {}) + addressing.setdefault("store", "mysql+oss") + addressing["sourceUrl"] = self._source_url(mc["endpoint"], bool(mc["secure"]), bucket, prefix) + sp_for_db["addressing"] = addressing + manifest_json = json.dumps(sp_for_db, ensure_ascii=False) + created_at = datetime.datetime.utcfromtimestamp(int(now_ts)).strftime("%Y-%m-%d %H:%M:%S") + + with conn.cursor() as cur: + cur.execute( + f"INSERT INTO {self._table} " + f"(game_id, version_id, source_hash, manifest_json, addressing, created_at) " + f"VALUES (%s, %s, %s, %s, %s, %s)", + (gid, version_id, source_hash, manifest_json, addressing["sourceUrl"], created_at), + ) + print(f"[tier2-store] BackendStore.save 落库成功: id={gid} versionId={version_id} " + f"sourceHash={source_hash[:12]} files={written} bucket={bucket} prefix={prefix}", flush=True) + return {"id": gid, "versionId": version_id, "sourceHash": source_hash} + except Exception as e: # noqa: BLE001 —— 落库失败响亮记日志(审计点),抛给调用方(不静默成功) + print(f"[tier2-store] ❌ BackendStore.save 落库失败: id={gid} " + f"sourceHash={source_hash[:12]} {type(e).__name__}: {e}", flush=True) + raise + finally: + try: + conn.close() + except Exception: # noqa: BLE001 —— 关连接失败不致命(连接已用完) + pass + + def _put_files(self, client, bucket: str, prefix: str, file_list: list[dict]) -> int: + """把源文件全文逐个 put 进 MinIO(object key = prefix + 工程内相对路径)。返回成功写入数。 + + 路径安全复用 _is_safe_rel_path(禁 '..' / 绝对路径,免把对象写到非预期前缀); + content 非 str 的项跳过(与 LocalFsStore 同口径)。单文件 put 失败响亮抛(落库要完整,不容残缺版本)。 + """ + import io # noqa: PLC0415 —— 仅 put 时把 content 包成字节流,延迟 import + + written = 0 + for f in (file_list or []): + rel = (f or {}).get("path") or "" + content = (f or {}).get("content") + if not _is_safe_rel_path(rel) or not isinstance(content, str): + continue + data = content.encode("utf-8") + object_name = f"{prefix}{rel}" + # put_object(bucket, object_name, data_stream, length):MinIO SDK 需流 + 显式长度。 + client.put_object(bucket, object_name, io.BytesIO(data), length=len(data), + content_type="text/plain; charset=utf-8") + written += 1 + return written def fetch(self, id: str, version_id: Optional[str] = None) -> Optional[dict]: - # Java 后端填(读 MySQL manifest + OSS 源文件回填 content);tier2 内不实连后端。 - raise NotImplementedError( - "BackendStore.fetch 待 Java 后端实现(读 MySQL manifest + OSS 源文件回填 fileTree content,见类注释);" - "tier2 spike 期用 LocalFsStore。") + """按 id(+ 可选 versionId)取回源工程重建(后端实现:读 MySQL manifest + MinIO 源文件回填)。 + + version_id=None → 取该 game_id 最新版本(created_at 最大);给定 → 取该历史版本(F1 版本寻址)。 + 取回时据 (game_id, version_id) 拼 MinIO key 拉源文件全文,回填进 fileTree[].content, + 使返回值可直接重新构建(与 LocalFsStore.fetch 同口径,含 missingContent 标注)。 + + Returns: + source_project(含各文件 content)或 None(无此 id / 无此版本)。读写失败响亮抛。 + """ + conn = self._connect_mysql() + try: + with conn.cursor() as cur: + if version_id is None: + cur.execute( + f"SELECT version_id, manifest_json FROM {self._table} " + f"WHERE game_id=%s ORDER BY created_at DESC, id DESC LIMIT 1", + (id,), + ) + else: + cur.execute( + f"SELECT version_id, manifest_json FROM {self._table} " + f"WHERE game_id=%s AND version_id=%s LIMIT 1", + (id, version_id), + ) + row = cur.fetchone() + if not row: + print(f"[tier2-store] BackendStore.fetch 未命中: id={id} " + f"versionId={version_id or '(latest)'}", flush=True) + return None + vid = row.get("version_id") + try: + sp = json.loads(row.get("manifest_json") or "{}") + except Exception as e: # noqa: BLE001 —— manifest 损坏 → 抛(取回不能返回半残数据冒充成功) + print(f"[tier2-store] ❌ BackendStore.fetch manifest 解析失败: id={id} versionId={vid} " + f"{type(e).__name__}: {e}", flush=True) + raise + finally: + try: + conn.close() + except Exception: # noqa: BLE001 + pass + + # ── 据 (game_id, version_id) 拼 MinIO key 前缀,把源文件全文回填进 fileTree[].content ── + client, bucket = self._connect_minio() + prefix = f"{id}/{vid}/" + for item in (sp.get("fileTree") or []): + rel = (item or {}).get("path") or "" + if not _is_safe_rel_path(rel): + item["missingContent"] = "路径不安全(含 '..' / 绝对路径),跳过取回" + continue + object_name = f"{prefix}{rel}" + try: + resp = client.get_object(bucket, object_name) + try: + item["content"] = resp.read().decode("utf-8") + finally: + # MinIO get_object 返回的 HTTPResponse 必须 close+release_conn(否则连接泄漏)。 + resp.close() + resp.release_conn() + except Exception as e: # noqa: BLE001 —— 单文件取回失败不致命:标注缺内容,继续(尽量重建) + item["missingContent"] = f"{type(e).__name__}: {e}" + print(f"[tier2-store] BackendStore.fetch 命中: id={id} versionId={vid} " + f"fileTree={len(sp.get('fileTree') or [])} bucket={bucket}", flush=True) + return sp -# ── 默认 store 工厂(收口处取实现的单一入口;spike 期恒回 LocalFsStore)── +# ── 默认 store 工厂(收口处取实现的单一入口;按 env TIER2_STORE 切换实现)── def default_store() -> SourceProjectStore: - """返回当前阶段的默认落库实现。 + """返回当前阶段的默认落库实现(按 env TIER2_STORE 切换)。 - spike 期 = LocalFsStore(本地真落,6c6g 可往返自检);产线接后端时,这里据配置/注入切到 BackendStore - (一处切换,run 收口与测试都经本工厂取实现,不在主链硬编码具体类)。 + - TIER2_STORE=backend → BackendStore(manifest 落 MySQL、源文件落 MinIO;部署时切)。 + - 其余(含未设)→ LocalFsStore(本地真落 _store/,6c6g 可往返自检;spike 期默认)。 + + 一处切换:run 收口(persist_source_project)与测试都经本工厂取实现,不在主链硬编码具体类。 + 口径对齐 service/infra_config(env 单点覆盖)与 AGENTS §6 条款 9(内网决策进配置/env,不靠散落硬编码)。 """ + if (os.environ.get("TIER2_STORE") or "").strip().lower() == "backend": + print("[tier2-store] default_store → BackendStore(TIER2_STORE=backend;MySQL+MinIO 落库)", flush=True) + return BackendStore() return LocalFsStore()