From bb406cd354a9872e3919fa918456a6a1bdbd572e Mon Sep 17 00:00:00 2001 From: dongzi Date: Wed, 15 Jul 2026 17:38:02 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E4=BA=8C=E9=98=B6=E6=AE=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../serialization-schema-check-config.yaml | 2 +- docs/CI集成说明.md | 2 +- docs/MQ序列化结构检测方案.md | 56 +-- docs/实施方案.md | 18 +- docs/配置说明.md | 29 +- .../cache/analyze/SchemaCheckAnalyzer.java | 34 +- .../cache/config/CheckerConfig.java | 6 +- .../cache/config/ConfigLoader.java | 2 +- .../cache/detector/MqReadHintDetector.java | 365 ++++++++++++++++++ .../cache/detector/MqWritePointDetector.java | 267 +++++++++++-- src/main/resources/default-config.yaml | 11 +- .../com/codechecker/cache/MqScenarioTest.java | 62 +++ .../cache/config/ConfigLoaderTest.java | 4 +- .../detector/MqReadHintDetectorTest.java | 51 +++ .../detector/MqWritePointDetectorTest.java | 109 ++++++ .../mq/kafka-patrol/CheckItemDetailVoNew.txt | 6 + .../mq/kafka-patrol/CheckItemDetailVoOld.txt | 7 + .../mq/kafka-patrol/PatrolService.txt | 14 + .../mq/kafka-record/PatrolNotifyListener.txt | 10 + .../mq/kafka-record/PatrolNotifyProducer.txt | 19 + .../mq/kafka-record/PatrolNotifyVo.txt | 6 + .../fixtures/mq/rocket-im/DutyImNotice.txt | 7 + .../mq/rocket-im/DutyImNoticeConsumer.txt | 10 + .../mq/rocket-im/DutyImNoticeProducer.txt | 28 ++ .../mq/rocket-wallet/CapitalMqConstants.txt | 6 + .../mq/rocket-wallet/WalletDeductProducer.txt | 11 + .../mq/rocket-wallet/WalletDeductReqNew.txt | 6 + .../mq/rocket-wallet/WalletDeductReqOld.txt | 7 + 28 files changed, 1076 insertions(+), 79 deletions(-) create mode 100644 src/main/java/com/codechecker/cache/detector/MqReadHintDetector.java create mode 100644 src/test/java/com/codechecker/cache/detector/MqReadHintDetectorTest.java create mode 100644 src/test/resources/fixtures/mq/kafka-patrol/CheckItemDetailVoNew.txt create mode 100644 src/test/resources/fixtures/mq/kafka-patrol/CheckItemDetailVoOld.txt create mode 100644 src/test/resources/fixtures/mq/kafka-patrol/PatrolService.txt create mode 100644 src/test/resources/fixtures/mq/kafka-record/PatrolNotifyListener.txt create mode 100644 src/test/resources/fixtures/mq/kafka-record/PatrolNotifyProducer.txt create mode 100644 src/test/resources/fixtures/mq/kafka-record/PatrolNotifyVo.txt create mode 100644 src/test/resources/fixtures/mq/rocket-im/DutyImNotice.txt create mode 100644 src/test/resources/fixtures/mq/rocket-im/DutyImNoticeConsumer.txt create mode 100644 src/test/resources/fixtures/mq/rocket-im/DutyImNoticeProducer.txt create mode 100644 src/test/resources/fixtures/mq/rocket-wallet/CapitalMqConstants.txt create mode 100644 src/test/resources/fixtures/mq/rocket-wallet/WalletDeductProducer.txt create mode 100644 src/test/resources/fixtures/mq/rocket-wallet/WalletDeductReqNew.txt create mode 100644 src/test/resources/fixtures/mq/rocket-wallet/WalletDeductReqOld.txt diff --git a/.gitea/config/serialization-schema-check-config.yaml b/.gitea/config/serialization-schema-check-config.yaml index fade21e..f85b7cf 100644 --- a/.gitea/config/serialization-schema-check-config.yaml +++ b/.gitea/config/serialization-schema-check-config.yaml @@ -4,7 +4,7 @@ # 说明: # - 本配置文件为业务覆盖配置,会与 jar 内 default-config.yaml 深度合并 # - 未声明的项沿用工具内置默认值(忽略规则、检测模式等) -# - 当前覆盖 Redis 缓存 + MQ(RocketMQ/Kafka)生产侧投递检测 +# - 当前覆盖 Redis 缓存 + MQ(RocketMQ/Kafka)生产侧投递与 MQ-R 读侧补强 # 总开关 true-执行检测 false-跳过检测(流水线直接通过,不发通知) enabled: true diff --git a/docs/CI集成说明.md b/docs/CI集成说明.md index ab09844..8157012 100644 --- a/docs/CI集成说明.md +++ b/docs/CI集成说明.md @@ -144,7 +144,7 @@ com/codechecker/serialization-schema-checker/1.0.0/ | 大量 Redis 误报 | 锁/计数器未过滤 | 补充 ignore.key_patterns | | 大量 MQ 误报 | 压测/临时 topic | 补充 ignore.mq_destinations | | commit 数显示为 1(实际多个) | 浅克隆下 `rev-list` 看不到中间提交 | 已修复:优先事件 `commits` 长度;并 deepen 到 before 为祖先 | -| 漏报(模式/模块) | W0x/MQ 未开 / include_modules 过窄 | 确认 W01~W05 与 mq_patterns;检查模块过滤 | +| 漏报(模式/模块) | W0x/MQ 未开 / include_modules 过窄 | 确认 W01~W05 与 mq_patterns(含 MQ03~05、MQ-K03/K04);检查 `mq_read_hints_enabled` 与模块过滤 | | 类型展开不完整 | 类型在依赖 jar 中 | 补充 `manual_mappings.value_type` | --- diff --git a/docs/MQ序列化结构检测方案.md b/docs/MQ序列化结构检测方案.md index 6342ce6..ea5f350 100644 --- a/docs/MQ序列化结构检测方案.md +++ b/docs/MQ序列化结构检测方案.md @@ -1,8 +1,8 @@ # MQ 消息体序列化结构变更检测 — 方案 -> 版本:v0.3 +> 版本:v0.4 > 日期:2026-07-15 -> 状态:**Phase M1 已实现**(MQ01/MQ02 + MQ-K01/MQ-K02);M2/M3 未启动 +> 状态:**Phase M1 + M2 已实现**(生产侧 MQ01~05 / MQ-K01~K04 + 读侧 MQ-R) > 关联:复用 `serialization-schema-checker` 的 Schema 提取、Diff、企微通知与 CI 框架 > 业务样本仓:`jnpf-java-cloud`(**RocketMQ + Kafka**) @@ -37,7 +37,7 @@ 在 push 时静态分析 **消息体类型的序列化 Schema** 是否相对对比区间发生变更,覆盖 **RocketMQ + Kafka**,并复用现有企微通知 / notify|block 能力。 -### 1.3 非目标(本方案首版) +### 1.3 非目标 - 不连接真实 Broker,不拉取积压消息做运行时校验 - 不解析依赖 jar 内消息类型(仅本仓 `src/main/java`) @@ -56,8 +56,9 @@ | `rocketMQTemplate.syncSend(dest, dto)` | 高(如钱包扣费) | **纳入** | | `asyncSend` / `syncSendOrderly` 等 | 中 | **纳入** | | `convertAndSend` | 视封装而定 | **纳入** | -| 先 `JSON.toJSONString` 再发 String | 较低 | unwrap 后取类型 | -| 只发 `String` / `byte[]` / `MessageExt` | 有 | **默认忽略** | +| `MessageBuilder.withPayload` 再 send | 中(如 IM 延时消息) | **纳入**(MQ04) | +| 先 `JSON.toJSONString` 再发 String | 较低 | unwrap 后取类型(MQ05) | +| 只发 `String` / `byte[]` / 无泛型 `Message` | 有 | **默认忽略** | 样本(资金钱包): @@ -74,6 +75,8 @@ public class WalletDeductConsumer implements RocketMQListener { |------|----------|------| | `kafkaTemplate.send(topic, dto)` | 中(值班食安项等) | **纳入** | | `kafkaTemplate.send(topic, List)` | 中(巡店食安项列表) | **纳入**(rootArray) | +| `kafkaTemplate.send(ProducerRecord)` | 中 | **纳入**(MQ-K03) | +| 先 JSON 再 `send(topic, json)` | 较低 | unwrap(MQ-K04) | | `@KafkaListener` + `String` + `JSONObject.parseObject(..., Xxx.class)` | 有(数据分析中差评) | **读侧补强 MQ-R** | | `KafkaTopicUtil` 仅创建 Topic | 有(租户) | **忽略**(无消息体) | @@ -118,6 +121,8 @@ public void handleMessage(String message) { | RocketMQ | `RocketMQListener`、`onMessage(T)` + `@RocketMQMessageListener` | | Kafka | `@KafkaListener` 方法参数类型;或方法内 `parseObject/parseArray(..., Xxx.class)`(与现有 W06 共享解析能力) | +开关:`detection.mq_read_hints_enabled`(默认 **true**)。 + --- ## 3. 方案总览 @@ -135,8 +140,8 @@ public void handleMessage(String message) { ```text Git Diff → 变更 Java 文件 ├─ Redis:W01~W05 + W06 ← 已有 - └─ MQ:RocketMQ(MQ01~)+ Kafka(MQ-K*) - + Listener / parse 补强(MQ-R) ← 本方案 + └─ MQ:RocketMQ(MQ01~05)+ Kafka(MQ-K01~K04) + + Listener / parse 补强(MQ-R) ← 已实现 ↓ 同一套 TypeSchema / SchemaDiffer / Skeleton / WeCom ``` @@ -161,7 +166,7 @@ Git Diff → 变更 Java 文件 | MQ01 | `rocketMQTemplate.syncSend(dest, payload, …)` | dest、payload 类型 | | MQ02 | `asyncSend` / `syncSendOrderly` / `sendOneWay` 等 | 同上 | | MQ03 | `convertAndSend(dest, payload)` | 同上 | -| MQ04 | `MessageBuilder.withPayload(obj)` 再 send | payload 类型 | +| MQ04 | `MessageBuilder.withPayload(obj)` 再 send,或 `Message` | payload / `T` | | MQ05 | 先 `toJSONString`/`getObjectToString` 再 send String | unwrap 后类型 | ### 4.2 生产侧 — Kafka @@ -170,7 +175,7 @@ Git Diff → 变更 Java 文件 |---------|------------|------| | MQ-K01 | `kafkaTemplate.send(topic, payload)` | topic、payload 类型 | | MQ-K02 | `kafkaTemplate.send(topic, key, payload)` | 同上(忽略分区 key) | -| MQ-K03 | `send(ProducerRecord)` / `ListenableFuture` 封装若可解析 | topic + value 类型 | +| MQ-K03 | `send(ProducerRecord)` | topic + value 类型 | | MQ-K04 | 先 JSON 序列化为 String 再 `send(topic, json)` | unwrap 后类型 | Payload 为 `List` / `Collection` 时标记 **rootArray**,骨架为 JSON 数组(与 Redis List 一致)。 @@ -187,7 +192,7 @@ Payload 为 `List` / `Collection` 时标记 **rootArray**,骨架为 JSON |---------|------|------| | MQ-R01 | `RocketMQListener` / `@RocketMQMessageListener` | 补强同 destination 生产点 | | MQ-R02 | `@KafkaListener` + 参数类型 `T`(非 String) | 补强同 topic | -| MQ-R03 | Listener 内 `parseObject`/`parseArray(..., Xxx.class)` | 补强(可复用 W06 检测器) | +| MQ-R03 | Listener 内 `parseObject`/`parseArray(..., Xxx.class)` | 补强(复用 parse AST) | 开关:`detection.mq_read_hints_enabled`(默认 true)。 @@ -220,7 +225,6 @@ Payload 为 `List` / `Collection` 时标记 **rootArray**,骨架为 JSON - 删除字段橙 `warning`;新增绿 `info` - destination 未解析时灰色提示 -- 可选后缀:「请评估消费积压与兼容反序列化」 「通道」字段用于区分中间件;若模板求简,可省略通道仅靠 Topic 形态区分。 @@ -272,37 +276,36 @@ manual_mappings: ## 7. 分阶段交付 -### Phase M1 — MVP(RocketMQ + Kafka 基础投递) +### Phase M1 — MVP(RocketMQ + Kafka 基础投递)✅ | 任务 | 说明 | |------|------| | MQ01/MQ02 | RocketMQ `syncSend` / `asyncSend` | | MQ-K01/MQ-K02 | Kafka `send(topic, payload)` / 三参 send | | destination 推断 | 字面量、常量、拼接;Kafka 动态 topic → `*` | -| Schema Diff + 骨架通知 | 复用 ReportBuilder;可选「通道」行 | +| Schema Diff + 骨架通知 | 复用 ReportBuilder;「通道」行 | | 夹具 | `fixtures/mq/rocket-wallet/`、`fixtures/mq/kafka-patrol/` | | 配置 | `mq_patterns`(含 MQ-K*)、`ignore.mq_destinations` | +| List 根数组骨架 | `List` 等(与 M2 需求合并交付) | **验收**: 1. 删 `WalletDeductReq` 字段 → 企微出现 RocketMQ Topic 骨架变更 2. 删 `CheckItemDetailVo` 字段 → 企微出现 Kafka Topic 骨架变更 -### Phase M2 — 增强 +### Phase M2 — 增强 ✅ | 任务 | 说明 | |------|------| -| MQ03~MQ05、MQ-K03/K04 | convertAndSend、MessageBuilder、JSON 字符串发送、ProducerRecord | +| MQ03~MQ05、MQ-K03/K04 | convertAndSend、MessageBuilder/`Message`、JSON 字符串发送、ProducerRecord | | MQ-R01~R03 | RocketMQ Listener + Kafka `@KafkaListener` / parse 补强 | -| List 根数组骨架 | 巡店 `List` 等 | +| 夹具 | `fixtures/mq/rocket-im/`、`fixtures/mq/kafka-record/` | -### Phase M3 — 运营 +**验收**: -| 任务 | 说明 | -|------|------| -| 积压风险提示文案 | 统一提示评估消费 lag / 积压 | -| 更多夹具 | 值班 Kafka、中差评 Listener、IM Favorite 等 | -| 分 webhook / 标题前缀 | 缓存 vs MQ 可选拆分 | +1. `MessageBuilder.withPayload(DutyImNotice)` + `asyncSend` → 命中 MQ04,改 VO 字段可告警 +2. `kafkaTemplate.send(ProducerRecord)` → 命中 MQ-K03 +3. Listener / `parseObject` 可补强同 Topic 弱类型生产点 --- @@ -312,10 +315,11 @@ manual_mappings: |------|------| | Kafka topic 按租户动态拼接 | 归一 `prefix:*`;`manual_mappings` | | RocketMQ / Kafka 混用同一 VO | 各投递点独立告警(符合预期) | -| Listener 收 String、parse 在方法深处 | MQ-R03 + 复用 W06 AST | +| Listener 收 String、parse 在方法深处 | MQ-R03 + parse AST | | 生产/消费跨模块对不齐 | 同仓索引 + destination 对齐;失败则仅写侧 | | Jackson / Fastjson / Kafka JsonSerializer 细节差 | 首版字段名;必要时方言或 mapping | | 只改消费未改生产类型 | 不告警(工具职责是消息体 Schema) | +| `include_modules` 过窄 | 观察期扩大模块或置空全仓 | --- @@ -334,6 +338,6 @@ manual_mappings: ## 10. 下一步 1. ~~按 Phase M1 开发~~ **已完成**(MQ01/02 + MQ-K01/K02) -2. Phase M2:MQ03~05、MQ-K03/K04、MQ-R 读侧补强 -3. 业务仓验收 Topic:`capital-topic:WALLET_DEDUCT`(RocketMQ)、巡店/值班 Kafka topic -4. 观察期按需扩大 `include_modules` \ No newline at end of file +2. ~~按 Phase M2 开发~~ **已完成**(MQ03~05、MQ-K03/K04、MQ-R) +3. 业务仓验收:`MessageBuilder` IM Topic、巡店/值班 Kafka、`capital-topic:WALLET_DEDUCT` +4. 观察期按需扩大 `include_modules` diff --git a/docs/实施方案.md b/docs/实施方案.md index d97b515..f69293d 100644 --- a/docs/实施方案.md +++ b/docs/实施方案.md @@ -535,20 +535,20 @@ notify: | 误报反馈 | `suppressions` 按写入点 / change_types 精细忽略 | | 更多业务场景覆盖 | 考勤、文件下载进度等 | -### Phase 4 — MQ 消息体结构检测(方案已落地,含 Kafka,开发待启) +### Phase 4 — MQ 消息体结构检测(已落地,含 Kafka) 业务仓同时存在: -- **RocketMQ**:`RocketMQTemplate.syncSend(topic:tag, dto)` / `RocketMQListener` -- **Kafka**:`KafkaTemplate.send(topic, vo|List)` / `@KafkaListener` + `parseObject` +- **RocketMQ**:`RocketMQTemplate.syncSend` / `asyncSend` / `MessageBuilder` / `RocketMQListener` +- **Kafka**:`KafkaTemplate.send(topic|ProducerRecord, …)` / `@KafkaListener` + `parseObject` 消息体字段变更会导致积压旧消息反序列化失败,风险模型与 Redis 同类。 | 任务 | 说明 | 状态 | |------|------|------| -| 方案文档 | RocketMQ + Kafka 模式、destination、复用 Schema Diff/企微 | ✅ 见专用文档 | -| Phase M1 | RocketMQ syncSend + Kafka send + 骨架通知 + 夹具 | 待启 | -| Phase M2/M3 | Listener/parse 补强、更多投递形态、运营 | 待启 | +| 方案文档 | RocketMQ + Kafka 模式、destination、复用 Schema Diff/企微 | ✅ | +| Phase M1 | RocketMQ sync/async + Kafka send + 骨架通知 + 夹具 | ✅ | +| Phase M2 | MessageBuilder / convertAndSend / ProducerRecord / JSON 字符串 + MQ-R | ✅ | **专用方案:** [`docs/MQ序列化结构检测方案.md`](./MQ序列化结构检测方案.md) @@ -561,6 +561,8 @@ notify: - `SchemaDifferTest`:纯字段路径对比逻辑 - `JavaSchemaExtractorTest`:类字段展开、注解、内部类 - `RedisWritePointDetectorTest`:各种写入 AST 模式匹配 +- `MqWritePointDetectorTest`:RocketMQ/Kafka 投递(含 MQ04 MessageBuilder、MQ-K03 ProducerRecord) +- `MqReadHintDetectorTest`:MQ-R Listener / parseObject 补强 - `RedisKeyResolverTest`:常量、format、拼接推断 ### 11.2 样本夹具测试(fixtures) @@ -574,6 +576,10 @@ notify: | `fixtures/tenant/` | 包装结构变更(TenantVO → CacheEnvelope) | | `fixtures/lock/` | 锁/计数器/token 应被忽略 | | `fixtures/template/` | W04 Template 直写 | +| `fixtures/mq/rocket-wallet/` | RocketMQ syncSend(MQ01) | +| `fixtures/mq/rocket-im/` | MessageBuilder + asyncSend(MQ04)、convertAndSend(MQ03) | +| `fixtures/mq/kafka-patrol/` | Kafka List 根数组(MQ-K01) | +| `fixtures/mq/kafka-record/` | ProducerRecord(MQ-K03)、JSON 字符串 send(MQ-K04) | ### 11.3 端到端测试 diff --git a/docs/配置说明.md b/docs/配置说明.md index 23b2201..275d3ce 100644 --- a/docs/配置说明.md +++ b/docs/配置说明.md @@ -43,9 +43,9 @@ include_modules: 由 `redisCheck` 仓库维护,随 jar 发布,默认包含: - `detection.patterns`:**W01~W05**(JSON 字符串写入 + Template 直写 + Hash) -- `detection.mq_patterns`:**MQ01/MQ02、MQ-K01/MQ-K02**(RocketMQ / Kafka 生产侧) +- `detection.mq_patterns`:**MQ01~MQ05、MQ-K01~MQ-K04**(RocketMQ / Kafka 生产侧) - `detection.read_hints_enabled`:W06 读侧反序列化类型辅助(默认 true) -- `detection.mq_read_hints_enabled`:MQ 读侧补强(Phase M2,默认 false) +- `detection.mq_read_hints_enabled`:MQ-R 读侧 Listener / parse 补强(默认 true) - `ignore.key_patterns`(锁 / 计数器 / token) - `ignore.mq_destinations`(MQ destination 忽略,默认真空) - `detection.min_confidence`、`max_field_depth` @@ -113,18 +113,23 @@ detection: - W04 # redisTemplate 直写对象 - W05 # opsForHash().put - # MQ 生产侧投递(Phase M1) + # MQ 生产侧投递(Phase M1 + M2) mq_patterns: - MQ01 # rocketMQTemplate.syncSend - MQ02 # asyncSend / syncSendOrderly / sendOneWay + - MQ03 # convertAndSend + - MQ04 # MessageBuilder.withPayload / Message + - MQ05 # 先 toJSONString 再 send String - MQ-K01 # kafkaTemplate.send(topic, payload) - MQ-K02 # kafkaTemplate.send(topic, key, payload) + - MQ-K03 # kafkaTemplate.send(ProducerRecord) + - MQ-K04 # 先 JSON 序列化为 String 再 send # W06:读侧反序列化类型辅助(不产生独立告警) read_hints_enabled: true - # MQ-R 读侧补强(Phase M2,默认关闭) - mq_read_hints_enabled: false + # MQ-R 读侧补强(Listener / parseObject) + mq_read_hints_enabled: true # 类型推断最低置信度,低于此值标记为低置信度提示 min_confidence: 0.6 @@ -240,7 +245,7 @@ detection: - **不是写入模式**:不会单独因为「多了一处 parse」而告警 - 能关联到 `redis get(key)` 时,还可补强 unresolved key - 覆盖优先级:`manual_mappings` > W06 > 写侧 AST -- **仅补强 Redis 写入点**;MQ 投递点不走 W06(MQ-R 为 Phase M2) +- **仅补强 Redis 写入点**;MQ 投递点走 `mq_read_hints_enabled`(MQ-R) 关闭示例: @@ -249,12 +254,12 @@ detection: read_hints_enabled: false ``` -### 3.5.2 detection.mq_patterns / ignore.mq_destinations(MQ Phase M1) +### 3.5.2 detection.mq_patterns / ignore.mq_destinations(MQ) | 配置项 | 说明 | |--------|------| -| `mq_patterns` | `MQ01`/`MQ02`(RocketMQ)、`MQ-K01`/`MQ-K02`(Kafka);默认已开启 | -| `mq_read_hints_enabled` | MQ-R 读侧补强,**默认 false** | +| `mq_patterns` | RocketMQ:`MQ01`~`MQ05`;Kafka:`MQ-K01`~`MQ-K04`;默认已全部开启 | +| `mq_read_hints_enabled` | MQ-R 读侧补强(Listener / parse),**默认 true** | | `ignore.mq_destinations` | 忽略 destination(glob),如 `benchmark:*` | 企微 MQ 块:`- Topic --> ...`,并带 `> **通道**: RocketMQ|Kafka`。 @@ -375,10 +380,10 @@ include_modules: [] # 全仓 --- -## 7. MQ 消息体检测(Phase M1 已落地) +## 7. MQ 消息体检测(Phase M1 + M2 已落地) -同一 jar 默认启用 RocketMQ / Kafka **生产侧**投递检测(`detection.mq_patterns`),与 Redis 共用 Diff、骨架与企微模板。 +同一 jar 默认启用 RocketMQ / Kafka **生产侧**投递检测(`detection.mq_patterns`)与 **MQ-R** 读侧补强,与 Redis 共用 Diff、骨架与企微模板。 - 配置说明见 §3.5.2 -- 读侧 MQ-R、`MessageBuilder` / `ProducerRecord` 等见方案 Phase M2 +- 含 `MessageBuilder` / `convertAndSend` / `ProducerRecord` / JSON 字符串发送、Listener/`parseObject` 补强 - 完整设计:[MQ序列化结构检测方案.md](./MQ序列化结构检测方案.md) \ No newline at end of file diff --git a/src/main/java/com/codechecker/cache/analyze/SchemaCheckAnalyzer.java b/src/main/java/com/codechecker/cache/analyze/SchemaCheckAnalyzer.java index 96565d1..fed6b78 100644 --- a/src/main/java/com/codechecker/cache/analyze/SchemaCheckAnalyzer.java +++ b/src/main/java/com/codechecker/cache/analyze/SchemaCheckAnalyzer.java @@ -3,6 +3,7 @@ package com.codechecker.cache.analyze; import com.codechecker.cache.config.CheckerConfig; import com.codechecker.cache.detector.CacheReadHint; import com.codechecker.cache.detector.CacheReadHintDetector; +import com.codechecker.cache.detector.MqReadHintDetector; import com.codechecker.cache.detector.MqWritePointDetector; import com.codechecker.cache.detector.RedisWritePointDetector; import com.codechecker.cache.detector.WritePoint; @@ -105,6 +106,9 @@ public class SchemaCheckAnalyzer { SchemaDiffer differ = new SchemaDiffer(); SkeletonJsonRenderer skeletonRenderer = new SkeletonJsonRenderer(); + List mqHintsNew = collectMqReadHints(newContents, newIndex); + List mqHintsOld = collectMqReadHints(oldContents, oldIndex); + List allChanges = new ArrayList<>(); Map keyChanges = new LinkedHashMap<>(); @@ -129,6 +133,8 @@ public class SchemaCheckAnalyzer { if (oldContent != null) { applyReadHints(oldWps, path, oldContent, oldIndex); } + applyMqReadHints(newWps, mqHintsNew); + applyMqReadHints(oldWps, mqHintsOld); newWps.forEach(this::applyManualMappings); oldWps.forEach(this::applyManualMappings); @@ -474,12 +480,38 @@ public class SchemaCheckAnalyzer { } for (WritePoint wp : writePoints) { if (wp.isMq()) { - continue; // Redis W06 不补强 MQ 投递点(MQ-R 为 Phase M2) + continue; // MQ 投递点由 MQ-R 补强 } enrichWritePointFromHints(wp, hints); } } + private List collectMqReadHints(Map contents, SourceIndex index) { + if (!config.getDetection().isMqReadHintsEnabled() || contents == null || contents.isEmpty()) { + return Collections.emptyList(); + } + MqReadHintDetector detector = new MqReadHintDetector(index); + List all = new ArrayList<>(); + for (Map.Entry e : contents.entrySet()) { + all.addAll(detector.detect(e.getKey(), e.getValue())); + } + return all; + } + + /** MQ-R:按 destination 匹配,补强低置信度 / 缺类型的 MQ 投递点。 */ + private void applyMqReadHints(List writePoints, List hints) { + if (!config.getDetection().isMqReadHintsEnabled() + || writePoints == null || writePoints.isEmpty() + || hints == null || hints.isEmpty()) { + return; + } + for (WritePoint wp : writePoints) { + if (wp.isMq()) { + enrichWritePointFromHints(wp, hints); + } + } + } + private void enrichWritePointFromHints(WritePoint wp, List hints) { boolean needType = wp.getResolvedValueType() == null || wp.getResolvedValueType().isEmpty() || wp.getConfidence() < config.getDetection().getMinConfidence(); diff --git a/src/main/java/com/codechecker/cache/config/CheckerConfig.java b/src/main/java/com/codechecker/cache/config/CheckerConfig.java index fb6821c..210dd37 100644 --- a/src/main/java/com/codechecker/cache/config/CheckerConfig.java +++ b/src/main/java/com/codechecker/cache/config/CheckerConfig.java @@ -118,14 +118,14 @@ public class CheckerConfig { public static class Detection { private List patterns = new ArrayList<>(); - /** MQ 投递检测模式:MQ01/MQ02/MQ-K01/MQ-K02… */ + /** MQ 投递检测模式:MQ01~MQ05、MQ-K01~MQ-K04 */ private List mqPatterns = new ArrayList<>(); private double minConfidence = 0.6; private int maxFieldDepth = 8; /** W06:是否启用读侧反序列化类型辅助补强 */ private boolean readHintsEnabled = true; - /** MQ-R:读侧 Listener / parse 补强(Phase M2;M1 默认 false) */ - private boolean mqReadHintsEnabled = false; + /** MQ-R:读侧 Listener / parse 补强 */ + private boolean mqReadHintsEnabled = true; public List getPatterns() { return patterns; diff --git a/src/main/java/com/codechecker/cache/config/ConfigLoader.java b/src/main/java/com/codechecker/cache/config/ConfigLoader.java index af732c5..5f60bea 100644 --- a/src/main/java/com/codechecker/cache/config/ConfigLoader.java +++ b/src/main/java/com/codechecker/cache/config/ConfigLoader.java @@ -103,7 +103,7 @@ public final class ConfigLoader { d.setMinConfidence(dbl(detection, "min_confidence", 0.6)); d.setMaxFieldDepth((int) lng(detection, "max_field_depth", 8)); d.setReadHintsEnabled(bool(detection, "read_hints_enabled", true)); - d.setMqReadHintsEnabled(bool(detection, "mq_read_hints_enabled", false)); + d.setMqReadHintsEnabled(bool(detection, "mq_read_hints_enabled", true)); Map severityOverrides = asMap(map.get("severity_overrides")); Map so = new LinkedHashMap<>(); diff --git a/src/main/java/com/codechecker/cache/detector/MqReadHintDetector.java b/src/main/java/com/codechecker/cache/detector/MqReadHintDetector.java new file mode 100644 index 0000000..086ae4b --- /dev/null +++ b/src/main/java/com/codechecker/cache/detector/MqReadHintDetector.java @@ -0,0 +1,365 @@ +package com.codechecker.cache.detector; + +import com.codechecker.cache.key.RedisKeyResolver; +import com.codechecker.cache.schema.SourceIndex; +import com.github.javaparser.StaticJavaParser; +import com.github.javaparser.ast.CompilationUnit; +import com.github.javaparser.ast.body.ClassOrInterfaceDeclaration; +import com.github.javaparser.ast.body.MethodDeclaration; +import com.github.javaparser.ast.body.Parameter; +import com.github.javaparser.ast.expr.AnnotationExpr; +import com.github.javaparser.ast.expr.ArrayInitializerExpr; +import com.github.javaparser.ast.expr.ClassExpr; +import com.github.javaparser.ast.expr.Expression; +import com.github.javaparser.ast.expr.MarkerAnnotationExpr; +import com.github.javaparser.ast.expr.MemberValuePair; +import com.github.javaparser.ast.expr.MethodCallExpr; +import com.github.javaparser.ast.expr.NormalAnnotationExpr; +import com.github.javaparser.ast.expr.SingleMemberAnnotationExpr; +import com.github.javaparser.ast.type.ClassOrInterfaceType; +import com.github.javaparser.ast.type.Type; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.HashSet; +import java.util.List; +import java.util.Optional; +import java.util.Set; + +/** + * MQ-R:消费侧类型提示,补强同 destination 生产点的 value 类型(不单独告警)。 + *
    + *
  • MQ-R01:{@code RocketMQListener} + {@code @RocketMQMessageListener}
  • + *
  • MQ-R02:{@code @KafkaListener} 非 String 参数类型
  • + *
  • MQ-R03:Listener 内 {@code parseObject/parseArray(..., Xxx.class)}
  • + *
+ */ +public class MqReadHintDetector { + + private static final Set OBJECT_PARSE = new HashSet<>(Arrays.asList( + "parseObject", "parse", "getJsonToBean", "toJavaObject", "readValue")); + private static final Set ARRAY_PARSE = new HashSet<>(Arrays.asList( + "parseArray", "getJsonToList", "parseArrayObject")); + private static final Set SKIP_PARAM_TYPES = new HashSet<>(Arrays.asList( + "String", "byte", "Byte", "ConsumerRecord", "MessageExt", "Message")); + + private final SourceIndex index; + private final RedisKeyResolver keyResolver; + + public MqReadHintDetector(SourceIndex index) { + this.index = index; + this.keyResolver = new RedisKeyResolver(index); + } + + public List detect(String filePath, String content) { + List result = new ArrayList<>(); + if (content == null || content.isEmpty()) { + return result; + } + CompilationUnit cu; + try { + cu = StaticJavaParser.parse(content); + } catch (RuntimeException e) { + return result; + } + for (ClassOrInterfaceDeclaration clazz : cu.findAll(ClassOrInterfaceDeclaration.class)) { + result.addAll(detectRocketListener(clazz, filePath)); + for (MethodDeclaration md : clazz.getMethods()) { + result.addAll(detectKafkaListener(md, clazz, filePath)); + result.addAll(detectParseInListener(md, clazz, filePath)); + } + } + return result; + } + + /** MQ-R01 */ + private List detectRocketListener(ClassOrInterfaceDeclaration clazz, String filePath) { + List result = new ArrayList<>(); + Optional listenerType = findRocketListenerType(clazz); + if (!listenerType.isPresent()) { + return result; + } + AnnotationExpr ann = findAnnotation(clazz.getAnnotations(), "RocketMQMessageListener"); + if (ann == null) { + return result; + } + String dest = resolveRocketDestination(ann, clazz); + InferredType payload = resolveTypeArg(listenerType.get(), index.get( + clazz.getFullyQualifiedName().orElse(clazz.getNameAsString()))); + if (payload.fqn == null) { + return result; + } + CacheReadHint hint = baseHint(filePath, clazz, "", dest, payload); + hint.setConfidence(0.85); + result.add(hint); + return result; + } + + /** MQ-R02 */ + private List detectKafkaListener(MethodDeclaration md, + ClassOrInterfaceDeclaration clazz, + String filePath) { + List result = new ArrayList<>(); + AnnotationExpr ann = findAnnotation(md.getAnnotations(), "KafkaListener"); + if (ann == null) { + return result; + } + List topics = resolveKafkaTopics(ann, clazz); + if (topics.isEmpty()) { + return result; + } + SourceIndex.IndexedType context = index.get( + clazz.getFullyQualifiedName().orElse(clazz.getNameAsString())); + InferredType payload = null; + for (Parameter p : md.getParameters()) { + InferredType t = resolveParamPayload(p.getType(), context); + if (t != null && t.fqn != null) { + payload = t; + break; + } + } + if (payload == null) { + return result; + } + for (String topic : topics) { + CacheReadHint hint = baseHint(filePath, clazz, md.getNameAsString(), topic, payload); + hint.setConfidence(0.85); + result.add(hint); + } + return result; + } + + /** MQ-R03:在已标注 KafkaListener 的方法,或 RocketMQListener 类内 parse */ + private List detectParseInListener(MethodDeclaration md, + ClassOrInterfaceDeclaration clazz, + String filePath) { + List result = new ArrayList<>(); + boolean kafka = findAnnotation(md.getAnnotations(), "KafkaListener") != null; + boolean rocket = findRocketListenerType(clazz).isPresent(); + if (!kafka && !rocket) { + return result; + } + List destinations = new ArrayList<>(); + if (kafka) { + destinations.addAll(resolveKafkaTopics( + findAnnotation(md.getAnnotations(), "KafkaListener"), clazz)); + } + if (rocket) { + AnnotationExpr ann = findAnnotation(clazz.getAnnotations(), "RocketMQMessageListener"); + if (ann != null) { + String dest = resolveRocketDestination(ann, clazz); + if (dest != null && !dest.isEmpty()) { + destinations.add(dest); + } + } + } + if (destinations.isEmpty()) { + destinations.add(null); + } + + SourceIndex.IndexedType context = index.get( + clazz.getFullyQualifiedName().orElse(clazz.getNameAsString())); + for (MethodCallExpr mce : md.findAll(MethodCallExpr.class)) { + String name = mce.getNameAsString(); + boolean array = ARRAY_PARSE.contains(name); + boolean object = OBJECT_PARSE.contains(name); + if (!array && !object || mce.getArguments().size() < 2) { + continue; + } + Expression classArg = mce.getArgument(1); + if (!(classArg instanceof ClassExpr)) { + continue; + } + Type type = ((ClassExpr) classArg).getType(); + if (!(type instanceof ClassOrInterfaceType)) { + continue; + } + String fqn = index.resolveFqn(((ClassOrInterfaceType) type).getNameWithScope(), context); + if (fqn == null) { + fqn = index.resolveFqn(((ClassOrInterfaceType) type).getNameAsString(), context); + } + if (fqn == null) { + continue; + } + InferredType payload = new InferredType(fqn, array); + for (String dest : destinations) { + CacheReadHint hint = baseHint(filePath, clazz, md.getNameAsString(), dest, payload); + hint.setLineNumber(mce.getBegin().map(p -> p.line).orElse(0)); + hint.setConfidence(dest == null ? 0.7 : 0.8); + result.add(hint); + } + } + return result; + } + + private CacheReadHint baseHint(String filePath, ClassOrInterfaceDeclaration clazz, + String method, String dest, InferredType payload) { + CacheReadHint hint = new CacheReadHint(); + hint.setFilePath(filePath); + hint.setLineNumber(clazz.getBegin().map(p -> p.line).orElse(0)); + hint.setEnclosingClass(clazz.getFullyQualifiedName().orElse(clazz.getNameAsString())); + hint.setEnclosingMethod(method); + hint.setResolvedKeyPattern(dest); + hint.setResolvedValueType(payload.fqn); + hint.setRootArray(payload.isArray); + return hint; + } + + private Optional findRocketListenerType(ClassOrInterfaceDeclaration clazz) { + for (ClassOrInterfaceType t : clazz.getImplementedTypes()) { + String n = t.getNameAsString(); + if ("RocketMQListener".equals(n) || "RocketMQReplyListener".equals(n)) { + return Optional.of(t); + } + } + return Optional.empty(); + } + + private InferredType resolveParamPayload(Type type, SourceIndex.IndexedType context) { + if (!(type instanceof ClassOrInterfaceType)) { + return null; + } + ClassOrInterfaceType cit = (ClassOrInterfaceType) type; + String simple = cit.getNameAsString(); + if (SKIP_PARAM_TYPES.contains(simple) && !"ConsumerRecord".equals(simple)) { + return null; + } + if ("ConsumerRecord".equals(simple) || "List".equals(simple) || "ArrayList".equals(simple)) { + Optional last = cit.getTypeArguments() + .filter(a -> !a.isEmpty()) + .map(a -> a.get(a.size() - 1)); + if (!last.isPresent()) { + return null; + } + InferredType inner = resolveTypeArgFromType(last.get(), context); + if (inner != null && ("List".equals(simple) || "ArrayList".equals(simple))) { + return new InferredType(inner.fqn, true); + } + return inner; + } + return resolveTypeArgFromType(type, context); + } + + private InferredType resolveTypeArg(ClassOrInterfaceType listenerType, SourceIndex.IndexedType context) { + Optional arg = listenerType.getTypeArguments().filter(a -> !a.isEmpty()).map(a -> a.get(0)); + if (!arg.isPresent()) { + return new InferredType(null, false); + } + return resolveTypeArgFromType(arg.get(), context); + } + + private InferredType resolveTypeArgFromType(Type type, SourceIndex.IndexedType context) { + if (!(type instanceof ClassOrInterfaceType)) { + return new InferredType(null, false); + } + ClassOrInterfaceType cit = (ClassOrInterfaceType) type; + if ("List".equals(cit.getNameAsString()) || "ArrayList".equals(cit.getNameAsString())) { + Optional el = cit.getTypeArguments().filter(a -> !a.isEmpty()).map(a -> a.get(0)); + if (el.isPresent() && el.get() instanceof ClassOrInterfaceType) { + return new InferredType(resolveFqn((ClassOrInterfaceType) el.get(), context), true); + } + return new InferredType(null, true); + } + String simple = cit.getNameAsString(); + if (SKIP_PARAM_TYPES.contains(simple)) { + return new InferredType(null, false); + } + return new InferredType(resolveFqn(cit, context), false); + } + + private String resolveFqn(ClassOrInterfaceType cit, SourceIndex.IndexedType context) { + String fqn = index.resolveFqn(cit.getNameWithScope(), context); + if (fqn == null) { + fqn = index.resolveFqn(cit.getNameAsString(), context); + } + return fqn; + } + + private String resolveRocketDestination(AnnotationExpr ann, ClassOrInterfaceDeclaration clazz) { + SourceIndex.IndexedType context = index.get( + clazz.getFullyQualifiedName().orElse(clazz.getNameAsString())); + Expression topicExpr = annotationValue(ann, "topic"); + if (topicExpr == null) { + return null; + } + String topicPattern = keyResolver.resolve(topicExpr, clazz, context); + Expression tagExpr = annotationValue(ann, "selectorExpression"); + if (tagExpr != null) { + String tagPattern = keyResolver.resolve(tagExpr, clazz, context); + if (tagPattern != null && !tagPattern.isEmpty() + && !"*".equals(tagPattern) && !tagPattern.contains("||")) { + return topicPattern + ":" + tagPattern; + } + } + return topicPattern; + } + + private List resolveKafkaTopics(AnnotationExpr ann, ClassOrInterfaceDeclaration clazz) { + List result = new ArrayList<>(); + if (ann == null) { + return result; + } + SourceIndex.IndexedType context = index.get( + clazz.getFullyQualifiedName().orElse(clazz.getNameAsString())); + Expression topicsExpr = annotationValue(ann, "topics"); + if (topicsExpr == null && ann instanceof SingleMemberAnnotationExpr) { + topicsExpr = ((SingleMemberAnnotationExpr) ann).getMemberValue(); + } + if (topicsExpr == null) { + return result; + } + List items = new ArrayList<>(); + if (topicsExpr instanceof ArrayInitializerExpr) { + items.addAll(((ArrayInitializerExpr) topicsExpr).getValues()); + } else { + items.add(topicsExpr); + } + for (Expression item : items) { + String resolved = keyResolver.resolve(item, clazz, context); + if (resolved != null && !resolved.isEmpty()) { + result.add(resolved); + } + } + return result; + } + + private static Expression annotationValue(AnnotationExpr ann, String name) { + if (ann instanceof NormalAnnotationExpr) { + for (MemberValuePair pair : ((NormalAnnotationExpr) ann).getPairs()) { + if (name.equals(pair.getNameAsString())) { + return pair.getValue(); + } + } + } + return null; + } + + private static AnnotationExpr findAnnotation(List annotations, String simpleName) { + for (AnnotationExpr ann : annotations) { + String n; + if (ann instanceof MarkerAnnotationExpr) { + n = ((MarkerAnnotationExpr) ann).getNameAsString(); + } else if (ann instanceof SingleMemberAnnotationExpr) { + n = ((SingleMemberAnnotationExpr) ann).getNameAsString(); + } else if (ann instanceof NormalAnnotationExpr) { + n = ((NormalAnnotationExpr) ann).getNameAsString(); + } else { + n = ann.getNameAsString(); + } + if (n.equals(simpleName) || n.endsWith("." + simpleName)) { + return ann; + } + } + return null; + } + + private static final class InferredType { + final String fqn; + final boolean isArray; + + InferredType(String fqn, boolean isArray) { + this.fqn = fqn; + this.isArray = isArray; + } + } +} diff --git a/src/main/java/com/codechecker/cache/detector/MqWritePointDetector.java b/src/main/java/com/codechecker/cache/detector/MqWritePointDetector.java index 86d6fe1..14c6341 100644 --- a/src/main/java/com/codechecker/cache/detector/MqWritePointDetector.java +++ b/src/main/java/com/codechecker/cache/detector/MqWritePointDetector.java @@ -32,7 +32,12 @@ import java.util.Optional; import java.util.Set; /** - * 检测 RocketMQ / Kafka 生产侧投递点(Phase M1:MQ01/MQ02、MQ-K01/MQ-K02)。 + * 检测 RocketMQ / Kafka 生产侧投递点(Phase M1 + M2)。 + *
    + *
  • M1:MQ01/MQ02、MQ-K01/MQ-K02
  • + *
  • M2:MQ03 convertAndSend、MQ04 MessageBuilder/`Message<T>`、MQ05 JSON 字符串、 + * MQ-K03 ProducerRecord、MQ-K04 JSON 字符串
  • + *
*/ public class MqWritePointDetector { @@ -41,12 +46,15 @@ public class MqWritePointDetector { private static final Set ROCKET_SYNC = new HashSet<>(Arrays.asList("syncSend")); private static final Set ROCKET_ASYNC = new HashSet<>(Arrays.asList( "asyncSend", "syncSendOrderly", "sendOneWay", "asyncSendOrderly")); + private static final Set ROCKET_CONVERT = new HashSet<>(Arrays.asList("convertAndSend")); private static final Set TRIVIAL_VALUE_CALLS = new HashSet<>(Arrays.asList( "randomUUID", "toString", "valueOf")); private static final Set COLLECTION_SIMPLE = new HashSet<>(Arrays.asList( "List", "ArrayList", "LinkedList", "Set", "HashSet", "Collection")); - private static final Set IGNORE_PAYLOAD_TYPES = new HashSet<>(Arrays.asList( - "String", "byte", "Byte", "MessageExt", "Message", "ProducerRecord")); + private static final Set ENVELOPE_TYPES = new HashSet<>(Arrays.asList( + "Message", "MessageExt", "ProducerRecord")); + private static final Set IGNORE_BARE_TYPES = new HashSet<>(Arrays.asList( + "String", "byte", "Byte")); private final SourceIndex index; private final Set enabledPatterns; @@ -83,19 +91,19 @@ public class MqWritePointDetector { private WritePoint tryDetectRocketMq(MethodCallExpr mce, String filePath) { String method = mce.getNameAsString(); - String pattern = null; + String basePattern = null; if (ROCKET_SYNC.contains(method) && enabledPatterns.contains("MQ01")) { - pattern = "MQ01"; + basePattern = "MQ01"; } else if (ROCKET_ASYNC.contains(method) && enabledPatterns.contains("MQ02")) { - pattern = "MQ02"; + basePattern = "MQ02"; + } else if (ROCKET_CONVERT.contains(method) && enabledPatterns.contains("MQ03")) { + basePattern = "MQ03"; } - if (pattern == null || !isRocketMqScope(mce) || mce.getArguments().size() < 2) { + if (basePattern == null || !isRocketMqScope(mce) || mce.getArguments().size() < 2) { return null; } - Expression destArg = mce.getArgument(0); - Expression payloadArg = mce.getArgument(1); - return buildMqWritePoint(mce, filePath, pattern, WritePoint.CHANNEL_ROCKETMQ, - destArg, payloadArg); + return buildMqWritePoint(mce, filePath, basePattern, WritePoint.CHANNEL_ROCKETMQ, + mce.getArgument(0), mce.getArgument(1), "MQ04", "MQ05"); } private WritePoint tryDetectKafka(MethodCallExpr mce, String filePath) { @@ -103,32 +111,80 @@ public class MqWritePointDetector { return null; } int argc = mce.getArguments().size(); - String pattern; + if (argc == 1 && enabledPatterns.contains("MQ-K03")) { + return buildFromProducerRecord(mce, filePath); + } + String basePattern; Expression topicArg; Expression payloadArg; if (argc == 2 && enabledPatterns.contains("MQ-K01")) { - pattern = "MQ-K01"; + basePattern = "MQ-K01"; topicArg = mce.getArgument(0); payloadArg = mce.getArgument(1); } else if (argc >= 3 && enabledPatterns.contains("MQ-K02")) { - pattern = "MQ-K02"; + basePattern = "MQ-K02"; topicArg = mce.getArgument(0); payloadArg = mce.getArgument(argc - 1); } else { return null; } - return buildMqWritePoint(mce, filePath, pattern, WritePoint.CHANNEL_KAFKA, - topicArg, payloadArg); + return buildMqWritePoint(mce, filePath, basePattern, WritePoint.CHANNEL_KAFKA, + topicArg, payloadArg, null, "MQ-K04"); } - private WritePoint buildMqWritePoint(MethodCallExpr mce, String filePath, String pattern, - String channel, Expression destArg, Expression payloadArg) { - Expression serialized = unwrapSerializer(payloadArg); - Expression typeExpr = serialized != null ? serialized : payloadArg; - if (serialized == null && isTrivialValue(payloadArg)) { + private WritePoint buildFromProducerRecord(MethodCallExpr mce, String filePath) { + Expression recordArg = mce.getArgument(0); + ObjectCreationExpr creation = findProducerRecordCreation(recordArg, mce); + if (creation == null || creation.getArguments().size() < 2) { return null; } - if (isBareStringOrBytesPayload(typeExpr, mce)) { + Expression topicArg = creation.getArgument(0); + Expression payloadArg = creation.getArgument(creation.getArguments().size() - 1); + WritePoint wp = buildMqWritePoint(mce, filePath, "MQ-K03", WritePoint.CHANNEL_KAFKA, + topicArg, payloadArg, null, "MQ-K04"); + if (wp != null) { + wp.setValueExpression(recordArg.toString()); + } + return wp; + } + + private ObjectCreationExpr findProducerRecordCreation(Expression expr, MethodCallExpr contextCall) { + if (expr instanceof ObjectCreationExpr) { + ObjectCreationExpr oce = (ObjectCreationExpr) expr; + if ("ProducerRecord".equals(oce.getType().getNameAsString())) { + return oce; + } + return null; + } + if (expr instanceof NameExpr) { + String name = ((NameExpr) expr).getNameAsString(); + Optional callable = contextCall.findAncestor(CallableDeclaration.class); + if (!callable.isPresent()) { + return null; + } + for (VariableDeclarator var : callable.get().findAll(VariableDeclarator.class)) { + if (!var.getNameAsString().equals(name) || !var.getInitializer().isPresent()) { + continue; + } + Expression init = var.getInitializer().get(); + if (init instanceof ObjectCreationExpr + && "ProducerRecord".equals(((ObjectCreationExpr) init).getType().getNameAsString())) { + return (ObjectCreationExpr) init; + } + } + } + return null; + } + + private WritePoint buildMqWritePoint(MethodCallExpr mce, String filePath, String basePattern, + String channel, Expression destArg, Expression payloadArg, + String messagePattern, String jsonPattern) { + PayloadResolution resolved = resolvePayload(payloadArg, mce, messagePattern, jsonPattern); + if (resolved == null) { + return null; + } + String pattern = resolved.overridePattern != null ? resolved.overridePattern : basePattern; + if (!enabledPatterns.contains(pattern)) { return null; } @@ -146,7 +202,7 @@ public class MqWritePointDetector { .findAncestor(ClassOrInterfaceDeclaration.class).orElse(null); wp.setResolvedKeyPattern(keyResolver.resolve(destArg, enclosingDecl, context)); - InferredType inferred = inferType(typeExpr, mce, context); + InferredType inferred = inferType(resolved.typeExpr, mce, context); if (inferred != null) { wp.setResolvedValueType(inferred.fqn); wp.setRootArray(inferred.isArray); @@ -157,6 +213,150 @@ public class MqWritePointDetector { return wp; } + /** + * 解析 payload:Message/MessageBuilder(MQ04)、JSON 字符串(MQ05/MQ-K04)、直传对象。 + * 无法解析的 Message/纯 String 返回 null(忽略)。 + */ + private PayloadResolution resolvePayload(Expression payloadArg, MethodCallExpr mce, + String messagePattern, String jsonPattern) { + if (payloadArg == null) { + return null; + } + // MQ04:Message<T> 泛型 / MessageBuilder.withPayload + if (messagePattern != null && enabledPatterns.contains(messagePattern)) { + Expression fromMessage = unwrapMessagePayload(payloadArg, mce); + if (fromMessage != null) { + return new PayloadResolution(fromMessage, messagePattern); + } + if (isEnvelopeBare(payloadArg, mce)) { + return null; + } + } else if (isEnvelopeBare(payloadArg, mce)) { + return null; + } + + // 内联 JSON 序列化 + Expression serialized = unwrapSerializer(payloadArg); + if (serialized != null) { + if (jsonPattern != null && enabledPatterns.contains(jsonPattern)) { + return new PayloadResolution(serialized, jsonPattern); + } + return new PayloadResolution(serialized, null); + } + + // 局部 String = toJSONString(...) + if (jsonPattern != null && enabledPatterns.contains(jsonPattern)) { + Expression fromLocalJson = unwrapJsonStringVar(payloadArg, mce); + if (fromLocalJson != null) { + return new PayloadResolution(fromLocalJson, jsonPattern); + } + } + + if (isTrivialValue(payloadArg) || isBareStringOrBytesPayload(payloadArg, mce)) { + return null; + } + return new PayloadResolution(payloadArg, null); + } + + private Expression unwrapMessagePayload(Expression payloadArg, MethodCallExpr mce) { + Expression fromBuilder = unwrapMessageBuilderPayload(payloadArg, mce); + if (fromBuilder != null) { + return fromBuilder; + } + // Message<DutyImNotice> → DutyImNotice(无 withPayload 也可) + if (payloadArg instanceof NameExpr) { + Type declared = findVariableType(((NameExpr) payloadArg).getNameAsString(), mce); + if (declared instanceof ClassOrInterfaceType) { + ClassOrInterfaceType cit = (ClassOrInterfaceType) declared; + if (ENVELOPE_TYPES.contains(cit.getNameAsString())) { + Optional payloadType = cit.getTypeArguments() + .filter(a -> !a.isEmpty()) + .map(a -> a.get(a.size() - 1)); + if (payloadType.isPresent() && payloadType.get() instanceof ClassOrInterfaceType) { + // 用伪 ObjectCreation 表达类型不便;改为在 infer 前用 NameExpr 不够 + // 返回一个标记:借助 CastExpr 包装类型信息 + return new CastExpr(payloadType.get(), payloadArg); + } + } + } + } + return null; + } + + private Expression unwrapMessageBuilderPayload(Expression payloadArg, MethodCallExpr mce) { + Expression chain = payloadArg; + if (payloadArg instanceof NameExpr) { + String name = ((NameExpr) payloadArg).getNameAsString(); + Optional callable = mce.findAncestor(CallableDeclaration.class); + if (callable.isPresent()) { + for (VariableDeclarator var : callable.get().findAll(VariableDeclarator.class)) { + if (var.getNameAsString().equals(name) && var.getInitializer().isPresent()) { + chain = var.getInitializer().get(); + break; + } + } + } + } + return findWithPayloadArg(chain); + } + + private Expression findWithPayloadArg(Expression expr) { + Expression current = expr; + int guard = 0; + while (current instanceof MethodCallExpr && guard++ < 16) { + MethodCallExpr call = (MethodCallExpr) current; + if ("withPayload".equals(call.getNameAsString()) && !call.getArguments().isEmpty()) { + return call.getArgument(0); + } + if (!call.getScope().isPresent()) { + break; + } + current = call.getScope().get(); + } + return null; + } + + private Expression unwrapJsonStringVar(Expression payloadArg, MethodCallExpr mce) { + if (!(payloadArg instanceof NameExpr)) { + return null; + } + String name = ((NameExpr) payloadArg).getNameAsString(); + Type declared = findVariableType(name, mce); + if (declared != null) { + String simple = declared.isClassOrInterfaceType() + ? declared.asClassOrInterfaceType().getNameAsString() + : declared.asString(); + if (!"String".equals(simple)) { + return null; + } + } + Optional callable = mce.findAncestor(CallableDeclaration.class); + if (!callable.isPresent()) { + return null; + } + for (VariableDeclarator var : callable.get().findAll(VariableDeclarator.class)) { + if (!var.getNameAsString().equals(name) || !var.getInitializer().isPresent()) { + continue; + } + return unwrapSerializer(var.getInitializer().get()); + } + return null; + } + + private boolean isEnvelopeBare(Expression typeExpr, MethodCallExpr contextCall) { + if (!(typeExpr instanceof NameExpr)) { + return false; + } + Type declared = findVariableType(((NameExpr) typeExpr).getNameAsString(), contextCall); + if (declared == null) { + return false; + } + String simple = declared.isClassOrInterfaceType() + ? declared.asClassOrInterfaceType().getNameAsString() + : declared.asString(); + return ENVELOPE_TYPES.contains(simple); + } + private boolean isRocketMqScope(MethodCallExpr mce) { String scope = mce.getScope().map(Expression::toString).orElse("").toLowerCase(); return scope.contains("rocketmqtemplate") || scope.contains("rocketmq"); @@ -177,7 +377,7 @@ public class MqWritePointDetector { String simple = declared.isClassOrInterfaceType() ? declared.asClassOrInterfaceType().getNameAsString() : declared.asString(); - if (IGNORE_PAYLOAD_TYPES.contains(simple) || "byte[]".equals(declared.asString())) { + if (IGNORE_BARE_TYPES.contains(simple) || "byte[]".equals(declared.asString())) { return true; } } @@ -306,6 +506,15 @@ public class MqWritePointDetector { } return new InferredType(null, true); } + if (ENVELOPE_TYPES.contains(simple)) { + Optional payload = cit.getTypeArguments() + .filter(a -> !a.isEmpty()) + .map(a -> a.get(a.size() - 1)); + if (payload.isPresent()) { + return resolveTypeNode(payload.get(), context); + } + return new InferredType(null, false); + } return new InferredType(resolveFqn(cit, context), false); } @@ -317,6 +526,16 @@ public class MqWritePointDetector { return fqn; } + private static final class PayloadResolution { + final Expression typeExpr; + final String overridePattern; + + PayloadResolution(Expression typeExpr, String overridePattern) { + this.typeExpr = typeExpr; + this.overridePattern = overridePattern; + } + } + private static final class InferredType { final String fqn; final boolean isArray; diff --git a/src/main/resources/default-config.yaml b/src/main/resources/default-config.yaml index b4a0389..1c0ce7b 100644 --- a/src/main/resources/default-config.yaml +++ b/src/main/resources/default-config.yaml @@ -46,16 +46,21 @@ detection: - W03 # stringRedisTemplate.opsForValue().set(key, JsonUtil.getObjectToString(x), ...) - W04 # redisTemplate.opsForValue().set(key, obj, ...) - W05 # redisTemplate.opsForHash().put(key, field, obj) - # MQ 生产侧投递检测(Phase M1) + # MQ 生产侧投递检测(Phase M1 + M2) mq_patterns: - MQ01 # rocketMQTemplate.syncSend(dest, payload) - MQ02 # asyncSend / syncSendOrderly / sendOneWay + - MQ03 # convertAndSend(dest, payload) + - MQ04 # MessageBuilder.withPayload / Message + - MQ05 # 先 toJSONString 再 send String - MQ-K01 # kafkaTemplate.send(topic, payload) - MQ-K02 # kafkaTemplate.send(topic, key, payload) + - MQ-K03 # kafkaTemplate.send(ProducerRecord) + - MQ-K04 # 先 JSON 序列化为 String 再 send # W06 读侧辅助:用 parseObject / getJsonToBean 等补强写入点 value 类型(非写入模式) read_hints_enabled: true - # MQ-R 读侧补强(Phase M2,默认关闭) - mq_read_hints_enabled: false + # MQ-R 读侧补强(Listener / parseObject) + mq_read_hints_enabled: true # 类型推断最低置信度,低于此值标记为低置信度提示 min_confidence: 0.6 # 字段展开最大深度(防止循环引用) diff --git a/src/test/java/com/codechecker/cache/MqScenarioTest.java b/src/test/java/com/codechecker/cache/MqScenarioTest.java index c81f609..32a814b 100644 --- a/src/test/java/com/codechecker/cache/MqScenarioTest.java +++ b/src/test/java/com/codechecker/cache/MqScenarioTest.java @@ -1,5 +1,7 @@ package com.codechecker.cache; +import com.codechecker.cache.detector.CacheReadHint; +import com.codechecker.cache.detector.MqReadHintDetector; import com.codechecker.cache.detector.MqWritePointDetector; import com.codechecker.cache.detector.WritePoint; import com.codechecker.cache.diff.ChangeType; @@ -19,6 +21,7 @@ import java.util.Set; import java.util.stream.Collectors; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertTrue; class MqScenarioTest { @@ -120,4 +123,63 @@ class MqScenarioTest { assertTrue(md.contains("> **通道**: Kafka"), md); assertTrue(md.contains("> **类型**: List"), md); } + + @Test + void messageBuilderProducerResolvesDutyImNotice() { + String notice = TestSupport.fixture("fixtures/mq/rocket-im/DutyImNotice.txt"); + String producer = TestSupport.fixture("fixtures/mq/rocket-im/DutyImNoticeProducer.txt"); + + SourceIndex index = new SourceIndex(); + index.addSource(notice); + index.addSource(producer); + + Set patterns = new HashSet<>(); + patterns.add("MQ02"); + patterns.add("MQ04"); + WritePoint wp = new MqWritePointDetector(index, patterns) + .detect("DutyImNoticeProducer.java", producer).stream() + .filter(w -> "MQ04".equals(w.getPattern())) + .findFirst() + .orElseThrow(() -> new AssertionError("应命中 MQ04")); + assertEquals("duty-im-notice-topic", wp.getResolvedKeyPattern()); + assertTrue(wp.getResolvedValueType().endsWith("DutyImNotice")); + } + + @Test + void mqReadHintEnrichesWritePointByDestination() { + String notice = TestSupport.fixture("fixtures/mq/rocket-im/DutyImNotice.txt"); + String producer = TestSupport.fixture("fixtures/mq/rocket-im/DutyImNoticeProducer.txt"); + String consumer = TestSupport.fixture("fixtures/mq/rocket-im/DutyImNoticeConsumer.txt"); + + SourceIndex index = new SourceIndex(); + index.addSource(notice); + index.addSource(producer); + index.addSource(consumer); + + List hints = new MqReadHintDetector(index) + .detect("DutyImNoticeConsumer.java", consumer); + assertFalse(hints.isEmpty()); + + Set patterns = new HashSet<>(); + patterns.add("MQ02"); + patterns.add("MQ04"); + WritePoint wp = new MqWritePointDetector(index, patterns) + .detect("DutyImNoticeProducer.java", producer).stream() + .filter(w -> "MQ04".equals(w.getPattern())) + .findFirst() + .orElseThrow(() -> new AssertionError("应命中 MQ04")); + + // 模拟无泛型 Message 导致类型缺失时,MQ-R 可按 destination 补强 + wp.setResolvedValueType(null); + wp.setConfidence(0.4); + CacheReadHint hint = hints.get(0); + assertEquals(wp.getResolvedKeyPattern(), hint.getResolvedKeyPattern()); + + if (hint.getResolvedValueType() != null) { + wp.setResolvedValueType(hint.getResolvedValueType()); + wp.setConfidence(Math.max(wp.getConfidence(), hint.getConfidence())); + } + assertTrue(wp.getResolvedValueType().endsWith("DutyImNotice")); + assertTrue(wp.getConfidence() >= 0.8); + } } diff --git a/src/test/java/com/codechecker/cache/config/ConfigLoaderTest.java b/src/test/java/com/codechecker/cache/config/ConfigLoaderTest.java index 4b9e9cd..e5b5fbc 100644 --- a/src/test/java/com/codechecker/cache/config/ConfigLoaderTest.java +++ b/src/test/java/com/codechecker/cache/config/ConfigLoaderTest.java @@ -21,7 +21,9 @@ class ConfigLoaderTest { assertTrue(config.getDetection().getPatterns().contains("W01")); assertTrue(config.getDetection().getMqPatterns().contains("MQ01")); assertTrue(config.getDetection().getMqPatterns().contains("MQ-K01")); - assertFalse(config.getDetection().isMqReadHintsEnabled()); + assertTrue(config.getDetection().getMqPatterns().contains("MQ04")); + assertTrue(config.getDetection().getMqPatterns().contains("MQ-K03")); + assertTrue(config.getDetection().isMqReadHintsEnabled()); assertFalse(config.isScanTestSources()); } diff --git a/src/test/java/com/codechecker/cache/detector/MqReadHintDetectorTest.java b/src/test/java/com/codechecker/cache/detector/MqReadHintDetectorTest.java new file mode 100644 index 0000000..6efb224 --- /dev/null +++ b/src/test/java/com/codechecker/cache/detector/MqReadHintDetectorTest.java @@ -0,0 +1,51 @@ +package com.codechecker.cache.detector; + +import com.codechecker.cache.TestSupport; +import com.codechecker.cache.schema.SourceIndex; +import org.junit.jupiter.api.Test; + +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class MqReadHintDetectorTest { + + @Test + void detectsRocketMqListenerGenericType() { + String notice = TestSupport.fixture("fixtures/mq/rocket-im/DutyImNotice.txt"); + String producer = TestSupport.fixture("fixtures/mq/rocket-im/DutyImNoticeProducer.txt"); + String consumer = TestSupport.fixture("fixtures/mq/rocket-im/DutyImNoticeConsumer.txt"); + + SourceIndex index = new SourceIndex(); + index.addSource(notice); + index.addSource(producer); + index.addSource(consumer); + + List hints = new MqReadHintDetector(index) + .detect("DutyImNoticeConsumer.java", consumer); + assertFalse(hints.isEmpty()); + CacheReadHint hint = hints.get(0); + assertTrue(hint.getResolvedValueType().endsWith("DutyImNotice")); + assertEquals("duty-im-notice-topic", hint.getResolvedKeyPattern()); + assertFalse(hint.isRootArray()); + } + + @Test + void detectsKafkaListenerParseObject() { + String vo = TestSupport.fixture("fixtures/mq/kafka-record/PatrolNotifyVo.txt"); + String listener = TestSupport.fixture("fixtures/mq/kafka-record/PatrolNotifyListener.txt"); + + SourceIndex index = new SourceIndex(); + index.addSource(vo); + index.addSource(listener); + + List hints = new MqReadHintDetector(index) + .detect("PatrolNotifyListener.java", listener); + assertTrue(hints.stream().anyMatch(h -> + h.getResolvedValueType() != null + && h.getResolvedValueType().endsWith("PatrolNotifyVo") + && "patrol-notify-topic".equals(h.getResolvedKeyPattern()))); + } +} diff --git a/src/test/java/com/codechecker/cache/detector/MqWritePointDetectorTest.java b/src/test/java/com/codechecker/cache/detector/MqWritePointDetectorTest.java index d87b86d..84cc8c4 100644 --- a/src/test/java/com/codechecker/cache/detector/MqWritePointDetectorTest.java +++ b/src/test/java/com/codechecker/cache/detector/MqWritePointDetectorTest.java @@ -28,6 +28,9 @@ class MqWritePointDetectorTest { Set patterns = new HashSet<>(); patterns.add("MQ01"); patterns.add("MQ02"); + patterns.add("MQ03"); + patterns.add("MQ04"); + patterns.add("MQ05"); List wps = new MqWritePointDetector(index, patterns) .detect("WalletDeductProducer.java", producer); @@ -63,4 +66,110 @@ class MqWritePointDetectorTest { assertTrue(wp.getResolvedValueType().endsWith("CheckItemDetailVo")); assertTrue(wp.isRootArray()); } + + @Test + void detectsMessageBuilderAsyncSendAsMq04() { + String notice = TestSupport.fixture("fixtures/mq/rocket-im/DutyImNotice.txt"); + String producer = TestSupport.fixture("fixtures/mq/rocket-im/DutyImNoticeProducer.txt"); + + SourceIndex index = new SourceIndex(); + index.addSource(notice); + index.addSource(producer); + + Set patterns = allRocketPatterns(); + List wps = new MqWritePointDetector(index, patterns) + .detect("DutyImNoticeProducer.java", producer); + + WritePoint mq04 = findByPattern(wps, "MQ04"); + assertEquals(WritePoint.CHANNEL_ROCKETMQ, mq04.getChannel()); + assertEquals("duty-im-notice-topic", mq04.getResolvedKeyPattern()); + assertTrue(mq04.getResolvedValueType().endsWith("DutyImNotice"), mq04.getResolvedValueType()); + } + + @Test + void detectsConvertAndSendAsMq03() { + String notice = TestSupport.fixture("fixtures/mq/rocket-im/DutyImNotice.txt"); + String producer = TestSupport.fixture("fixtures/mq/rocket-im/DutyImNoticeProducer.txt"); + SourceIndex index = new SourceIndex(); + index.addSource(notice); + index.addSource(producer); + + WritePoint mq03 = findByPattern( + new MqWritePointDetector(index, allRocketPatterns()) + .detect("DutyImNoticeProducer.java", producer), + "MQ03"); + assertTrue(mq03.getResolvedValueType().endsWith("DutyImNotice")); + } + + @Test + void detectsJsonStringSyncSendAsMq05() { + String notice = TestSupport.fixture("fixtures/mq/rocket-im/DutyImNotice.txt"); + String producer = TestSupport.fixture("fixtures/mq/rocket-im/DutyImNoticeProducer.txt"); + SourceIndex index = new SourceIndex(); + index.addSource(notice); + index.addSource(producer); + + WritePoint mq05 = findByPattern( + new MqWritePointDetector(index, allRocketPatterns()) + .detect("DutyImNoticeProducer.java", producer), + "MQ05"); + assertTrue(mq05.getResolvedValueType().endsWith("DutyImNotice")); + } + + @Test + void detectsProducerRecordSendAsMqK03() { + String vo = TestSupport.fixture("fixtures/mq/kafka-record/PatrolNotifyVo.txt"); + String producer = TestSupport.fixture("fixtures/mq/kafka-record/PatrolNotifyProducer.txt"); + SourceIndex index = new SourceIndex(); + index.addSource(vo); + index.addSource(producer); + + Set patterns = allKafkaPatterns(); + WritePoint wp = findByPattern( + new MqWritePointDetector(index, patterns).detect("PatrolNotifyProducer.java", producer), + "MQ-K03"); + assertEquals("patrol-notify-topic", wp.getResolvedKeyPattern()); + assertTrue(wp.getResolvedValueType().endsWith("PatrolNotifyVo")); + } + + @Test + void detectsKafkaJsonStringAsMqK04() { + String vo = TestSupport.fixture("fixtures/mq/kafka-record/PatrolNotifyVo.txt"); + String producer = TestSupport.fixture("fixtures/mq/kafka-record/PatrolNotifyProducer.txt"); + SourceIndex index = new SourceIndex(); + index.addSource(vo); + index.addSource(producer); + + WritePoint wp = findByPattern( + new MqWritePointDetector(index, allKafkaPatterns()) + .detect("PatrolNotifyProducer.java", producer), + "MQ-K04"); + assertTrue(wp.getResolvedValueType().endsWith("PatrolNotifyVo")); + } + + private static Set allRocketPatterns() { + Set patterns = new HashSet<>(); + patterns.add("MQ01"); + patterns.add("MQ02"); + patterns.add("MQ03"); + patterns.add("MQ04"); + patterns.add("MQ05"); + return patterns; + } + + private static Set allKafkaPatterns() { + Set patterns = new HashSet<>(); + patterns.add("MQ-K01"); + patterns.add("MQ-K02"); + patterns.add("MQ-K03"); + patterns.add("MQ-K04"); + return patterns; + } + + private static WritePoint findByPattern(List wps, String pattern) { + return wps.stream() + .filter(w -> pattern.equals(w.getPattern())) + .findFirst() + .orElseThrow(() -> new AssertionError("缺少模式 " + pattern + ",实际: " + wps)); + } } diff --git a/src/test/resources/fixtures/mq/kafka-patrol/CheckItemDetailVoNew.txt b/src/test/resources/fixtures/mq/kafka-patrol/CheckItemDetailVoNew.txt new file mode 100644 index 0000000..683786f --- /dev/null +++ b/src/test/resources/fixtures/mq/kafka-patrol/CheckItemDetailVoNew.txt @@ -0,0 +1,6 @@ +package demo.kafka; + +public class CheckItemDetailVo { + private String itemId; + private String itemName; +} diff --git a/src/test/resources/fixtures/mq/kafka-patrol/CheckItemDetailVoOld.txt b/src/test/resources/fixtures/mq/kafka-patrol/CheckItemDetailVoOld.txt new file mode 100644 index 0000000..76d7911 --- /dev/null +++ b/src/test/resources/fixtures/mq/kafka-patrol/CheckItemDetailVoOld.txt @@ -0,0 +1,7 @@ +package demo.kafka; + +public class CheckItemDetailVo { + private String itemId; + private String itemName; + private Integer score; +} diff --git a/src/test/resources/fixtures/mq/kafka-patrol/PatrolService.txt b/src/test/resources/fixtures/mq/kafka-patrol/PatrolService.txt new file mode 100644 index 0000000..0d7ed78 --- /dev/null +++ b/src/test/resources/fixtures/mq/kafka-patrol/PatrolService.txt @@ -0,0 +1,14 @@ +package demo.kafka; + +import java.util.List; +import org.springframework.kafka.core.KafkaTemplate; + +public class PatrolService { + private KafkaTemplate kafkaTemplate; + private static final String TOPIC = "patrol-store-food-safe:%s"; + + public void sendFoodSafeData(String tenantId, List thousandsData) { + String topic = String.format(TOPIC, tenantId); + kafkaTemplate.send(topic, thousandsData); + } +} diff --git a/src/test/resources/fixtures/mq/kafka-record/PatrolNotifyListener.txt b/src/test/resources/fixtures/mq/kafka-record/PatrolNotifyListener.txt new file mode 100644 index 0000000..b7729ff --- /dev/null +++ b/src/test/resources/fixtures/mq/kafka-record/PatrolNotifyListener.txt @@ -0,0 +1,10 @@ +package demo.kafka; + +import org.springframework.kafka.annotation.KafkaListener; + +public class PatrolNotifyListener { + @KafkaListener(topics = "patrol-notify-topic", groupId = "patrol-group") + public void handle(String message) { + PatrolNotifyVo vo = JSONObject.parseObject(message, PatrolNotifyVo.class); + } +} diff --git a/src/test/resources/fixtures/mq/kafka-record/PatrolNotifyProducer.txt b/src/test/resources/fixtures/mq/kafka-record/PatrolNotifyProducer.txt new file mode 100644 index 0000000..c7403ce --- /dev/null +++ b/src/test/resources/fixtures/mq/kafka-record/PatrolNotifyProducer.txt @@ -0,0 +1,19 @@ +package demo.kafka; + +import org.apache.kafka.clients.producer.ProducerRecord; +import org.springframework.kafka.core.KafkaTemplate; + +public class PatrolNotifyProducer { + private KafkaTemplate kafkaTemplate; + private static final String TOPIC = "patrol-notify-topic"; + + public void sendRecord(PatrolNotifyVo vo) { + ProducerRecord record = new ProducerRecord<>(TOPIC, vo); + kafkaTemplate.send(record); + } + + public void sendJson(PatrolNotifyVo vo) { + String json = JSON.toJSONString(vo); + kafkaTemplate.send(TOPIC, json); + } +} diff --git a/src/test/resources/fixtures/mq/kafka-record/PatrolNotifyVo.txt b/src/test/resources/fixtures/mq/kafka-record/PatrolNotifyVo.txt new file mode 100644 index 0000000..9f3fbaa --- /dev/null +++ b/src/test/resources/fixtures/mq/kafka-record/PatrolNotifyVo.txt @@ -0,0 +1,6 @@ +package demo.kafka; + +public class PatrolNotifyVo { + private String storeId; + private String status; +} diff --git a/src/test/resources/fixtures/mq/rocket-im/DutyImNotice.txt b/src/test/resources/fixtures/mq/rocket-im/DutyImNotice.txt new file mode 100644 index 0000000..fb559ea --- /dev/null +++ b/src/test/resources/fixtures/mq/rocket-im/DutyImNotice.txt @@ -0,0 +1,7 @@ +package demo.mq; + +public class DutyImNotice { + private String tenantId; + private String id; + private String content; +} diff --git a/src/test/resources/fixtures/mq/rocket-im/DutyImNoticeConsumer.txt b/src/test/resources/fixtures/mq/rocket-im/DutyImNoticeConsumer.txt new file mode 100644 index 0000000..0cc49d8 --- /dev/null +++ b/src/test/resources/fixtures/mq/rocket-im/DutyImNoticeConsumer.txt @@ -0,0 +1,10 @@ +package demo.mq; + +import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; +import org.apache.rocketmq.spring.core.RocketMQListener; + +@RocketMQMessageListener(topic = DutyImNoticeProducer.MQ_TOPIC, consumerGroup = "duty-im-group") +public class DutyImNoticeConsumer implements RocketMQListener { + public void onMessage(DutyImNotice message) { + } +} diff --git a/src/test/resources/fixtures/mq/rocket-im/DutyImNoticeProducer.txt b/src/test/resources/fixtures/mq/rocket-im/DutyImNoticeProducer.txt new file mode 100644 index 0000000..a12e154 --- /dev/null +++ b/src/test/resources/fixtures/mq/rocket-im/DutyImNoticeProducer.txt @@ -0,0 +1,28 @@ +package demo.mq; + +import org.apache.rocketmq.spring.core.RocketMQTemplate; +import org.springframework.messaging.Message; +import org.springframework.messaging.support.MessageBuilder; + +public class DutyImNoticeProducer { + private RocketMQTemplate rocketMQTemplate; + public static final String MQ_TOPIC = "duty-im-notice-topic"; + + public void addMq(DutyImNotice message, long timestamp) { + String uniqueKey = message.getTenantId() + "_" + message.getId(); + Message mqMessage = MessageBuilder.withPayload(message) + .setHeader("KEYS", uniqueKey) + .setHeader("executeTime", timestamp) + .build(); + rocketMQTemplate.asyncSend(MQ_TOPIC, mqMessage); + } + + public void convertSend(DutyImNotice message) { + rocketMQTemplate.convertAndSend(MQ_TOPIC, message); + } + + public void sendJson(DutyImNotice message) { + String body = JSON.toJSONString(message); + rocketMQTemplate.syncSend(MQ_TOPIC, body); + } +} diff --git a/src/test/resources/fixtures/mq/rocket-wallet/CapitalMqConstants.txt b/src/test/resources/fixtures/mq/rocket-wallet/CapitalMqConstants.txt new file mode 100644 index 0000000..62da126 --- /dev/null +++ b/src/test/resources/fixtures/mq/rocket-wallet/CapitalMqConstants.txt @@ -0,0 +1,6 @@ +package demo.mq; + +public final class CapitalMqConstants { + public static final String TOPIC = "capital-topic"; + public static final String TAG_WALLET_DEDUCT = "WALLET_DEDUCT"; +} diff --git a/src/test/resources/fixtures/mq/rocket-wallet/WalletDeductProducer.txt b/src/test/resources/fixtures/mq/rocket-wallet/WalletDeductProducer.txt new file mode 100644 index 0000000..d0f4958 --- /dev/null +++ b/src/test/resources/fixtures/mq/rocket-wallet/WalletDeductProducer.txt @@ -0,0 +1,11 @@ +package demo.mq; + +import org.apache.rocketmq.spring.core.RocketMQTemplate; + +public class WalletDeductProducer { + private RocketMQTemplate rocketMQTemplate; + + public void send(WalletDeductReq req) { + rocketMQTemplate.syncSend(CapitalMqConstants.TOPIC + ":" + CapitalMqConstants.TAG_WALLET_DEDUCT, req); + } +} diff --git a/src/test/resources/fixtures/mq/rocket-wallet/WalletDeductReqNew.txt b/src/test/resources/fixtures/mq/rocket-wallet/WalletDeductReqNew.txt new file mode 100644 index 0000000..d11152c --- /dev/null +++ b/src/test/resources/fixtures/mq/rocket-wallet/WalletDeductReqNew.txt @@ -0,0 +1,6 @@ +package demo.mq; + +public class WalletDeductReq { + private String walletId; + private Long amount; +} diff --git a/src/test/resources/fixtures/mq/rocket-wallet/WalletDeductReqOld.txt b/src/test/resources/fixtures/mq/rocket-wallet/WalletDeductReqOld.txt new file mode 100644 index 0000000..c45f06a --- /dev/null +++ b/src/test/resources/fixtures/mq/rocket-wallet/WalletDeductReqOld.txt @@ -0,0 +1,7 @@ +package demo.mq; + +public class WalletDeductReq { + private String walletId; + private Long amount; + private String remark; +}