From 7b346420af61f307d41413b32ba711cb74563931 Mon Sep 17 00:00:00 2001 From: lingniu Date: Mon, 29 Jun 2026 16:35:02 +0800 Subject: [PATCH] fix: preserve jt808 registration metadata --- .../vehicle-ingest-tdengine-verification.md | 6 +++ .../jt808/inbound/Jt808ChannelHandler.java | 18 +++++++ .../jt808/mapper/Jt808EventMapper.java | 18 ++++++- .../inbound/Jt808ChannelHandlerTest.java | 47 +++++++++++++++++++ .../jt808/mapper/Jt808EventMapperTest.java | 24 ++++++++++ 5 files changed, 112 insertions(+), 1 deletion(-) diff --git a/docs/operations/vehicle-ingest-tdengine-verification.md b/docs/operations/vehicle-ingest-tdengine-verification.md index 71eff355..53c80d58 100644 --- a/docs/operations/vehicle-ingest-tdengine-verification.md +++ b/docs/operations/vehicle-ingest-tdengine-verification.md @@ -209,6 +209,7 @@ JT808 真实转发链路已验证: - `vehicle_locations` 中 `protocol='JT808'` 的位置行数已超过 `1900`。 - `telemetry_fields` 中 `protocol='JT808'` 的字段行数已超过 `15000`。 - 最新样例包含终端号 `13079963310`、消息 ID `512`、`archive://2026/06/29/JT808/...` 原始帧引用。 +- 合成 0x0100 注册帧终端号 `13079969999` 已验证:TCP `808` 返回 `0x8100` 注册 ACK,`raw_frames.metadata_json` 包含 `jt808.register.province`、`jt808.register.city`、`jt808.register.maker`、`jt808.register.deviceType`、`jt808.register.deviceId`、`jt808.register.plateColor`、`jt808.register.plate`,对应 raw archive 文件存在。 复查 SQL: @@ -222,6 +223,11 @@ SELECT event_time, vehicle_key, phone, message_id, raw_uri WHERE protocol = 'JT808' ORDER BY event_time DESC LIMIT 5; +SELECT ts, frame_id, phone, message_id, metadata_json, raw_uri + FROM raw_frames + WHERE protocol = 'JT808' AND phone = '13079969999' AND message_id = 256 + ORDER BY ts DESC + LIMIT 3; ``` 复查 API: diff --git a/modules/protocols/protocol-jt808/src/main/java/com/lingniu/ingest/protocol/jt808/inbound/Jt808ChannelHandler.java b/modules/protocols/protocol-jt808/src/main/java/com/lingniu/ingest/protocol/jt808/inbound/Jt808ChannelHandler.java index a84e8202..294dbc85 100644 --- a/modules/protocols/protocol-jt808/src/main/java/com/lingniu/ingest/protocol/jt808/inbound/Jt808ChannelHandler.java +++ b/modules/protocols/protocol-jt808/src/main/java/com/lingniu/ingest/protocol/jt808/inbound/Jt808ChannelHandler.java @@ -163,6 +163,7 @@ public class Jt808ChannelHandler extends SimpleChannelInboundHandler { } meta.put("seq", Integer.toString(msg.header().serialNo())); meta.put("peer", addr(ctx)); + addRegistrationMetadata(meta, msg.body()); RawFrame rf = new RawFrame( ProtocolId.JT808, msg.header().messageId(), @@ -309,6 +310,23 @@ public class Jt808ChannelHandler extends SimpleChannelInboundHandler { return body instanceof Jt808Body.Register reg ? reg.plate() : ""; } + private static void addRegistrationMetadata(Map meta, Jt808Body body) { + if (!(body instanceof Jt808Body.Register reg)) { + return; + } + meta.put("jt808.register.province", Integer.toString(reg.province())); + meta.put("jt808.register.city", Integer.toString(reg.city())); + meta.put("jt808.register.maker", nullToEmpty(reg.maker())); + meta.put("jt808.register.deviceType", nullToEmpty(reg.deviceType())); + meta.put("jt808.register.deviceId", nullToEmpty(reg.deviceId())); + meta.put("jt808.register.plateColor", Integer.toString(reg.plateColor())); + meta.put("jt808.register.plate", nullToEmpty(reg.plate())); + } + + private static String nullToEmpty(String value) { + return value == null ? "" : value; + } + private void ackTerminal(ChannelHandlerContext ctx, String phone, int ackSerial, int ackMsgId, int result) { var cmd = Jt808Commands.platformAck(ackSerial, ackMsgId, result); byte[] frame = Jt808FrameEncoder.encode(cmd.messageId(), phone, channelRegistry.nextSerial(phone), cmd.body()); diff --git a/modules/protocols/protocol-jt808/src/main/java/com/lingniu/ingest/protocol/jt808/mapper/Jt808EventMapper.java b/modules/protocols/protocol-jt808/src/main/java/com/lingniu/ingest/protocol/jt808/mapper/Jt808EventMapper.java index 1ca5fdea..93578108 100644 --- a/modules/protocols/protocol-jt808/src/main/java/com/lingniu/ingest/protocol/jt808/mapper/Jt808EventMapper.java +++ b/modules/protocols/protocol-jt808/src/main/java/com/lingniu/ingest/protocol/jt808/mapper/Jt808EventMapper.java @@ -52,7 +52,7 @@ public final class Jt808EventMapper implements EventMapper { case Jt808Body.Register reg -> List.of(new VehicleEvent.Login( UUID.randomUUID().toString(), - vin, ProtocolId.JT808, ingestTime, ingestTime, null, meta, + vin, ProtocolId.JT808, ingestTime, ingestTime, null, registerMeta(meta, reg), reg.deviceId(), header.version().name())); case Jt808Body.Auth auth -> List.of(new VehicleEvent.Login( @@ -285,6 +285,18 @@ public final class Jt808EventMapper implements EventMapper { return Map.copyOf(copy); } + private static Map registerMeta(Map meta, Jt808Body.Register reg) { + java.util.HashMap copy = new java.util.HashMap<>(meta); + copy.put("jt808.register.province", Integer.toString(reg.province())); + copy.put("jt808.register.city", Integer.toString(reg.city())); + copy.put("jt808.register.maker", nullToEmpty(reg.maker())); + copy.put("jt808.register.deviceType", nullToEmpty(reg.deviceType())); + copy.put("jt808.register.deviceId", nullToEmpty(reg.deviceId())); + copy.put("jt808.register.plateColor", Integer.toString(reg.plateColor())); + copy.put("jt808.register.plate", nullToEmpty(reg.plate())); + return Map.copyOf(copy); + } + private static Map withMeta(Map meta, String key, String value) { java.util.HashMap copy = new java.util.HashMap<>(meta); copy.put(key, value == null ? "" : value); @@ -350,6 +362,10 @@ public final class Jt808EventMapper implements EventMapper { return body instanceof Jt808Body.Register reg ? reg.plate() : ""; } + private static String nullToEmpty(String value) { + return value == null ? "" : value; + } + private record IdentityResolution(VehicleIdentity identity, String errorMessage) { private static IdentityResolution ok(VehicleIdentity identity) { return new IdentityResolution(identity, null); diff --git a/modules/protocols/protocol-jt808/src/test/java/com/lingniu/ingest/protocol/jt808/inbound/Jt808ChannelHandlerTest.java b/modules/protocols/protocol-jt808/src/test/java/com/lingniu/ingest/protocol/jt808/inbound/Jt808ChannelHandlerTest.java index ede2811d..2fc838c0 100644 --- a/modules/protocols/protocol-jt808/src/test/java/com/lingniu/ingest/protocol/jt808/inbound/Jt808ChannelHandlerTest.java +++ b/modules/protocols/protocol-jt808/src/test/java/com/lingniu/ingest/protocol/jt808/inbound/Jt808ChannelHandlerTest.java @@ -237,6 +237,53 @@ class Jt808ChannelHandlerTest { eventBus.close(); } + @Test + void registerArchivePreservesRegistrationFieldsInMetadata() throws Exception { + RecordingSink sink = new RecordingSink(1); + DisruptorEventBus eventBus = new DisruptorEventBus(1024, "blocking", List.of(sink)); + AsyncBatchExecutor batchExecutor = new AsyncBatchExecutor(eventBus::publish); + Dispatcher dispatcher = new Dispatcher( + new HandlerRegistry(), + new InterceptorChain(List.of()), + new HandlerInvoker(), + eventBus, + batchExecutor); + InMemoryVehicleIdentityService identity = new InMemoryVehicleIdentityService(); + Jt808ChannelHandler handler = new Jt808ChannelHandler( + new Jt808MessageDecoder(new BodyParserRegistry(List.of(new RegisterBodyParser()))), + dispatcher, + new InMemorySessionStore(), + identity, + new Jt808ChannelRegistry(), + new Jt808PendingRequests()); + EmbeddedChannel channel = new EmbeddedChannel(handler); + byte[] frame = buildFrame( + Jt808MessageId.TERMINAL_REGISTER, + "123456789012", + 1, + buildRegisterBody("DEV808", "B80808")); + + channel.writeInbound(frame); + + assertThat(sink.await()).isTrue(); + assertThat(sink.events).singleElement().satisfies(event -> { + assertThat(event).isInstanceOf(VehicleEvent.RawArchive.class); + VehicleEvent.RawArchive archive = (VehicleEvent.RawArchive) event; + assertThat(archive.command()).isEqualTo(Jt808MessageId.TERMINAL_REGISTER); + assertThat(archive.metadata()) + .containsEntry("jt808.register.province", "44") + .containsEntry("jt808.register.city", "4401") + .containsEntry("jt808.register.maker", "MAKER") + .containsEntry("jt808.register.deviceType", "TYPE-A") + .containsEntry("jt808.register.deviceId", "DEV808") + .containsEntry("jt808.register.plateColor", "1") + .containsEntry("jt808.register.plate", "B80808"); + }); + + batchExecutor.close(); + eventBus.close(); + } + @Test void malformedFrameIsArchivedAndDispatchedAsPassthrough() throws Exception { RecordingSink sink = new RecordingSink(2); diff --git a/modules/protocols/protocol-jt808/src/test/java/com/lingniu/ingest/protocol/jt808/mapper/Jt808EventMapperTest.java b/modules/protocols/protocol-jt808/src/test/java/com/lingniu/ingest/protocol/jt808/mapper/Jt808EventMapperTest.java index d0341046..87f362db 100644 --- a/modules/protocols/protocol-jt808/src/test/java/com/lingniu/ingest/protocol/jt808/mapper/Jt808EventMapperTest.java +++ b/modules/protocols/protocol-jt808/src/test/java/com/lingniu/ingest/protocol/jt808/mapper/Jt808EventMapperTest.java @@ -133,6 +133,30 @@ class Jt808EventMapperTest { .first().isInstanceOf(VehicleEvent.Heartbeat.class); } + @Test + void registerBodyPreservesRegistrationFieldsInMetadata() { + var header = new Jt808Header( + Jt808MessageId.TERMINAL_REGISTER, 37, 0, false, + Jt808Header.ProtocolVersion.V2013, "123456789012", 1, 0, 0); + var body = new Jt808Body.Register( + 44, 1, "G7", "MODEL-X", "DEV123456", 2, "粤B12345"); + + List events = mapper.toEvents(new Jt808Message(header, body)); + + assertThat(events).singleElement() + .satisfies(event -> { + assertThat(event).isInstanceOf(VehicleEvent.Login.class); + assertThat(event.metadata()) + .containsEntry("jt808.register.province", "44") + .containsEntry("jt808.register.city", "1") + .containsEntry("jt808.register.maker", "G7") + .containsEntry("jt808.register.deviceType", "MODEL-X") + .containsEntry("jt808.register.deviceId", "DEV123456") + .containsEntry("jt808.register.plateColor", "2") + .containsEntry("jt808.register.plate", "粤B12345"); + }); + } + @Test void unregisterBodyProducesLogoutEvent() { var header = new Jt808Header(