feat(aigc-3b-B): 派发面真 worker 服务壳 + HMAC 服务间签名(关 @PermitAll 裸缺口)

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) <noreply@anthropic.com>
This commit is contained in:
zizi 2026-06-14 15:03:53 +00:00
parent ce9085df6b
commit d59093f84c
6 changed files with 551 additions and 21 deletions

View File

@ -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<CallbackSignatureVerifier> 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-A2026-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 stagingworker 与后端
* mini-desktop localhost 3b-B worker / 未来外部接入前必须补服务间签名校验
* HMAC/共享密钥或 mTLS与上方 /dify/callback Dify 接入补签名TODO 同治上线前若本路由
* 可被外网路由到须在网关层显式拦截或加签名门否则任意外部请求可伪造回调驱动落包安全红线
* TenantIgnoreworker 无租户上下文 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}而非直接 VOHMAC 是逐字节
* 摘要必须对 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<Boolean> 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-BHMAC-SHA256 服务间签名X-Callback-Signature错签 401 拒。")
public CommonResult<Boolean> 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 401bodyLen={}, 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 三表写链建版本/组包/落包/回填含受理置 RUNNINGUPDATE/INSERT updater/creator
// null 直出撞 NOT NULLDataIntegrityViolation任务卡 RUNNINGRBAC /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()));

View File

@ -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());
}
/**
* 生成执行器§5tick 触发面 + 执行体构造时执行 §5.2 阈值不等式校验误配置 fail-fast
*

View File

@ -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 建版本/组包/落包入 feed3b-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
*

View File

@ -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-SHA256hex 小写置头 {@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 造梦AIP3 派发面·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() 返回 falsecontroller 据此回退 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=签名匹配或验签未启用时恒 truefalse=签名缺失/不匹配/算法异常调用方应 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-SHA256hex 小写
*
* @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);
}
}

View File

@ -177,9 +177,12 @@ aigc:
api-key: ${NEWAPI_KEY:} # 密钥只走环境变量占位符,严禁写真值入 repo
# P3 W-G1 派发面3b-A2026-06-14generic 模板派发外置 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:

View File

@ -0,0 +1,300 @@
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""service.py —— W-G1 真生成 worker 的 HTTP 服务壳P3 派发面·3b-B2026-06-14
职责只加服务壳run.py / agent_loop 逐行不改 L2 agentic studiodesign 展开一句话 + 自产 gatespec
便宜模型写 GameHostFactory esbuild __GameBundle 九门真玩 harness 失败回喂重试包成 HTTP 服务
POST /generate 执行器WorkerDispatchClient投递 §6.1 jobjob_id/brief/model/gameId/templateId/callback/traceId
立即回 202投递握手成功执行器不同步长挂任务留 RUNNING 等回调
后台线程串行跑 run_studio真烧 new-api 便宜模型生成
过九门 game-runtime/games/_wg1-gen/<game_id>/bundle.iife.js 全文作 engineBundle
组后端回调入参 DifyCallbackReqVOtraceId/status/templateId/gameConfig 占位/engineBundle
HMAC-SHA256 对回调原始字节签名置头 X-Callback-Signature3b-B 安全红线
POST /admin-api/aigc/dify/callback-internal 后端 handleCallback 建版本/组包/落包入 feed
3b-A stub workerwg1_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须逐字节一致
串行一次一 jobrun_studio 用固定 serve 端口 4320 + CDP 9222serve-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 造梦AIP3 派发面·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.pathimport runworker/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.pyGEN_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):
"""带前缀日志stdoutflush 便于 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() # 串行锁:一次一 jobserve 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-indict
: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.2engine 路游戏逻辑全在 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 把任务置 FAILEDfailureReason 供排障
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=Truefinally 释放保证一次一 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_studiodesign 展开一句话 + 自产 gatespecgate-H→ 代码/打包/九门真玩 → 失败回喂重试。
# play_spec 传空 {}gatespec 由设计 agent 自产并自动合入(不再手写 brief play-spec
# async → 本线程内 asyncio.runserve 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一次一 jobrun_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}/generateCtrl-C 退出)")
try:
server.serve_forever()
except KeyboardInterrupt:
log("收到中断,退出")
server.shutdown()
if __name__ == "__main__":
main()