diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/framework/metrics/GenMetrics.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/framework/metrics/GenMetrics.java new file mode 100644 index 00000000..f91bb3ec --- /dev/null +++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/framework/metrics/GenMetrics.java @@ -0,0 +1,257 @@ +package com.wanxiang.huijing.game.module.aigc.framework.metrics; + +import com.wanxiang.huijing.game.module.aigc.controller.admin.task.vo.DifyCallbackReqVO; +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.AigcTaskStatusEnum; +import com.wanxiang.huijing.game.module.aigc.enums.FailureReasonEnum; +import com.wanxiang.huijing.framework.tenant.core.util.TenantUtils; +import io.micrometer.core.instrument.DistributionSummary; +import io.micrometer.core.instrument.Gauge; +import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.Timer; +import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; + +import java.util.Map; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * 生成业务语义埋点(观测体系阶段四 · 面五) + * + *
把 MVP 真正要盯的几个生成指标以 micrometer 形式暴露给 {@code /actuator/prometheus}:生成成功率、单次耗时、 + * 单次成本、九门逐门失败、在飞队列深度。这些值此前只落在产物目录的 json 或回调的 trace 字段里,没有一处能实时看、 + * 能告警——本组件是把它们变成可采集指标的收口。 + * + *
开关与装配(旁路铁律):整组埋点由 {@code huijing.aigc.metrics.enabled} 控制,默认关—— + * {@link ConditionalOnProperty} 保证关时本组件根本不装配(无 Bean、无 @Scheduled、无 meter 注册), + * 对生成主链零字节影响。生产 profile 显式置 true 才开。开启的前提是「面一」已给单体引入 actuator + micrometer、 + * 提供 {@link MeterRegistry} bean;本组件只消费该 bean、不负责装 actuator。故未做面一的环境保持开关关即可。 + * + *
埋点绝不咬主链(硬红线):{@link #recordTerminal} 全程吞异常,调用方({@code DifyCallbackServiceImpl} + * 的提交后 best-effort 段)再兜一层,双保险确保任何埋点故障都只留痕、绝不外抛去影响回调的成败判定或后续通知/回填。 + * + *
Micrometer 点名 → Prometheus 名(PrometheusMeterRegistry 自动转 snake_case 并补后缀): + *
覆盖全部终态——成功、业务失败、以及超时收尸(stale RUNNING 收尸为 failed+failureReason=timeout)均经 + * 回调唯一写入路径到达提交后段,故此处一个挂点即全覆盖,无需另挂 watchdog。 + * + *
全程吞异常:埋点是纯旁路,任何解析/注册异常只留痕、绝不外抛(旁路铁律)。
+ *
+ * @param reqVO 回调入参(取 trace 载荷:wallS / cost.totalRmb / sevenGateVerdict.guards)
+ * @param task 提交后按 traceId 回查到的任务(取终态 status + failureReason;可能为 null=未定位到)
+ */
+ public void recordTerminal(DifyCallbackReqVO reqVO, AigcTaskDO task) {
+ try {
+ recordTaskTerminal(task);
+ Map guards 在便宜档主路与 tier2 两条生成线都逐门产出(result_out / service._extract_trace 同源),故本项可做到真逐门,
+ * 不必降级为合成分。门判真值解析复用 {@code ReadinessScorer.guardPass} 同款口径:H_progress 是嵌套对象
+ * {@code {pass,checks,latch,...}} 取其 .pass,其余门为裸 bool。只有明确「未过」的门才 +1。
+ */
+ @SuppressWarnings("unchecked")
+ private void recordGateFails(Map 为何不把 Gauge 直接绑 {@code countGlobalInflight}:直接绑会让每次 Prometheus scrape(以及 micrometer
+ * 内部 gauge 轮询)都打一次全表 COUNT,DB 压力与 scrape 频率耦合、不可控。改为定时刷进 {@link #queueDepth}
+ * AtomicInteger,把查询频率钉死在固定周期,scrape 时只读内存零 DB——符合「勿每次 scrape 打全表重查」。
+ * 复用既有控制平面计数 {@code countGlobalInflight}(D12 门②背压同一口径),不新增查询。
+ *
+ * 跨租户:定时线程无租户上下文,须 {@code TenantUtils.executeIgnore} 包裹取全局在飞口径
+ * (与执行器 watchdog 收尸扫描同款)。best-effort:刷新失败只留痕、保留上次值,绝不外抛。
+ */
+ @Scheduled(fixedRateString = "${huijing.aigc.metrics.queue-depth-refresh-ms:15000}")
+ public void refreshQueueDepth() {
+ try {
+ // 用 lambda(非方法引用)对齐执行器既有 executeIgnore(Callable) 用法,消除 Runnable/Callable 重载歧义。
+ Long inflight = TenantUtils.executeIgnore(() -> aigcTaskMapper.countGlobalInflight());
+ queueDepth.set(inflight == null ? 0 : inflight.intValue());
+ } catch (Exception e) {
+ log.warn("[GenMetrics] 在飞队列深度刷新失败(保留上次值,不阻断)", e);
+ }
+ }
+
+}
diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/callback/DifyCallbackServiceImpl.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/callback/DifyCallbackServiceImpl.java
index 8b44a228..c9700036 100644
--- a/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/callback/DifyCallbackServiceImpl.java
+++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/main/java/com/wanxiang/huijing/game/module/aigc/service/callback/DifyCallbackServiceImpl.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.AigcTaskStatusEnum;
+import com.wanxiang.huijing.game.module.aigc.framework.metrics.GenMetrics;
import com.wanxiang.huijing.game.module.community.api.CommunityNotifyApi;
import com.wanxiang.huijing.game.module.studio.api.SourceProjectApi;
import com.wanxiang.huijing.game.module.studio.dto.SourceProjectLandReqDTO;
@@ -13,6 +14,7 @@ import com.wanxiang.huijing.framework.common.exception.enums.GlobalErrorCodeCons
import com.wanxiang.huijing.framework.common.pojo.CommonResult;
import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
@@ -85,6 +87,16 @@ public class DifyCallbackServiceImpl implements DifyCallbackService {
@Resource
private SourceProjectApi sourceProjectApi;
+ /**
+ * 生成业务埋点组件(观测阶段四 · 面五;可选注入)。
+ *
+ * 可选(required=false)的原因:埋点由开关 {@code huijing.aigc.metrics.enabled} 控制,默认关——关时
+ * {@link GenMetrics} 整个不装配({@code @ConditionalOnProperty}),此处注入为 null。故本类在提交后段调用前必判空
+ * ({@link #recordGenMetricsQuietly}),关时天然旁路、零埋点副作用(旁路铁律)。开时(生产 profile)才在席。
+ */
+ @Autowired(required = false)
+ private GenMetrics genMetrics;
+
@Override
public Boolean handleCallback(DifyCallbackReqVO reqVO) {
log.info("[handleCallback] 收到 Dify 回调 traceId={}, status={}, templateId={}",
@@ -144,6 +156,10 @@ public class DifyCallbackServiceImpl implements DifyCallbackService {
private void postCommitBestEffort(DifyCallbackReqVO reqVO, SourceLanding landing) {
try {
AigcTaskDO task = aigcTaskMapper.selectByTraceId(reqVO.getTraceId());
+ // ===== 面五 观测埋点(终态收口):post-commit best-effort,独立兜底、与下方通知/回填解耦 =====
+ // 置于 realSuccess 判定之前,使成功/失败/超时收尸(failed+timeout)三类终态都被计入;
+ // genMetrics 关时为 null → 天然旁路;开时其内部再吞一层异常,双保险确保埋点绝不咬主链、不跳过下方源态回填。
+ recordGenMetricsQuietly(reqVO, task);
boolean realSuccess = task != null
&& Objects.equals(task.getStatus(), AigcTaskStatusEnum.SUCCEEDED.getStatus())
&& task.getVersionId() != null;
@@ -165,6 +181,28 @@ public class DifyCallbackServiceImpl implements DifyCallbackService {
}
}
+ /**
+ * 面五 观测埋点静默调用(提交后 best-effort):把终态计数/耗时/成本/门失败交给 {@link GenMetrics}。
+ *
+ * 双保险不咬主链:① 开关关时 {@code genMetrics} 为 null → 直接旁路,零副作用;② 开时即便 {@code recordTerminal}
+ * 意外抛错,本方法 catch 吞掉——若不兜,异常会冒进 {@code postCommitBestEffort} 外层 catch、连带跳过其后的通知与
+ * 源态回填(那才是真正咬主链)。故此处必须独立吞异常,不能省。
+ *
+ * @param reqVO 回调入参(取 trace 载荷)
+ * @param task 提交后回查到的任务(取终态 status/failureReason;可能为 null)
+ */
+ private void recordGenMetricsQuietly(DifyCallbackReqVO reqVO, AigcTaskDO task) {
+ if (genMetrics == null) {
+ return; // 埋点开关关 / 组件未装配 → 零埋点副作用(旁路铁律)
+ }
+ try {
+ genMetrics.recordTerminal(reqVO, task);
+ } catch (Exception e) {
+ log.error("[handleCallback] 生成业务埋点失败(best-effort,不影响主链、不跳过源态回填)traceId={}",
+ reqVO.getTraceId(), e);
+ }
+ }
+
/**
* 本次落源结果(U2/B8 §5.2):持源行 ID + gameId + sourceHash,供「建包成功回填 version_id」「建包失败标孤儿」复用。
* landing 为 null(reqVO 无 sourceProject / 落源失败)→ 后续两步均旁路(best-effort,不阻断主链)。
diff --git a/game-cloud/game-module-aigc/game-module-aigc-server/src/test/java/com/wanxiang/huijing/game/module/aigc/framework/metrics/GenMetricsTest.java b/game-cloud/game-module-aigc/game-module-aigc-server/src/test/java/com/wanxiang/huijing/game/module/aigc/framework/metrics/GenMetricsTest.java
new file mode 100644
index 00000000..bf25a00d
--- /dev/null
+++ b/game-cloud/game-module-aigc/game-module-aigc-server/src/test/java/com/wanxiang/huijing/game/module/aigc/framework/metrics/GenMetricsTest.java
@@ -0,0 +1,260 @@
+package com.wanxiang.huijing.game.module.aigc.framework.metrics;
+
+import com.wanxiang.huijing.game.module.aigc.controller.admin.task.vo.DifyCallbackReqVO;
+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.AigcTaskStatusEnum;
+import com.wanxiang.huijing.game.module.aigc.service.callback.DifyCallbackServiceImpl;
+import com.wanxiang.huijing.game.module.aigc.service.callback.DifyCallbackTxService;
+import com.wanxiang.huijing.game.module.community.api.CommunityNotifyApi;
+import io.micrometer.core.instrument.MeterRegistry;
+import io.micrometer.core.instrument.simple.SimpleMeterRegistry;
+import org.junit.jupiter.api.Test;
+import org.springframework.boot.test.context.runner.ApplicationContextRunner;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import java.util.LinkedHashMap;
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * {@link GenMetrics}(观测阶段四 · 面五 生成业务埋点)单元测试(纯 Mockito + SimpleMeterRegistry,不依赖 DB/Spring 全量上下文)
+ *
+ * 覆盖三条硬要求:① 开关开时终态各 status 正确递增 gen_task_total(含超时还原 timed_out)+ 耗时/成本/门失败真落表;
+ * ② 开关关时组件不装配(ApplicationContextRunner 断言 Bean 缺席);③ 埋点内部抛错不影响回调主流程返回与源态回填。
+ *
+ * @author 绘境AI(观测阶段四 · 面五)
+ */
+class GenMetricsTest {
+
+ // ============================== 用例①-a:gen_task_total{status} 按终态递增 ==============================
+
+ /**
+ * 三类终态各自递增对应 status 标签:succeeded / failed(非超时) / timed_out(failed+failureReason=timeout 还原)。
+ * 坐实成功率分母口径 = succeeded + failed + timed_out,且超时据真实 failureReason 还原、非编造。
+ */
+ @Test
+ void testGenTaskTotal_incrementsByTerminalStatus() {
+ SimpleMeterRegistry registry = new SimpleMeterRegistry();
+ GenMetrics metrics = new GenMetrics(registry, mock(AigcTaskMapper.class));
+
+ metrics.recordTerminal(reqWithTrace(true, 1), taskOf(AigcTaskStatusEnum.SUCCEEDED, null));
+ metrics.recordTerminal(reqWithTrace(false, 3), taskOf(AigcTaskStatusEnum.FAILED, "llm_error"));
+ metrics.recordTerminal(reqWithTrace(false, 3), taskOf(AigcTaskStatusEnum.FAILED, "timeout"));
+
+ assertEquals1(registry.get(GenMetrics.M_TASK).tag("status", "succeeded").counter().count());
+ assertEquals1(registry.get(GenMetrics.M_TASK).tag("status", "failed").counter().count());
+ // 超时:status 落 failed(3)、failureReason=timeout → 还原为 timed_out 标签
+ assertEquals1(registry.get(GenMetrics.M_TASK).tag("status", "timed_out").counter().count());
+ }
+
+ // ============================== 用例①-b:gen_duration_seconds / llm_cost 真落表 ==============================
+
+ /**
+ * 单次终态回调 → gen_duration_seconds 记入一笔(≈trace.wallS 秒)、llm_cost 记入一笔(≈trace.cost.totalRmb)。
+ */
+ @Test
+ void testDurationAndCost_recordedFromTrace() {
+ SimpleMeterRegistry registry = new SimpleMeterRegistry();
+ GenMetrics metrics = new GenMetrics(registry, mock(AigcTaskMapper.class));
+
+ metrics.recordTerminal(reqWithTrace(true, 1), taskOf(AigcTaskStatusEnum.SUCCEEDED, null));
+
+ // 耗时:count=1,总时长≈wallS(23.4s)
+ assertTrue(registry.get(GenMetrics.M_DURATION).timer().count() == 1L, "gen_duration 应记入一笔");
+ double sec = registry.get(GenMetrics.M_DURATION).timer().totalTime(TimeUnit.SECONDS);
+ assertTrue(Math.abs(sec - 23.4) < 0.1, "gen_duration 总时长应≈wallS=23.4s,实际=" + sec);
+
+ // 成本:count=1,总额≈totalRmb(0.031)
+ assertTrue(registry.get(GenMetrics.M_LLM_COST).summary().count() == 1L, "llm_cost 应记入一笔");
+ double rmb = registry.get(GenMetrics.M_LLM_COST).summary().totalAmount();
+ assertTrue(Math.abs(rmb - 0.031) < 1e-6, "llm_cost 总额应≈totalRmb=0.031,实际=" + rmb);
+ }
+
+ // ============================== 用例①-c:gen_gate_fail_total{gate} 逐门失败 ==============================
+
+ /**
+ * guards = {C_frame:true(过), H_progress:{pass:false}(嵌套·未过), E_live:false(裸·未过)}
+ * → 仅两个未过门各 +1;过门不计;嵌套 H_progress 取 .pass 判定(复用 ReadinessScorer 口径)。
+ */
+ @Test
+ void testGenGateFail_perFailingGate() {
+ SimpleMeterRegistry registry = new SimpleMeterRegistry();
+ GenMetrics metrics = new GenMetrics(registry, mock(AigcTaskMapper.class));
+
+ DifyCallbackReqVO reqVO = new DifyCallbackReqVO();
+ reqVO.setTraceId("t-gate");
+ Map
+ *
+ */
+ private String resolveStatusLabel(AigcTaskDO task) {
+ AigcTaskStatusEnum st = AigcTaskStatusEnum.of(task.getStatus());
+ if (st == null) {
+ return null; // 非法状态值:不计(避免脏标签污染指标基数)
+ }
+ if (st == AigcTaskStatusEnum.FAILED
+ && FailureReasonEnum.TIMEOUT.getReason().equals(task.getFailureReason())) {
+ return "timed_out";
+ }
+ return st.name().toLowerCase();
+ }
+
+ /** gen_duration_seconds:trace.wallS(秒,worker 已算)→ Timer.record。缺失/非数/负值跳过。 */
+ private void recordDuration(Map