From d1748fcc2f1d61b7ff06f6fa14ef2de5be3e7821 Mon Sep 17 00:00:00 2001 From: lingniu Date: Mon, 29 Jun 2026 18:47:19 +0800 Subject: [PATCH] fix: archive gb32960 decode failures --- .../inbound/Gb32960ChannelHandler.java | 17 ++++++++++ .../Gb32960ChannelHandlerAckBoundaryTest.java | 31 +++++++++++++++++++ 2 files changed, 48 insertions(+) diff --git a/modules/protocols/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandler.java b/modules/protocols/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandler.java index 96b2bcc5..8285f6aa 100644 --- a/modules/protocols/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandler.java +++ b/modules/protocols/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandler.java @@ -109,6 +109,7 @@ public class Gb32960ChannelHandler extends SimpleChannelInboundHandler { } catch (Exception e) { // 解码失败只丢当前帧,不主动断开连接。真实设备偶发脏字节时,FrameDecoder 会继续寻找下一帧头。 log.warn("[gb32960] decode failed peer={} len={}", addr(ctx), frame.length, e); + dispatchMalformed(ctx, frame, e.getMessage()); return; } @@ -224,6 +225,22 @@ public class Gb32960ChannelHandler extends SimpleChannelInboundHandler { Instant.now()); } + private void dispatchMalformed(ChannelHandlerContext ctx, byte[] frame, String errorMessage) { + Map sourceMeta = new HashMap<>(4); + sourceMeta.put("vin", "unknown"); + sourceMeta.put("peer", addr(ctx)); + sourceMeta.put("parseError", "true"); + sourceMeta.put("parseErrorMessage", errorMessage == null ? "" : errorMessage); + dispatcher.dispatch(new RawFrame( + ProtocolId.GB32960, + 0, + 0, + frame, + frame == null ? new byte[0] : frame, + sourceMeta, + Instant.now())); + } + private static boolean requiresDurableAck(CommandType cmd) { return switch (cmd) { case VEHICLE_LOGIN, diff --git a/modules/protocols/protocol-gb32960/src/test/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandlerAckBoundaryTest.java b/modules/protocols/protocol-gb32960/src/test/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandlerAckBoundaryTest.java index fe83a700..d2ce859d 100644 --- a/modules/protocols/protocol-gb32960/src/test/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandlerAckBoundaryTest.java +++ b/modules/protocols/protocol-gb32960/src/test/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandlerAckBoundaryTest.java @@ -97,6 +97,29 @@ class Gb32960ChannelHandlerAckBoundaryTest { } } + @Test + void decodeFailureIsArchivedWithParseErrorMetadata() { + try (DispatchHarness dispatch = new DispatchHarness()) { + EmbeddedChannel channel = new EmbeddedChannel(handlerWith(dispatch.dispatcher())); + byte[] malformed = new byte[]{0x23, 0x23, 0x02}; + + assertThat(channel.writeInbound(malformed)).isFalse(); + + dispatch.awaitPublishCount(1); + VehicleEvent event = dispatch.event(0); + assertThat(event).isInstanceOf(VehicleEvent.RawArchive.class); + VehicleEvent.RawArchive raw = (VehicleEvent.RawArchive) event; + assertThat(raw.source()).isEqualTo(com.lingniu.ingest.api.ProtocolId.GB32960); + assertThat(raw.command()).isZero(); + assertThat(raw.rawBytes()).containsExactly(malformed); + assertThat(raw.metadata()) + .containsEntry("vin", "unknown") + .containsEntry("parseError", "true") + .containsKey("parseErrorMessage"); + channel.finishAndReleaseAll(); + } + } + private static ByteBuf readOutboundAfterRunningTasks(EmbeddedChannel channel) { long deadline = System.currentTimeMillis() + 3000; while (System.currentTimeMillis() < deadline) { @@ -190,6 +213,12 @@ class Gb32960ChannelHandlerAckBoundaryTest { throw new AssertionError("expected " + count + " sink publishes, actual=" + sink.futures.size()); } + VehicleEvent event(int index) { + synchronized (sink.futures) { + return sink.events.get(index); + } + } + @Override public void close() { batchExecutor.close(); @@ -199,6 +228,7 @@ class Gb32960ChannelHandlerAckBoundaryTest { private static final class ControlledSink implements EventSink { private final List> futures = new ArrayList<>(); + private final List events = new ArrayList<>(); @Override public String name() { @@ -209,6 +239,7 @@ class Gb32960ChannelHandlerAckBoundaryTest { public CompletableFuture publish(VehicleEvent event) { CompletableFuture future = new CompletableFuture<>(); synchronized (futures) { + events.add(event); futures.add(future); } return future;