feat(tier2): P4 服务化核心——Agent Service(create_app)+ 第二装载真落库(MySQL+MinIO)
【U1 Agent Service】service/{app,infra_config,bootstrap}.py + config/infra.yaml:
经 agentscope 2.0.2 官方 create_app 把 tier2 引擎接成 Agent Service 壳——复用 studio 的 agent 装配零件
(roles/九工具/四熔断+trace/BYPASS/历史压缩)接到 create_app 扩展点(extra_agent_tools/extra_agent_middlewares),
不重写生成逻辑。Redis(mini-infra)做 storage(db0)+MessageBus(db1)。session 三路(新建/续接/加载工程迭代)经 REST(bootstrap.py)。
诚实差距:有界resume/工作室多agent/L2-L3/落库=循环外编排,服务态由控制面承接(已标 followup)。
【U2 真落库】store.py BackendStore 真实现(替 P3 NotImplementedError)+ schema SQL:
manifest→MySQL(幂等键 game_id+source_hash)+源文件全文→MinIO(bucket tier2-src);fetch 回填重建。
pymysql+minio 惰性 import;TIER2_STORE=backend 切换。
全 tier2/ 内、零碰 Tier0/1、py_compile+惰性import实证+forbidden-import 通过。真部署(连基建+uvicorn+smoke)待 mini-desktop。
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
78ff9f5d86
commit
ce491314c8
58
tier2/config/infra.yaml
Normal file
58
tier2/config/infra.yaml
Normal file
@ -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__<SECTION>__<KEY>(全大写、双下划线分隔)——临时压一个值不必改文件
|
||||
# (例:换 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/<game_id>/<version_id>/<工程内相对路径>(见凭据档)。与 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)
|
||||
33
tier2/config/schema/tier2_source_project_version.sql
Normal file
33
tier2/config/schema/tier2_source_project_version.sql
Normal file
@ -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 = <game_id>/<version_id>/<工程内相对路径>,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 第二装载落库)';
|
||||
@ -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 依赖内。
|
||||
|
||||
15
tier2/gen-worker/service/__init__.py
Normal file
15
tier2/gen-worker/service/__init__.py
Normal file
@ -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)。
|
||||
"""
|
||||
278
tier2/gen-worker/service/app.py
Normal file
278
tier2/gen-worker/service/app.py
Normal file
@ -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/<game_id>,与 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/<game_id>(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/<session_id>),与 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()
|
||||
262
tier2/gen-worker/service/bootstrap.py
Normal file
262
tier2/gen-worker/service/bootstrap.py
Normal file
@ -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: <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>(按 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}
|
||||
251
tier2/gen-worker/service/infra_config.py
Normal file
251
tier2/gen-worker/service/infra_config.py
Normal file
@ -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__<SECTION>__<KEY>(全大写、双下划线分隔)——临时压一个值不必改文件;
|
||||
② 外部 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>(双下划线分隔 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] 三级回落 + 转型自测通过。")
|
||||
@ -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 形态 = <game_id>/<version_id>/<工程内相对路径>(bucket 由 infra.yaml minio.bucket 决定,
|
||||
# 默认 tier2-src)。注:bucket 名不进 key(MinIO/S3 的 object key 不含 bucket 前缀),与凭据档
|
||||
# "key=tier2-src/<game_id>/..." 的写法差异 = 那处把 bucket 也写进了人读路径示意,实际 SDK put/get 只传
|
||||
# bucket + 不含 bucket 的 object key,故此处 key 从 <game_id>/ 起。
|
||||
|
||||
|
||||
# ── 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/<game_id>/<version_id>/<相对路径>;
|
||||
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 = <game_id>/<version_id>/<工程内相对路径>。
|
||||
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)://<endpoint>/<bucket>/<game_id>/<version_id>/(指向该版本所有源文件的前缀)。
|
||||
仅作 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()
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user