From d59093f84cb511ae04a7f049592763548c419191 Mon Sep 17 00:00:00 2001 From: zizi Date: Sun, 14 Jun 2026 15:03:53 +0000 Subject: [PATCH] =?UTF-8?q?feat(aigc-3b-B):=20=E6=B4=BE=E5=8F=91=E9=9D=A2?= =?UTF-8?q?=E7=9C=9F=20worker=20=E6=9C=8D=E5=8A=A1=E5=A3=B3=20+=20HMAC=20?= =?UTF-8?q?=E6=9C=8D=E5=8A=A1=E9=97=B4=E7=AD=BE=E5=90=8D(=E5=85=B3=20@Perm?= =?UTF-8?q?itAll=20=E8=A3=B8=E7=BC=BA=E5=8F=A3)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit P3 W-G1 派发面 3b-B:把 3b-A 的 stub worker 升级为真生成 worker 服务壳, 并补 HMAC-SHA256 服务间签名关掉 /dify/callback-internal 的裸 @PermitAll 安全缺口。 后端安全面(评审验证): - CallbackSignatureVerifier(新):HMAC-SHA256 对回调原始字节验签,MessageDigest.isEqual 常数时间比对(抗时序),JDK 原生零新依赖,空密钥=回退 3b-A 仅内网兜底(灰度/回滚兼容) - AigcExecutorProperties:加 callbackSecret 字段(来源 aigc.executor.callback-secret) - AigcExecutorConfiguration:注册 Verifier @Bean(随 executor.enabled 装配) - AdminAigcTaskController.difyCallbackInternal:接 @RequestBody String rawBody 先 HMAC 验签 (错签真 HTTP 401),再用同字节反序列化为 VO,补 traceId/status 非空(绕 @Valid 手动校验); 系统身份 id=0 注入 + finally clearContext 保留(部署窗口#3 同根因) - application-staging.yaml:加 callback-secret 占位符 + worker-url 指真 worker :9401 worker 服务壳(语法+契约对账验证): - wg1/gen-worker/worker/service.py(新):http.server 服务壳包 L2 run_studio (design 展开一句话→自产 gatespec→九门真玩),run.py/agent_loop 逐行不改; POST /generate 收 §6.1 job 立即回 202→后台串行真生成→读 bundle.iife.js→ HMAC 签名(逐字节对齐 Java 验签:json.dumps ensure_ascii=False 一次算定既签既发)→ POST 回调入参驱动落包入 feed 契约对账:执行器 dispatchGeneric job 含 brief/gameId/traceId/templateId ↔ service.py 读取一致; HMAC 两侧密钥(CALLBACK_SECRET/AIGC_CALLBACK_SECRET)须逐字节一致。 e2e(worker 起+HMAC双向证+一句话真生成入feed真玩)待 mini-desktop 跑。 Co-Authored-By: Claude Opus 4.8 (1M context) --- .../admin/task/AdminAigcTaskController.java | 95 ++++-- .../config/AigcExecutorConfiguration.java | 16 + .../executor/AigcExecutorProperties.java | 18 ++ .../executor/CallbackSignatureVerifier.java | 138 ++++++++ .../main/resources/application-staging.yaml | 5 +- wg1/gen-worker/worker/service.py | 300 ++++++++++++++++++ 6 files changed, 551 insertions(+), 21 deletions(-) create mode 100644 game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/service/executor/CallbackSignatureVerifier.java create mode 100644 wg1/gen-worker/worker/service.py diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/controller/admin/task/AdminAigcTaskController.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/controller/admin/task/AdminAigcTaskController.java index 67c5ddbe..a3c4ef7e 100644 --- a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/controller/admin/task/AdminAigcTaskController.java +++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/controller/admin/task/AdminAigcTaskController.java @@ -6,18 +6,22 @@ import cn.wanxiang.game.module.aigc.controller.app.task.vo.AigcTaskRespVO; import cn.wanxiang.game.module.aigc.convert.task.AigcTaskConvert; import cn.wanxiang.game.module.aigc.dal.dataobject.task.AigcTaskDO; import cn.wanxiang.game.module.aigc.service.callback.DifyCallbackService; +import cn.wanxiang.game.module.aigc.service.executor.CallbackSignatureVerifier; import cn.wanxiang.game.module.aigc.service.task.AigcTaskService; import cn.iocoder.yudao.framework.common.enums.UserTypeEnum; import cn.iocoder.yudao.framework.common.pojo.CommonResult; import cn.iocoder.yudao.framework.common.pojo.PageResult; import cn.iocoder.yudao.framework.security.core.LoginUser; import cn.iocoder.yudao.framework.tenant.core.aop.TenantIgnore; +import com.fasterxml.jackson.databind.ObjectMapper; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.tags.Tag; import jakarta.annotation.Resource; import jakarta.annotation.security.PermitAll; +import jakarta.servlet.http.HttpServletResponse; import jakarta.validation.Valid; import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.ObjectProvider; import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.security.authentication.UsernamePasswordAuthenticationToken; import org.springframework.security.core.context.SecurityContextHolder; @@ -26,6 +30,7 @@ import org.springframework.web.bind.annotation.*; import java.util.Collections; +import static cn.iocoder.yudao.framework.common.exception.enums.GlobalErrorCodeConstants.UNAUTHORIZED; import static cn.iocoder.yudao.framework.common.pojo.CommonResult.success; /** @@ -49,6 +54,17 @@ public class AdminAigcTaskController { @Resource private DifyCallbackService difyCallbackService; + /** 回调原始 body 反序列化用(验签后再解析为 VO;实例无共享可变状态,线程安全) */ + private static final ObjectMapper CALLBACK_MAPPER = new ObjectMapper(); + + /** + * 回调签名校验器(P3 W-G1·3b-B;软依赖):随 aigc.executor.enabled 装配。 + * 用 ObjectProvider 软注入——executor 未启用时本 Bean 缺席(内网回调本无 worker 来源), + * difyCallbackInternal 回退 3b-A「仅内网可达」兜底;启用且密钥非空 → 强制 HMAC 验签。 + */ + @Resource + private ObjectProvider signatureVerifierProvider; + @GetMapping("/task/page") @Operation(summary = "生成任务全量分页", description = "跨创作者排障,按状态/模板/traceId 检索(T-TEL-14)") @PreAuthorize("@ss.hasPermission('aigc:task:query')") @@ -70,39 +86,78 @@ public class AdminAigcTaskController { } /** - * 内网生成 worker 回调(P3 W-G1 派发面·3b-A,2026-06-14)。 + * 内网生成 worker 回调(P3 W-G1 派发面·3b-A 起;3b-B 补 HMAC 服务间签名)。 * * 【为何另开本路由(§6.5 auth 孤儿解法)】P3 generic 走「派发外置 worker」范式:worker 是【内部服务】、 * 无 RBAC 登录态 token,无法调上方带 {@code @PreAuthorize('aigc:dify:callback')} 的 /dify/callback。 * 故按本仓既有「外部/内部服务回调」范式(SmsCallbackController 的 {@code @PermitAll}、pay-notify 回调同款) * 另开本 {@code @PermitAll} 内网回调子路由,复用同一 {@link DifyCallbackService#handleCallback} 唯一写入路径 - * (建版本/组包/落包/回填全不变),把鉴权问题收敛为「仅内网可达」。worker 把 §6.1 result-out 经 - * {@link DifyCallbackReqVO}(含 engineBundle 字段)POST 到本端点驱动任务终态。 + * (建版本/组包/落包/回填全不变)。 * - * 【安全边界(诚实分界)】本路由跳过 admin RBAC,仅恃【内网不可外达】兜底(MVP staging:worker 与后端 - * 同 mini-desktop localhost)。⚠️ 3b-B 真 worker / 未来外部接入前,必须补【服务间签名校验】 - * (HMAC/共享密钥或 mTLS),与上方 /dify/callback 的「真 Dify 接入补签名」TODO 同治;上线前若本路由 - * 可被外网路由到,须在网关层显式拦截或加签名门,否则任意外部请求可伪造回调驱动落包(安全红线)。 - * 【TenantIgnore】worker 无租户上下文;与 SmsCallbackController 同款 {@code @TenantIgnore},任务定位经 - * traceId 全局唯一(handleCallback 内 selectByTraceId 不依赖租户上下文)。 + * 【3b-B 安全升级(关掉裸 @PermitAll 缺口)】3b-A 仅恃「内网不可外达」兜底——任意可达端点都能伪造回调 + * 驱动落包(安全红线)。3b-B 补 HMAC-SHA256 服务间签名:worker 用共享密钥(aigc.executor.callback-secret) + * 对【回调原始 body 字节】算签名置头 {@code X-Callback-Signature},本端口以同密钥重算【常数时间】比对, + * 不匹配即 HTTP 401 拒——升级为「内网且持密钥」双因子。 + * 【验签对原始字节(关键)】本方法接 {@code @RequestBody String rawBody}(而非直接 VO)——HMAC 是逐字节 + * 摘要,必须对 worker 实际发送、未经反序列化重整的原始字节验签,再用 {@link #CALLBACK_MAPPER} 解析为 VO, + * 否则字段顺序/空白差异必然误拒。 + * 【向后兼容/回滚】密钥未配(验签器空串)或 executor 未启用(验签器 Bean 缺席)→ 回退 3b-A 仅内网行为, + * 不阻断(灰度/回滚态);正式 3b-B 必须配非空密钥才算关掉裸缺口。 + * 【TenantIgnore + 系统身份】worker 无租户上下文 + 无 RBAC 登录态:与 SmsCallbackController 同款 @TenantIgnore; + * 并自注入系统身份 id=0(否则三表写链 updater/creator 撞 NOT NULL,部署窗口#3 同根因),finally 必清理。 + * + * @param rawBody 回调原始请求体(验签所用、与 worker 签名字节同源;随后反序列化为 DifyCallbackReqVO) + * @param signatureHeader 请求头 {@code X-Callback-Signature}(worker 算的 HMAC-SHA256 hex;缺失为 null) + * @param response 用于验签失败时直接写 HTTP 401 状态(便于 worker/排障侧以状态码判被拒) */ @PostMapping("/dify/callback-internal") @PermitAll @TenantIgnore - @Operation(summary = "内网生成 worker 完成回调(P3 派发面·@PermitAll 内网)", - description = "P3 generic 派发外置 worker 后,worker 把含 engineBundle 的 §6.1 result-out 经本内网免鉴权" - + "子路由 POST 回来,复用 handleCallback 唯一写入路径建版本/组包/落包。仅内网可达;真 worker/外部接入前补签名校验。") - public CommonResult difyCallbackInternal(@Valid @RequestBody DifyCallbackReqVO reqVO) { - // 复用唯一写入路径(与 /dify/callback 完全同款 handleCallback,仅鉴权门不同):worker 回调携 engineBundle, - // 经 DifyCallbackTxService.resolveEngineBundleText 取值后 put 进 GamePackage.engineBundle 落包入 feed。 - log.info("[aigc-callback-internal] 收到内网 worker 回调 traceId={}, status={}, engineBundleLen={}", + @Operation(summary = "内网生成 worker 完成回调(P3 派发面·@PermitAll 内网 + 3b-B HMAC 验签)", + description = "P3 generic 派发外置 worker 后,worker 把含 engineBundle 的 §6.1 result-out 经本内网子路由 POST 回来," + + "复用 handleCallback 唯一写入路径建版本/组包/落包。3b-B:HMAC-SHA256 服务间签名(X-Callback-Signature),错签 401 拒。") + public CommonResult difyCallbackInternal(@RequestBody String rawBody, + @RequestHeader(value = CallbackSignatureVerifier.SIGNATURE_HEADER, required = false) String signatureHeader, + HttpServletResponse response) { + // ===== 步骤①:HMAC 服务间签名校验(3b-B 安全红线;对原始字节验签)===== + CallbackSignatureVerifier verifier = signatureVerifierProvider.getIfAvailable(); + if (verifier != null && verifier.enabled()) { + if (!verifier.verify(rawBody, signatureHeader)) { + // 错签/缺签:HTTP 401 拒(关掉裸 @PermitAll 缺口;不落库、不驱动任何状态机) + response.setStatus(HttpServletResponse.SC_UNAUTHORIZED); + log.warn("[aigc-callback-internal] 回调签名校验未过,拒绝(HTTP 401)bodyLen={}, hasSig={}", + rawBody == null ? 0 : rawBody.length(), signatureHeader != null && !signatureHeader.isBlank()); + return CommonResult.error(UNAUTHORIZED.getCode(), "回调签名校验未过"); + } + } + // else:验签器缺席(executor 未启用)或密钥空串(验签关闭)→ 回退 3b-A 仅内网兜底(启动已 warn) + + // ===== 步骤②:验签通过后反序列化为 VO(用同一份原始字节解析)===== + DifyCallbackReqVO reqVO; + try { + reqVO = CALLBACK_MAPPER.readValue(rawBody, DifyCallbackReqVO.class); + } catch (Exception e) { + // body 非合法 JSON / 字段类型不符:400 语义(业务码),不驱动状态机 + log.warn("[aigc-callback-internal] 回调 body 解析失败,拒绝 bodyLen={}", rawBody == null ? 0 : rawBody.length(), e); + response.setStatus(HttpServletResponse.SC_BAD_REQUEST); + return CommonResult.error(400, "回调 body 解析失败"); + } + // 复用 /dify/callback 的入参非空约束(traceId/status @NotBlank):手动反序列化绕过了 @Valid,此处显式校验 + if (reqVO.getTraceId() == null || reqVO.getTraceId().isBlank() + || reqVO.getStatus() == null || reqVO.getStatus().isBlank()) { + log.warn("[aigc-callback-internal] 回调缺 traceId/status,拒绝"); + response.setStatus(HttpServletResponse.SC_BAD_REQUEST); + return CommonResult.error(400, "回调缺 traceId/status"); + } + + // ===== 步骤③:复用唯一写入路径(与 /dify/callback 完全同款 handleCallback,仅鉴权门不同)===== + log.info("[aigc-callback-internal] 收到内网 worker 回调(验签通过)traceId={}, status={}, engineBundleLen={}", reqVO.getTraceId(), reqVO.getStatus(), reqVO.getEngineBundle() == null ? 0 : reqVO.getEngineBundle().length()); // 【根因·与部署窗口#3 同源】本路由 @PermitAll、无 RBAC 登录态 → DefaultDBFieldHandler 取不到 userId → - // handleCallback 三表写链(建版本/组包/落包/回填,含受理置 RUNNING)UPDATE/INSERT 的 updater/creator - // 带 null 直出撞 NOT NULL(DataIntegrityViolation,任务卡 RUNNING)。RBAC 的 /dify/callback 由登录管理员 - // 提供身份故无此问题;本内网路由须自注入【系统身份 id=0】(审计列落 "0",镜像 - // AigcGenerateExecutor.executeWithSystemIdentity 同款,finally 必清理防 web 线程池复用污染)。 + // handleCallback 三表写链 UPDATE/INSERT 的 updater/creator 带 null 直出撞 NOT NULL(任务卡 RUNNING)。 + // 故须自注入【系统身份 id=0】(审计列落 "0",镜像 AigcGenerateExecutor.executeWithSystemIdentity, + // finally 必清理防 web 线程池复用污染)。 LoginUser systemUser = new LoginUser().setId(0L).setUserType(UserTypeEnum.ADMIN.getValue()); SecurityContextHolder.getContext().setAuthentication( new UsernamePasswordAuthenticationToken(systemUser, null, Collections.emptyList())); diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/framework/executor/config/AigcExecutorConfiguration.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/framework/executor/config/AigcExecutorConfiguration.java index 27fc30dc..366efd15 100644 --- a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/framework/executor/config/AigcExecutorConfiguration.java +++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/framework/executor/config/AigcExecutorConfiguration.java @@ -4,6 +4,7 @@ import cn.wanxiang.game.module.aigc.dal.mysql.task.AigcTaskMapper; import cn.wanxiang.game.module.aigc.service.callback.DifyCallbackService; import cn.wanxiang.game.module.aigc.service.executor.AigcExecutorProperties; import cn.wanxiang.game.module.aigc.service.executor.AigcGenerateExecutor; +import cn.wanxiang.game.module.aigc.service.executor.CallbackSignatureVerifier; import cn.wanxiang.game.module.aigc.service.executor.ExecutorLlmClient; import cn.wanxiang.game.module.aigc.service.executor.GameConfigSchemaValidator; import cn.wanxiang.game.module.aigc.service.executor.PromptResourceLoader; @@ -85,6 +86,21 @@ public class AigcExecutorConfiguration { return new WorkerDispatchClient(properties); } + /** + * 回调签名校验器(P3 W-G1·3b-B 安全红线:HMAC-SHA256 关掉 /dify/callback-internal 裸 @PermitAll 缺口) + * + * 随 aigc.executor.enabled 总开关一同装配(executor 启用才有 worker 派发、才有内网回调来源)。 + * AdminAigcTaskController 以 ObjectProvider 软注入本 Bean:存在且密钥非空 → 强制验签;executor 未启用 + * (本 Bean 缺席)→ 内网回调本无 worker 来源,controller 回退 3b-A 仅内网兜底(灰度/回滚态)。 + * + * @param properties 执行器配置(callback-secret 共享密钥) + * @return 校验器(密钥空串=验签关闭,启动 warn 提示) + */ + @Bean + public CallbackSignatureVerifier aigcCallbackSignatureVerifier(AigcExecutorProperties properties) { + return new CallbackSignatureVerifier(properties.getCallbackSecret()); + } + /** * 生成执行器(§5:tick 触发面 + 执行体;构造时执行 §5.2 阈值不等式校验,误配置 fail-fast) * diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/service/executor/AigcExecutorProperties.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/service/executor/AigcExecutorProperties.java index 51b1b45d..edc07c07 100644 --- a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/service/executor/AigcExecutorProperties.java +++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/service/executor/AigcExecutorProperties.java @@ -132,6 +132,24 @@ public class AigcExecutorProperties { */ private String callbackUrl = "http://localhost:48080/admin-api/aigc/dify/callback-internal"; + /** + * 服务间回调签名共享密钥(P3 W-G1·3b-B 安全红线:关掉 /dify/callback-internal「裸 @PermitAll」缺口)。 + * + * 【为何加(3b-A 标的安全红线)】3b-A 的内网回调子路由仅恃「内网不可外达」兜底——任意能路由到该端点的 + * 请求都能伪造回调驱动落包(DifyCallbackTxService 建版本/组包/落包入 feed)。3b-B 真 worker 上线,补 + * HMAC-SHA256 服务间签名:worker 用本密钥对回调原始 body 算签名置头 {@code X-Callback-Signature}, + * 后端入口重算比对(常数时间),不匹配 401 拒——把「仅内网可达」升级为「内网且持密钥」双因子。 + * + * 【取值约定】内网占位真值可入库(与项目「内网地址/密码可入库」口径一致),生产经环境变量覆盖更稳: + * staging application-staging.yaml 占位符 {@code ${AIGC_CALLBACK_SECRET:...}},worker 侧 .env + * {@code CALLBACK_SECRET} 必须与之【逐字节一致】,否则签名永不匹配。 + * + * 【空串语义(向后兼容/灰度回滚)】空串 = 验签关闭(回到 3b-A「裸 @PermitAll 仅内网」行为)—— + * 仅用于回滚或密钥未配阶段;正式 3b-B 必须配非空密钥才算「关掉裸缺口」。进程内执行器回调路 + * (AigcGenerateExecutor.invokeCallback 直接调 handleCallback、不走 HTTP)不受本签名门影响。 + */ + private String callbackSecret = ""; + /** * 阈值关系铁律校验(HJ-AGENT-LOOP-EXEC-002 §5.2 精确式;执行器装配时调用,误配置 fail-fast) * diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/service/executor/CallbackSignatureVerifier.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/service/executor/CallbackSignatureVerifier.java new file mode 100644 index 00000000..9e182b6f --- /dev/null +++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/cn/wanxiang/game/module/aigc/service/executor/CallbackSignatureVerifier.java @@ -0,0 +1,138 @@ +package cn.wanxiang.game.module.aigc.service.executor; + +import lombok.extern.slf4j.Slf4j; + +import javax.crypto.Mac; +import javax.crypto.spec.SecretKeySpec; +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; + +/** + * 服务间回调签名校验器(P3 W-G1·3b-B 安全红线:HMAC-SHA256 关掉 /dify/callback-internal 裸 @PermitAll 缺口) + * + * 【职责】对内网生成 worker 回调做服务间签名校验:worker 用共享密钥对回调【原始请求体字节】算 + * HMAC-SHA256(hex 小写)置头 {@code X-Callback-Signature},本校验器以同密钥重算并【常数时间】比对, + * 不匹配即判失败。把 3b-A 的「仅内网可达」升级为「内网且持密钥」双因子,杜绝任意可达端点伪造回调落包。 + * + * 【为何对「原始字节」而非反序列化后的 VO 算签】HMAC 是逐字节摘要——若后端先反序列化成 VO 再重新序列化 + * 去验签,字段顺序/空白/Unicode 转义任一差异都会让摘要不同(必然误拒)。故 controller 必须接「原始 body + * 字符串」、用同一份字节验签后再反序列化(worker 侧亦保证「签名用字节 == HTTP 发送字节」同源)。 + * + * 【为何用 JDK 原生 javax.crypto.Mac 而非 hutool】零新依赖、零模块 pom 改动、跨平台稳定;比对用 + * {@link MessageDigest#isEqual}(JDK 实现为常数时间,抗时序侧信道,避免逐字符短路比较泄漏匹配前缀长度)。 + * + * 【Bean 注册铁律(§3)】本类经 {@code AigcExecutorConfiguration} @Bean 注册(随 aigc.executor.enabled + * 总开关一同装配),类上不标 @Component/@Service,防组件扫描绕过开关。 + * + * @author 造梦AI(P3 派发面·3b-B 服务间签名) + */ +@Slf4j +public class CallbackSignatureVerifier { + + /** HMAC 算法名(JDK 标准名) */ + private static final String HMAC_ALGO = "HmacSHA256"; + + /** 回调签名头名(worker 置此头,后端读此头比对) */ + public static final String SIGNATURE_HEADER = "X-Callback-Signature"; + + /** 共享密钥(来自 aigc.executor.callback-secret;空串=验签关闭,仅回滚/未配阶段用) */ + private final String secret; + + /** + * 构造(经 AigcExecutorConfiguration @Bean 注册) + * + * @param secret 共享密钥(空串=验签关闭:enabled() 返回 false,controller 据此回退 3b-A 仅内网行为) + */ + public CallbackSignatureVerifier(String secret) { + this.secret = secret == null ? "" : secret; + // 启动留痕(不打印密钥本身,仅打印是否启用 + 长度,便于排障「为何 401/为何未验签」) + if (this.secret.isEmpty()) { + log.warn("[callback-sign] 回调签名校验【未启用】(aigc.executor.callback-secret 为空)——" + + "/dify/callback-internal 回退 3b-A「仅内网可达」兜底;3b-B 正式态应配非空密钥关掉裸缺口"); + } else { + log.info("[callback-sign] 回调签名校验【已启用】(HMAC-SHA256,密钥长度={});" + + "worker 须以同密钥对原始 body 算签置头 {}", this.secret.length(), SIGNATURE_HEADER); + } + } + + /** + * 是否启用验签(密钥非空即启用) + * + * @return true=已配密钥、强制验签;false=未配密钥、回退仅内网兜底 + */ + public boolean enabled() { + return !secret.isEmpty(); + } + + /** + * 校验回调签名(对原始请求体字节重算 HMAC-SHA256 并常数时间比对) + * + * @param rawBody 回调原始请求体(与 worker 签名所用、HTTP 发送的字节【同源】;UTF-8 编码) + * @param signatureHeader 请求头 {@code X-Callback-Signature} 值(worker 算的 hex 小写签名;缺失为 null) + * @return true=签名匹配(或验签未启用时恒 true);false=签名缺失/不匹配/算法异常(调用方应 401 拒) + */ + public boolean verify(String rawBody, String signatureHeader) { + if (!enabled()) { + // 未配密钥:验签关闭,放行(回退 3b-A 仅内网兜底;启动已 warn 提示) + return true; + } + if (signatureHeader == null || signatureHeader.isBlank()) { + // 启用验签但请求未带签名头 → 拒(可观测:便于排查 worker 漏置头 / 外部伪造无签名) + log.warn("[callback-sign] 回调缺签名头 {}(验签已启用),拒绝", SIGNATURE_HEADER); + return false; + } + String expected = hmacSha256Hex(rawBody == null ? "" : rawBody); + if (expected == null) { + // 算法异常(理论不可达:HmacSHA256 是 JDK 标配)→ 安全起见判失败 + return false; + } + // 常数时间比对(MessageDigest.isEqual 抗时序侧信道):按字节比,需同字符集 + boolean ok = MessageDigest.isEqual( + expected.getBytes(StandardCharsets.UTF_8), + signatureHeader.trim().getBytes(StandardCharsets.UTF_8)); + if (!ok) { + // 不匹配留痕:打印期望/实际的【前 8 位】辅助排查密钥不一致/body 被中间篡改,不泄露完整签名 + log.warn("[callback-sign] 回调签名不匹配,拒绝:expected(前8)={}, actual(前8)={}", + safePrefix(expected), safePrefix(signatureHeader.trim())); + } + return ok; + } + + /** + * 计算 HMAC-SHA256(hex 小写) + * + * @param data 待签数据(UTF-8) + * @return hex 小写签名;算法异常返回 null(理论不可达) + */ + private String hmacSha256Hex(String data) { + try { + Mac mac = Mac.getInstance(HMAC_ALGO); + mac.init(new SecretKeySpec(secret.getBytes(StandardCharsets.UTF_8), HMAC_ALGO)); + byte[] raw = mac.doFinal(data.getBytes(StandardCharsets.UTF_8)); + StringBuilder sb = new StringBuilder(raw.length * 2); + for (byte b : raw) { + sb.append(Character.forDigit((b >> 4) & 0xF, 16)); + sb.append(Character.forDigit(b & 0xF, 16)); + } + return sb.toString(); + } catch (Exception e) { + // HmacSHA256/密钥初始化异常(理论不可达):留痕后返回 null,调用方判验签失败 + log.error("[callback-sign] HMAC-SHA256 计算异常(理论不可达,请排查 JDK 加密提供方)", e); + return null; + } + } + + /** + * 取签名前 8 位(排障用,不泄露完整签名) + * + * @param sig 签名 + * @return 前 8 位(不足则原样) + */ + private static String safePrefix(String sig) { + if (sig == null) { + return ""; + } + return sig.length() <= 8 ? sig : sig.substring(0, 8); + } + +} diff --git a/game-cloud/yudao-server/src/main/resources/application-staging.yaml b/game-cloud/yudao-server/src/main/resources/application-staging.yaml index 63c95b81..68c6966c 100644 --- a/game-cloud/yudao-server/src/main/resources/application-staging.yaml +++ b/game-cloud/yudao-server/src/main/resources/application-staging.yaml @@ -177,9 +177,12 @@ aigc: api-key: ${NEWAPI_KEY:} # 密钥只走环境变量占位符,严禁写真值入 repo # P3 W-G1 派发面(3b-A,2026-06-14):generic 模板派发外置 worker 生成 engineBundle 入 feed worker-url: ${AIGC_WORKER_URL:} # 外置生成 worker 端点(空=派发面无端点→generic 任务判 failed; - # 3b-A 经 .env 注入 stub worker 地址,如 http://localhost:9301/wg1/generate) + # 3b-B 经 .env 注入真 worker 地址,如 http://localhost:9401/generate) # callback-url 默认即内网回调子路由(/admin-api/aigc/dify/callback-internal,@PermitAll);如需覆盖再放开: # callback-url: ${AIGC_CALLBACK_URL:http://localhost:48080/admin-api/aigc/dify/callback-internal} + # P3 W-G1·3b-B 安全红线:HMAC-SHA256 服务间回调签名共享密钥(关掉 /dify/callback-internal 裸 @PermitAll 缺口)。 + # worker 侧 .env CALLBACK_SECRET 必须与之逐字节一致;空=验签关闭(回滚 3b-A 仅内网)。内网占位真值可入库。 + callback-secret: ${AIGC_CALLBACK_SECRET:} --- #################### passport 鉴权(真实鉴权与匿名玩家 HJ-PASSPORT-EXEC-001)#################### wanxiang: diff --git a/wg1/gen-worker/worker/service.py b/wg1/gen-worker/worker/service.py new file mode 100644 index 00000000..295e1f0c --- /dev/null +++ b/wg1/gen-worker/worker/service.py @@ -0,0 +1,300 @@ +#!/usr/bin/env python3 +# -*- coding: utf-8 -*- +"""service.py —— W-G1 真生成 worker 的 HTTP 服务壳(P3 派发面·3b-B,2026-06-14)。 + +【职责(只加服务壳,run.py / agent_loop 逐行不改)】把 L2 agentic studio(design 展开一句话 + 自产 gatespec → + 便宜模型写 GameHostFactory → esbuild 打 __GameBundle → 九门真玩 harness → 失败回喂重试)包成 HTTP 服务: + + POST /generate ← 执行器(WorkerDispatchClient)投递 §6.1 job(job_id/brief/model/gameId/templateId/callback/traceId) + ↓ 立即回 202(投递握手成功,执行器不同步长挂;任务留 RUNNING 等回调) + 后台线程串行跑 run_studio(真烧 new-api 便宜模型生成)→ + ↓ 过九门 → 读 game-runtime/games/_wg1-gen//bundle.iife.js 全文作 engineBundle + 组后端回调入参 DifyCallbackReqVO(traceId/status/templateId/gameConfig 占位/engineBundle)→ + ↓ HMAC-SHA256 对回调原始字节签名置头 X-Callback-Signature(3b-B 安全红线) + POST 回 /admin-api/aigc/dify/callback-internal → 后端 handleCallback 建版本/组包/落包入 feed。 + +【与 3b-A stub worker(wg1_stub_worker.py)的关系】结构同款(http.server + 后台线程异步回调),仅两处升级: + ① 固定 pong bundle → 真 run_studio 真生成(design 展开一句话 + 九门真玩); + ② 回调 POST 加 HMAC 服务间签名(关掉 /dify/callback-internal 裸 @PermitAll 缺口,与后端 CallbackSignatureVerifier 对账)。 + +【HMAC 逐字节对齐铁律(与 Java CallbackSignatureVerifier 对账)】 + - 签名所用字节 == HTTP 发送字节 == json.dumps(payload, ensure_ascii=False).encode("utf-8")(一次算、原样发); + - Content-Type 显式带 charset=utf-8,使 Spring @RequestBody String 以 UTF-8 解码 → rawBody.getBytes(UTF_8) 还原同字节; + - hexdigest() 为 hex 小写,与 Java Character.forDigit 小写十六进制一致;密钥两侧(CALLBACK_SECRET / AIGC_CALLBACK_SECRET)须逐字节一致。 + +【串行(一次一 job)】run_studio 用固定 serve 端口 4320 + CDP 9222(serve-and-play.sh),不可并发;故全局 Lock, + 忙时回 409(执行器据此知 worker 占用,可重投或等)。 + +【运行(mini-desktop,须 cwd=repo 根 + worker venv + .env 配 new-api key/CALLBACK_SECRET)】 + cd /root/games-development-ai + CALLBACK_SECRET=<与后端 AIGC_CALLBACK_SECRET 同值> \ + wg1/gen-worker/.venv/bin/python wg1/gen-worker/worker/service.py --port 9401 + 执行器 application-staging.yaml 配 aigc.executor.worker-url=http://localhost:9401/generate + aigc.executor.callback-secret=<同值>。 + +@author 造梦AI(P3 派发面·3b-B 真 worker 服务壳) +""" + +import argparse +import asyncio +import hashlib +import hmac +import json +import os +import sys +import threading +import urllib.error +import urllib.request +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path + +# ── 把 worker/ 挂上 sys.path:import run(worker/run.py)+ from agent_loop.studio import run_studio ── +# 与 run.py / studio.py 的「直跑兜底」同款(python worker/service.py → sys.path[0]=worker/)。 +WORKER_DIR = Path(__file__).resolve().parent +if str(WORKER_DIR) not in sys.path: + sys.path.insert(0, str(WORKER_DIR)) + +import run # noqa: E402 worker/run.py(GEN_DIR/REPO_ROOT 常量 + scaffold/build/play) +from agent_loop.studio import run_studio # noqa: E402 L2 agentic 闭环(design 展一句话→九门) +from agent_loop import config # noqa: E402 models.yaml 阶段(默认 default_stage) + +# ---------- 默认常量(mini-desktop staging 口径)---------- +DEFAULT_PORT = 9401 +# 后端内网免鉴权回调子路由(AdminAigcTaskController.difyCallbackInternal,@PermitAll + 3b-B HMAC 验签) +DEFAULT_CALLBACK = "http://localhost:48080/admin-api/aigc/dify/callback-internal" +# meta 段标题最大长(与后端 META_TITLE 口径一致,截断防超) +META_TITLE_MAX = 20 + + +def log(msg): + """带前缀日志(stdout,flush 便于 nohup/systemd 落盘排障)。""" + print(f"[wg1-worker] {msg}", flush=True) + + +class WorkerState: + """进程级状态:回调 URL + HMAC 密钥 + 生成阶段 + 串行锁。""" + + def __init__(self, callback_url, secret, stage): + self.callback_url = callback_url # 后端内网回调子路由 + self.secret = secret or "" # HMAC 共享密钥(空=不签,仅 stub/回滚态) + self.stage = stage # models.yaml 阶段(None=config 默认) + self.lock = threading.Lock() # 串行锁:一次一 job(serve 4320/CDP 9222 不可并发) + self.busy = False # 占用标记(忙时 do_POST 回 409) + + +def _derive_title(brief): + """从一句话 brief 派生 meta 标题(截断到 META_TITLE_MAX,防超后端 meta 段约束)。""" + if not brief: + return "AI 生成游戏" + return brief.strip()[:META_TITLE_MAX] + + +def build_callback_payload(job, status, bundle_text, failure_reason=None): + """据 job + 生成结果组后端回调入参(DifyCallbackReqVO 形态)。 + + :param job: 收到的 §6.1 job-in(dict) + :param status: 'succeeded' | 'failed'(由 run_studio result['pass'] 派生) + :param bundle_text: engineBundle 真值(成功路 = bundle.iife.js 全文;失败路 = None) + :param failure_reason: 失败原因(失败路填,供后端落 failureReason 排障) + :return: 回调入参 dict + """ + trace_id = job.get("traceId") or job.get("job_id") # traceId 贯穿全链(与 job_id 同值) + template_id = job.get("templateId") or "generic" + brief = job.get("brief") or "" + payload = { + "traceId": trace_id, + "status": status, + "templateId": template_id, + # gameConfig 最小合法占位(§5.2):engine 路游戏逻辑全在 engineBundle 内,gameConfig 仍不可空 + # (GamePackage required + additionalProperties:false);故回最小合法占位。 + "gameConfig": { + "templateId": template_id, # 须与包级 templateId 一致(validateSucceededPayload 一致性校验) + "title": _derive_title(brief), # 供 meta 段派生 cover/summary + "theme": "generic", # 占位主题 + "engineDriven": True, # 标记本款走引擎 bundle 路(非填参 GameConfig) + }, + "assets": [], + } + if status == "succeeded": + # 成功路携 bundle 全文 → 后端 resolveEngineBundleText 取值落包入 feed + payload["engineBundle"] = bundle_text + else: + # 失败路:后端据 status=failed 把任务置 FAILED;failureReason 供排障 + payload["failureReason"] = (failure_reason or "generation_failed")[:500] + return payload + + +def post_callback(state, payload): + """把回调入参 POST 回后端内网回调子路由,并按 3b-B 红线加 HMAC 签名。 + + 【HMAC 对账(关键)】data 一次算定 → 既算签名又作 body 发送(同字节),charset=utf-8 使后端同字节还原。 + + :param state: WorkerState(含 callback_url + secret) + :param payload: 回调入参 dict + :return: (http_status, body_text) + """ + data = json.dumps(payload, ensure_ascii=False).encode("utf-8") # 签名字节 == 发送字节(铁律) + headers = { + "Content-Type": "application/json; charset=utf-8", # 显式 charset → Spring 以 UTF-8 解 @RequestBody String + "tenant-id": "0", # 兜底(@TenantIgnore 已忽略,带上无害) + } + if state.secret: + # HMAC-SHA256(secret, raw bytes) → hex 小写 → 头(与 Java CallbackSignatureVerifier.verify 对账) + sig = hmac.new(state.secret.encode("utf-8"), data, hashlib.sha256).hexdigest() + headers["X-Callback-Signature"] = sig + req = urllib.request.Request(state.callback_url, data=data, method="POST", headers=headers) + try: + with urllib.request.urlopen(req, timeout=30) as resp: + return resp.status, resp.read().decode("utf-8", "replace") + except urllib.error.HTTPError as e: + body = e.read().decode("utf-8", "replace") if e.fp else "" + return e.code, body + except Exception as e: # 网络异常:返回 -1 + 异常文本(不抛,日志记真因) + return -1, f"{type(e).__name__}: {e}" + + +def handle_job(state, job): + """后台线程:真生成(run_studio)→ 读 bundle → 组回调入参(含 HMAC)→ POST 回调。 + + 串行:进入时已持 state.lock 并置 busy=True;finally 释放(保证一次一 job + 不漏锁)。 + """ + trace_id = job.get("traceId") or job.get("job_id") + num_game_id = job.get("gameId") + brief = job.get("brief") or "" + # studio 落盘目录名(文件系统安全):gen-<数字 gameId>;与后端 job 的数字 gameId 解耦(后者是组包元数据) + game_id = f"gen-{num_game_id}" + log(f"开始真生成 trace_id={trace_id}, gameId={num_game_id}, game_dir={game_id}, brief={brief[:40]!r}") + status = "failed" + bundle_text = None + failure_reason = None + try: + # run_studio:design 展开一句话 + 自产 gatespec(gate-H)→ 代码/打包/九门真玩 → 失败回喂重试。 + # play_spec 传空 {}:gatespec 由设计 agent 自产并自动合入(不再手写 brief play-spec)。 + # async → 本线程内 asyncio.run(serve 4320 / CDP 9222 串行,全局锁已保证不并发)。 + result = asyncio.run(run_studio(game_id, brief, play_spec={}, stage=state.stage)) + passed = bool(result.get("pass")) + cost = (result.get("cost") or {}).get("total_rmb") + wall = result.get("wall_s") + log(f"run_studio 完成 trace_id={trace_id}, pass={passed}, repairs={result.get('repairs')}, " + f"wall={wall}s, ¥={cost}") + if passed: + # 读真生成 bundle 全文(绝对路径,run.GEN_DIR 已是绝对) + bundle_path = run.GEN_DIR / game_id / "bundle.iife.js" + if bundle_path.is_file(): + bundle_text = bundle_path.read_text(encoding="utf-8") + if "__GameBundle" not in bundle_text: + # 出厂自检:缺全局名宿主取不到工厂 → 判失败(不落坏包入 feed) + status = "failed" + failure_reason = "bundle_missing_global_name" + log(f"❌ bundle 缺 __GameBundle 全局名:{bundle_path}") + else: + status = "succeeded" + log(f"✅ bundle 已读 {len(bundle_text)} 字符(含 __GameBundle ✓)→ {bundle_path}") + else: + status = "failed" + failure_reason = "bundle_file_missing" + log(f"❌ pass 但 bundle 文件缺失:{bundle_path}") + else: + failure_reason = "nine_gate_failed" + except Exception as e: + failure_reason = f"{type(e).__name__}: {e}" + log(f"❌ 真生成异常 trace_id={trace_id}: {failure_reason}") + finally: + # 无论成败都回调(成功落包 / 失败置 FAILED),并释放串行锁 + try: + payload = build_callback_payload(job, status, bundle_text, failure_reason) + http_status, body = post_callback(state, payload) + log(f"回调完成 trace_id={trace_id}, status={status}, callbackHttp={http_status}, body={body[:200]}") + except Exception as e: + log(f"❌ 回调 POST 异常 trace_id={trace_id}: {type(e).__name__}: {e}") + finally: + state.busy = False + try: + state.lock.release() + except RuntimeError: + pass # 幂等:未持锁时释放无害 + + +def make_handler(state): + """工厂:绑定 state 的 HTTP handler 类。""" + + class Handler(BaseHTTPRequestHandler): + def log_message(self, fmt, *args): + return # 静默默认访问日志(用自定义 log) + + def _send_json(self, code, obj): + data = json.dumps(obj, ensure_ascii=False).encode("utf-8") + self.send_response(code) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(data))) + self.end_headers() + self.wfile.write(data) + + def do_GET(self): + # 健康探针:GET /health → 200(含 busy 标记,便于编排判 worker 占用) + if self.path.rstrip("/") in ("/health", "/wg1/health"): + self._send_json(200, {"ok": True, "busy": state.busy}) + else: + self._send_json(404, {"error": "not found"}) + + def do_POST(self): + # 派发入口:POST /generate(兼容 /wg1/generate)← 执行器投递 §6.1 job + if self.path.rstrip("/") not in ("/generate", "/wg1/generate"): + self._send_json(404, {"error": "not found", "path": self.path}) + return + length = int(self.headers.get("Content-Length", "0") or "0") + raw = self.rfile.read(length) if length > 0 else b"" + try: + job = json.loads(raw.decode("utf-8")) if raw else {} + except Exception as e: + log(f"job 解析失败:{e}") + self._send_json(400, {"accepted": False, "error": "bad job json"}) + return + trace_id = job.get("traceId") or job.get("job_id") + # 串行抢锁:忙则 409(一次一 job;run_studio 占用 serve 4320/CDP 9222 不可并发) + if not state.lock.acquire(blocking=False): + log(f"worker 占用中,拒 job trace_id={trace_id}(409)") + self._send_json(409, {"accepted": False, "error": "worker busy", "traceId": trace_id}) + return + state.busy = True + log(f"收到 job trace_id={trace_id}, templateId={job.get('templateId')}, gameId={job.get('gameId')}") + # 投递握手:立即回 202,后台线程真生成 + 回调(持锁,handle_job 的 finally 释放) + threading.Thread(target=handle_job, args=(state, job), daemon=True).start() + self._send_json(202, {"accepted": True, "job_id": job.get("job_id"), "traceId": trace_id}) + + return Handler + + +def main(): + parser = argparse.ArgumentParser(description="WG1 真生成 worker 服务壳(P3 派发面·3b-B)") + parser.add_argument("--port", type=int, default=DEFAULT_PORT, help=f"监听端口(默认 {DEFAULT_PORT})") + parser.add_argument("--callback", default=DEFAULT_CALLBACK, help="后端内网回调子路由 URL") + parser.add_argument("--stage", default=None, help="models.yaml 阶段(默认 default_stage)") + args = parser.parse_args() + + # HMAC 密钥:环境变量 CALLBACK_SECRET(须与后端 AIGC_CALLBACK_SECRET 逐字节一致;空=不签,仅回滚态) + secret = os.environ.get("CALLBACK_SECRET", "") + # 阶段:默认走 config 的 default_stage + stage = args.stage + if stage is None: + _cfg, default_stage = config.load_config() + stage = default_stage + + if secret: + log(f"HMAC 服务间签名【已启用】(密钥长度={len(secret)})→ 头 X-Callback-Signature") + else: + log("⚠️ HMAC 签名【未启用】(CALLBACK_SECRET 为空)——仅回滚/桩态用;3b-B 正式态须配密钥") + log(f"回调子路由:{args.callback}") + log(f"生成阶段:{stage}(models.yaml)") + + state = WorkerState(args.callback, secret, stage) + server = ThreadingHTTPServer(("0.0.0.0", args.port), make_handler(state)) + log(f"真生成 worker 已启动 → http://0.0.0.0:{args.port}/generate(Ctrl-C 退出)") + try: + server.serve_forever() + except KeyboardInterrupt: + log("收到中断,退出") + server.shutdown() + + +if __name__ == "__main__": + main()