diff --git a/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/cn/wanxiang/game/module/telemetry/service/event/EventIngestServiceImpl.java b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/cn/wanxiang/game/module/telemetry/service/event/EventIngestServiceImpl.java index 85d06ab1..e2b22773 100644 --- a/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/cn/wanxiang/game/module/telemetry/service/event/EventIngestServiceImpl.java +++ b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/cn/wanxiang/game/module/telemetry/service/event/EventIngestServiceImpl.java @@ -12,10 +12,14 @@ import cn.wanxiang.game.module.telemetry.enums.TelemetryEventEnum; import cn.wanxiang.game.module.telemetry.service.quality.AnomalyFilterService; import cn.wanxiang.game.module.feed.api.FeedApi; import cn.wanxiang.game.module.feed.dto.FeedRankUpsertReqDTO; +import cn.iocoder.yudao.framework.common.enums.UserTypeEnum; import cn.iocoder.yudao.framework.common.util.json.JsonUtils; import cn.iocoder.yudao.framework.common.util.servlet.ServletUtils; +import cn.iocoder.yudao.framework.security.core.LoginUser; import jakarta.annotation.Resource; import lombok.extern.slf4j.Slf4j; +import org.springframework.security.authentication.UsernamePasswordAuthenticationToken; +import org.springframework.security.core.context.SecurityContextHolder; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import org.springframework.util.StringUtils; @@ -25,6 +29,7 @@ import java.math.RoundingMode; import java.time.Instant; import java.time.LocalDate; import java.time.ZoneId; +import java.util.Collections; import java.util.Map; import java.util.Set; import java.util.UUID; @@ -103,32 +108,44 @@ public class EventIngestServiceImpl implements EventIngestService { @Override @Transactional(rollbackFor = Exception.class) // §3.5 C5 同步写:落库+聚合+回灌+置已聚合在同一本地事务,任一步失败整批回滚 public EventBatchResultVO ingestBatch(EventBatchReqVO reqVO, Long userId) { - // 批量上限:超限整体拒绝(统一错误码出口,前端据此分批重传) - if (reqVO.getBatch().size() > MAX_BATCH_SIZE) { - throw exception(TELEMETRY_BATCH_SIZE_EXCEEDED); - } - // 本次批量受理回执 traceId(便于排障;优先复用首条信封 traceId,缺失则生成) - String batchTraceId = resolveBatchTraceId(reqVO); - - int accepted = 0; - int rejected = 0; - // 部分成功语义:逐条轻量校验,通过则同步写(落库+聚合+回灌),失败计 rejected(不中断整批) - for (EnvelopeReqVO envelope : reqVO.getBatch()) { - if (!isEnvelopeValid(envelope)) { - rejected++; - continue; + // 2026-06-10 e2e 二轮实证(P0-1 补口):匿名上报是多级联写(事件 insert → stat 原子累加 → feed rank 回灌 + // → markAggregated UPDATE),其中 markAggregated 的 updateById 运行时实证仍把 updater=null 写进 + // UPDATE SET(与 M-b 三因叠加同形态)。逐写点 set 兜底易漏级联(feed 模块回灌同线程同雷),故在入口 + // 对匿名请求注入系统身份(照 AigcGenerateExecutor.executeWithSystemIdentity 范本的 Web 线程版), + // 使全链 DefaultDBFieldHandler 填充 creator/updater="0";finally 清理防 Web 线程上下文泄漏。 + boolean injected = injectSystemIdentityIfAnonymous(userId); + try { + // 批量上限:超限整体拒绝(统一错误码出口,前端据此分批重传) + if (reqVO.getBatch().size() > MAX_BATCH_SIZE) { + throw exception(TELEMETRY_BATCH_SIZE_EXCEEDED); } - ingestOne(envelope, userId); - accepted++; - } - // 关键链路日志:受理成功记 INFO,含回执 traceId 与计数,便于全链路排障 - log.info("[ingestBatch] 批量同步写完成 batchTraceId={}, accepted={}, rejected={}", batchTraceId, accepted, rejected); + // 本次批量受理回执 traceId(便于排障;优先复用首条信封 traceId,缺失则生成) + String batchTraceId = resolveBatchTraceId(reqVO); - EventBatchResultVO result = new EventBatchResultVO(); - result.setAccepted(accepted); - result.setRejected(rejected); - result.setTraceId(batchTraceId); - return result; + int accepted = 0; + int rejected = 0; + // 部分成功语义:逐条轻量校验,通过则同步写(落库+聚合+回灌),失败计 rejected(不中断整批) + for (EnvelopeReqVO envelope : reqVO.getBatch()) { + if (!isEnvelopeValid(envelope)) { + rejected++; + continue; + } + ingestOne(envelope, userId); + accepted++; + } + // 关键链路日志:受理成功记 INFO,含回执 traceId 与计数,便于全链路排障 + log.info("[ingestBatch] 批量同步写完成 batchTraceId={}, accepted={}, rejected={}", batchTraceId, accepted, rejected); + + EventBatchResultVO result = new EventBatchResultVO(); + result.setAccepted(accepted); + result.setRejected(rejected); + result.setTraceId(batchTraceId); + return result; + } finally { + if (injected) { + SecurityContextHolder.clearContext(); + } + } } @Override @@ -144,8 +161,40 @@ public class EventIngestServiceImpl implements EventIngestService { reqVO == null ? null : reqVO.getTraceId()); return Boolean.TRUE; } - ingestOne(reqVO, userId); - return Boolean.TRUE; + // 匿名注入系统身份:与 ingestBatch 同口径(级联写审计列兜底),finally 清理 + boolean injected = injectSystemIdentityIfAnonymous(userId); + try { + ingestOne(reqVO, userId); + return Boolean.TRUE; + } finally { + if (injected) { + SecurityContextHolder.clearContext(); + } + } + } + + /** + * 匿名上报注入系统身份(2026-06-10 e2e 二轮 P0-1 补口,范本=AigcGenerateExecutor.executeWithSystemIdentity)。 + * + * 仅当 userId==null(@PermitAll 匿名请求)时注入 LoginUser(id=0),使本次同步写全级联 + * (事件 insert / game_stat 原子累加 / FeedApi.upsertRank 回灌 / markAggregated UPDATE)的 + * MyBatis-Plus 审计列填充统一取到 "0"(与 M-b 执行器系统身份同口径);登录态请求保持原样不注入。 + * 身份仅用于审计列填充:本链业务归属一律以信封 user_id/anon_id 承载,不读登录上下文,注入无业务副作用。 + * + * @param userId 控制器解析的登录用户编号(匿名为 null) + * @return 是否已注入(调用方 finally 中据此 clearContext,防 Web 线程安全上下文污染) + */ + private boolean injectSystemIdentityIfAnonymous(Long userId) { + if (userId != null) { + return false; + } + // MVP 单租户:tenant_id 恒 1(与遥测表默认一致);ADMIN 类型仅为构造合法 LoginUser,无权限放大(不经过授权判定) + LoginUser systemUser = new LoginUser().setId(0L) + .setUserType(UserTypeEnum.ADMIN.getValue()) + .setTenantId(1L); + SecurityContextHolder.getContext().setAuthentication( + new UsernamePasswordAuthenticationToken(systemUser, null, Collections.emptyList())); + return true; } // ============================== 同步写核心(§3.5 C5)============================== @@ -338,7 +387,8 @@ public class EventIngestServiceImpl implements EventIngestService { // DefaultDBFieldHandler.insertFill 仅在 getLoginUserId()!=null 时填 creator/updater // (见 yudao-spring-boot-starter-mybatis .../handler/DefaultDBFieldHandler.java:38-43),匿名两列保持 null 撞 NOT NULL → 500。 // 归属维度已用 user_id/anon_id 承载,审计列匿名时显式兜底 "0"(与系统身份 id=0 同口径,显式值优先于 handler)。 - // 注:后续 markAggregated 的 updateById 仅更新 process_status,其 updater=null 不进 UPDATE SET(MP 跳过 null 字段),不覆盖本行 "0"。 + // 注:markAggregated 的 updateById 运行时实证仍会把 updater=null 写进 UPDATE SET(e2e 二轮复测逮到,「MP 跳过 + // null 字段」的静态推断被证伪)——级联写兜底已上移到入口 injectSystemIdentityIfAnonymous,本处 DO 级显式值保留作纵深防御。 if (eventDO.getUserId() == null) { eventDO.setCreator("0"); eventDO.setUpdater("0");