From 4eaaa88e9c2ea7143dc6b7684d9896187e18a434 Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 18:22:51 +0800 Subject: [PATCH] chore: remove xinda ingestion modules --- CHANGELOG.md | 2 +- DECISIONS.md | 17 +- README.md | 5 +- docs/module-data-flow.html | 2 +- .../plans/2026-06-23-gb32960-service-split.md | 4 +- ...-23-gb32960-production-readiness-design.md | 2 +- docs/target-architecture.md | 7 +- .../historyapp/MavenModuleProfileTest.java | 28 +- modules/apps/xinda-push-app/pom.xml | 79 ---- .../xindapushapp/XindaPushApplication.java | 21 - .../src/main/resources/application.yml | 111 ----- .../XindaPushAppCompositionTest.java | 69 --- .../XindaPushAppDefaultsTest.java | 63 --- .../com/lingniu/ingest/api/ProtocolId.java | 3 +- .../ingest/api/event/VehicleEvent.java | 2 +- .../core/dispatcher/HandlerRegistryTest.java | 6 +- .../InMemoryVehicleIdentityService.java | 48 ++ modules/inbound/inbound-xinda-push/pom.xml | 118 ----- .../ingest/inbound/xinda/XindaPushClient.java | 413 ------------------ .../inbound/xinda/XindaPushEventMapper.java | 413 ------------------ .../inbound/xinda/XindaPushMessage.java | 32 -- .../xinda/XindaPushRealtimeHandler.java | 37 -- .../config/XindaPushAutoConfiguration.java | 52 --- .../xinda/config/XindaPushProperties.java | 45 -- ...ot.autoconfigure.AutoConfiguration.imports | 1 - .../inbound/xinda/XindaPushClientTest.java | 272 ------------ .../xinda/XindaPushEventMapperTest.java | 295 ------------- .../xinda/XindaPushRealtimeHandlerTest.java | 30 -- .../InMemoryVehicleIdentityService.java | 48 ++ .../InMemoryVehicleIdentityService.java | 48 ++ .../InMemoryVehicleIdentityService.java | 48 ++ .../LocationHistoryControllerTest.java | 2 +- .../RawFrameHistoryControllerTest.java | 2 +- pom.xml | 39 -- 34 files changed, 234 insertions(+), 2130 deletions(-) delete mode 100644 modules/apps/xinda-push-app/pom.xml delete mode 100644 modules/apps/xinda-push-app/src/main/java/com/lingniu/ingest/xindapushapp/XindaPushApplication.java delete mode 100644 modules/apps/xinda-push-app/src/main/resources/application.yml delete mode 100644 modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppCompositionTest.java delete mode 100644 modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppDefaultsTest.java create mode 100644 modules/inbound/inbound-mqtt/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java delete mode 100644 modules/inbound/inbound-xinda-push/pom.xml delete mode 100644 modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushClient.java delete mode 100644 modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushEventMapper.java delete mode 100644 modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushMessage.java delete mode 100644 modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushRealtimeHandler.java delete mode 100644 modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/config/XindaPushAutoConfiguration.java delete mode 100644 modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/config/XindaPushProperties.java delete mode 100644 modules/inbound/inbound-xinda-push/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports delete mode 100644 modules/inbound/inbound-xinda-push/src/test/java/com/lingniu/ingest/inbound/xinda/XindaPushClientTest.java delete mode 100644 modules/inbound/inbound-xinda-push/src/test/java/com/lingniu/ingest/inbound/xinda/XindaPushEventMapperTest.java delete mode 100644 modules/inbound/inbound-xinda-push/src/test/java/com/lingniu/ingest/inbound/xinda/XindaPushRealtimeHandlerTest.java create mode 100644 modules/protocols/protocol-jsatl12/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java create mode 100644 modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java create mode 100644 modules/protocols/protocol-jt808/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java diff --git a/CHANGELOG.md b/CHANGELOG.md index 1d3b030e..e85b9c04 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -54,7 +54,7 @@ V2016/V2025 双版本信息体 parser、平台登入鉴权、VIN 白名单、读空闲告警。 - **JT/T 808 / JT/T 1078 / JSATL12 接入**(`protocol-jt808` / `protocol-jt1078` / `protocol-jsatl12`)。 -- **MQTT 入站**与**信达平台 push 接入**(`inbound-mqtt` / `inbound-xinda-push`)。 +- **MQTT 入站**(`inbound-mqtt`)。 - **Disruptor 事件总线**(`ingest-core`):高吞吐解耦解码与下游 sink。 - **Sink**:本地归档(`sink-archive`)、Kafka(`sink-kafka`,含 protobuf envelope)。 - **会话状态**(`session-core`)。 diff --git a/DECISIONS.md b/DECISIONS.md index a63c9f9b..7acb356f 100644 --- a/DECISIONS.md +++ b/DECISIONS.md @@ -20,16 +20,11 @@ - **Decision**: 拆出独立 `command-gateway` 模块,复用 `session-core`;本次重构同步交付 - **Consequences**: HTTP 对外接口(原 JT808Controller / JT1078Controller)迁移到此模块 -## ADR-004 信达 Push:Legacy only +## ADR-004 信达 Push:Removed - **Status**: Superseded by ADR-011 - **Context**: 旧实现硬编码凭证、无重连、0300/0401 空实现 -- **Decision**: - 1. 新模块 `inbound-xinda-push` 隔离第三方库 `gps-push-client` - 2. 配置外置(yml + env) - 3. 自动重连 + Resilience4j 熔断 - 4. 0200/0300/0401 **仅做协议解析**,转换为 `VehicleEvent` 后直接投 Kafka - 5. **业务处理全部下沉到下游消费者**,本服务不再实现报警/透传业务逻辑 -- **Consequences**: 信达 Push 现在不属于默认生产面;模块仅保留在 `legacy-xinda` profile,后续优化优先投向 GB32960、JT808、宇通 MQTT、历史和统计链路。 +- **Decision**: 删除信达 Push 源码、Maven profile、私仓依赖和部署入口。 +- **Consequences**: 信达 Push 后续已从源码和构建面删除,优化优先投向 GB32960、JT808、宇通 MQTT、历史和统计链路。 ## ADR-005 JT1078 / JSATL12 进第一批 - **Status**: Superseded by ADR-012 @@ -66,12 +61,12 @@ ## ADR-011 默认生产面:三协议接入 + 历史 + 统计 - **Status**: Accepted -- **Context**: 生产接入范围已经从 GB32960 三应用拆分,扩展为三条活跃接入链路和两个消费应用;信达 Push 已废弃,不应再出现在默认构建、部署或优化目标里。 +- **Context**: 生产接入范围已经从 GB32960 三应用拆分,扩展为三条活跃接入链路和两个消费应用;信达 Push 已废弃并删除,不应再出现在构建、部署或优化目标里。 - **Decision**: 1. 默认生产面包含 GB32960、JT808、Yutong MQTT、vehicle-history-app、vehicle-analytics-app。 - 2. Xinda Push 仅保留在 `legacy-xinda` profile,不进入默认 Maven reactor、Woodpecker 镜像发布和历史消费绑定。 + 2. Xinda Push 源码、Maven profile、Woodpecker 镜像发布和历史消费绑定全部移除。 3. `command-gateway` 和 JT1078 继续作为可选能力,通过 `optional-command-gateway` profile 显式启用。 -- **Consequences**: 默认构建和部署保持精简;后续性能、可靠性、字段解析和存储优化都优先服务三条活跃协议链路。 +- **Consequences**: 默认构建和部署保持精简;后续性能、可靠性、字段解析和存储优化都服务三条活跃协议链路。 ## ADR-012 JSATL12 附件上传:Optional only - **Status**: Accepted diff --git a/README.md b/README.md index 0fe276a1..8e2976b6 100644 --- a/README.md +++ b/README.md @@ -45,8 +45,7 @@ lingniu-vehicle-ingest/ │ │ ├── protocol-jt1078/ JT/T 1078(optional-command-gateway profile;808 信令按需桥接 + 常用下行信令编码 + TCP/UDP RTP 媒体流分段归档 + 事件 metadata 内部 VIN + 坏 RTP 统一 RawArchive/Passthrough + 归档失败兜底 + SIM 身份映射) │ │ └── protocol-jsatl12/ 苏标主动安全报警附件(optional-attachments profile) │ ├── inbound/ -│ │ ├── inbound-mqtt/ MQTT 接入(endpoint 生命周期 + profile 注册扩展 + 统一身份映射 + PEM 双向 TLS + 未知 profile/解析失败/profile异常/连接订阅失败兜底 + 统一 UNKNOWN 身份 metadata) -│ │ └── inbound-xinda-push/ 信达 Push 接入(废弃兼容模块,仅 -Plegacy-xinda 显式构建) +│ │ └── inbound-mqtt/ MQTT 接入(endpoint 生命周期 + profile 注册扩展 + 统一身份映射 + PEM 双向 TLS + 未知 profile/解析失败/profile异常/连接订阅失败兜底 + 统一 UNKNOWN 身份 metadata) │ ├── sinks/ │ │ ├── sink-kafka/ Kafka producer + Protobuf Envelope │ │ └── sink-archive/ 原始报文冷存 @@ -75,7 +74,7 @@ mvn -pl :gb32960-ingest-app,:jt808-ingest-app,:yutong-mqtt-app,:vehicle-history- 生产服务只在 ECS/Portainer 上运行;本机只用于源码检查、Maven 构建和单元测试。部署、健康检查和真实流量验证步骤见 `docs/operations/gb32960-service-split-runbook.md` 与 `docs/operations/vehicle-ingest-tdengine-verification.md`。 -信达 Push 已废弃,不参与默认 Maven reactor、CI、Portainer 部署和历史消费链路;如需兼容排查,可显式使用 `-Plegacy-xinda` 构建。 +信达 Push 已废弃并删除源码,不参与 Maven reactor、CI、Portainer 部署和历史消费链路。 JSATL12 附件上传不在默认 Maven reactor 中;需要时使用 `-Poptional-attachments` 构建。 diff --git a/docs/module-data-flow.html b/docs/module-data-flow.html index 3220eba3..377abbc0 100644 --- a/docs/module-data-flow.html +++ b/docs/module-data-flow.html @@ -265,7 +265,7 @@
接入方式 - Netty TCP、MQTT、后续可扩展更多 Inbound Adapter;Xinda Push 不在默认生产面。 + Netty TCP、MQTT、后续可扩展更多 Inbound Adapter;信达 Push 源码已删除。
统一模型 diff --git a/docs/superpowers/plans/2026-06-23-gb32960-service-split.md b/docs/superpowers/plans/2026-06-23-gb32960-service-split.md index 84f02c7a..f538d55f 100644 --- a/docs/superpowers/plans/2026-06-23-gb32960-service-split.md +++ b/docs/superpowers/plans/2026-06-23-gb32960-service-split.md @@ -13,7 +13,9 @@ ## Current-State Notes - Existing all-in-one entrypoint: `modules/apps/bootstrap-all/src/main/java/com/lingniu/ingest/bootstrap/IngestApplication.java`. -- Existing all-in-one module: `modules/apps/bootstrap-all/pom.xml` imports protocols, sinks, history, state, stat, command gateway, MQTT, Xinda, and JT modules. +- Existing all-in-one module at the time of this plan imported protocols, + sinks, history, state, stat, command gateway, MQTT, and JT modules. Xinda was + later removed from source and build configuration. - Existing GB32960 handler: `modules/protocols/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandler.java`. - Existing handler currently writes report ACK before `dispatcher.dispatch(rf)`. - Existing `DisruptorEventBus` publishes to sinks asynchronously and does not wait for sink completion. diff --git a/docs/superpowers/specs/2026-06-23-gb32960-production-readiness-design.md b/docs/superpowers/specs/2026-06-23-gb32960-production-readiness-design.md index 23aaafbb..0f023d5f 100644 --- a/docs/superpowers/specs/2026-06-23-gb32960-production-readiness-design.md +++ b/docs/superpowers/specs/2026-06-23-gb32960-production-readiness-design.md @@ -2,7 +2,7 @@ > Use `docs/target-architecture.md` and > `docs/superpowers/specs/2026-06-29-vehicle-ingest-redesign.md` for the > current production architecture: TDengine is the default hot history store, -> `event-file-store` is optional compatibility only, and Xinda Push is legacy. +> `event-file-store` and Xinda Push have been removed from the production codebase. # GB32960 Production Readiness Design diff --git a/docs/target-architecture.md b/docs/target-architecture.md index f2d44c91..2f636c4a 100644 --- a/docs/target-architecture.md +++ b/docs/target-architecture.md @@ -17,9 +17,8 @@ Active production protocols: | JT/T 808 | `jt808-ingest-app` | `vehicle.event.jt808.v1` | `vehicle.raw.jt808.v1` | | Yutong MQTT | `yutong-mqtt-app` | `vehicle.event.mqtt-yutong.v1` | `vehicle.raw.mqtt-yutong.v1` | -Xinda Push is legacy and is not part of the default production surface. -It may remain in source for reference, but new optimization work should target -the active protocols above. +Xinda Push source, Maven profile, deployment bindings, and history consumers +are removed. New optimization work should target the active protocols above. JSATL12 attachment upload is also outside the default production reactor. Keep it available through the `optional-attachments` profile until an attachment app @@ -187,7 +186,7 @@ production deployment. - Active production deployment contains GB32960, JT808, Yutong MQTT, vehicle-history, and vehicle-analytics apps. -- Xinda Push is absent from default deployment and history consumer bindings. +- Xinda Push source and deployment bindings are removed. - JSATL12 is absent from the default Maven reactor and available through `optional-attachments`. - vehicle-state-service is absent from the default Maven reactor and available diff --git a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/MavenModuleProfileTest.java b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/MavenModuleProfileTest.java index 6ccedcd6..b1759268 100644 --- a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/MavenModuleProfileTest.java +++ b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/MavenModuleProfileTest.java @@ -119,27 +119,34 @@ class MavenModuleProfileTest { } @Test - void xindaIsLegacyOnlyAndOutsideDefaultProductionReactor() throws Exception { + void xindaModulesAreRemovedFromBuildSurface() throws Exception { Document pom = rootPom(); + String rootPomText = Files.readString(repositoryRoot().resolve("pom.xml")); assertThat(defaultModules(pom)) .doesNotContain("modules/inbound/inbound-xinda-push") .doesNotContain("modules/apps/xinda-push-app"); - assertThat(profileModules(pom, "legacy-xinda")) - .containsExactly( - "modules/inbound/inbound-xinda-push", - "modules/apps/xinda-push-app"); + .as("xinda is deleted, not hidden behind an opt-in profile") + .isEmpty(); + assertThat(Files.exists(repositoryRoot().resolve("modules/inbound/inbound-xinda-push"))) + .isFalse(); + assertThat(Files.exists(repositoryRoot().resolve("modules/apps/xinda-push-app"))) + .isFalse(); + assertThat(rootPomText) + .doesNotContain("legacy-xinda") + .doesNotContain("inbound-xinda-push") + .doesNotContain("xinda-push-app"); } @Test - void privateLingniuRepositoryIsScopedToLegacyXindaProfile() throws Exception { + void privateLingniuRepositoryIsRemovedWithXinda() throws Exception { Document pom = rootPom(); assertThat(repositoryIds(pom.getDocumentElement())) .doesNotContain("lingniu-nexus"); assertThat(profileRepositoryIds(pom, "legacy-xinda")) - .containsExactly("lingniu-nexus"); + .isEmpty(); } @Test @@ -167,7 +174,7 @@ class MavenModuleProfileTest { assertThat(profileDependencyManagementArtifacts(pom, "optional-event-file-store")) .isEmpty(); assertThat(profileDependencyManagementArtifacts(pom, "legacy-xinda")) - .containsExactly("inbound-xinda-push", "xinda-push-app"); + .isEmpty(); } @Test @@ -472,12 +479,12 @@ class MavenModuleProfileTest { String decisions = Files.readString(repositoryRoot().resolve("DECISIONS.md")); assertThat(decisions) - .contains("## ADR-004 信达 Push:Legacy only") + .contains("## ADR-004 信达 Push:Removed") .contains("## ADR-006 部署形态:GB32960 三应用拆分\n- **Status**: Superseded by ADR-011") .contains("## ADR-011 默认生产面:三协议接入 + 历史 + 统计") .contains("## ADR-012 JSATL12 附件上传:Optional only") .contains("GB32960、JT808、Yutong MQTT、vehicle-history-app、vehicle-analytics-app") - .contains("Xinda Push 仅保留在 `legacy-xinda` profile") + .contains("Xinda Push 源码、Maven profile、Woodpecker 镜像发布和历史消费绑定全部移除") .contains("JSATL12 仅保留在 `optional-attachments` profile") .contains("vehicle-state-service 仅保留在 `optional-latest-state` profile") .contains("raw-archive-store 原型已删除") @@ -486,6 +493,7 @@ class MavenModuleProfileTest { .doesNotContain("event-file-store") .doesNotContain("optional-event-file-store") .doesNotContain("optional-raw-archive-store") + .doesNotContain("legacy-xinda") .doesNotContain("## ADR-004 信达 Push:保留模块但彻底重写\n- **Status**: Accepted") .doesNotContain("CI/CD 分别构建三个镜像") .doesNotContain("Spring Boot 3.4.x"); diff --git a/modules/apps/xinda-push-app/pom.xml b/modules/apps/xinda-push-app/pom.xml deleted file mode 100644 index 052910c6..00000000 --- a/modules/apps/xinda-push-app/pom.xml +++ /dev/null @@ -1,79 +0,0 @@ - - - 4.0.0 - - com.lingniu.ingest - lingniu-vehicle-ingest - 0.1.0-SNAPSHOT - ../../../pom.xml - - - xinda-push-app - xinda-push-app - Xinda Push ingress runtime: platform push client, frame parse, identity mapping, Kafka production, and raw archive. - - - - com.lingniu.ingest - ingest-core - - - com.lingniu.ingest - observability - - - com.lingniu.ingest - vehicle-identity - - - com.lingniu.ingest - inbound-xinda-push - - - com.lingniu.ingest - sink-kafka - - - com.lingniu.ingest - sink-archive - - - org.springframework.boot - spring-boot-starter-web - - - com.alibaba.cloud - spring-cloud-starter-alibaba-nacos-config - - - org.springdoc - springdoc-openapi-starter-webmvc-ui - - - org.springframework.boot - spring-boot-starter-test - test - - - - - xinda-push-app - - - org.springframework.boot - spring-boot-maven-plugin - - --sun-misc-unsafe-memory-access=allow - - - - - repackage - build-info - - - - - - - diff --git a/modules/apps/xinda-push-app/src/main/java/com/lingniu/ingest/xindapushapp/XindaPushApplication.java b/modules/apps/xinda-push-app/src/main/java/com/lingniu/ingest/xindapushapp/XindaPushApplication.java deleted file mode 100644 index 7d558c3d..00000000 --- a/modules/apps/xinda-push-app/src/main/java/com/lingniu/ingest/xindapushapp/XindaPushApplication.java +++ /dev/null @@ -1,21 +0,0 @@ -package com.lingniu.ingest.xindapushapp; - -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; - -@SpringBootApplication(scanBasePackages = "com.lingniu.ingest") -public class XindaPushApplication { - - private static final Logger log = LoggerFactory.getLogger(XindaPushApplication.class); - - public static void main(String[] args) { - try { - SpringApplication.run(XindaPushApplication.class, args); - } catch (Throwable t) { - log.error("Xinda Push ingest application failed to start, forcing JVM exit", t); - System.exit(1); - } - } -} diff --git a/modules/apps/xinda-push-app/src/main/resources/application.yml b/modules/apps/xinda-push-app/src/main/resources/application.yml deleted file mode 100644 index 3da42635..00000000 --- a/modules/apps/xinda-push-app/src/main/resources/application.yml +++ /dev/null @@ -1,111 +0,0 @@ -spring: - application: - name: xinda-push-app - config: - import: - - optional:nacos:${spring.application.name}.${NACOS_CONFIG_FILE_EXTENSION:yml}?group=${NACOS_GROUP:DEFAULT_GROUP}&refreshEnabled=${NACOS_REFRESH_ENABLED:true} - cloud: - nacos: - config: - enabled: ${NACOS_CONFIG_ENABLED:true} - server-addr: ${NACOS_SERVER_ADDR:127.0.0.1:8848} - namespace: ${NACOS_NAMESPACE:} - group: ${NACOS_GROUP:DEFAULT_GROUP} - file-extension: ${NACOS_CONFIG_FILE_EXTENSION:yml} - username: ${NACOS_USERNAME:} - password: ${NACOS_PASSWORD:} - threads: - virtual: - enabled: true - -server: - port: ${HTTP_PORT:20600} - -springdoc: - swagger-ui: - path: /swagger-ui.html - api-docs: - path: /v3/api-docs - -lingniu: - ingest: - xinda-push: - enabled: ${XINDA_PUSH_ENABLED:false} - host: ${XINDA_PUSH_HOST:} - port: ${XINDA_PUSH_PORT:10100} - username: ${XINDA_PUSH_USERNAME:} - password: ${XINDA_PUSH_PASSWORD:} - client-desc: ${XINDA_PUSH_CLIENT_DESC:lingniu-ingest} - subscribe-msg-ids: - - ${XINDA_PUSH_MSG_ID_LOCATION:0200} - - ${XINDA_PUSH_MSG_ID_ALARM:0300} - - ${XINDA_PUSH_MSG_ID_PASSTHROUGH:0401} - reconnect-interval-sec: ${XINDA_PUSH_RECONNECT_INTERVAL_SEC:5} - pipeline: - disruptor: - ring-buffer-size: ${PIPELINE_RING_BUFFER_SIZE:131072} - wait-strategy: ${PIPELINE_WAIT_STRATEGY:yielding} - producer-type: multi - dedup: - enabled: true - cache-size: 200000 - ttl-seconds: 600 - rate-limit: - per-vin-qps: 50 - identity: - store: ${VEHICLE_IDENTITY_STORE:mysql} - mysql: - jdbc-url: ${VEHICLE_IDENTITY_MYSQL_JDBC_URL:jdbc:mysql://127.0.0.1:3306/lingniu_vehicle?useUnicode=true&characterEncoding=utf8&useSSL=false&serverTimezone=Asia/Shanghai} - username: ${VEHICLE_IDENTITY_MYSQL_USERNAME:root} - password: ${VEHICLE_IDENTITY_MYSQL_PASSWORD:} - table-name: ${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_bindings} - initialize-schema: ${VEHICLE_IDENTITY_MYSQL_INITIALIZE_SCHEMA:true} - sink: - kafka: - enabled: true - bootstrap-servers: ${KAFKA_BROKERS:114.55.58.251:9092} - compression-type: zstd - linger-ms: 20 - batch-size: 65536 - acks: all - enable-idempotence: true - node-id: ${KAFKA_NODE_ID:xinda-push-local} - topics: - realtime: ${KAFKA_TOPIC_XINDA_PUSH_EVENT:vehicle.event.xinda-push.v1} - location: ${KAFKA_TOPIC_XINDA_PUSH_EVENT:vehicle.event.xinda-push.v1} - alarm: ${KAFKA_TOPIC_XINDA_PUSH_EVENT:vehicle.event.xinda-push.v1} - session: ${KAFKA_TOPIC_XINDA_PUSH_EVENT:vehicle.event.xinda-push.v1} - media-meta: ${KAFKA_TOPIC_MEDIA_META:vehicle.media.meta.v1} - raw-archive: ${KAFKA_TOPIC_XINDA_PUSH_RAW:vehicle.raw.xinda-push.v1} - dlq: ${KAFKA_TOPIC_XINDA_PUSH_DLQ:vehicle.dlq.xinda-push.v1} - consumer: - enabled: false - archive: - enabled: ${SINK_ARCHIVE_ENABLED:true} - type: local - path: ${SINK_ARCHIVE_PATH:./archive/} - event-file-store: - enabled: false - event-history: - enabled: false - vehicle-state: - enabled: false - vehicle-stat: - enabled: false - -management: - endpoints: - web: - exposure: - include: health,info,metrics,prometheus,env - health: - redis: - enabled: ${MANAGEMENT_HEALTH_REDIS_ENABLED:false} - metrics: - tags: - application: xinda-push-app - -logging: - level: - root: INFO - com.lingniu.ingest: INFO diff --git a/modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppCompositionTest.java b/modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppCompositionTest.java deleted file mode 100644 index 34d7e6f1..00000000 --- a/modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppCompositionTest.java +++ /dev/null @@ -1,69 +0,0 @@ -package com.lingniu.ingest.xindapushapp; - -import com.lingniu.ingest.core.config.IngestCoreAutoConfiguration; -import com.lingniu.ingest.identity.config.VehicleIdentityAutoConfiguration; -import com.lingniu.ingest.inbound.xinda.XindaPushClient; -import com.lingniu.ingest.inbound.xinda.XindaPushEventMapper; -import com.lingniu.ingest.inbound.xinda.XindaPushRealtimeHandler; -import com.lingniu.ingest.inbound.xinda.config.XindaPushAutoConfiguration; -import com.lingniu.ingest.sink.archive.ArchiveStore; -import com.lingniu.ingest.sink.archive.RawArchiveEventSink; -import com.lingniu.ingest.sink.archive.config.SinkArchiveAutoConfiguration; -import com.lingniu.ingest.sink.kafka.KafkaEnvelopeDeadLetterSink; -import com.lingniu.ingest.sink.kafka.KafkaEventSink; -import com.lingniu.ingest.sink.kafka.KafkaSinkAutoConfiguration; -import org.apache.kafka.clients.producer.KafkaProducer; -import org.junit.jupiter.api.Test; -import org.springframework.boot.autoconfigure.AutoConfigurations; -import org.springframework.boot.test.context.runner.ApplicationContextRunner; - -import static org.assertj.core.api.Assertions.assertThat; -import static org.mockito.Mockito.mock; - -class XindaPushAppCompositionTest { - - private final ApplicationContextRunner contextRunner = new ApplicationContextRunner() - .withConfiguration(AutoConfigurations.of( - IngestCoreAutoConfiguration.class, - VehicleIdentityAutoConfiguration.class, - KafkaSinkAutoConfiguration.class, - SinkArchiveAutoConfiguration.class, - XindaPushAutoConfiguration.class)) - .withAllowBeanDefinitionOverriding(true) - .withBean("kafkaProducer", KafkaProducer.class, XindaPushAppCompositionTest::kafkaProducer) - .withPropertyValues( - "lingniu.ingest.xinda-push.enabled=true", - "lingniu.ingest.xinda-push.host=127.0.0.1", - "lingniu.ingest.xinda-push.port=10100", - "lingniu.ingest.xinda-push.username=test", - "lingniu.ingest.xinda-push.password=test", - "lingniu.ingest.identity.store=mysql", - "lingniu.ingest.identity.mysql.initialize-schema=false", - "lingniu.ingest.sink.kafka.enabled=true", - "lingniu.ingest.sink.kafka.bootstrap-servers=localhost:9092", - "lingniu.ingest.sink.kafka.consumer.enabled=false", - "lingniu.ingest.sink.archive.enabled=true", - "lingniu.ingest.sink.archive.path=target/test-archive-xinda", - "lingniu.ingest.event-file-store.enabled=false", - "lingniu.ingest.event-history.enabled=false", - "lingniu.ingest.vehicle-state.enabled=false", - "lingniu.ingest.vehicle-stat.enabled=false"); - - @Test - void createsXindaPushIngressAndKafkaSinksWithoutHistoryStorage() { - contextRunner.run(context -> { - assertThat(context).hasSingleBean(XindaPushClient.class); - assertThat(context).hasSingleBean(XindaPushEventMapper.class); - assertThat(context).hasSingleBean(XindaPushRealtimeHandler.class); - assertThat(context).hasSingleBean(KafkaEventSink.class); - assertThat(context).hasSingleBean(KafkaEnvelopeDeadLetterSink.class); - assertThat(context).hasSingleBean(ArchiveStore.class); - assertThat(context).hasSingleBean(RawArchiveEventSink.class); - }); - } - - @SuppressWarnings("unchecked") - private static KafkaProducer kafkaProducer() { - return mock(KafkaProducer.class); - } -} diff --git a/modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppDefaultsTest.java b/modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppDefaultsTest.java deleted file mode 100644 index e8d118fe..00000000 --- a/modules/apps/xinda-push-app/src/test/java/com/lingniu/ingest/xindapushapp/XindaPushAppDefaultsTest.java +++ /dev/null @@ -1,63 +0,0 @@ -package com.lingniu.ingest.xindapushapp; - -import org.junit.jupiter.api.Test; -import org.springframework.beans.factory.config.YamlPropertiesFactoryBean; -import org.springframework.core.io.ClassPathResource; - -import java.io.IOException; -import java.nio.charset.StandardCharsets; -import java.util.Properties; - -import static org.assertj.core.api.Assertions.assertThat; - -class XindaPushAppDefaultsTest { - - @Test - void applicationDefaultsKeepXindaPushAsKafkaProducerOnly() throws IOException { - Properties properties = applicationProperties(); - - assertThat(properties) - .containsEntry("spring.application.name", "xinda-push-app") - .containsEntry( - "spring.config.import[0]", - "optional:nacos:${spring.application.name}.${NACOS_CONFIG_FILE_EXTENSION:yml}?group=${NACOS_GROUP:DEFAULT_GROUP}&refreshEnabled=${NACOS_REFRESH_ENABLED:true}") - .containsEntry("spring.cloud.nacos.config.server-addr", "${NACOS_SERVER_ADDR:127.0.0.1:8848}") - .containsEntry("server.port", "${HTTP_PORT:20600}") - .containsEntry("lingniu.ingest.xinda-push.enabled", "${XINDA_PUSH_ENABLED:false}") - .containsEntry("lingniu.ingest.xinda-push.host", "${XINDA_PUSH_HOST:}") - .containsEntry("lingniu.ingest.xinda-push.port", "${XINDA_PUSH_PORT:10100}") - .containsEntry("lingniu.ingest.xinda-push.username", "${XINDA_PUSH_USERNAME:}") - .containsEntry("lingniu.ingest.xinda-push.subscribe-msg-ids[0]", "${XINDA_PUSH_MSG_ID_LOCATION:0200}") - .containsEntry("lingniu.ingest.sink.kafka.enabled", true) - .containsEntry("lingniu.ingest.sink.kafka.consumer.enabled", false) - .containsEntry("lingniu.ingest.sink.kafka.topics.realtime", "${KAFKA_TOPIC_XINDA_PUSH_EVENT:vehicle.event.xinda-push.v1}") - .containsEntry("lingniu.ingest.sink.kafka.topics.raw-archive", "${KAFKA_TOPIC_XINDA_PUSH_RAW:vehicle.raw.xinda-push.v1}") - .containsEntry("lingniu.ingest.sink.kafka.topics.dlq", "${KAFKA_TOPIC_XINDA_PUSH_DLQ:vehicle.dlq.xinda-push.v1}") - .containsEntry("lingniu.ingest.identity.store", "${VEHICLE_IDENTITY_STORE:mysql}") - .containsEntry("lingniu.ingest.identity.mysql.table-name", - "${VEHICLE_IDENTITY_MYSQL_TABLE:vehicle_identity_bindings}") - .containsEntry("lingniu.ingest.sink.archive.enabled", "${SINK_ARCHIVE_ENABLED:true}") - .containsEntry("lingniu.ingest.event-file-store.enabled", false) - .containsEntry("lingniu.ingest.event-history.enabled", false) - .containsEntry("lingniu.ingest.vehicle-state.enabled", false) - .containsEntry("lingniu.ingest.vehicle-stat.enabled", false); - - assertThat(applicationYaml()) - .doesNotContain("KAFKA_ENABLED"); - } - - private static Properties applicationProperties() { - YamlPropertiesFactoryBean factory = new YamlPropertiesFactoryBean(); - factory.setResources(new ClassPathResource("application.yml")); - - Properties properties = factory.getObject(); - - assertThat(properties).isNotNull(); - return properties; - } - - private static String applicationYaml() throws IOException { - return new String(new ClassPathResource("application.yml").getInputStream().readAllBytes(), - StandardCharsets.UTF_8); - } -} diff --git a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/ProtocolId.java b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/ProtocolId.java index 53b8ec2a..8814ac33 100644 --- a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/ProtocolId.java +++ b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/ProtocolId.java @@ -10,6 +10,5 @@ public enum ProtocolId { JT808, JT1078, JSATL12, - MQTT_YUTONG, - XINDA_PUSH + MQTT_YUTONG } diff --git a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/VehicleEvent.java b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/VehicleEvent.java index d6ee000e..d4f4e90a 100644 --- a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/VehicleEvent.java +++ b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/VehicleEvent.java @@ -53,7 +53,7 @@ public sealed interface VehicleEvent RealtimePayload payload ) implements VehicleEvent {} - /** 纯位置事件(808 的 0x0200 / Xinda 0200 映射到此)。 */ + /** 纯位置事件(808 的 0x0200 等映射到此)。 */ record Location( String eventId, String vin, diff --git a/modules/core/ingest-core/src/test/java/com/lingniu/ingest/core/dispatcher/HandlerRegistryTest.java b/modules/core/ingest-core/src/test/java/com/lingniu/ingest/core/dispatcher/HandlerRegistryTest.java index 03cd2f27..1eaa8f72 100644 --- a/modules/core/ingest-core/src/test/java/com/lingniu/ingest/core/dispatcher/HandlerRegistryTest.java +++ b/modules/core/ingest-core/src/test/java/com/lingniu/ingest/core/dispatcher/HandlerRegistryTest.java @@ -21,15 +21,15 @@ class HandlerRegistryTest { registry.register(exactDef); registry.register(wildcardDef); - assertThat(registry.resolve(ProtocolId.XINDA_PUSH, 0x0200, 0)) + assertThat(registry.resolve(ProtocolId.JT808, 0x0200, 0)) .containsExactly(exactDef); - assertThat(registry.resolve(ProtocolId.XINDA_PUSH, 0x0500, 0)) + assertThat(registry.resolve(ProtocolId.JT808, 0x0500, 0)) .containsExactly(wildcardDef); } private static HandlerDefinition definition(int command, String desc, Object bean, Method method) { return new HandlerDefinition( - ProtocolId.XINDA_PUSH, + ProtocolId.JT808, command, 0, desc, diff --git a/modules/inbound/inbound-mqtt/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java b/modules/inbound/inbound-mqtt/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java new file mode 100644 index 00000000..f30866f6 --- /dev/null +++ b/modules/inbound/inbound-mqtt/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java @@ -0,0 +1,48 @@ +package com.lingniu.ingest.identity; + +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; + +public final class InMemoryVehicleIdentityService implements VehicleIdentityResolver, VehicleIdentityRegistry { + + private final List bindings = new CopyOnWriteArrayList<>(); + + @Override + public VehicleIdentity resolve(VehicleIdentityLookup lookup) { + if (lookup == null) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.UNKNOWN); + } + if (!lookup.vin().isBlank() && !"unknown".equalsIgnoreCase(lookup.vin())) { + return new VehicleIdentity(lookup.vin(), true, VehicleIdentitySource.EXPLICIT_VIN); + } + for (VehicleIdentityBinding binding : bindings) { + if (binding.protocol() != lookup.protocol()) { + continue; + } + if (!lookup.phone().isBlank() && lookup.phone().equals(binding.phone())) { + return new VehicleIdentity(binding.vin(), true, VehicleIdentitySource.BOUND_PHONE); + } + if (!lookup.deviceId().isBlank() && lookup.deviceId().equals(binding.deviceId())) { + return new VehicleIdentity(binding.vin(), true, VehicleIdentitySource.BOUND_DEVICE_ID); + } + if (!lookup.plate().isBlank() && lookup.plate().equals(binding.plate())) { + return new VehicleIdentity(binding.vin(), true, VehicleIdentitySource.BOUND_PLATE); + } + } + if (!lookup.phone().isBlank()) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.FALLBACK_PHONE); + } + if (!lookup.deviceId().isBlank()) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.FALLBACK_DEVICE_ID); + } + if (!lookup.plate().isBlank()) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.FALLBACK_PLATE); + } + return new VehicleIdentity("unknown", false, VehicleIdentitySource.UNKNOWN); + } + + @Override + public void bind(VehicleIdentityBinding binding) { + bindings.add(binding); + } +} diff --git a/modules/inbound/inbound-xinda-push/pom.xml b/modules/inbound/inbound-xinda-push/pom.xml deleted file mode 100644 index 27e97eb3..00000000 --- a/modules/inbound/inbound-xinda-push/pom.xml +++ /dev/null @@ -1,118 +0,0 @@ - - - 4.0.0 - - com.lingniu.ingest - lingniu-vehicle-ingest - 0.1.0-SNAPSHOT - ../../../pom.xml - - inbound-xinda-push - inbound-xinda-push - - 信达平台推送接入。基于私仓 org.lingniu:gps-push-client:1.0 封装的 - com.gps31.push.netty.PushClient 重写,保留信达真实帧格式, - 登录/心跳/订阅/位置/报警 全程对接成功后把事件直接投到 Kafka。 - - - - - com.lingniu.ingest - ingest-api - - - com.lingniu.ingest - ingest-core - - - com.lingniu.ingest - vehicle-identity - - - org.springframework.boot - spring-boot-starter - - - org.springframework.boot - spring-boot-configuration-processor - true - - - io.netty - netty-all - - - io.netty - netty-common - - - io.netty - netty-buffer - - - io.netty - netty-transport - - - io.netty - netty-resolver - - - io.netty - netty-codec - - - io.netty - netty-handler - - - com.fasterxml.jackson.core - jackson-databind - - - - - org.lingniu - gps-push-client - 1.0 - - - - com.alibaba - fastjson - 1.2.83 - - - - org.apache.commons - commons-collections4 - 4.4 - - - - org.apache.commons - commons-lang3 - - - - commons-codec - commons-codec - - - - org.junit.jupiter - junit-jupiter - test - - - org.assertj - assertj-core - test - - - org.mockito - mockito-core - test - - - diff --git a/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushClient.java b/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushClient.java deleted file mode 100644 index 9715b9a1..00000000 --- a/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushClient.java +++ /dev/null @@ -1,413 +0,0 @@ -package com.lingniu.ingest.inbound.xinda; - -import com.alibaba.fastjson.JSON; -import com.alibaba.fastjson.JSONObject; -import com.gps31.push.netty.PushClient; -import com.gps31.push.netty.PushMsg; -import com.gps31.push.netty.client.TcpClient; -import com.lingniu.ingest.api.ProtocolId; -import com.lingniu.ingest.api.event.RawArchiveKeys; -import com.lingniu.ingest.api.pipeline.RawFrame; -import com.lingniu.ingest.core.dispatcher.Dispatcher; -import com.lingniu.ingest.identity.InMemoryVehicleIdentityService; -import com.lingniu.ingest.identity.VehicleIdentity; -import com.lingniu.ingest.identity.VehicleIdentityLookup; -import com.lingniu.ingest.identity.VehicleIdentityResolver; -import com.lingniu.ingest.identity.VehicleIdentitySource; -import com.lingniu.ingest.inbound.xinda.config.XindaPushProperties; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import java.nio.charset.StandardCharsets; -import java.time.Instant; -import java.util.HashMap; -import java.util.Map; - -/** - * 信达 Push 接入客户端。基于 {@code org.lingniu:gps-push-client:1.0} 的 {@link PushClient} - * 实现,完全使用信达真实协议帧格式(帧头 / 长度 / cmd / json / 帧尾 + 序列号 + 心跳)。 - * - *

行为: - *

    - *
  • {@link #start()} 调用基类 {@code TcpClient.start()} 建立连接并自动重连 - *
  • 基类在 {@code channelActive} 自动发送 {@code 0x0001 登录} 帧,携带 userName / pwd / csTag / desc - *
  • 登录成功收到应答后(由基类内部处理)会自动订阅 {@code subMsgIds}(默认空,需显式调用 setSubMsgIds) - *
  • {@link #messageReceived(TcpClient, PushMsg)} 为业务入口:cmd 字符串区分消息类型 - *
  • 位置 {@code 0200} / 报警 {@code 0300} / 透传 {@code 0401} → RawFrame → Dispatcher - *
- * - *

类继承基类 {@code PushClient},不再需要手工处理帧编解码、心跳、重连、订阅等机制。 - */ -public class XindaPushClient extends PushClient { - - private static final Logger log = LoggerFactory.getLogger(XindaPushClient.class); - - private final XindaPushProperties props; - private final VehicleIdentityResolver identityResolver; - private final Dispatcher dispatcher; - - public XindaPushClient(XindaPushProperties props, Dispatcher dispatcher) { - this(props, new InMemoryVehicleIdentityService(), dispatcher); - } - - public XindaPushClient(XindaPushProperties props, - VehicleIdentityResolver identityResolver, - Dispatcher dispatcher) { - super(); - this.props = props; - this.identityResolver = identityResolver == null ? new InMemoryVehicleIdentityService() : identityResolver; - this.dispatcher = dispatcher; - this.isLog = false; - } - - /** - * 启动客户端:填入 host/port/user/pwd/desc,配置订阅消息 ID,交给基类建立连接。 - */ - public void start() { - if (props.getHost() == null || props.getHost().isBlank()) { - publishOperationalError("startup", "host not configured"); - return; - } - if (props.getUsername() == null || props.getPassword() == null) { - publishOperationalError("startup", "credentials missing"); - return; - } - setHost(props.getHost()); - setPort(props.getPort()); - setUserName(props.getUsername()); - setPwd(props.getPassword()); - setDesc(props.getClientDesc()); - // 订阅消息 ID:信达基类 PushClient.setSubMsgIds 用管道符 "|" 作为分隔符 - // (内部调用 StringUtil.splitStr(s, "|")) - if (props.getSubscribeMsgIds() != null && !props.getSubscribeMsgIds().isEmpty()) { - setSubMsgIds(String.join("|", props.getSubscribeMsgIds())); - } - // 基类 PushClient 构造器默认 isLog=false / isLogBytes=false; - // 这里尊重默认值,不再覆盖为 true,避免 <<>>Log接收 每条帧都被 ERROR 级打印。 - // 真要打开 wire-level 诊断时,临时改为 setLog(true) 即可。 - try { - log.info("[xinda-push] starting host={} port={} user={} subMsgIds={}", - getHost(), getPort(), getUserName(), getSubMsgIds()); - super.start(); - } catch (Exception e) { - log.error("[xinda-push] start failed", e); - publishOperationalError("startup", e.getMessage() == null ? e.getClass().getSimpleName() : e.getMessage()); - } - } - - /** - * Spring 生命周期停止钩子。 - */ - public void stop() { - try { - destroy(); - log.info("[xinda-push] stopped"); - } catch (Exception e) { - log.warn("[xinda-push] stop error", e); - } - } - - // ======================================================================== - // 业务回调:收到完整的 PushMsg 后分派事件 - // ======================================================================== - - @Override - public void messageReceived(TcpClient client, PushMsg msg) throws Exception { - super.messageReceived(client, msg); - String cmd = msg.getCmd(); - if (cmd == null) return; - try { - // 信达只把业务消息派发到 Dispatcher;登录/心跳/订阅 ACK 只影响连接状态和日志。 - switch (cmd) { - case "8001" -> handleLoginAck(msg); - case "8002" -> log.debug("[xinda-push] heartbeat ack"); - case "8003" -> log.info("[xinda-push] subscribe ack json={}", msg.getJson()); - case "0200", "0300", "0401" -> dispatchBusinessMessage(msg, false); - default -> { - log.debug("[xinda-push] rx unknown business cmd={} json={}", cmd, msg.getJson()); - dispatchBusinessMessage(msg, true); - } - } - } catch (Exception e) { - log.warn("[xinda-push] dispatch failed cmd={} json={}", cmd, msg.getJson(), e); - // 推送回调异常也生成错误 RawFrame,保证原始 JSON 不丢。 - publishMessageError(msg, e.getMessage() == null ? e.getClass().getSimpleName() : e.getMessage()); - } - } - - /** - * 收到 8001 登录应答后处理:若 {@code rspResult=0} 则主动发送 0003 订阅命令。 - * - *

信达订阅命令的 body 格式(从旧 {@code PushClientHandler} 和 8003 错误提示反推得到): - *

{@code
-     * {
-     *   "seq":     "1",                             // 序列号
-     *   "action":  "add",                           // 操作: add 订阅 / del 取消
-     *   "msgIds":  "[\"0200\",\"0300\",\"0401\"]"  // JSON 数组字符串(注意是字符串,不是数组)
-     * }
-     * }
- *

- * 其中 {@code msgIds} 的值来自基类 {@link PushClient#getSubCmdSet()},内容由 - * {@link PushClient#setSubMsgIds(String)} 填充(分隔符是 {@code |})。 - */ - private void handleLoginAck(PushMsg msg) throws Exception { - String json = msg.getJson() == null ? "" : msg.getJson(); - log.info("[xinda-push] login ack json={}", json); - JSONObject parsed; - try { - parsed = JSON.parseObject(json); - } catch (Exception e) { - publishOperationalError("login", "login ack parse failed: " - + (e.getMessage() == null ? e.getClass().getSimpleName() : e.getMessage())); - return; - } - String rsp = parsed == null ? null : parsed.getString("rspResult"); - if (!"0".equals(rsp)) { - log.error("[xinda-push] login REJECTED rspResult={}", rsp); - publishOperationalError("login", "login rejected rspResult=" + (rsp == null ? "" : rsp)); - return; - } - Map body = new java.util.HashMap<>(); - body.put("seq", "1"); - body.put("action", "add"); - body.put("msgIds", JSON.toJSONString(getSubCmdSet())); - PushMsg subscribe = getInstance("0003", body); - sendMsg(subscribe); - log.info("[xinda-push] subscribe sent action=add msgIds={}", getSubCmdSet()); - } - - @Override - public void channelActive(TcpClient client) throws Exception { - super.channelActive(client); - log.info("[xinda-push] channel active, login frame sent"); - } - - @Override - public void channelInactive(TcpClient client) throws Exception { - super.channelInactive(client); - log.warn("[xinda-push] channel inactive, will auto-reconnect"); - } - - private void dispatchBusinessMessage(PushMsg msg, boolean unknownCommand) { - String json = msg.getJson() == null ? "" : msg.getJson(); - XindaPushMessage payload = new XindaPushMessage(msg.getCmd(), json); - byte[] raw = json.getBytes(StandardCharsets.UTF_8); - JsonParseResult parsed = parseForMetadata(json); - JSONObject identityFields = parsed.json(); - IdentityResolution identity = resolveIdentity(identityFields); - Map meta = new HashMap<>(); - // 保留外部字段和内部 VIN 的双轨元数据,方便定位绑定错误或第三方字段变更。 - meta.put("source", "xinda-push"); - meta.put("cmd", payload.command()); - putRawFields(meta, identityFields); - meta.put("vin", identity.identity().vin()); - meta.put("externalVin", firstNonBlank(identityFields, "vin")); - meta.put("phone", phone(identityFields)); - meta.put("deviceId", deviceId(identityFields)); - meta.put("plateNo", plateNo(identityFields)); - meta.put("identityResolved", Boolean.toString(identity.identity().resolved())); - meta.put("identitySource", identity.identity().source().name()); - if (identity.errorMessage() != null) { - meta.put("identityError", "true"); - meta.put("identityErrorMessage", identity.errorMessage()); - } - if (unknownCommand) { - meta.put("unknownCommand", "true"); - } - if (payload.malformedCommand()) { - meta.put("malformedCommand", "true"); - } - if (parsed.parseError()) { - meta.put("parseError", "true"); - meta.put("parseErrorMessage", parsed.errorMessage()); - } - meta.put(RawArchiveKeys.META_PARSED_JSON, parsedJson(payload.command(), json, parsed)); - dispatcher.dispatch(new RawFrame( - ProtocolId.XINDA_PUSH, - payload.commandCode(), - 0, - payload, - raw, - meta, - Instant.now())); - } - - private void publishOperationalError(String phase, String reason) { - log.error("[xinda-push] operational error phase={} reason={}", phase, reason); - JSONObject json = new JSONObject(); - json.put("operationalError", true); - json.put("phase", phase == null ? "" : phase); - json.put("reason", reason == null ? "" : reason); - XindaPushMessage payload = new XindaPushMessage("0000", json.toJSONString()); - byte[] raw = payload.json().getBytes(StandardCharsets.UTF_8); - Map meta = new HashMap<>(); - meta.put("source", "xinda-push"); - meta.put("cmd", payload.command()); - meta.put("vin", "unknown"); - meta.put("identityResolved", "false"); - meta.put("identitySource", "UNKNOWN"); - meta.put("operationalError", "true"); - meta.put("phase", phase == null ? "" : phase); - meta.put("reason", reason == null ? "" : reason); - meta.put(RawArchiveKeys.META_PARSED_JSON, - parsedJson(payload.command(), payload.json(), new JsonParseResult(json, false, ""))); - dispatcher.dispatch(new RawFrame( - ProtocolId.XINDA_PUSH, - payload.commandCode(), - 0, - payload, - raw, - meta, - Instant.now())); - } - - private void publishMessageError(PushMsg msg, String reason) { - String cmd = msg == null || msg.getCmd() == null ? "0000" : msg.getCmd(); - String json = msg == null || msg.getJson() == null ? "" : msg.getJson(); - XindaPushMessage payload = new XindaPushMessage(cmd, json); - byte[] raw = json.getBytes(StandardCharsets.UTF_8); - Map meta = new HashMap<>(); - meta.put("source", "xinda-push"); - meta.put("cmd", cmd); - meta.put("vin", "unknown"); - meta.put("identityResolved", "false"); - meta.put("identitySource", "UNKNOWN"); - meta.put("operationalError", "true"); - meta.put("phase", "message"); - meta.put("reason", reason == null ? "" : reason); - JsonParseResult parsed = parseForMetadata(json); - meta.put(RawArchiveKeys.META_PARSED_JSON, parsedJson(cmd, json, parsed)); - dispatcher.dispatch(new RawFrame( - ProtocolId.XINDA_PUSH, - 0, - 0, - payload, - raw, - meta, - Instant.now())); - } - - // ======================================================================== - // JSON helpers - // ======================================================================== - - private static String firstNonBlank(JSONObject node, String... keys) { - if (node == null) return ""; - for (String k : keys) { - Object v = node.get(k); - if (v != null && !v.toString().isBlank()) return v.toString(); - } - return ""; - } - - private IdentityResolution resolveIdentity(JSONObject json) { - try { - return new IdentityResolution(identityResolver.resolve(new VehicleIdentityLookup( - ProtocolId.XINDA_PUSH, - firstNonBlank(json, "vin"), - phone(json), - deviceId(json), - plateNo(json))), null); - } catch (RuntimeException e) { - return new IdentityResolution( - new VehicleIdentity("unknown", false, VehicleIdentitySource.UNKNOWN), - e.getMessage() == null ? "" : e.getMessage()); - } - } - - private static String plateNo(JSONObject node) { - return firstNonBlank(node, "plateNo", "carNo", "carName"); - } - - private static String phone(JSONObject node) { - return firstNonBlank(node, "sim", "phone", "mobile", "simCode", "tmnKey"); - } - - private static String deviceId(JSONObject node) { - return firstNonBlank(node, "deviceId", "terminalId", "imei", "tmnKey", "tmnCode"); - } - - private static void putRawFields(Map meta, JSONObject json) { - if (meta == null || json == null || json.isEmpty()) { - return; - } - for (Map.Entry entry : json.entrySet()) { - if (entry.getKey() == null || entry.getKey().isBlank() || entry.getValue() == null) { - continue; - } - Object value = entry.getValue(); - if (value instanceof JSONObject object) { - for (Map.Entry nested : object.entrySet()) { - if (nested.getKey() != null && !nested.getKey().isBlank() && nested.getValue() != null) { - meta.put("raw." + entry.getKey() + "." + nested.getKey(), nested.getValue().toString()); - } - } - } else { - String text = value.toString(); - JSONObject nestedJson = parseObjectOrNull(text); - if (nestedJson != null) { - for (Map.Entry nested : nestedJson.entrySet()) { - if (nested.getKey() != null && !nested.getKey().isBlank() && nested.getValue() != null) { - meta.put("raw." + entry.getKey() + "." + nested.getKey(), nested.getValue().toString()); - } - } - } else { - meta.put("raw." + entry.getKey(), text); - } - } - } - } - - private static JSONObject parseObjectOrNull(String value) { - if (value == null || value.isBlank() || !value.trim().startsWith("{")) { - return null; - } - try { - return JSON.parseObject(value); - } catch (Exception ignored) { - return null; - } - } - - private static JSONObject parseQuietly(String rawJson) { - if (rawJson == null || rawJson.isBlank()) { - return new JSONObject(); - } - try { - JSONObject parsed = JSON.parseObject(rawJson); - return parsed == null ? new JSONObject() : parsed; - } catch (Exception ignored) { - return new JSONObject(); - } - } - - private static JsonParseResult parseForMetadata(String rawJson) { - if (rawJson == null || rawJson.isBlank()) { - return new JsonParseResult(new JSONObject(), false, ""); - } - try { - JSONObject parsed = JSON.parseObject(rawJson); - return new JsonParseResult(parsed == null ? new JSONObject() : parsed, false, ""); - } catch (Exception e) { - String message = e.getMessage() == null ? e.getClass().getSimpleName() : e.getMessage(); - return new JsonParseResult(new JSONObject(), true, message); - } - } - - private static String parsedJson(String command, String rawJson, JsonParseResult parsed) { - JSONObject root = new JSONObject(true); - root.put("command", command == null ? "" : command); - root.put("parseError", parsed != null && parsed.parseError()); - if (parsed != null && parsed.parseError()) { - root.put("parseErrorMessage", parsed.errorMessage() == null ? "" : parsed.errorMessage()); - root.put("rawJson", rawJson == null ? "" : rawJson); - } else { - root.put("body", parsed == null || parsed.json() == null ? new JSONObject() : parsed.json()); - } - return root.toJSONString(); - } - - private record JsonParseResult(JSONObject json, boolean parseError, String errorMessage) {} - - private record IdentityResolution(VehicleIdentity identity, String errorMessage) {} -} diff --git a/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushEventMapper.java b/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushEventMapper.java deleted file mode 100644 index e9002546..00000000 --- a/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushEventMapper.java +++ /dev/null @@ -1,413 +0,0 @@ -package com.lingniu.ingest.inbound.xinda; - -import com.alibaba.fastjson.JSON; -import com.alibaba.fastjson.JSONObject; -import com.lingniu.ingest.api.ProtocolId; -import com.lingniu.ingest.api.event.AlarmPayload; -import com.lingniu.ingest.api.event.LocationPayload; -import com.lingniu.ingest.api.event.VehicleEvent; -import com.lingniu.ingest.api.spi.EventMapper; -import com.lingniu.ingest.identity.InMemoryVehicleIdentityService; -import com.lingniu.ingest.identity.VehicleIdentity; -import com.lingniu.ingest.identity.VehicleIdentityLookup; -import com.lingniu.ingest.identity.VehicleIdentityRegistry; -import com.lingniu.ingest.identity.VehicleIdentityResolver; -import com.lingniu.ingest.identity.VehicleIdentitySource; -import com.lingniu.ingest.identity.VehicleRegistrationBinding; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import java.nio.charset.StandardCharsets; -import java.time.Instant; -import java.util.LinkedHashSet; -import java.util.List; -import java.util.Map; -import java.util.UUID; - -public final class XindaPushEventMapper implements EventMapper { - - private static final Logger log = LoggerFactory.getLogger(XindaPushEventMapper.class); - private static final VehicleIdentityRegistry NOOP_REGISTRY = binding -> {}; - - private final VehicleIdentityResolver identityResolver; - private final VehicleIdentityRegistry identityRegistry; - - public XindaPushEventMapper() { - this(new InMemoryVehicleIdentityService()); - } - - public XindaPushEventMapper(VehicleIdentityResolver identityResolver) { - this.identityResolver = identityResolver == null ? new InMemoryVehicleIdentityService() : identityResolver; - this.identityRegistry = this.identityResolver instanceof VehicleIdentityRegistry registry ? registry : NOOP_REGISTRY; - } - - public XindaPushEventMapper(VehicleIdentityResolver identityResolver, VehicleIdentityRegistry identityRegistry) { - this.identityResolver = identityResolver == null ? new InMemoryVehicleIdentityService() : identityResolver; - this.identityRegistry = identityRegistry == null ? NOOP_REGISTRY : identityRegistry; - } - - @Override - public List toEvents(XindaPushMessage message) { - if (message == null || message.command().isBlank()) { - return List.of(); - } - JSONObject json; - try { - json = parse(message.json()); - } catch (Exception e) { - // JSON 解析失败不丢弃,转 passthrough 保留 rawJson,后续可通过历史库回查。 - return List.of(passthrough(message, new JSONObject(), true, false, - e.getMessage() == null ? e.getClass().getSimpleName() : e.getMessage())); - } - if (json.getBooleanValue("operationalError")) { - // 接入侧运行错误也统一作为 passthrough 事件进入管线,便于监控和审计。 - return List.of(passthrough(message, json, false, false)); - } - return switch (message.command()) { - case "0200" -> List.of(location(message, json)); - case "0300" -> List.of(alarm(message, json)); - case "0401" -> List.of(passthrough(message, json, false, false)); - default -> List.of(passthrough(message, json, false, true)); - }; - } - - private VehicleEvent.Location location(XindaPushMessage message, JSONObject json) { - IdentityResolution identity = resolveIdentity(json); - double lon = normalizeCoord(asDouble(json, "lng", "lon", "longitude")); - double lat = normalizeCoord(asDouble(json, "lat", "latitude")); - double speed = asDouble(json, "speed"); - Double totalMileageKm = optionalDouble(json, - "totalMileage", "totalMileageKm", "mileage", "mileageKm", "total_mileage"); - double direction = asDouble(json, "drct", "direction"); - long alarmFlag = (long) asDouble(json, "alarmStts", "alarmFlag"); - long statusFlag = (long) asDouble(json, "statusStts", "statusFlag"); - Instant now = Instant.now(); - return new VehicleEvent.Location( - UUID.randomUUID().toString(), - identity.vin(), - ProtocolId.XINDA_PUSH, - now, - now, - null, - metadata(message, identity, false, false), - new LocationPayload(lon, lat, 0.0, speed, direction, alarmFlag, statusFlag, totalMileageKm)); - } - - private VehicleEvent.Alarm alarm(XindaPushMessage message, JSONObject json) { - IdentityResolution identity = resolveIdentity(json); - double lon = normalizeCoord(asDouble(json, "lng", "lon", "longitude")); - double lat = normalizeCoord(asDouble(json, "lat", "latitude")); - String alarmType = firstNonBlank(json, "alarmType", "alarmCode", "type", "warnType"); - String alarmName = firstNonBlank(json, "alarmName", "alarmDesc", "alarmContent", "content", "warnName"); - int alarmTypeCode = alarmTypeCode(alarmType); - boolean hydrogenLeak = isHydrogenLeak(alarmType, alarmName, firstNonBlank(json, "safetyCategory")); - AlarmPayload.AlarmLevel level = alarmLevel(asInt(json, "alarmLevel", "level"), hydrogenLeak); - Instant now = Instant.now(); - return new VehicleEvent.Alarm( - UUID.randomUUID().toString(), - identity.vin(), - ProtocolId.XINDA_PUSH, - now, - now, - null, - alarmMetadata(message, identity, hydrogenLeak), - new AlarmPayload( - level, - alarmTypeCode, - alarmType.isBlank() ? alarmName : alarmType, - faultCodes(json, alarmType, alarmName), - activeBits(json, alarmType, alarmName, hydrogenLeak), - lon, - lat, - hydrogenLeak ? AlarmPayload.SafetyCategory.HYDROGEN_LEAK : AlarmPayload.SafetyCategory.GENERAL, - hydrogenLeak, - hydrogenLeak ? AlarmPayload.HydrogenLeakLevel.CRITICAL : AlarmPayload.HydrogenLeakLevel.NONE, - hydrogenLeak)); - } - - private VehicleEvent.Passthrough passthrough(XindaPushMessage message, - JSONObject json, - boolean parseError, - boolean unknownCommand) { - return passthrough(message, json, parseError, unknownCommand, ""); - } - - private VehicleEvent.Passthrough passthrough(XindaPushMessage message, - JSONObject json, - boolean parseError, - boolean unknownCommand, - String parseErrorMessage) { - IdentityResolution identity = resolveIdentity(json); - byte[] raw = message.json().getBytes(StandardCharsets.UTF_8); - Instant now = Instant.now(); - return new VehicleEvent.Passthrough( - UUID.randomUUID().toString(), - identity.vin(), - ProtocolId.XINDA_PUSH, - now, - now, - null, - metadata(message, identity, parseError, unknownCommand, parseErrorMessage), - message.commandCode(), - raw); - } - - private IdentityResolution resolveIdentity(JSONObject json) { - try { - IdentityResolution resolution = IdentityResolution.ok(identityResolver.resolve(new VehicleIdentityLookup( - ProtocolId.XINDA_PUSH, - firstNonBlank(json, "vin"), - phone(json), - deviceId(json), - plateNo(json)))); - registerIdentity(json, resolution); - return resolution; - } catch (RuntimeException e) { - return new IdentityResolution( - new VehicleIdentity("unknown", false, VehicleIdentitySource.UNKNOWN), - e.getMessage() == null ? "" : e.getMessage()); - } - } - - private static Map metadata(XindaPushMessage message, - IdentityResolution identity, - boolean parseError, - boolean unknownCommand) { - return metadata(message, identity, parseError, unknownCommand, ""); - } - - private static Map metadata(XindaPushMessage message, - IdentityResolution identity, - boolean parseError, - boolean unknownCommand, - String parseErrorMessage) { - JSONObject json = parseQuietly(message.json()); - Map out = new java.util.LinkedHashMap<>(); - out.put("source", "xinda-push"); - out.put("cmd", message.command()); - out.put("vin", normalizedVin(identity)); - out.put("externalVin", firstNonBlank(json, "vin")); - out.put("phone", phone(json)); - out.put("deviceId", deviceId(json)); - out.put("plateNo", plateNo(json)); - out.put("identityResolved", Boolean.toString(identity.identity().resolved())); - out.put("identitySource", identity.identity().source().name()); - if (identity.errorMessage() != null) { - out.put("identityError", "true"); - out.put("identityErrorMessage", identity.errorMessage()); - } - out.put("parseError", Boolean.toString(parseError)); - if (parseError) { - out.put("parseErrorMessage", parseErrorMessage == null ? "" : parseErrorMessage); - } - out.put("unknownCommand", Boolean.toString(unknownCommand)); - out.put("malformedCommand", Boolean.toString(message.malformedCommand())); - if (json.getBooleanValue("operationalError")) { - out.put("operationalError", "true"); - out.put("phase", firstNonBlank(json, "phase")); - out.put("reason", firstNonBlank(json, "reason")); - } - return Map.copyOf(out); - } - - private static String normalizedVin(IdentityResolution identity) { - if (identity == null || identity.vin() == null || identity.vin().isBlank()) { - return "unknown"; - } - return identity.vin(); - } - - private static Map alarmMetadata(XindaPushMessage message, - IdentityResolution identity, - boolean hydrogenLeak) { - JSONObject json = parseQuietly(message.json()); - Map out = new java.util.LinkedHashMap<>(metadata(message, identity, false, false)); - out.put("alarmType", firstNonBlank(json, "alarmType", "alarmCode", "type", "warnType")); - out.put("alarmName", firstNonBlank(json, "alarmName", "alarmDesc", "alarmContent", "content", "warnName")); - out.put("alarmLevel", firstNonBlank(json, "alarmLevel", "level")); - out.put("hydrogenLeakDetected", Boolean.toString(hydrogenLeak)); - if (hydrogenLeak) { - out.put("safetyCategory", AlarmPayload.SafetyCategory.HYDROGEN_LEAK.name()); - out.put("hydrogenLeakLevel", AlarmPayload.HydrogenLeakLevel.CRITICAL.name()); - out.put("hydrogenLeakActionRequired", "true"); - } - return Map.copyOf(out); - } - - private static JSONObject parse(String raw) { - if (raw == null || raw.isBlank()) { - return new JSONObject(); - } - return JSON.parseObject(raw); - } - - private static JSONObject parseQuietly(String raw) { - try { - return parse(raw); - } catch (Exception ignored) { - return new JSONObject(); - } - } - - private static String firstNonBlank(JSONObject node, String... keys) { - if (node == null) return ""; - for (String k : keys) { - Object v = node.get(k); - if (v != null && !v.toString().isBlank()) return v.toString(); - } - return ""; - } - - private static String plateNo(JSONObject node) { - return firstNonBlank(node, "plateNo", "carNo", "carName"); - } - - private static String phone(JSONObject node) { - return firstNonBlank(node, "sim", "phone", "mobile", "simCode", "tmnKey"); - } - - private static String deviceId(JSONObject node) { - return firstNonBlank(node, "deviceId", "terminalId", "imei", "tmnKey", "tmnCode"); - } - - private void registerIdentity(JSONObject json, IdentityResolution identity) { - String phone = phone(json); - String deviceId = deviceId(json); - String plate = plateNo(json); - if (phone.isBlank() && deviceId.isBlank() && plate.isBlank()) { - return; - } - String vin = identity != null && identity.identity().resolved() ? identity.identity().vin() : "unknown"; - try { - identityRegistry.register(new VehicleRegistrationBinding( - ProtocolId.XINDA_PUSH, - vin, - phone, - deviceId, - plate, - null, - null, - "", - "", - optionalInteger(json, "plateColor", "carColor"))); - } catch (RuntimeException e) { - log.warn("[xinda-push] vehicle identity registration failed plate={} phone={} deviceId={}", - plate, phone, deviceId, e); - } - } - - private static double asDouble(JSONObject node, String... keys) { - Double value = optionalDouble(node, keys); - return value == null ? 0.0 : value; - } - - private static Double optionalDouble(JSONObject node, String... keys) { - if (node == null) return null; - for (String k : keys) { - Object v = node.get(k); - if (v instanceof Number n) return n.doubleValue(); - if (v != null) { - try { - return Double.parseDouble(v.toString()); - } catch (NumberFormatException ignored) { - } - } - } - return null; - } - - private static int asInt(JSONObject node, String... keys) { - return (int) asDouble(node, keys); - } - - private static Integer optionalInteger(JSONObject node, String... keys) { - Double value = optionalDouble(node, keys); - return value == null ? null : value.intValue(); - } - - private static double normalizeCoord(double value) { - if (Math.abs(value) > 1000) { - return value / 1_000_000.0; - } - return value; - } - - private static boolean isHydrogenLeak(String... values) { - for (String value : values) { - if (value == null || value.isBlank()) { - continue; - } - String v = value.trim().toLowerCase(java.util.Locale.ROOT); - if (v.contains("hydrogen_leak") - || v.contains("h2_leak") - || v.contains("hydrogenleak") - || v.contains("氢气泄") - || v.contains("氢泄") - || v.contains("漏氢")) { - return true; - } - } - return false; - } - - private static AlarmPayload.AlarmLevel alarmLevel(int code, boolean hydrogenLeak) { - if (hydrogenLeak || code >= 3) { - return AlarmPayload.AlarmLevel.CRITICAL; - } - if (code == 2) { - return AlarmPayload.AlarmLevel.MAJOR; - } - if (code == 1) { - return AlarmPayload.AlarmLevel.MINOR; - } - return AlarmPayload.AlarmLevel.INFO; - } - - private static int alarmTypeCode(String alarmType) { - if (alarmType == null || alarmType.isBlank()) { - return 0; - } - try { - return Integer.parseInt(alarmType.trim()); - } catch (NumberFormatException ignored) { - return Math.abs(alarmType.trim().hashCode()); - } - } - - private static List faultCodes(JSONObject json, String alarmType, String alarmName) { - LinkedHashSet codes = new LinkedHashSet<>(); - addIfPresent(codes, alarmType); - addIfPresent(codes, firstNonBlank(json, "faultCode", "faultCodes", "alarmCode")); - addIfPresent(codes, alarmName); - return List.copyOf(codes); - } - - private static java.util.Set activeBits(JSONObject json, - String alarmType, - String alarmName, - boolean hydrogenLeak) { - LinkedHashSet bits = new LinkedHashSet<>(); - addIfPresent(bits, alarmType); - addIfPresent(bits, alarmName); - if (hydrogenLeak) { - bits.add("HYDROGEN_LEAK"); - } - return java.util.Set.copyOf(bits); - } - - private static void addIfPresent(java.util.Set out, String value) { - if (value != null && !value.isBlank()) { - out.add(value.trim()); - } - } - - private record IdentityResolution(VehicleIdentity identity, String errorMessage) { - private static IdentityResolution ok(VehicleIdentity identity) { - return new IdentityResolution(identity, null); - } - - private String vin() { - return identity.vin(); - } - } -} diff --git a/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushMessage.java b/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushMessage.java deleted file mode 100644 index 49ee818a..00000000 --- a/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushMessage.java +++ /dev/null @@ -1,32 +0,0 @@ -package com.lingniu.ingest.inbound.xinda; - -public record XindaPushMessage(String command, String json) { - - public XindaPushMessage { - command = command == null ? "" : command.trim(); - json = json == null ? "" : json; - } - - public int commandCode() { - if (command.isBlank()) { - return 0; - } - try { - return Integer.parseInt(command, 16); - } catch (NumberFormatException e) { - return 0; - } - } - - public boolean malformedCommand() { - if (command.isBlank()) { - return false; - } - try { - Integer.parseInt(command, 16); - return false; - } catch (NumberFormatException e) { - return true; - } - } -} diff --git a/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushRealtimeHandler.java b/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushRealtimeHandler.java deleted file mode 100644 index 680fb293..00000000 --- a/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/XindaPushRealtimeHandler.java +++ /dev/null @@ -1,37 +0,0 @@ -package com.lingniu.ingest.inbound.xinda; - -import com.lingniu.ingest.api.ProtocolId; -import com.lingniu.ingest.api.annotation.EventEmit; -import com.lingniu.ingest.api.annotation.MessageMapping; -import com.lingniu.ingest.api.annotation.ProtocolHandler; -import com.lingniu.ingest.api.event.VehicleEvent; - -import java.util.List; - -@ProtocolHandler(protocol = ProtocolId.XINDA_PUSH) -public final class XindaPushRealtimeHandler { - - private final XindaPushEventMapper mapper; - - public XindaPushRealtimeHandler(XindaPushEventMapper mapper) { - this.mapper = mapper; - } - - @MessageMapping(command = 0x0200, desc = "信达定位") - @EventEmit(VehicleEvent.Location.class) - public List onLocation(XindaPushMessage message) { - return mapper.toEvents(message); - } - - @MessageMapping(command = 0x0300, desc = "信达报警") - @EventEmit(VehicleEvent.Alarm.class) - public List onAlarm(XindaPushMessage message) { - return mapper.toEvents(message); - } - - @MessageMapping(desc = "信达透传/未知业务消息兜底") - @EventEmit(VehicleEvent.Passthrough.class) - public List onPassthrough(XindaPushMessage message) { - return mapper.toEvents(message); - } -} diff --git a/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/config/XindaPushAutoConfiguration.java b/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/config/XindaPushAutoConfiguration.java deleted file mode 100644 index e8fd63d0..00000000 --- a/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/config/XindaPushAutoConfiguration.java +++ /dev/null @@ -1,52 +0,0 @@ -package com.lingniu.ingest.inbound.xinda.config; - -import com.lingniu.ingest.core.dispatcher.Dispatcher; -import com.lingniu.ingest.identity.VehicleIdentityRegistry; -import com.lingniu.ingest.identity.VehicleIdentityResolver; -import com.lingniu.ingest.inbound.xinda.XindaPushClient; -import com.lingniu.ingest.inbound.xinda.XindaPushEventMapper; -import com.lingniu.ingest.inbound.xinda.XindaPushRealtimeHandler; -import org.springframework.boot.ApplicationArguments; -import org.springframework.boot.ApplicationRunner; -import org.springframework.boot.autoconfigure.AutoConfiguration; -import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; -import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; -import org.springframework.boot.context.properties.EnableConfigurationProperties; -import org.springframework.context.annotation.Bean; - -@AutoConfiguration -@EnableConfigurationProperties(XindaPushProperties.class) -@ConditionalOnProperty(prefix = "lingniu.ingest.xinda-push", name = "enabled", havingValue = "true") -public class XindaPushAutoConfiguration { - - @Bean(destroyMethod = "stop") - @ConditionalOnMissingBean - public XindaPushClient xindaPushClient(XindaPushProperties props, - VehicleIdentityResolver identityResolver, - Dispatcher dispatcher) { - return new XindaPushClient(props, identityResolver, dispatcher); - } - - @Bean - @ConditionalOnMissingBean - public XindaPushEventMapper xindaPushEventMapper(VehicleIdentityResolver identityResolver, - VehicleIdentityRegistry identityRegistry) { - return new XindaPushEventMapper(identityResolver, identityRegistry); - } - - @Bean - @ConditionalOnMissingBean - public XindaPushRealtimeHandler xindaPushRealtimeHandler(XindaPushEventMapper mapper) { - return new XindaPushRealtimeHandler(mapper); - } - - @Bean - public ApplicationRunner xindaPushStartupRunner(XindaPushClient client) { - return new ApplicationRunner() { - @Override - public void run(ApplicationArguments args) { - client.start(); - } - }; - } -} diff --git a/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/config/XindaPushProperties.java b/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/config/XindaPushProperties.java deleted file mode 100644 index 92838d03..00000000 --- a/modules/inbound/inbound-xinda-push/src/main/java/com/lingniu/ingest/inbound/xinda/config/XindaPushProperties.java +++ /dev/null @@ -1,45 +0,0 @@ -package com.lingniu.ingest.inbound.xinda.config; - -import org.springframework.boot.context.properties.ConfigurationProperties; - -import java.util.ArrayList; -import java.util.List; - -/** - * 信达 Push 接入配置。所有敏感字段均可通过环境变量注入,彻底取代旧代码里的硬编码。 - * - *

该入口属于第三方推送协议,与 GB32960 TCP 接收链路并列。生产排查 32960 时, - * 应先确认本开关是否关闭,避免把信达推送事件误认为 32960 原始包。 - */ -@ConfigurationProperties(prefix = "lingniu.ingest.xinda-push") -public class XindaPushProperties { - - /** 总开关。关闭时不会向信达平台建连/订阅。 */ - private boolean enabled = false; - private String host; - private int port = 10100; - private String username; - private String password; - private String clientDesc = "lingniu-ingest"; - /** 默认订阅常见定位/轨迹类消息,新增业务消息需要先在 mapper 中确认字段映射。 */ - private List subscribeMsgIds = new ArrayList<>(List.of("0200", "0300", "0401")); - /** 断线重连间隔(秒)。 */ - private int reconnectIntervalSec = 5; - - public boolean isEnabled() { return enabled; } - public void setEnabled(boolean enabled) { this.enabled = enabled; } - public String getHost() { return host; } - public void setHost(String host) { this.host = host; } - public int getPort() { return port; } - public void setPort(int port) { this.port = port; } - public String getUsername() { return username; } - public void setUsername(String username) { this.username = username; } - public String getPassword() { return password; } - public void setPassword(String password) { this.password = password; } - public String getClientDesc() { return clientDesc; } - public void setClientDesc(String clientDesc) { this.clientDesc = clientDesc; } - public List getSubscribeMsgIds() { return subscribeMsgIds; } - public void setSubscribeMsgIds(List subscribeMsgIds) { this.subscribeMsgIds = subscribeMsgIds; } - public int getReconnectIntervalSec() { return reconnectIntervalSec; } - public void setReconnectIntervalSec(int reconnectIntervalSec) { this.reconnectIntervalSec = reconnectIntervalSec; } -} diff --git a/modules/inbound/inbound-xinda-push/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports b/modules/inbound/inbound-xinda-push/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports deleted file mode 100644 index 43286331..00000000 --- a/modules/inbound/inbound-xinda-push/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports +++ /dev/null @@ -1 +0,0 @@ -com.lingniu.ingest.inbound.xinda.config.XindaPushAutoConfiguration diff --git a/modules/inbound/inbound-xinda-push/src/test/java/com/lingniu/ingest/inbound/xinda/XindaPushClientTest.java b/modules/inbound/inbound-xinda-push/src/test/java/com/lingniu/ingest/inbound/xinda/XindaPushClientTest.java deleted file mode 100644 index a299d2fd..00000000 --- a/modules/inbound/inbound-xinda-push/src/test/java/com/lingniu/ingest/inbound/xinda/XindaPushClientTest.java +++ /dev/null @@ -1,272 +0,0 @@ -package com.lingniu.ingest.inbound.xinda; - -import com.gps31.push.netty.PushMsg; -import com.lingniu.ingest.api.ProtocolId; -import com.lingniu.ingest.api.pipeline.RawFrame; -import com.lingniu.ingest.core.dispatcher.Dispatcher; -import com.lingniu.ingest.identity.InMemoryVehicleIdentityService; -import com.lingniu.ingest.identity.VehicleIdentityBinding; -import com.lingniu.ingest.inbound.xinda.config.XindaPushProperties; -import org.junit.jupiter.api.Test; -import org.mockito.ArgumentCaptor; - -import java.nio.charset.StandardCharsets; - -import static org.assertj.core.api.Assertions.assertThat; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.verify; - -class XindaPushClientTest { - - @Test - void missingHostIsDispatchedAsOperationalPassthroughCandidate() { - Dispatcher dispatcher = mock(Dispatcher.class); - XindaPushProperties props = new XindaPushProperties(); - props.setUsername("user"); - props.setPassword("pwd"); - XindaPushClient client = new XindaPushClient(props, dispatcher); - - client.start(); - - ArgumentCaptor captor = ArgumentCaptor.forClass(RawFrame.class); - verify(dispatcher).dispatch(captor.capture()); - RawFrame frame = captor.getValue(); - assertThat(frame.protocolId()).isEqualTo(ProtocolId.XINDA_PUSH); - assertThat(frame.command()).isZero(); - assertThat(frame.payload()).isInstanceOf(XindaPushMessage.class); - assertThat(frame.rawBytes()).asString(StandardCharsets.UTF_8) - .contains("\"operationalError\":true") - .contains("\"phase\":\"startup\"") - .contains("\"reason\":\"host not configured\""); - assertThat(frame.sourceMeta()) - .containsEntry("source", "xinda-push") - .containsEntry("operationalError", "true") - .containsEntry("phase", "startup") - .containsEntry("reason", "host not configured") - .containsEntry("vin", "unknown") - .containsEntry("identityResolved", "false") - .containsEntry("identitySource", "UNKNOWN"); - } - - @Test - void unknownBusinessCommandIsDispatchedForArchiveAndPassthrough() throws Exception { - Dispatcher dispatcher = mock(Dispatcher.class); - XindaPushClient client = new XindaPushClient(new XindaPushProperties(), dispatcher); - PushMsg msg = new PushMsg(); - msg.init("0500", "{\"vin\":\"LNVIN000000000303\",\"vendor\":\"extension\"}"); - - client.messageReceived(null, msg); - - ArgumentCaptor captor = ArgumentCaptor.forClass(RawFrame.class); - verify(dispatcher).dispatch(captor.capture()); - RawFrame frame = captor.getValue(); - assertThat(frame.protocolId()).isEqualTo(ProtocolId.XINDA_PUSH); - assertThat(frame.command()).isEqualTo(0x0500); - assertThat(frame.payload()).isEqualTo(new XindaPushMessage("0500", - "{\"vin\":\"LNVIN000000000303\",\"vendor\":\"extension\"}")); - assertThat(frame.rawBytes()).containsExactly( - "{\"vin\":\"LNVIN000000000303\",\"vendor\":\"extension\"}".getBytes(StandardCharsets.UTF_8)); - assertThat(frame.sourceMeta()) - .containsEntry("source", "xinda-push") - .containsEntry("cmd", "0500") - .containsEntry("vin", "LNVIN000000000303") - .containsEntry("unknownCommand", "true"); - } - - @Test - void businessMessageRawFrameUsesResolvedVehicleIdentity() throws Exception { - Dispatcher dispatcher = mock(Dispatcher.class); - InMemoryVehicleIdentityService identities = new InMemoryVehicleIdentityService(); - identities.bind(new VehicleIdentityBinding( - ProtocolId.XINDA_PUSH, - "LNVIN000000XINDA1", - "", - "XD-DEVICE-001", - "粤B12345")); - XindaPushClient client = new XindaPushClient(new XindaPushProperties(), identities, dispatcher); - PushMsg msg = new PushMsg(); - msg.init("0200", "{\"deviceId\":\"XD-DEVICE-001\",\"plateNo\":\"粤B12345\",\"lng\":113.1,\"lat\":23.1}"); - - client.messageReceived(null, msg); - - ArgumentCaptor captor = ArgumentCaptor.forClass(RawFrame.class); - verify(dispatcher).dispatch(captor.capture()); - RawFrame frame = captor.getValue(); - assertThat(frame.sourceMeta()) - .containsEntry("vin", "LNVIN000000XINDA1") - .containsEntry("deviceId", "XD-DEVICE-001") - .containsEntry("plateNo", "粤B12345") - .containsEntry("identityResolved", "true") - .containsEntry("identitySource", "BOUND_DEVICE_ID"); - } - - @Test - void businessMessageRawFrameResolvesIdentityByCarNamePlateAlias() throws Exception { - Dispatcher dispatcher = mock(Dispatcher.class); - InMemoryVehicleIdentityService identities = new InMemoryVehicleIdentityService(); - identities.bind(new VehicleIdentityBinding( - ProtocolId.XINDA_PUSH, - "LB9A32A24P0LS1270", - "", - "", - "粤BD12345")); - XindaPushClient client = new XindaPushClient(new XindaPushProperties(), identities, dispatcher); - PushMsg msg = new PushMsg(); - msg.init("0200", "{\"carName\":\"粤BD12345\",\"tmnKey\":\"XD-TMN-001\",\"lng\":113.1,\"lat\":23.1}"); - - client.messageReceived(null, msg); - - ArgumentCaptor captor = ArgumentCaptor.forClass(RawFrame.class); - verify(dispatcher).dispatch(captor.capture()); - RawFrame frame = captor.getValue(); - assertThat(frame.sourceMeta()) - .containsEntry("vin", "LB9A32A24P0LS1270") - .containsEntry("plateNo", "粤BD12345") - .containsEntry("deviceId", "XD-TMN-001") - .containsEntry("raw.carName", "粤BD12345") - .containsEntry("raw.tmnKey", "XD-TMN-001") - .containsEntry("raw.lng", "113.1") - .containsEntry("identityResolved", "true") - .containsEntry("identitySource", "BOUND_PLATE"); - } - - @Test - void malformedCommandIsStillDispatchedForArchiveAndPassthrough() throws Exception { - Dispatcher dispatcher = mock(Dispatcher.class); - XindaPushClient client = new XindaPushClient(new XindaPushProperties(), dispatcher); - PushMsg msg = new PushMsg(); - msg.init("bad-cmd", "{\"vin\":\"LNVIN000000000404\",\"vendor\":\"broken\"}"); - - client.messageReceived(null, msg); - - ArgumentCaptor captor = ArgumentCaptor.forClass(RawFrame.class); - verify(dispatcher).dispatch(captor.capture()); - RawFrame frame = captor.getValue(); - assertThat(frame.protocolId()).isEqualTo(ProtocolId.XINDA_PUSH); - assertThat(frame.command()).isZero(); - assertThat(frame.payload()).isEqualTo(new XindaPushMessage("bad-cmd", - "{\"vin\":\"LNVIN000000000404\",\"vendor\":\"broken\"}")); - assertThat(frame.rawBytes()).containsExactly( - "{\"vin\":\"LNVIN000000000404\",\"vendor\":\"broken\"}".getBytes(StandardCharsets.UTF_8)); - assertThat(frame.sourceMeta()) - .containsEntry("source", "xinda-push") - .containsEntry("cmd", "bad-cmd") - .containsEntry("vin", "LNVIN000000000404") - .containsEntry("unknownCommand", "true") - .containsEntry("malformedCommand", "true"); - } - - @Test - void malformedJsonBusinessMessageMarksRawFrameForArchiveDiagnostics() throws Exception { - Dispatcher dispatcher = mock(Dispatcher.class); - XindaPushClient client = new XindaPushClient(new XindaPushProperties(), dispatcher); - PushMsg msg = new PushMsg(); - msg.init("0200", "{bad-json"); - - client.messageReceived(null, msg); - - ArgumentCaptor captor = ArgumentCaptor.forClass(RawFrame.class); - verify(dispatcher).dispatch(captor.capture()); - RawFrame frame = captor.getValue(); - assertThat(frame.protocolId()).isEqualTo(ProtocolId.XINDA_PUSH); - assertThat(frame.command()).isEqualTo(0x0200); - assertThat(frame.rawBytes()).containsExactly("{bad-json".getBytes(StandardCharsets.UTF_8)); - assertThat(frame.sourceMeta()) - .containsEntry("source", "xinda-push") - .containsEntry("cmd", "0200") - .containsEntry("vin", "unknown") - .containsEntry("identityResolved", "false") - .containsEntry("identitySource", "UNKNOWN") - .containsEntry("parseError", "true") - .containsKey("parseErrorMessage"); - assertThat(frame.sourceMeta().get("parseErrorMessage")).contains("expect"); - } - - @Test - void identityResolverFailureStillDispatchesBusinessMessageWithUnknownIdentity() throws Exception { - Dispatcher dispatcher = mock(Dispatcher.class); - XindaPushClient client = new XindaPushClient( - new XindaPushProperties(), - lookup -> { - throw new IllegalStateException("identity backend unavailable"); - }, - dispatcher); - PushMsg msg = new PushMsg(); - msg.init("0200", "{\"deviceId\":\"XD-DEVICE-001\",\"lng\":113.1,\"lat\":23.1}"); - - client.messageReceived(null, msg); - - ArgumentCaptor captor = ArgumentCaptor.forClass(RawFrame.class); - verify(dispatcher).dispatch(captor.capture()); - RawFrame frame = captor.getValue(); - assertThat(frame.protocolId()).isEqualTo(ProtocolId.XINDA_PUSH); - assertThat(frame.command()).isEqualTo(0x0200); - assertThat(frame.rawBytes()).containsExactly( - "{\"deviceId\":\"XD-DEVICE-001\",\"lng\":113.1,\"lat\":23.1}".getBytes(StandardCharsets.UTF_8)); - assertThat(frame.sourceMeta()) - .containsEntry("source", "xinda-push") - .containsEntry("cmd", "0200") - .containsEntry("vin", "unknown") - .containsEntry("identityResolved", "false") - .containsEntry("identitySource", "UNKNOWN") - .containsEntry("identityError", "true") - .containsEntry("identityErrorMessage", "identity backend unavailable"); - assertThat(frame.sourceMeta()).doesNotContainKeys("operationalError", "phase", "reason"); - } - - @Test - void rejectedLoginAckIsDispatchedAsOperationalPassthroughCandidate() throws Exception { - Dispatcher dispatcher = mock(Dispatcher.class); - XindaPushClient client = new XindaPushClient(new XindaPushProperties(), dispatcher); - PushMsg msg = new PushMsg(); - msg.init("8001", "{\"rspResult\":\"1\",\"message\":\"bad credentials\"}"); - - client.messageReceived(null, msg); - - ArgumentCaptor captor = ArgumentCaptor.forClass(RawFrame.class); - verify(dispatcher).dispatch(captor.capture()); - RawFrame frame = captor.getValue(); - assertThat(frame.protocolId()).isEqualTo(ProtocolId.XINDA_PUSH); - assertThat(frame.command()).isZero(); - assertThat(frame.rawBytes()).asString(StandardCharsets.UTF_8) - .contains("\"operationalError\":true") - .contains("\"phase\":\"login\"") - .contains("login rejected rspResult=1"); - assertThat(frame.sourceMeta()) - .containsEntry("source", "xinda-push") - .containsEntry("operationalError", "true") - .containsEntry("phase", "login") - .containsEntry("vin", "unknown") - .containsEntry("identityResolved", "false") - .containsEntry("identitySource", "UNKNOWN"); - assertThat(frame.sourceMeta().get("reason")).contains("login rejected rspResult=1"); - } - - @Test - void malformedLoginAckIsDispatchedAsOperationalPassthroughCandidate() throws Exception { - Dispatcher dispatcher = mock(Dispatcher.class); - XindaPushClient client = new XindaPushClient(new XindaPushProperties(), dispatcher); - PushMsg msg = new PushMsg(); - msg.init("8001", "{bad-json"); - - client.messageReceived(null, msg); - - ArgumentCaptor captor = ArgumentCaptor.forClass(RawFrame.class); - verify(dispatcher).dispatch(captor.capture()); - RawFrame frame = captor.getValue(); - assertThat(frame.protocolId()).isEqualTo(ProtocolId.XINDA_PUSH); - assertThat(frame.command()).isZero(); - assertThat(frame.rawBytes()).asString(StandardCharsets.UTF_8) - .contains("\"operationalError\":true") - .contains("\"phase\":\"login\"") - .contains("login ack parse failed"); - assertThat(frame.sourceMeta()) - .containsEntry("source", "xinda-push") - .containsEntry("operationalError", "true") - .containsEntry("phase", "login") - .containsEntry("vin", "unknown") - .containsEntry("identityResolved", "false") - .containsEntry("identitySource", "UNKNOWN"); - assertThat(frame.sourceMeta().get("reason")).contains("login ack parse failed"); - } -} diff --git a/modules/inbound/inbound-xinda-push/src/test/java/com/lingniu/ingest/inbound/xinda/XindaPushEventMapperTest.java b/modules/inbound/inbound-xinda-push/src/test/java/com/lingniu/ingest/inbound/xinda/XindaPushEventMapperTest.java deleted file mode 100644 index 6d8802c4..00000000 --- a/modules/inbound/inbound-xinda-push/src/test/java/com/lingniu/ingest/inbound/xinda/XindaPushEventMapperTest.java +++ /dev/null @@ -1,295 +0,0 @@ -package com.lingniu.ingest.inbound.xinda; - -import com.lingniu.ingest.api.ProtocolId; -import com.lingniu.ingest.api.event.AlarmPayload; -import com.lingniu.ingest.api.event.VehicleEvent; -import com.lingniu.ingest.identity.InMemoryVehicleIdentityService; -import com.lingniu.ingest.identity.VehicleIdentityBinding; -import com.lingniu.ingest.identity.VehicleIdentityRegistry; -import com.lingniu.ingest.identity.VehicleRegistrationBinding; -import org.junit.jupiter.api.Test; - -import java.util.List; - -import static org.assertj.core.api.Assertions.assertThat; - -class XindaPushEventMapperTest { - - private final InMemoryVehicleIdentityService identity = new InMemoryVehicleIdentityService(); - private final XindaPushEventMapper mapper = new XindaPushEventMapper(identity); - - @Test - void maps0200ToLocationEvent() { - XindaPushMessage message = new XindaPushMessage("0200", - "{\"vin\":\"LNVIN000000000101\",\"sim\":\"13800138000\",\"deviceId\":\"XD-101\",\"plateNo\":\"粤B10101\",\"lng\":116397128,\"lat\":39916527,\"speed\":52.3,\"drct\":90}"); - - List events = mapper.toEvents(message); - - assertThat(events).hasSize(1); - assertThat(events.get(0)).isInstanceOf(VehicleEvent.Location.class); - VehicleEvent.Location loc = (VehicleEvent.Location) events.get(0); - assertThat(loc.source()).isEqualTo(ProtocolId.XINDA_PUSH); - assertThat(loc.vin()).isEqualTo("LNVIN000000000101"); - assertThat(loc.payload().longitude()).isEqualTo(116.397128); - assertThat(loc.payload().latitude()).isEqualTo(39.916527); - assertThat(loc.metadata()) - .containsEntry("vin", "LNVIN000000000101") - .containsEntry("externalVin", "LNVIN000000000101") - .containsEntry("phone", "13800138000") - .containsEntry("deviceId", "XD-101") - .containsEntry("plateNo", "粤B10101") - .doesNotContainKey("rawJson"); - } - - @Test - void maps0200MileageToInternalLocationPayload() { - XindaPushMessage message = new XindaPushMessage("0200", - "{\"vin\":\"LNVIN000000000101\",\"lng\":116397128,\"lat\":39916527,\"totalMileage\":12345.6}"); - - VehicleEvent.Location loc = (VehicleEvent.Location) mapper.toEvents(message).getFirst(); - - assertThat(loc.payload().totalMileageKm()).isEqualTo(12345.6); - } - - @Test - void mapsAlarmAndPassthroughEvents() { - assertThat(mapper.toEvents(new XindaPushMessage("0300", "{\"plateNo\":\"粤B12345\"}"))) - .singleElement() - .isInstanceOf(VehicleEvent.Alarm.class) - .extracting(VehicleEvent::source) - .isEqualTo(ProtocolId.XINDA_PUSH); - - assertThat(mapper.toEvents(new XindaPushMessage("0401", "{\"plateNo\":\"粤B12345\"}"))) - .singleElement() - .isInstanceOf(VehicleEvent.Passthrough.class); - } - - @Test - void mapsHydrogenLeakAlarmToCriticalAlarmEvent() { - identity.bind(new VehicleIdentityBinding( - ProtocolId.XINDA_PUSH, "LNVIN000000H2LEAK", "", "XD-H2-001", "")); - XindaPushMessage message = new XindaPushMessage("0300", """ - { - "deviceId": "XD-H2-001", - "alarmType": "HYDROGEN_LEAK", - "alarmName": "氢气泄漏", - "alarmLevel": 3, - "lng": 116397128, - "lat": 39916527 - } - """); - - List events = mapper.toEvents(message); - - assertThat(events).singleElement().satisfies(event -> { - assertThat(event).isInstanceOf(VehicleEvent.Alarm.class); - VehicleEvent.Alarm alarm = (VehicleEvent.Alarm) event; - assertThat(alarm.vin()).isEqualTo("LNVIN000000H2LEAK"); - assertThat(alarm.payload().level()).isEqualTo(AlarmPayload.AlarmLevel.CRITICAL); - assertThat(alarm.payload().safetyCategory()).isEqualTo(AlarmPayload.SafetyCategory.HYDROGEN_LEAK); - assertThat(alarm.payload().hydrogenLeakDetected()).isTrue(); - assertThat(alarm.payload().hydrogenLeakLevel()).isEqualTo(AlarmPayload.HydrogenLeakLevel.CRITICAL); - assertThat(alarm.payload().hydrogenLeakActionRequired()).isTrue(); - assertThat(alarm.payload().longitude()).isEqualTo(116.397128); - assertThat(alarm.payload().latitude()).isEqualTo(39.916527); - assertThat(alarm.metadata()) - .containsEntry("alarmType", "HYDROGEN_LEAK") - .containsEntry("alarmName", "氢气泄漏") - .containsEntry("hydrogenLeakDetected", "true"); - }); - } - - @Test - void resolvesVehicleIdentityByPlate() { - identity.bind(new VehicleIdentityBinding( - ProtocolId.XINDA_PUSH, "LNVIN000000000202", "", "", "粤B12345")); - - List events = mapper.toEvents(new XindaPushMessage("0200", - "{\"plateNo\":\"粤B12345\",\"lng\":116397128,\"lat\":39916527}")); - - assertThat(events).singleElement() - .satisfies(event -> { - assertThat(event.vin()).isEqualTo("LNVIN000000000202"); - assertThat(event.metadata()) - .containsEntry("vin", "LNVIN000000000202") - .containsEntry("identityResolved", "true"); - }); - } - - @Test - void resolvesVehicleIdentityByCarNamePlateFromXindaPayload() { - identity.bind(new VehicleIdentityBinding( - ProtocolId.XINDA_PUSH, "LNVIN000000000203", "", "", "浙F00780F")); - - List events = mapper.toEvents(new XindaPushMessage("0200", - "{\"carName\":\"浙F00780F\",\"tmnKey\":\"13307811342\",\"simCode\":\"8986062582008805648\"," - + "\"lng\":121.075223,\"lat\":30.583753}")); - - assertThat(events).singleElement() - .satisfies(event -> { - assertThat(event.vin()).isEqualTo("LNVIN000000000203"); - assertThat(event.metadata()) - .containsEntry("vin", "LNVIN000000000203") - .containsEntry("plateNo", "浙F00780F") - .containsEntry("phone", "8986062582008805648") - .containsEntry("deviceId", "13307811342") - .containsEntry("identityResolved", "true"); - }); - } - - @Test - void registersXindaIdentityFromCarNamePayload() { - RecordingIdentityRegistry registry = new RecordingIdentityRegistry(); - identity.bind(new VehicleIdentityBinding( - ProtocolId.XINDA_PUSH, "LA9GG68L8PBAF4776", "", "", "浙F00780F")); - XindaPushEventMapper mapper = new XindaPushEventMapper(identity, registry); - - mapper.toEvents(new XindaPushMessage("0200", - "{\"carName\":\"浙F00780F\",\"tmnKey\":\"13307811342\",\"simCode\":\"8986062582008805648\"," - + "\"carColor\":\"93\",\"lng\":121.075223,\"lat\":30.583753}")); - - assertThat(registry.registration) - .isNotNull() - .satisfies(registration -> { - assertThat(registration.protocol()).isEqualTo(ProtocolId.XINDA_PUSH); - assertThat(registration.vin()).isEqualTo("LA9GG68L8PBAF4776"); - assertThat(registration.phone()).isEqualTo("8986062582008805648"); - assertThat(registration.deviceId()).isEqualTo("13307811342"); - assertThat(registration.plate()).isEqualTo("浙F00780F"); - assertThat(registration.plateColor()).isEqualTo(93); - }); - } - - @Test - void registersUnknownVinWhenXindaIdentityIsNotResolved() { - RecordingIdentityRegistry registry = new RecordingIdentityRegistry(); - XindaPushEventMapper mapper = new XindaPushEventMapper(identity, registry); - - mapper.toEvents(new XindaPushMessage("0200", - "{\"carName\":\"浙F00780F\",\"tmnKey\":\"13307811342\"," - + "\"lng\":121.075223,\"lat\":30.583753}")); - - assertThat(registry.registration) - .isNotNull() - .satisfies(registration -> { - assertThat(registration.protocol()).isEqualTo(ProtocolId.XINDA_PUSH); - assertThat(registration.vin()).isEqualTo("unknown"); - assertThat(registration.phone()).isEqualTo("13307811342"); - assertThat(registration.deviceId()).isEqualTo("13307811342"); - assertThat(registration.plate()).isEqualTo("浙F00780F"); - }); - } - - @Test - void identityResolverFailureStillProducesLocationEventWithUnknownVin() { - XindaPushEventMapper mapperWithFailingIdentity = new XindaPushEventMapper(lookup -> { - throw new IllegalStateException("identity backend unavailable"); - }); - - List events = mapperWithFailingIdentity.toEvents(new XindaPushMessage("0200", - "{\"deviceId\":\"XD-DEVICE-001\",\"lng\":116397128,\"lat\":39916527,\"speed\":52.3}")); - - assertThat(events).singleElement() - .satisfies(event -> { - assertThat(event).isInstanceOf(VehicleEvent.Location.class); - assertThat(event.vin()).isEqualTo("unknown"); - assertThat(event.metadata()) - .containsEntry("vin", "unknown") - .containsEntry("deviceId", "XD-DEVICE-001") - .containsEntry("identityResolved", "false") - .containsEntry("identitySource", "UNKNOWN") - .containsEntry("identityError", "true") - .containsEntry("identityErrorMessage", "identity backend unavailable"); - }); - } - - @Test - void malformedJsonFallsBackToPassthroughEvent() { - List events = mapper.toEvents(new XindaPushMessage("0200", "{bad-json")); - - assertThat(events).singleElement() - .satisfies(event -> { - assertThat(event).isInstanceOf(VehicleEvent.Passthrough.class); - VehicleEvent.Passthrough passthrough = (VehicleEvent.Passthrough) event; - assertThat(passthrough.data()).containsExactly("{bad-json".getBytes(java.nio.charset.StandardCharsets.UTF_8)); - assertThat(event.metadata()) - .containsEntry("vin", "unknown") - .containsEntry("parseError", "true") - .containsKey("parseErrorMessage"); - assertThat(event.metadata().get("parseErrorMessage")).contains("expect"); - }); - } - - @Test - void unknownBusinessCommandFallsBackToPassthroughEvent() { - List events = mapper.toEvents(new XindaPushMessage("0500", - "{\"vin\":\"LNVIN000000000303\",\"vendor\":\"extension\"}")); - - assertThat(events).singleElement() - .satisfies(event -> { - assertThat(event).isInstanceOf(VehicleEvent.Passthrough.class); - assertThat(event.source()).isEqualTo(ProtocolId.XINDA_PUSH); - assertThat(event.vin()).isEqualTo("LNVIN000000000303"); - assertThat(event.metadata()) - .containsEntry("vin", "LNVIN000000000303") - .containsEntry("cmd", "0500") - .containsEntry("unknownCommand", "true") - .containsEntry("parseError", "false") - .containsEntry("externalVin", "LNVIN000000000303"); - }); - } - - @Test - void malformedCommandFallsBackToPassthroughEventWithMetadata() { - List events = mapper.toEvents(new XindaPushMessage("bad-cmd", - "{\"vin\":\"LNVIN000000000404\",\"vendor\":\"broken\"}")); - - assertThat(events).singleElement() - .satisfies(event -> { - assertThat(event).isInstanceOf(VehicleEvent.Passthrough.class); - assertThat(event.source()).isEqualTo(ProtocolId.XINDA_PUSH); - assertThat(event.vin()).isEqualTo("LNVIN000000000404"); - assertThat(((VehicleEvent.Passthrough) event).passthroughType()).isZero(); - assertThat(event.metadata()) - .containsEntry("vin", "LNVIN000000000404") - .containsEntry("cmd", "bad-cmd") - .containsEntry("unknownCommand", "true") - .containsEntry("malformedCommand", "true") - .containsEntry("parseError", "false"); - }); - } - - @Test - void operationalFailureMessageBecomesPassthroughEventWithFailureMetadata() { - List events = mapper.toEvents(new XindaPushMessage("0000", - "{\"operationalError\":true,\"phase\":\"startup\",\"reason\":\"host not configured\"}")); - - assertThat(events).singleElement() - .satisfies(event -> { - assertThat(event).isInstanceOf(VehicleEvent.Passthrough.class); - assertThat(event.source()).isEqualTo(ProtocolId.XINDA_PUSH); - assertThat(event.vin()).isEqualTo("unknown"); - assertThat(event.metadata()) - .containsEntry("vin", "unknown") - .containsEntry("cmd", "0000") - .containsEntry("operationalError", "true") - .containsEntry("phase", "startup") - .containsEntry("reason", "host not configured") - .containsEntry("unknownCommand", "false") - .containsEntry("parseError", "false"); - }); - } - - private static final class RecordingIdentityRegistry implements VehicleIdentityRegistry { - private VehicleRegistrationBinding registration; - - @Override - public void bind(VehicleIdentityBinding binding) { - } - - @Override - public void register(VehicleRegistrationBinding registration) { - this.registration = registration; - } - } -} diff --git a/modules/inbound/inbound-xinda-push/src/test/java/com/lingniu/ingest/inbound/xinda/XindaPushRealtimeHandlerTest.java b/modules/inbound/inbound-xinda-push/src/test/java/com/lingniu/ingest/inbound/xinda/XindaPushRealtimeHandlerTest.java deleted file mode 100644 index 740e769e..00000000 --- a/modules/inbound/inbound-xinda-push/src/test/java/com/lingniu/ingest/inbound/xinda/XindaPushRealtimeHandlerTest.java +++ /dev/null @@ -1,30 +0,0 @@ -package com.lingniu.ingest.inbound.xinda; - -import com.lingniu.ingest.api.annotation.EventEmit; -import com.lingniu.ingest.api.annotation.MessageMapping; -import com.lingniu.ingest.api.event.VehicleEvent; -import org.junit.jupiter.api.Test; - -import java.lang.reflect.Method; - -import static org.assertj.core.api.Assertions.assertThat; - -class XindaPushRealtimeHandlerTest { - - @Test - void declares0300AsAlarmHandler() throws Exception { - Method method = XindaPushRealtimeHandler.class.getMethod("onAlarm", XindaPushMessage.class); - - assertThat(method.getAnnotation(MessageMapping.class).command()).containsExactly(0x0300); - assertThat(method.getAnnotation(EventEmit.class).value()).containsExactly(VehicleEvent.Alarm.class); - } - - @Test - void alarmHandlerEmitsAlarmEvent() { - XindaPushRealtimeHandler handler = new XindaPushRealtimeHandler(new XindaPushEventMapper()); - - assertThat(handler.onAlarm(new XindaPushMessage("0300", "{\"alarmType\":\"HYDROGEN_LEAK\"}"))) - .singleElement() - .isInstanceOf(VehicleEvent.Alarm.class); - } -} diff --git a/modules/protocols/protocol-jsatl12/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java b/modules/protocols/protocol-jsatl12/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java new file mode 100644 index 00000000..f30866f6 --- /dev/null +++ b/modules/protocols/protocol-jsatl12/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java @@ -0,0 +1,48 @@ +package com.lingniu.ingest.identity; + +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; + +public final class InMemoryVehicleIdentityService implements VehicleIdentityResolver, VehicleIdentityRegistry { + + private final List bindings = new CopyOnWriteArrayList<>(); + + @Override + public VehicleIdentity resolve(VehicleIdentityLookup lookup) { + if (lookup == null) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.UNKNOWN); + } + if (!lookup.vin().isBlank() && !"unknown".equalsIgnoreCase(lookup.vin())) { + return new VehicleIdentity(lookup.vin(), true, VehicleIdentitySource.EXPLICIT_VIN); + } + for (VehicleIdentityBinding binding : bindings) { + if (binding.protocol() != lookup.protocol()) { + continue; + } + if (!lookup.phone().isBlank() && lookup.phone().equals(binding.phone())) { + return new VehicleIdentity(binding.vin(), true, VehicleIdentitySource.BOUND_PHONE); + } + if (!lookup.deviceId().isBlank() && lookup.deviceId().equals(binding.deviceId())) { + return new VehicleIdentity(binding.vin(), true, VehicleIdentitySource.BOUND_DEVICE_ID); + } + if (!lookup.plate().isBlank() && lookup.plate().equals(binding.plate())) { + return new VehicleIdentity(binding.vin(), true, VehicleIdentitySource.BOUND_PLATE); + } + } + if (!lookup.phone().isBlank()) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.FALLBACK_PHONE); + } + if (!lookup.deviceId().isBlank()) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.FALLBACK_DEVICE_ID); + } + if (!lookup.plate().isBlank()) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.FALLBACK_PLATE); + } + return new VehicleIdentity("unknown", false, VehicleIdentitySource.UNKNOWN); + } + + @Override + public void bind(VehicleIdentityBinding binding) { + bindings.add(binding); + } +} diff --git a/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java b/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java new file mode 100644 index 00000000..f30866f6 --- /dev/null +++ b/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java @@ -0,0 +1,48 @@ +package com.lingniu.ingest.identity; + +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; + +public final class InMemoryVehicleIdentityService implements VehicleIdentityResolver, VehicleIdentityRegistry { + + private final List bindings = new CopyOnWriteArrayList<>(); + + @Override + public VehicleIdentity resolve(VehicleIdentityLookup lookup) { + if (lookup == null) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.UNKNOWN); + } + if (!lookup.vin().isBlank() && !"unknown".equalsIgnoreCase(lookup.vin())) { + return new VehicleIdentity(lookup.vin(), true, VehicleIdentitySource.EXPLICIT_VIN); + } + for (VehicleIdentityBinding binding : bindings) { + if (binding.protocol() != lookup.protocol()) { + continue; + } + if (!lookup.phone().isBlank() && lookup.phone().equals(binding.phone())) { + return new VehicleIdentity(binding.vin(), true, VehicleIdentitySource.BOUND_PHONE); + } + if (!lookup.deviceId().isBlank() && lookup.deviceId().equals(binding.deviceId())) { + return new VehicleIdentity(binding.vin(), true, VehicleIdentitySource.BOUND_DEVICE_ID); + } + if (!lookup.plate().isBlank() && lookup.plate().equals(binding.plate())) { + return new VehicleIdentity(binding.vin(), true, VehicleIdentitySource.BOUND_PLATE); + } + } + if (!lookup.phone().isBlank()) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.FALLBACK_PHONE); + } + if (!lookup.deviceId().isBlank()) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.FALLBACK_DEVICE_ID); + } + if (!lookup.plate().isBlank()) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.FALLBACK_PLATE); + } + return new VehicleIdentity("unknown", false, VehicleIdentitySource.UNKNOWN); + } + + @Override + public void bind(VehicleIdentityBinding binding) { + bindings.add(binding); + } +} diff --git a/modules/protocols/protocol-jt808/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java b/modules/protocols/protocol-jt808/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java new file mode 100644 index 00000000..f30866f6 --- /dev/null +++ b/modules/protocols/protocol-jt808/src/test/java/com/lingniu/ingest/identity/InMemoryVehicleIdentityService.java @@ -0,0 +1,48 @@ +package com.lingniu.ingest.identity; + +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; + +public final class InMemoryVehicleIdentityService implements VehicleIdentityResolver, VehicleIdentityRegistry { + + private final List bindings = new CopyOnWriteArrayList<>(); + + @Override + public VehicleIdentity resolve(VehicleIdentityLookup lookup) { + if (lookup == null) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.UNKNOWN); + } + if (!lookup.vin().isBlank() && !"unknown".equalsIgnoreCase(lookup.vin())) { + return new VehicleIdentity(lookup.vin(), true, VehicleIdentitySource.EXPLICIT_VIN); + } + for (VehicleIdentityBinding binding : bindings) { + if (binding.protocol() != lookup.protocol()) { + continue; + } + if (!lookup.phone().isBlank() && lookup.phone().equals(binding.phone())) { + return new VehicleIdentity(binding.vin(), true, VehicleIdentitySource.BOUND_PHONE); + } + if (!lookup.deviceId().isBlank() && lookup.deviceId().equals(binding.deviceId())) { + return new VehicleIdentity(binding.vin(), true, VehicleIdentitySource.BOUND_DEVICE_ID); + } + if (!lookup.plate().isBlank() && lookup.plate().equals(binding.plate())) { + return new VehicleIdentity(binding.vin(), true, VehicleIdentitySource.BOUND_PLATE); + } + } + if (!lookup.phone().isBlank()) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.FALLBACK_PHONE); + } + if (!lookup.deviceId().isBlank()) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.FALLBACK_DEVICE_ID); + } + if (!lookup.plate().isBlank()) { + return new VehicleIdentity("unknown", false, VehicleIdentitySource.FALLBACK_PLATE); + } + return new VehicleIdentity("unknown", false, VehicleIdentitySource.UNKNOWN); + } + + @Override + public void bind(VehicleIdentityBinding binding) { + bindings.add(binding); + } +} diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/LocationHistoryControllerTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/LocationHistoryControllerTest.java index f419c274..f572aa76 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/LocationHistoryControllerTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/LocationHistoryControllerTest.java @@ -70,7 +70,7 @@ class LocationHistoryControllerTest { } @Test - void rejectsLegacyXindaProtocolOnDefaultLocationApi() { + void rejectsRemovedXindaProtocolOnDefaultLocationApi() { AtomicReference captured = new AtomicReference<>(); LocationHistoryController controller = new LocationHistoryController( new StubReader(captured, new TdenginePage<>(List.of(), null))); diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/RawFrameHistoryControllerTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/RawFrameHistoryControllerTest.java index 75ad0b7f..9d59ade1 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/RawFrameHistoryControllerTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/RawFrameHistoryControllerTest.java @@ -202,7 +202,7 @@ class RawFrameHistoryControllerTest { } @Test - void rejectsLegacyXindaProtocolOnDefaultRawFrameApi() { + void rejectsRemovedXindaProtocolOnDefaultRawFrameApi() { RawFrameHistoryController controller = new RawFrameHistoryController( new CapturingReader(new TdenginePage<>(List.of(), Optional.empty()))); diff --git a/pom.xml b/pom.xml index c91a846f..bb95ea2a 100644 --- a/pom.xml +++ b/pom.xml @@ -91,45 +91,6 @@ - - legacy-xinda - - modules/inbound/inbound-xinda-push - modules/apps/xinda-push-app - - - - - lingniu-nexus - Lingniu Private Maven Repository - https://nexus.lnh2e.com/repository/maven-public/ - - true - - - true - - - - - - - com.lingniu.ingest - inbound-xinda-push - ${project.version} - - - com.lingniu.ingest - xinda-push-app - ${project.version} - - - -