职责边界(经 Opus 评审纠正,别把空体 admit() 当缺陷):本门统计「每秒入队请求数」,不是「在跑生成并发」。 - * 空体入口的 entry/exit 是纳秒级,thread-grade 统计不到真正在跑的生成——那发生在 admit() 返回之后的 RocketMQ - * 消费线程里。故「在跑生成并发≤15」由 C3 消费者 consumeThreadMax=15 保证,不由本门保证。 + * 空体入口的 entry/exit 是纳秒级,thread-grade 统计不到真正在跑的生成——真生成发生在 admit() 返回、入队、消费线程 + * 投递之后的 worker 侧(异步跑数十秒)。故在跑生成并发不由本门保证:C3 consumeThreadMax=15 只 cap dispatch 投递速率, + * 真「在跑生成并发≤15」由有界 worker 池保证(follow-up,不在 C3)。 * *
三层正交: *
真「在跑生成并发≤15」的硬上限由有界 worker 池保证(创始人 2026-06-26 决策:异步任务队列 + 有界 worker 池 + * N=4-6 起步压测调、≤15 硬上限),作 follow-up 工作单元实现、不在 C3。 * *
越限行为:Sentinel 抛 {@link BlockException} → {@link #onAdmissionBlocked} 优雅转成业务异常 * {@code AIGC_ADMISSION_RATE_LIMITED}「生成服务繁忙…请稍后重试」(经全局 GlobalExceptionHandler 转 CommonResult.error, @@ -43,7 +46,7 @@ public class GenAdmissionResource { * 通过则直接返回(调用方继续走 DB 三门 + 入队);被 block 则由 Sentinel 切面转调 {@link #onAdmissionBlocked}。 * *
空体是设计,不是缺陷:{@code FLOW_GRADE_QPS} 只统计「每秒通过本资源的次数」(与方法耗时无关),
- * 空体正合适——让 Sentinel 计入队 QPS 削洪峰。业务准入(配额)在调用方 DB 三门;在跑并发≤15 在 C3 consumeThreadMax,均不在此。
+ * 空体正合适——让 Sentinel 计入队 QPS 削洪峰。业务准入(配额)在调用方 DB 三门;C3 consumeThreadMax=15 cap dispatch 投递速率、真在跑并发≤15 属有界 worker 池(follow-up),均不在此。
*/
@SentinelResource(value = RESOURCE, blockHandler = "onAdmissionBlocked")
public void admit() {
diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/framework/executor/config/AigcExecutorConfiguration.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/framework/executor/config/AigcExecutorConfiguration.java
index 5f80b03a..79715643 100644
--- a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/framework/executor/config/AigcExecutorConfiguration.java
+++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/framework/executor/config/AigcExecutorConfiguration.java
@@ -9,6 +9,7 @@ import com.wanxiang.huijing.game.module.aigc.service.executor.ExecutorLlmClient;
import com.wanxiang.huijing.game.module.aigc.service.executor.GameConfigSchemaValidator;
import com.wanxiang.huijing.game.module.aigc.service.executor.PromptResourceLoader;
import com.wanxiang.huijing.game.module.aigc.service.executor.WorkerClassifyClient;
+import com.wanxiang.huijing.game.module.aigc.mq.GenTaskProducer;
import com.wanxiang.huijing.game.module.aigc.saa.SaaGraphDispatcher;
import com.wanxiang.huijing.game.module.aigc.service.executor.GenerationDispatcher;
import com.wanxiang.huijing.game.module.aigc.service.executor.WorkerDispatchClient;
@@ -163,6 +164,8 @@ public class AigcExecutorConfiguration {
* @param llmClient LLM 通道
* @param workerDispatchClient 外置生成 worker 派发通道(P3 generic 路)
* @param sourceProjectApiProvider 源项目反查 seam(A11 M1②:modify 路据 baseVersionId 反查 base 源注入 job.sourceProject;软取)
+ * @param genTaskProducerProvider gen 队列生产者 seam(C3:@Scheduled 兜底 tick 对 stale queued 重发信号交有界消费者;软取,
+ * 缺席→兜底退回 tick 直派,非破坏性。与本配置同 aigc.executor.enabled 条件,enabled 时恒在席)
* @return 执行器(@Scheduled 注解由全局调度后处理器消费,fixedDelay 单飞)
*/
@Bean
@@ -174,15 +177,20 @@ public class AigcExecutorConfiguration {
ExecutorLlmClient llmClient,
WorkerDispatchClient workerDispatchClient,
SaaGraphDispatcher saaGraphDispatcher,
- ObjectProvider 复用既有、不重写:把「CAS 认领({@code claimQueuedTask}) + 生成流水派发({@code runGenerationPipeline}
+ * → generic 走 {@code SaaGraphDispatcher}/http worker)」收敛为一次调用 {@link AigcGenerateExecutor#claimAndDispatchById}。
+ * 与执行器 {@code @Scheduled} tick 的 {@code processCandidate} 同一执行体、同一系统身份 + 租户上下文,只是触发面
+ * 从「批量扫描」换成「MQ 单 taskId 事件」。不复制三表写链、不重写派发逻辑。
+ *
+ * 幂等 + 派发失败落 failed(细节见 {@link AigcGenerateExecutor#claimAndDispatchById}):
+ * CAS 认领失败(被并发抢走/非 queued/已取消/任务不存在)→ 返回 false 幂等丢弃;认领成功但派发抛异常
+ * → 内部落 failed(不留悬空 running)后返回 true。故本方法不外抛,消费者据返回值只做日志、不重投。
+ *
+ * 装配:随 {@code aigc.executor.enabled=true} 装配(与执行器/生产者/消费者同一开关);执行器 Bean 由
+ * {@code AigcExecutorConfiguration} 同条件注册,故此处硬注入 {@link AigcGenerateExecutor} 恒可解析。
+ *
+ * @author 绘境AI
+ */
+@Slf4j
+@Service
+@ConditionalOnProperty(prefix = "aigc.executor", name = "enabled", havingValue = "true")
+public class AigcGenDispatchService {
+
+ /** 既有生成执行器(AigcExecutorConfiguration @Bean 注册,同 aigc.executor.enabled 条件)——复用其认领+派发执行体。 */
+ private final AigcGenerateExecutor generateExecutor;
+
+ public AigcGenDispatchService(AigcGenerateExecutor generateExecutor) {
+ this.generateExecutor = generateExecutor;
+ }
+
+ /**
+ * 认领并派发一个生成任务(MQ 消费者收到 taskId 信号后调用)。
+ *
+ * @param taskId 任务 ID(消息体轻信号)
+ * @return true=CAS 认领成功(已派发或已落终态);false=认领未命中=幂等丢弃(同 taskId 重复投递/被并发抢走/非 queued)
+ */
+ public boolean claimAndDispatch(Long taskId) {
+ // 委托既有执行器:CAS 认领 + 复用 runGenerationPipeline 派发;认领失败返 false,派发抛异常内部落 failed(不外抛)。
+ boolean claimed = generateExecutor.claimAndDispatchById(taskId);
+ log.debug("[aigc-mq] 认领+派发结论 taskId={} claimed={}", taskId, claimed);
+ return claimed;
+ }
+}
diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/mq/GenTaskConsumer.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/mq/GenTaskConsumer.java
new file mode 100644
index 00000000..8610b57a
--- /dev/null
+++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/mq/GenTaskConsumer.java
@@ -0,0 +1,74 @@
+package com.wanxiang.huijing.game.module.aigc.mq;
+
+import lombok.extern.slf4j.Slf4j;
+import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
+import org.apache.rocketmq.spring.core.RocketMQListener;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.stereotype.Component;
+
+/**
+ * 生成任务 MQ 消费者(切片一 阶段〇 C3)——有界并发派发。
+ *
+ * {@code consumeThreadMax=15} = dispatch(投递握手)并发速率上限,不是在跑生成并发上限:消费线程收到
+ * taskId 信号后 CAS 认领 → 组 §6.1 job 投递生成躯体(http worker POST 拿 ACK / saa 进程内交后台池),投递握手即
+ * 返回释放(sub-second),不在消费线程内长挂等生成结果——生成是 worker 异步跑数十秒、回调驱动终态(见
+ * {@link com.wanxiang.huijing.game.module.aigc.service.executor.GenerationDispatcher}「仅投递握手,不等生成结果」契约)。
+ * 故这 15 线程 cap 的是「每秒能投递多少 job」,而非「同时真跑多少生成」(投递后任务即留 RUNNING,真在跑的生成挂在 worker 侧)。
+ *
+ * 真「在跑生成并发≤15」由有界 worker 池保证(follow-up,不在 C3):符合创始人 2026-06-26 决策——异步任务
+ * 队列 + 有界 worker 池(N=4-6 起步、压测调,≤15 硬上限);该 in-flight 硬上限作独立 follow-up 工作单元实现,本 C3
+ * 只落队列机制。三层正交:C2 Sentinel 管入口 QPS「每秒进多少」、DB 三门管业务配额「每人/全局多少在飞」、本消费线程
+ * 池管投递速率「每秒投递多少 job」。多实例扩容后的集群级全局限流留放量后(Token Server)。
+ *
+ * 幂等(CAS 认领去重):收到 taskId 信号 → 经 {@link AigcGenDispatchService#claimAndDispatch} 用既有
+ * CAS {@code claimQueuedTask}(queued→running) 认领,认领成功才派发;认领失败(已被别的消费线程/兜底 tick 拿走、
+ * 非 queued、已取消、任务不存在)即幂等丢弃——同一 taskId 重复投递不会双跑。
+ *
+ * 不无限重投:{@code claimAndDispatch} 内部已处置派发异常(认领成功但派发抛异常 → 落 failed,不留悬空
+ * running),正常返回 true/false 不外抛;故 {@link #onMessage} 认领失败/毒消息只记日志、正常返回(ack),
+ * 不让 RocketMQ 无限重投。消费模型 = 默认 CLUSTERING(每条消息集群内单实例消费,即工作队列语义)。
+ *
+ * 装配:随 {@code aigc.executor.enabled=true} 装配(与执行器/生产者同一开关)——executor 关闭时不起消费者、
+ * 不连 NameServer(避免无执行器时空转/无谓连接)。
+ *
+ * @author 绘境AI
+ */
+@Slf4j
+@Component
+@ConditionalOnProperty(prefix = "aigc.executor", name = "enabled", havingValue = "true")
+@RocketMQMessageListener(
+ topic = GenTaskProducer.TOPIC,
+ consumerGroup = "aigc_gen_consumer",
+ consumeThreadMax = 15, // rocketmq-client 无界队列致 consumeThreadMax 不单独生效(线程数永不超 core),故与 consumeThreadNumber 同设 15 让投递并发上限真为 15
+ consumeThreadNumber = 15 // min=max=15:15 个投递线程常驻(投递握手 sub-second、空闲线程零成本);这是 dispatch 投递并发、非在跑生成并发(后者=有界 worker 池 follow-up)
+)
+public class GenTaskConsumer implements RocketMQListener 消息体只放 taskId(轻信号):任务真身/状态仍以 {@code game_aigc_task} 表为准——MQ 是投递触发、
+ * 不是数据搬运。消费侧据 taskId 用既有 CAS {@code claimQueuedTask}(queued→running) 认领,认领成功才派发,
+ * 同一 taskId 重复投递被 CAS 天然去重(不双跑)。
+ *
+ * best-effort(命门):发 MQ 失败不回滚入队(任务已落库 queued)、也不外抛——只告警留痕。
+ * 丢失的信号由执行器 {@code @Scheduled} 降级兜底轮询补偿捡起(stale queued 重发信号,见
+ * {@link com.wanxiang.huijing.game.module.aigc.service.executor.AigcGenerateExecutor} 兜底补偿路)。
+ * 故 MQ 抖动只影响「即时性」、不影响「不丢任务」。
+ *
+ * 装配条件:随 {@code aigc.executor.enabled=true} 装配(与执行器/消费者同一开关)——executor 关闭时
+ * 无消费者、无执行器,发信号无意义,故一并缺席;{@link com.wanxiang.huijing.game.module.aigc.service.task.AigcTaskServiceImpl}
+ * 以 {@code ObjectProvider} 软注入本 Bean,缺席时入队照常落库、仅不发信号(现行骨架行为)。
+ *
+ * @author 绘境AI
+ */
+@Slf4j
+@Component
+@ConditionalOnProperty(prefix = "aigc.executor", name = "enabled", havingValue = "true")
+public class GenTaskProducer {
+
+ /** gen 队列 topic(与 B2 broker 在 mini-infra 建的 topic 对齐;消费者 {@link GenTaskConsumer} 同名订阅)。 */
+ public static final String TOPIC = "AIGC_GEN_TOPIC";
+
+ /** RocketMQ 模板(rocketmq-spring 自动装配;name-server 见 application*.yaml → mini-infra 100.64.0.8:9876)。 */
+ private final RocketMQTemplate rocketMQTemplate;
+
+ public GenTaskProducer(RocketMQTemplate rocketMQTemplate) {
+ this.rocketMQTemplate = rocketMQTemplate;
+ }
+
+ /**
+ * 发「新任务」信号(同步发,拿发送结果落日志;失败只告警、不抛,由执行器兜底轮询补偿)。
+ *
+ * @param taskId 生成任务 ID(消息体轻信号;消费侧据此反查 game_aigc_task + CAS 认领)
+ */
+ public void publishNewTask(Long taskId) {
+ try {
+ // 同步发:payload 只放 taskId 字符串(轻信号)。syncSend(String, Object) 拿 SendResult 落痕便于排障。
+ var result = rocketMQTemplate.syncSend(TOPIC, String.valueOf(taskId));
+ log.info("[aigc-mq] 发送 gen 任务信号 topic={} taskId={} msgId={} status={}",
+ TOPIC, taskId, result != null ? result.getMsgId() : null,
+ result != null ? result.getSendStatus() : null);
+ } catch (Exception e) {
+ // best-effort:MQ 抖动不连累入队(任务已落库 queued),也不外抛——由 @Scheduled 兜底轮询补偿捡起。
+ log.warn("[aigc-mq] 发送 gen 信号失败(不回滚入队、不外抛,由执行器兜底轮询补偿)taskId={}", taskId, e);
+ }
+ }
+}
diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/executor/AigcExecutorProperties.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/executor/AigcExecutorProperties.java
index ae12006e..d6b6a97b 100644
--- a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/executor/AigcExecutorProperties.java
+++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/executor/AigcExecutorProperties.java
@@ -70,8 +70,12 @@ public class AigcExecutorProperties {
/** LLM 模型名(沿 C1 spike / M-a 批跑同款通道) */
private String llmModel = "MiniMax-M2.7";
- /** tick 轮询间隔毫秒(@Scheduled fixedDelay:上一 tick 跑完才调度下一 tick,进程级单飞) */
- private Long pollIntervalMs = 10000L;
+ /**
+ * tick 轮询间隔毫秒(@Scheduled fixedDelay:上一 tick 跑完才调度下一 tick,进程级单飞)。
+ * C3:MQ 事件驱动为主触发后,tick 降级为兜底补偿,周期由 10s 拉大至 60s(不再作主触发路径);
+ * 同时作 C3 兜底路「stale queued」判据的一个周期基准(createTime 早于「一个 poll 周期前」才重发信号)。
+ */
+ private Long pollIntervalMs = 60000L;
/** 每 tick 认领上限(灰度小流量;扫描按 id 升序 FIFO 公平) */
private Integer scanBatchSize = 2;
diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/executor/AigcGenerateExecutor.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/executor/AigcGenerateExecutor.java
index 2a468d5c..aeaf9d5c 100644
--- a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/executor/AigcGenerateExecutor.java
+++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/executor/AigcGenerateExecutor.java
@@ -4,6 +4,7 @@ import com.wanxiang.huijing.game.module.aigc.controller.admin.task.vo.DifyCallba
import com.wanxiang.huijing.game.module.aigc.dal.dataobject.task.AigcTaskDO;
import com.wanxiang.huijing.game.module.aigc.dal.mysql.task.AigcTaskMapper;
import com.wanxiang.huijing.game.module.aigc.enums.FailureReasonEnum;
+import com.wanxiang.huijing.game.module.aigc.mq.GenTaskProducer;
import com.wanxiang.huijing.game.module.aigc.service.callback.DifyCallbackService;
import com.wanxiang.huijing.game.module.studio.api.SourceProjectApi;
import com.wanxiang.huijing.game.module.studio.dto.SourceProjectFetchRespDTO;
@@ -42,6 +43,13 @@ import static com.wanxiang.huijing.game.module.aigc.enums.ErrorCodeConstants.AIG
* {@link DifyCallbackService#handleCallback}(复用 EXEC-001 §8.4 唯一写入路径,不绕过、不复制三表逻辑)。
* 从此真实用户 submitGenerate 不再永久停 queued(M2 整体收口第②段)。
*
+ * 【C3 触发面变更(切片一 阶段〇:RocketMQ 异步 gen 队列)】主触发从「@Scheduled DB 轮询」改为「MQ 事件驱动」:
+ * 入队即发 AIGC_GEN_TOPIC 信号 → 有界消费者(consumeThreadMax=15=dispatch 投递握手并发速率上限,非在跑生成并发;
+ * 真在跑并发≤15 由有界 worker 池保证=follow-up)→ 经
+ * {@link #claimAndDispatchById} 走本类同一 CAS 认领 + 生成流水(复用,不重写)。{@link #tick()} @Scheduled
+ * 降级为兜底补偿(周期拉大至 60s):① stale queued(信号丢)重发信号交消费者、② stale running 僵死收尸
+ * failed。派发失败落 failed 不留悬空 running({@link #claimAndDispatchById} catch 路)。
+ *
* 【调度选型依据(§3,亲核 2026-06-10)】Spring @Scheduled 薄轮询,不依赖 XXL-Job:
* staging xxl.job.enabled=false 且无 xxl-admin → HuijingXxlJobAutoConfiguration 整类不装配 →
* @XxlJob 触发面在 staging 永不运行(trade SettlementJob 即此死法,B4 烟测实证)。
@@ -107,6 +115,15 @@ public class AigcGenerateExecutor {
private final SourceProjectApi sourceProjectApi;
/** 墙钟(预算门用;单测注入假时钟验证 M5) */
private final Clock clock;
+ /**
+ * gen 队列生产者(C3 兜底补偿路专用,可为 null):MQ 为主触发后,{@code @Scheduled} tick 降级为兜底补偿,
+ * 对捞到的 stale queued 重发信号(交消费者重走认领+派发),不在 tick 线程直接派发——这是让 MQ 保持
+ * 唯一投递路径的投递可靠性补偿(信号丢了补一条、回到同一条消费→派发链路,不另起第二条并行派发路径),
+ * 与并发 cap 无关(真「在跑生成并发≤15」属有界 worker 池=follow-up,不由消费线程数或此重发决定)。
+ * null 时降级为直派(producer 未注入的旧构造/单测态 / 生产者 Bean 因故缺席):tick 退回既有
+ * 「扫描-认领-派发」行为,避免任务卡死(韧性降级,非主路径)。生产经 AigcExecutorConfiguration 注入非空。
+ */
+ private final GenTaskProducer genTaskProducer;
/** 自禁用连续 tick 计数(自检恢复即清零;告警节流判据) */
private int selfCheckFailTicks = 0;
@@ -155,19 +172,34 @@ public class AigcGenerateExecutor {
}
/**
- * 10 参构造(A11 切片三 M1② 收口;经 AigcExecutorConfiguration @Bean 调用):在 9 参基础上接
- * {@link SourceProjectApi} 反查 seam——modify 路据 baseVersionId 反查 base 源注入 §6.1 job 的 sourceProject 键
- * (供便宜档 HTTP worker 无 DB 消费)。按 {@code aigc.executor.dispatcher} 选派发通道——{@code saa} 且 saaDispatcher
- * 非空 → 进程内 SaaGraphDispatcher(形态①);否则 → http {@link WorkerDispatchClient}(默认,现行)。
- *
- * @param saaDispatcher SAA 进程内派发器(可空:dispatcher=http 或 saa 未装配时为 null → 回落 http)
- * @param sourceProjectApi 源项目反查 seam(A11 modify 取 base 源;可 null=取源旁路,不注入 sourceProject)
+ * 10 参构造(A11 切片三 M1② 收口;现行单测沿用此签名):委托 11 参构造,{@code genTaskProducer} 传 null →
+ * {@code @Scheduled} 兜底 tick 退回「直接认领+派发」(C3 前行为,非破坏性)。生产经 11 参构造注入非空生产者,
+ * 兜底 tick 改为「重发信号交消费者」(投递可靠性补偿,让 MQ 保持唯一投递路径;与并发 cap 无关)。
*/
public AigcGenerateExecutor(AigcExecutorProperties properties, AigcTaskMapper aigcTaskMapper,
DifyCallbackService difyCallbackService, PromptResourceLoader promptResourceLoader,
GameConfigSchemaValidator schemaValidator, ExecutorLlmClient llmClient,
WorkerDispatchClient workerDispatchClient, GenerationDispatcher saaDispatcher,
SourceProjectApi sourceProjectApi, Clock clock) {
+ // 委托 11 参构造,genTaskProducer 传 null → 兜底 tick 直派(C3 前行为,非破坏性;现行单测沿用此签名)。
+ this(properties, aigcTaskMapper, difyCallbackService, promptResourceLoader, schemaValidator,
+ llmClient, workerDispatchClient, saaDispatcher, sourceProjectApi, clock, null);
+ }
+
+ /**
+ * 11 参构造(切片一 阶段〇 C3 收口;经 AigcExecutorConfiguration @Bean 调用):在 10 参基础上接 gen 队列生产者
+ * {@link GenTaskProducer}——MQ 为主触发后,{@code @Scheduled} tick 降级为兜底补偿,对 stale queued 重发信号
+ * (交消费者重走认领+派发),不在 tick 线程直接派发(投递可靠性补偿,让 MQ 保持唯一投递路径;与并发 cap 无关)。
+ *
+ * @param saaDispatcher SAA 进程内派发器(可空:dispatcher=http 或 saa 未装配时为 null → 回落 http)
+ * @param sourceProjectApi 源项目反查 seam(A11 modify 取 base 源;可 null=取源旁路,不注入 sourceProject)
+ * @param genTaskProducer gen 队列生产者(C3 兜底 tick 重发信号用;可 null=兜底退回直派,非破坏性/单测态)
+ */
+ public AigcGenerateExecutor(AigcExecutorProperties properties, AigcTaskMapper aigcTaskMapper,
+ DifyCallbackService difyCallbackService, PromptResourceLoader promptResourceLoader,
+ GameConfigSchemaValidator schemaValidator, ExecutorLlmClient llmClient,
+ WorkerDispatchClient workerDispatchClient, GenerationDispatcher saaDispatcher,
+ SourceProjectApi sourceProjectApi, Clock clock, GenTaskProducer genTaskProducer) {
// §5.2 阈值关系铁律:stale×60 > budget + 273,装配期校验(误配置宁可启动失败也不带病收割在飞任务)
properties.validateThresholds();
this.properties = properties;
@@ -184,11 +216,14 @@ public class AigcGenerateExecutor {
log.info("[executor] 生成派发通道选定 dispatcher={}(saa装配={})= {}", properties.getDispatcher(),
saaDispatcher != null, useSaa ? "进程内SaaGraphDispatcher(形态①)" : "http WorkerDispatchClient(现行)");
this.clock = clock;
+ // C3:gen 队列生产者(兜底 tick 重发信号用);非空=MQ 主触发 + 兜底重发信号(投递可靠性补偿),null=兜底退回直派(单测/降级态)。
+ this.genTaskProducer = genTaskProducer;
// 启动三行自检日志(§12-① 冒烟观察面:执行器已启用 / prompt 资源版本 / key 配置态;严禁打印密钥本身)
- log.info("[executor-selfcheck] aigc 生成执行器已启用(poll={}ms, batch={}, budget={}s, stale={}min, maxAge={}h, templates={}, sourceSeam={})",
+ log.info("[executor-selfcheck] aigc 生成执行器已启用(poll={}ms, batch={}, budget={}s, stale={}min, maxAge={}h, templates={}, sourceSeam={}, 触发={})",
properties.getPollIntervalMs(), properties.getScanBatchSize(), properties.getTaskBudgetSeconds(),
properties.getStaleRunningMinutes(), properties.getMaxTaskAgeHours(), properties.getSupportedTemplates(),
- sourceProjectApi != null ? "在席" : "缺席(modify取源旁路)");
+ sourceProjectApi != null ? "在席" : "缺席(modify取源旁路)",
+ genTaskProducer != null ? "MQ主触发+tick兜底重发信号" : "tick直派(无MQ生产者/单测态)");
// HJ-MC-TPL-EXEC-001 §6.4:多模板装载后无单一版本,按 templateId 逐模板列出版本(就绪时)
log.info("[executor-selfcheck] prompt 资源版本 {}", promptResourceLoader.isReady()
? promptResourceLoader.describePromptVersions() : ("未就绪:" + promptResourceLoader.getNotReadyReason()));
@@ -196,11 +231,13 @@ public class AigcGenerateExecutor {
}
/**
- * 调度触发面(§5.1):fixedDelay 轮询(默认 10s,上一 tick 跑完才开下一 tick)+ 启动延迟 30s
- * (避开 app 启动期 Bean 装配/Flyway 迁移窗口)。XXL-Job 上线后本方法被 @XxlJob 触发面替换,
- * 执行体 {@link #runOneTick()} 原样复用(§3 演化路径)。
+ * 调度触发面(§5.1;C3 降级为兜底补偿):fixedDelay 轮询(默认 60s,周期已拉大——MQ 事件驱动为主触发后,
+ * 本 tick 不再作主触发路径,只当兜底补偿)+ 启动延迟 30s(避开 app 启动期 Bean 装配/Flyway 迁移窗口)。
+ * 兜底补偿捞两类卡住任务({@link #runOneTick()}):① stale queued(入队了但 MQ 信号丢)→ 有生产者则
+ * 重发信号交有界消费者认领、无生产者则直派兜底;② stale running 超阈值(进程崩溃/派发丢状态遗留)→ watchdog
+ * 收尸 failed。XXL-Job 上线后本方法被 @XxlJob 触发面替换,执行体原样复用(§3 演化路径)。
*/
- @Scheduled(fixedDelayString = "${aigc.executor.poll-interval-ms:10000}", initialDelay = 30_000)
+ @Scheduled(fixedDelayString = "${aigc.executor.poll-interval-ms:60000}", initialDelay = 30_000)
public void tick() {
try {
runOneTick();
@@ -225,8 +262,28 @@ public class AigcGenerateExecutor {
int claimed = 0;
int succeeded = 0;
int failed = 0;
+ int resignaled = 0; // C3 兜底补偿:对 stale queued 重发 MQ 信号的条数(有生产者的生产态走此路,不直派)
+ // C3 stale queued 判据:createTime 早于「一个 poll 周期前」才算信号大概率已丢/滞留,才重发;
+ // 新鲜 queued(刚入队、信号在途)不重发以减噪(重发即便发生也被消费侧 CAS 幂等去重,无害)。
+ LocalDateTime staleQueuedBefore = LocalDateTime.now().minus(Duration.ofMillis(properties.getPollIntervalMs()));
// ===== 步骤③:逐任务处理(tick 线程内串行;单任务异常不中断本 tick 其余任务)=====
for (AigcTaskDO task : candidates) {
+ // ===== C3 兜底补偿路(生产态:MQ 生产者在席)=====
+ // 只对 stale queued「重发信号」交消费者重走认领+派发,不在 tick 线程直接派发——这是让 MQ 保持唯一投递
+ // 路径的投递可靠性补偿(信号丢了补发一条、回到同一条消费→派发链路,不另起第二条并行派发路径);
+ // 与并发 cap 无关(真「在跑生成并发≤15」属有界 worker 池=follow-up,不由消费线程数或此重发决定)。
+ if (genTaskProducer != null) {
+ if (task.getCreateTime() != null && task.getCreateTime().isAfter(staleQueuedBefore)) {
+ continue; // 新鲜 queued(信号大概率在途),本轮不重发,减噪
+ }
+ // best-effort 重发(生产者内部吞异常不外抛);消费侧 CAS 幂等去重,重复信号无害。
+ genTaskProducer.publishNewTask(task.getId());
+ resignaled++;
+ log.info("[executor-tick] 兜底补偿:stale queued 重发 gen 信号(交消费者认领,不在 tick 直派)taskId={}, traceId={}",
+ task.getId(), task.getTraceId());
+ continue;
+ }
+ // ===== 无 MQ 生产者(单测/降级态):退回既有「直接认领+派发」(行为逐字不变,兜底不让任务卡死)=====
try {
// 系统身份 + 所属租户上下文内处理整任务(含 CAS 认领与回调写链——三表行必须落对租户,§5.1-3b;
// 系统身份保证审计列 updater 可填充,否则 UPDATE null 直出撞 NOT NULL,部署窗口#3 根因修复)
@@ -249,16 +306,78 @@ public class AigcGenerateExecutor {
task.getId(), task.getTraceId(), e);
}
}
- // ===== 步骤④:stale watchdog(进程崩溃恢复:上次认领后死掉的任务收尸为 failed(timeout))=====
+ // ===== 步骤④:stale watchdog(进程崩溃恢复 / 派发丢状态遗留:RUNNING 超阈值收尸为 failed(timeout))=====
+ // 这即 C3 兜底补偿的第②类「stale running 僵死任务」——认领成功后进程崩溃/派发失败未落状态的漏网,超时收尸。
int reaped = watchdogSweep();
- // tick 级汇总行(§5.5):有活动打 info;空转降为 debug 防 10s 级刷屏
- if (claimed + reaped > 0) {
- log.info("[executor-tick] tick 汇总:认领 {} / 成功 {} / 失败 {} / 收尸 {}", claimed, succeeded, failed, reaped);
+ // tick 级汇总行(§5.5):有活动打 info;空转降为 debug 防刷屏
+ if (claimed + reaped + resignaled > 0) {
+ log.info("[executor-tick] tick 汇总(兜底补偿):认领 {} / 成功 {} / 失败 {} / 重发信号 {} / 收尸 {}",
+ claimed, succeeded, failed, resignaled, reaped);
} else {
- log.debug("[executor-tick] tick 汇总:认领 0 / 成功 0 / 失败 0 / 收尸 0(空转)");
+ log.debug("[executor-tick] tick 汇总:认领 0 / 成功 0 / 失败 0 / 重发信号 0 / 收尸 0(空转)");
}
}
+ /**
+ * MQ 事件驱动的单任务认领+派发(切片一 阶段〇 C3)——RocketMQ 消费者 {@link com.wanxiang.huijing.game.module.aigc.mq.GenTaskConsumer}
+ * 收到 taskId 信号后经 {@link com.wanxiang.huijing.game.module.aigc.mq.AigcGenDispatchService} 调用本方法。与 @Scheduled tick 的
+ * {@link #processCandidate} 同一执行体(CAS 认领 + {@link #runGenerationPipeline} 派发)、同一系统身份 + 租户上下文,
+ * 只是触发面从「批量扫描」换成「单 taskId 事件」。
+ *
+ * 幂等(CAS 去重):既有 {@code claimQueuedTask} 仅 queued→running 才成功;认领成功才派发,认领失败
+ * (被并发消费线程/兜底 tick 抢走、非 queued、已取消、任务不存在)即幂等丢弃返 false——同一 taskId 重复投递不双跑。
+ *
+ * 派发失败落 failed(不留悬空 running,Opus 评审):CAS 成功后 {@link #runGenerationPipeline} 抛异常
+ * (派发/渲染/DB 抖动等)→ catch 里经 {@link #callbackFailed} 落 failed(llm_error 归因,与 dispatchGeneric 投递
+ * 失败同桶),绝不留悬空 running 等 watchdog。若 pipeline 内部已自行落终态后再抛,二次回调被终态拒重入(M7)幂等吞掉,安全。
+ *
+ * 就绪门:key/prompt/schema 未就绪({@link #isExecutorReady()})则不认领、返 false——任务留 queued,
+ * 由 @Scheduled 兜底稍后重发信号重试(不烧一次认领把任务判失败)。
+ *
+ * @param taskId 任务 ID(消息体轻信号)
+ * @return true=CAS 认领成功(已派发或已落终态);false=认领未命中/未就绪=幂等丢弃(不重投)
+ */
+ public boolean claimAndDispatchById(Long taskId) {
+ // 消费线程无租户上下文 → executeIgnore 跨租户按 id 取任务(与 tick 扫描同款绕租户过滤)。
+ AigcTaskDO task = TenantUtils.executeIgnore(() -> aigcTaskMapper.selectById(taskId));
+ if (task == null) {
+ // 任务不存在(脏信号/已被物理清理):幂等丢弃,不重投。
+ log.warn("[executor-mq] 任务不存在,幂等丢弃 taskId={}", taskId);
+ return false;
+ }
+ // 就绪门(只读、不变更状态,多消费线程并发安全;不复用 selfCheckPass 以免竞争其节流计数):
+ // 未就绪则不认领,任务留 queued 交兜底轮询稍后重发信号重试。
+ if (!isExecutorReady()) {
+ log.warn("[executor-mq] 执行器未就绪(api-key/prompt/schema),暂不认领(兜底轮询稍后重试)taskId={}, traceId={}",
+ task.getId(), task.getTraceId());
+ return false;
+ }
+ // 系统身份 + 所属租户上下文内:CAS 认领 → 复用生成流水派发;认领失败幂等丢弃,派发抛异常落 failed(不留悬空 running)。
+ Boolean claimed = executeWithSystemIdentity(task, () -> {
+ if (!aigcTaskMapper.claimQueuedTask(task.getId())) {
+ // CAS 未命中:被并发消费线程/兜底 tick 抢走,或已取消/非 queued —— 幂等丢弃。
+ log.info("[executor-mq] CAS 认领未命中(已被认领/非 queued),幂等丢弃 taskId={}, traceId={}",
+ task.getId(), task.getTraceId());
+ return Boolean.FALSE;
+ }
+ long queuedSeconds = task.getCreateTime() == null ? -1
+ : Duration.between(task.getCreateTime(), LocalDateTime.now()).getSeconds();
+ log.info("[executor-mq] 任务认领成功(MQ 触发)taskId={}, traceId={}, templateId={}, 排队时长={}s",
+ task.getId(), task.getTraceId(), task.getTemplateId(), queuedSeconds);
+ try {
+ // 复用既有生成流水(与 tick processCandidate 同一 runGenerationPipeline:generic 走派发面 SaaGraphDispatcher/http worker)。
+ runGenerationPipeline(task);
+ } catch (Exception e) {
+ // 认领成功但派发抛异常:落 failed(llm_error,与 dispatchGeneric 投递失败同桶),绝不留悬空 running(Opus 评审)。
+ log.error("[executor-mq] 认领成功但派发异常,落 failed(不留悬空 running)taskId={}, traceId={}",
+ task.getId(), task.getTraceId(), e);
+ callbackFailed(task, FailureReasonEnum.LLM_ERROR);
+ }
+ return Boolean.TRUE;
+ });
+ return Boolean.TRUE.equals(claimed);
+ }
+
/** 单测可观测口径:自禁用告警已发条数(验证 §8.3 节流正确性) */
public int getSelfCheckAlertCount() {
return selfCheckAlertCount;
@@ -294,6 +413,19 @@ public class AigcGenerateExecutor {
return false;
}
+ /**
+ * 就绪只读检查(C3 MQ 路专用):key 非空 + prompt/schema classpath 资源就位。
+ * 与 {@link #selfCheckPass()} 判据同源,但不变更任何状态(不碰 selfCheckFailTicks 节流计数)——
+ * 供多消费线程(consumeThreadMax=15)并发安全调用(selfCheckPass 的节流计数是 tick 单线程专用,多线程并发会数据竞争)。
+ *
+ * @return true=就绪可认领
+ */
+ private boolean isExecutorReady() {
+ return StringUtils.hasText(properties.getApiKey())
+ && promptResourceLoader.isReady()
+ && schemaValidator.isReady();
+ }
+
// ============================== 系统身份任务边界(部署窗口#3 根因修复) ==============================
/** 执行器系统用户编号:审计列 updater/creator 统一落 "0"(DefaultDBFieldHandler 以 getLoginUserId().toString() 填充) */
diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/task/AigcTaskServiceImpl.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/task/AigcTaskServiceImpl.java
index 7e1ac3e9..e8340ceb 100644
--- a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/task/AigcTaskServiceImpl.java
+++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/task/AigcTaskServiceImpl.java
@@ -9,6 +9,7 @@ import com.wanxiang.huijing.game.module.aigc.dal.mysql.task.AigcTaskMapper;
import com.wanxiang.huijing.game.module.aigc.enums.AigcTaskSourceEnum;
import com.wanxiang.huijing.game.module.aigc.enums.AigcTaskStatusEnum;
import com.wanxiang.huijing.game.module.aigc.enums.FailureReasonEnum;
+import com.wanxiang.huijing.game.module.aigc.mq.GenTaskProducer;
import com.wanxiang.huijing.game.module.aigc.service.executor.AigcTemplateConstants;
import com.wanxiang.huijing.framework.common.pojo.PageResult;
import com.wanxiang.huijing.module.infra.api.config.ConfigApi;
@@ -81,11 +82,20 @@ public class AigcTaskServiceImpl implements AigcTaskService {
* 入口 QPS 突发闸(Sentinel FLOW_GRADE_QPS,规则热源 Nacos;切片一 阶段〇 C2)。
* 与上方控制平面软注入件不同:GenAdmissionResource 是无条件 @Component(恒在席),故直接 @Resource 硬注入
* (同 aigcTaskMapper/playerApi),不走 ObjectProvider。职责 = 入队入口速率闸削洪峰,与 DB 三门(业务配额)
- * / C3 consumeThreadMax(在跑并发≤15)三层正交,详见 {@link GenAdmissionResource} 类注释。
+ * / C3 consumeThreadMax(dispatch 投递速率上限;真在跑并发≤15 属有界 worker 池=follow-up)三层正交,详见 {@link GenAdmissionResource} 类注释。
*/
@Resource
private com.wanxiang.huijing.game.module.aigc.admission.GenAdmissionResource genAdmissionResource;
+ /**
+ * gen 队列生产者软注入(切片一 阶段〇 C3):入队成功后发 AIGC_GEN_TOPIC 信号,触发有界消费者(consumeThreadMax=15)
+ * 据 taskId CAS 认领拉起生成——把「触发」从执行器 @Scheduled DB 轮询改为 MQ 事件驱动。
+ * 与执行器/消费者同 {@code aigc.executor.enabled} 条件装配 → 用 {@link ObjectProvider} 软注入:enabled=false
+ * (无执行器/无消费者)时缺席,入队照常落库、仅不发信号(现行骨架行为,任务留 queued)。best-effort:发信号失败不回滚入队。
+ */
+ @Resource
+ private ObjectProvider 须在入队落库后调用({@code task.getId()} 已由 insert 回填)。
+ *
+ * @param task 已落库的 queued 任务(id 已回填)
+ */
+ private void publishGenSignal(AigcTaskDO task) {
+ GenTaskProducer producer = genTaskProducerProvider.getIfAvailable();
+ if (producer == null) {
+ // 执行器/消费者未装配(aigc.executor.enabled=false):无消费方,发信号无意义,跳过;任务留 queued(现行骨架行为)。
+ log.debug("[aigc-mq] gen 队列生产者缺席(executor 未启用),跳过发信号 taskId={}", task.getId());
+ return;
+ }
+ // best-effort:publishNewTask 内部已吞异常(发失败只告警、不外抛、不回滚入队),执行器兜底轮询补偿丢失信号。
+ producer.publishNewTask(task.getId());
+ }
+
/**
* 降级门读取(执行版 §6.1):读 infra ConfigApi 键 aigc.generate.paused。
* fail-open on switch read:ConfigApi 缺席 / 读不到值 / 非 "true" / 异常 → 一律视为「未暂停」,
diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/test/java/com/wanxiang/huijing/game/module/aigc/mq/GenTaskConsumerTest.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/test/java/com/wanxiang/huijing/game/module/aigc/mq/GenTaskConsumerTest.java
new file mode 100644
index 00000000..3e9e5e1e
--- /dev/null
+++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/test/java/com/wanxiang/huijing/game/module/aigc/mq/GenTaskConsumerTest.java
@@ -0,0 +1,53 @@
+package com.wanxiang.huijing.game.module.aigc.mq;
+
+import com.wanxiang.huijing.framework.test.core.ut.BaseMockitoUnitTest;
+import org.junit.jupiter.api.Test;
+import org.mockito.InjectMocks;
+import org.mockito.Mock;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * {@link GenTaskConsumer} 单元测试(纯 Mockito,不触真执行器/DB)——把守 C3 消费者「幂等丢弃 + 不重投」契约。
+ *
+ * 覆盖:① {@code claimAndDispatch} 返 false(已被并发消费线程/兜底 tick 认领、非 queued)→ {@link GenTaskConsumer#onMessage}
+ * 正常返回、不抛(幂等丢弃,天然去重不双跑);② 返 true(认领成功)→ onMessage 不抛、taskId 正确解析下派;
+ * ③ 毒消息(taskId 非数字)→ onMessage 不抛且不调 dispatch(丢弃不重投,防重投风暴)。
+ *
+ * @author 绘境AI
+ */
+class GenTaskConsumerTest extends BaseMockitoUnitTest {
+
+ @Mock
+ private AigcGenDispatchService dispatchService;
+
+ @InjectMocks
+ private GenTaskConsumer consumer;
+
+ /** 幂等丢弃:claimAndDispatch 返 false → onMessage 正常返回、不抛(同 taskId 重复投递不双跑)。 */
+ @Test
+ void testOnMessage_claimMiss_idempotentDiscardNoThrow() {
+ when(dispatchService.claimAndDispatch(789L)).thenReturn(false);
+ assertDoesNotThrow(() -> consumer.onMessage("789"));
+ verify(dispatchService).claimAndDispatch(789L);
+ }
+
+ /** 认领成功:claimAndDispatch 返 true → onMessage 不抛,taskId 正确解析下派。 */
+ @Test
+ void testOnMessage_claimed_dispatchInvokedNoThrow() {
+ when(dispatchService.claimAndDispatch(100L)).thenReturn(true);
+ assertDoesNotThrow(() -> consumer.onMessage("100"));
+ verify(dispatchService).claimAndDispatch(100L);
+ }
+
+ /** 毒消息:taskId 非数字 → onMessage 不抛且不调 dispatch(丢弃不重投,防无限重投风暴)。 */
+ @Test
+ void testOnMessage_malformedTaskId_discardNoDispatchNoThrow() {
+ assertDoesNotThrow(() -> consumer.onMessage("not-a-number"));
+ verify(dispatchService, never()).claimAndDispatch(anyLong());
+ }
+}
diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/test/java/com/wanxiang/huijing/game/module/aigc/mq/GenTaskProducerTest.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/test/java/com/wanxiang/huijing/game/module/aigc/mq/GenTaskProducerTest.java
new file mode 100644
index 00000000..eaf1a9b2
--- /dev/null
+++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/test/java/com/wanxiang/huijing/game/module/aigc/mq/GenTaskProducerTest.java
@@ -0,0 +1,51 @@
+package com.wanxiang.huijing.game.module.aigc.mq;
+
+import com.wanxiang.huijing.framework.test.core.ut.BaseMockitoUnitTest;
+import org.apache.rocketmq.spring.core.RocketMQTemplate;
+import org.junit.jupiter.api.Test;
+import org.mockito.InjectMocks;
+import org.mockito.Mock;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.verify;
+
+/**
+ * {@link GenTaskProducer} 单元测试(纯 Mockito,不触真 MQ)——把守 C3 生产者 best-effort 契约。
+ *
+ * 覆盖:① 正常发 AIGC_GEN_TOPIC 信号、payload = taskId 字符串;② {@code syncSend} 抛异常时
+ * {@link GenTaskProducer#publishNewTask} 不外抛(best-effort:不回滚入队、由执行器兜底轮询补偿)。
+ *
+ * 注:{@code RocketMQTemplate.syncSend} 有 (String,Object) 与 (String,Message) 两个 2 参重载;本测用
+ * String payload 匹配器({@code eq("123")}/{@code anyString()})强制解析到 (String,Object) 重载,与生产代码
+ * {@code syncSend(TOPIC, String.valueOf(taskId))} 同一重载。
+ *
+ * @author 绘境AI
+ */
+class GenTaskProducerTest extends BaseMockitoUnitTest {
+
+ @Mock
+ private RocketMQTemplate rocketMQTemplate;
+
+ @InjectMocks
+ private GenTaskProducer producer;
+
+ /** 正常路:发到正确 topic,payload = taskId 字符串。 */
+ @Test
+ void testPublishNewTask_sendsSignalToTopicWithTaskId() {
+ assertDoesNotThrow(() -> producer.publishNewTask(123L));
+ // 验证发到 AIGC_GEN_TOPIC + payload="123"(未 stub 返回值=默认 null,publishNewTask 的 result!=null 守卫吞掉,不 NPE)
+ verify(rocketMQTemplate).syncSend(eq(GenTaskProducer.TOPIC), eq("123"));
+ }
+
+ /** best-effort(核心断言):syncSend 抛异常 → publishNewTask 吞异常不外抛(否则会连累入队事务/接口)。 */
+ @Test
+ void testPublishNewTask_bestEffort_swallowsSendException() {
+ doThrow(new RuntimeException("MQ down")).when(rocketMQTemplate).syncSend(anyString(), anyString());
+ // 关键:MQ 抖动不外抛(任务已落库 queued,丢失信号由执行器 @Scheduled 兜底轮询补偿)
+ assertDoesNotThrow(() -> producer.publishNewTask(456L));
+ verify(rocketMQTemplate).syncSend(eq(GenTaskProducer.TOPIC), eq("456"));
+ }
+}
diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/test/java/com/wanxiang/huijing/game/module/aigc/service/executor/AigcGenerateExecutorTest.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/test/java/com/wanxiang/huijing/game/module/aigc/service/executor/AigcGenerateExecutorTest.java
index 9adfe3b5..34c96464 100644
--- a/game-cloud/game-module-aigc/game-module-aigc-server/src/test/java/com/wanxiang/huijing/game/module/aigc/service/executor/AigcGenerateExecutorTest.java
+++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/test/java/com/wanxiang/huijing/game/module/aigc/service/executor/AigcGenerateExecutorTest.java
@@ -4,6 +4,7 @@ import com.wanxiang.huijing.game.module.aigc.controller.admin.task.vo.DifyCallba
import com.wanxiang.huijing.game.module.aigc.dal.dataobject.task.AigcTaskDO;
import com.wanxiang.huijing.game.module.aigc.dal.mysql.task.AigcTaskMapper;
import com.wanxiang.huijing.game.module.aigc.framework.executor.config.AigcExecutorConfiguration;
+import com.wanxiang.huijing.game.module.aigc.mq.GenTaskProducer;
import com.wanxiang.huijing.game.module.aigc.service.callback.DifyCallbackService;
import com.wanxiang.huijing.game.module.studio.api.SourceProjectApi;
import com.wanxiang.huijing.game.module.studio.dto.SourceProjectFetchRespDTO;
@@ -39,6 +40,7 @@ import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.mock;
@@ -832,4 +834,77 @@ class AigcGenerateExecutorTest extends BaseMockitoUnitTest {
verifyNoInteractions(sourceProjectApi);
}
+ // ============================== 用例27-29:C3 MQ 事件驱动认领+派发 & 兜底重发信号 ==============================
+
+ /**
+ * C3-A 幂等({@link AigcGenerateExecutor#claimAndDispatchById}):CAS 认领未命中(被并发消费线程/兜底 tick 抢走、
+ * 非 queued)→ 返 false,且不派发、不回调(同 taskId 重复投递不双跑=天然去重)。
+ */
+ @Test
+ void testClaimAndDispatchById_casMiss_idempotentNoDispatch() {
+ AigcTaskDO task = queuedTask(901L, "aigc-mq-miss");
+ when(aigcTaskMapper.selectById(901L)).thenReturn(task);
+ when(aigcTaskMapper.claimQueuedTask(901L)).thenReturn(false); // CAS 未命中(已被认领/非 queued)
+
+ boolean claimed = newExecutor(readyProps(), Clock.systemDefaultZone()).claimAndDispatchById(901L);
+
+ assertFalse(claimed); // 幂等丢弃
+ verifyNoInteractions(difyCallbackService); // 未派发、未回调(无终态写入)
+ }
+
+ /**
+ * C3-B 派发失败落 failed 不留悬空 running(Opus 评审命门):CAS 认领成功后派发通道抛异常 →
+ * {@link AigcGenerateExecutor#claimAndDispatchById} catch 里经唯一写入路径落 failed(llm_error),绝不留悬空 running。
+ */
+ @Test
+ void testClaimAndDispatchById_dispatchThrows_marksFailedNoOrphanRunning() {
+ AigcExecutorProperties props = genericProps();
+ props.setDispatcher("saa"); // 选进程内派发通道 = 下方注入的抛异常 dispatcher
+ ExecutorLlmClient llmClient = new ExecutorLlmClient(props, sender, millis -> { });
+ WorkerDispatchClient wdc = new WorkerDispatchClient(props, dispatchSender);
+ GenerationDispatcher boom = mock(GenerationDispatcher.class);
+ when(boom.dispatch(any())).thenThrow(new RuntimeException("dispatch boom")); // 派发抛异常(非返 false 的前置失败)
+ // 11 参构造:saaDispatcher=boom + dispatcher=saa → generationDispatcher=boom;producer=null(走 MQ 单任务路,非 tick)。
+ AigcGenerateExecutor exec = new AigcGenerateExecutor(props, aigcTaskMapper, difyCallbackService,
+ LOADER, VALIDATOR, llmClient, wdc, boom, null, Clock.systemDefaultZone(), null);
+ AigcTaskDO task = genericQueuedTask(902L, "aigc-mq-throw");
+ when(aigcTaskMapper.selectById(902L)).thenReturn(task);
+ when(aigcTaskMapper.claimQueuedTask(902L)).thenReturn(true); // CAS 认领成功
+
+ boolean claimed = exec.claimAndDispatchById(902L);
+
+ assertTrue(claimed); // 已认领(虽派发失败),不重投
+ // 关键:落 failed(llm_error) 走唯一写入路径 handleCallback,绝不留悬空 running 等 watchdog
+ ArgumentCaptor