From 42e8742b104738783f8412e1c6388250c807c318 Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 09:07:12 +0800 Subject: [PATCH] refactor: require identity resolver for jt1078 malformed rtp --- .../ingest/session/InMemorySessionStore.java | 4 +- .../lingniu/ingest/session/SessionStore.java | 2 +- .../StatelessVehicleIdentityResolver.java | 34 ----------- .../StatelessVehicleIdentityResolverTest.java | 57 ------------------- .../handler/Jt1078MalformedRtpHandler.java | 8 +-- .../Jt1078MalformedRtpEventPublisher.java | 10 +--- .../Jt1078MalformedRtpEventPublisherTest.java | 3 +- .../jt1078/Jt1078MalformedRtpHandlerTest.java | 14 +++++ .../config/Jt1078AutoConfigurationTest.java | 44 ++++++++++++++ 9 files changed, 68 insertions(+), 108 deletions(-) delete mode 100644 modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/StatelessVehicleIdentityResolver.java delete mode 100644 modules/core/vehicle-identity/src/test/java/com/lingniu/ingest/identity/StatelessVehicleIdentityResolverTest.java diff --git a/modules/core/session-core/src/main/java/com/lingniu/ingest/session/InMemorySessionStore.java b/modules/core/session-core/src/main/java/com/lingniu/ingest/session/InMemorySessionStore.java index 1b022b9d..ccdb80d1 100644 --- a/modules/core/session-core/src/main/java/com/lingniu/ingest/session/InMemorySessionStore.java +++ b/modules/core/session-core/src/main/java/com/lingniu/ingest/session/InMemorySessionStore.java @@ -10,10 +10,10 @@ import java.util.concurrent.ConcurrentMap; import java.util.function.UnaryOperator; /** - * 纯内存会话存储。支持按 sessionId / vin / phone 三种 key 查询。 + * 测试/本地用纯内存会话存储。支持按 sessionId / vin / phone 三种 key 查询。 * *

失活清理:基于 {@link Caffeine} 的 {@code expireAfterAccess 30 分钟}。 - * 生产环境若需多节点一致性,可用 Redis 实现替换此 Bean。 + * 生产自动配置拒绝 memory 模式,应使用 Redis 会话索引。 */ public final class InMemorySessionStore implements SessionStore { diff --git a/modules/core/session-core/src/main/java/com/lingniu/ingest/session/SessionStore.java b/modules/core/session-core/src/main/java/com/lingniu/ingest/session/SessionStore.java index 32c94941..38811206 100644 --- a/modules/core/session-core/src/main/java/com/lingniu/ingest/session/SessionStore.java +++ b/modules/core/session-core/src/main/java/com/lingniu/ingest/session/SessionStore.java @@ -4,7 +4,7 @@ import java.util.Optional; import java.util.function.UnaryOperator; /** - * 设备会话存储 SPI。默认实现是内存 + Caffeine,生产可替换为 Redis 兜底。 + * 设备会话存储 SPI。生产自动配置只创建 Redis 实现;测试可显式提供轻量实现。 */ public interface SessionStore { diff --git a/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/StatelessVehicleIdentityResolver.java b/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/StatelessVehicleIdentityResolver.java deleted file mode 100644 index ecaa45f2..00000000 --- a/modules/core/vehicle-identity/src/main/java/com/lingniu/ingest/identity/StatelessVehicleIdentityResolver.java +++ /dev/null @@ -1,34 +0,0 @@ -package com.lingniu.ingest.identity; - -/** - * Stateless production fallback used when a caller has no injected identity service. - * - *

Persistent binding belongs to {@link MySqlVehicleIdentityService}; this resolver only keeps - * ingestion flowing with explicit VIN or stable external identifiers. - */ -public final class StatelessVehicleIdentityResolver implements VehicleIdentityResolver { - - @Override - public VehicleIdentity resolve(VehicleIdentityLookup lookup) { - if (lookup == null) { - return unknown(); - } - if (!lookup.vin().isBlank()) { - return new VehicleIdentity(lookup.vin(), true, VehicleIdentitySource.EXPLICIT_VIN); - } - if (!lookup.deviceId().isBlank()) { - return new VehicleIdentity(lookup.deviceId(), false, VehicleIdentitySource.FALLBACK_DEVICE_ID); - } - if (!lookup.phone().isBlank()) { - return new VehicleIdentity(lookup.phone(), false, VehicleIdentitySource.FALLBACK_PHONE); - } - if (!lookup.plate().isBlank()) { - return new VehicleIdentity(lookup.plate(), false, VehicleIdentitySource.FALLBACK_PLATE); - } - return unknown(); - } - - private static VehicleIdentity unknown() { - return new VehicleIdentity("unknown", false, VehicleIdentitySource.UNKNOWN); - } -} diff --git a/modules/core/vehicle-identity/src/test/java/com/lingniu/ingest/identity/StatelessVehicleIdentityResolverTest.java b/modules/core/vehicle-identity/src/test/java/com/lingniu/ingest/identity/StatelessVehicleIdentityResolverTest.java deleted file mode 100644 index 5314cf19..00000000 --- a/modules/core/vehicle-identity/src/test/java/com/lingniu/ingest/identity/StatelessVehicleIdentityResolverTest.java +++ /dev/null @@ -1,57 +0,0 @@ -package com.lingniu.ingest.identity; - -import com.lingniu.ingest.api.ProtocolId; -import org.junit.jupiter.api.Test; - -import static org.assertj.core.api.Assertions.assertThat; - -class StatelessVehicleIdentityResolverTest { - - private final StatelessVehicleIdentityResolver resolver = new StatelessVehicleIdentityResolver(); - - @Test - void explicitVinWinsWithoutPersistingAnyBinding() { - VehicleIdentity identity = resolver.resolve(new VehicleIdentityLookup( - ProtocolId.MQTT_YUTONG, "LNVIN000000000001", "13900000000", "DEV001", "粤B12345")); - - assertThat(identity.vin()).isEqualTo("LNVIN000000000001"); - assertThat(identity.resolved()).isTrue(); - assertThat(identity.source()).isEqualTo(VehicleIdentitySource.EXPLICIT_VIN); - } - - @Test - void fallsBackToStableExternalIdentifiersWithoutMarkingResolved() { - assertThat(resolver.resolve(new VehicleIdentityLookup( - ProtocolId.JT808, "", "13900000001", "DEV001", "粤B12345"))) - .satisfies(identity -> { - assertThat(identity.vin()).isEqualTo("DEV001"); - assertThat(identity.resolved()).isFalse(); - assertThat(identity.source()).isEqualTo(VehicleIdentitySource.FALLBACK_DEVICE_ID); - }); - - assertThat(resolver.resolve(new VehicleIdentityLookup( - ProtocolId.JT808, "", "13900000001", "", "粤B12345"))) - .satisfies(identity -> { - assertThat(identity.vin()).isEqualTo("13900000001"); - assertThat(identity.resolved()).isFalse(); - assertThat(identity.source()).isEqualTo(VehicleIdentitySource.FALLBACK_PHONE); - }); - - assertThat(resolver.resolve(new VehicleIdentityLookup( - ProtocolId.JT808, "", "", "", "粤B12345"))) - .satisfies(identity -> { - assertThat(identity.vin()).isEqualTo("粤B12345"); - assertThat(identity.resolved()).isFalse(); - assertThat(identity.source()).isEqualTo(VehicleIdentitySource.FALLBACK_PLATE); - }); - } - - @Test - void returnsUnknownWhenNoIdentitySignalExists() { - VehicleIdentity identity = resolver.resolve(new VehicleIdentityLookup(ProtocolId.JT808, "", "", "", "")); - - assertThat(identity.vin()).isEqualTo("unknown"); - assertThat(identity.resolved()).isFalse(); - assertThat(identity.source()).isEqualTo(VehicleIdentitySource.UNKNOWN); - } -} diff --git a/modules/protocols/protocol-jt1078/src/main/java/com/lingniu/ingest/protocol/jt1078/handler/Jt1078MalformedRtpHandler.java b/modules/protocols/protocol-jt1078/src/main/java/com/lingniu/ingest/protocol/jt1078/handler/Jt1078MalformedRtpHandler.java index 06fbf4b9..5939384f 100644 --- a/modules/protocols/protocol-jt1078/src/main/java/com/lingniu/ingest/protocol/jt1078/handler/Jt1078MalformedRtpHandler.java +++ b/modules/protocols/protocol-jt1078/src/main/java/com/lingniu/ingest/protocol/jt1078/handler/Jt1078MalformedRtpHandler.java @@ -5,7 +5,6 @@ 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 com.lingniu.ingest.identity.StatelessVehicleIdentityResolver; import com.lingniu.ingest.identity.VehicleIdentity; import com.lingniu.ingest.identity.VehicleIdentityLookup; import com.lingniu.ingest.identity.VehicleIdentityResolver; @@ -17,6 +16,7 @@ import java.time.Instant; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Objects; import java.util.UUID; @ProtocolHandler(protocol = ProtocolId.JT1078) @@ -24,12 +24,8 @@ public final class Jt1078MalformedRtpHandler { private final VehicleIdentityResolver identityResolver; - public Jt1078MalformedRtpHandler() { - this(new StatelessVehicleIdentityResolver()); - } - public Jt1078MalformedRtpHandler(VehicleIdentityResolver identityResolver) { - this.identityResolver = identityResolver == null ? new StatelessVehicleIdentityResolver() : identityResolver; + this.identityResolver = Objects.requireNonNull(identityResolver, "identityResolver must not be null"); } @MessageMapping(command = Jt1078MediaMessageId.MALFORMED_RTP, desc = "JT1078 RTP 入口坏包兜底") diff --git a/modules/protocols/protocol-jt1078/src/main/java/com/lingniu/ingest/protocol/jt1078/media/Jt1078MalformedRtpEventPublisher.java b/modules/protocols/protocol-jt1078/src/main/java/com/lingniu/ingest/protocol/jt1078/media/Jt1078MalformedRtpEventPublisher.java index 5961d364..5ebff586 100644 --- a/modules/protocols/protocol-jt1078/src/main/java/com/lingniu/ingest/protocol/jt1078/media/Jt1078MalformedRtpEventPublisher.java +++ b/modules/protocols/protocol-jt1078/src/main/java/com/lingniu/ingest/protocol/jt1078/media/Jt1078MalformedRtpEventPublisher.java @@ -2,7 +2,6 @@ package com.lingniu.ingest.protocol.jt1078.media; import com.lingniu.ingest.api.ProtocolId; import com.lingniu.ingest.api.pipeline.RawFrame; -import com.lingniu.ingest.identity.StatelessVehicleIdentityResolver; import com.lingniu.ingest.identity.VehicleIdentity; import com.lingniu.ingest.identity.VehicleIdentityLookup; import com.lingniu.ingest.identity.VehicleIdentityResolver; @@ -11,6 +10,7 @@ import com.lingniu.ingest.identity.VehicleIdentitySource; import java.time.Instant; import java.util.HashMap; import java.util.Map; +import java.util.Objects; import java.util.function.Consumer; public final class Jt1078MalformedRtpEventPublisher { @@ -18,14 +18,10 @@ public final class Jt1078MalformedRtpEventPublisher { private final VehicleIdentityResolver identityResolver; private final Consumer dispatcher; - public Jt1078MalformedRtpEventPublisher(Consumer dispatcher) { - this(new StatelessVehicleIdentityResolver(), dispatcher); - } - public Jt1078MalformedRtpEventPublisher(VehicleIdentityResolver identityResolver, Consumer dispatcher) { - this.identityResolver = identityResolver == null ? new StatelessVehicleIdentityResolver() : identityResolver; - this.dispatcher = dispatcher; + this.identityResolver = Objects.requireNonNull(identityResolver, "identityResolver must not be null"); + this.dispatcher = Objects.requireNonNull(dispatcher, "dispatcher must not be null"); } public void publish(Jt1078MalformedRtpPacket malformed) { diff --git a/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/protocol/jt1078/Jt1078MalformedRtpEventPublisherTest.java b/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/protocol/jt1078/Jt1078MalformedRtpEventPublisherTest.java index ef064dc1..5e77dc33 100644 --- a/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/protocol/jt1078/Jt1078MalformedRtpEventPublisherTest.java +++ b/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/protocol/jt1078/Jt1078MalformedRtpEventPublisherTest.java @@ -19,7 +19,8 @@ class Jt1078MalformedRtpEventPublisherTest { @Test void malformedRtpPacketProducesRawFrameForUnifiedArchiveAndPassthrough() { ArrayList frames = new ArrayList<>(); - Jt1078MalformedRtpEventPublisher publisher = new Jt1078MalformedRtpEventPublisher(frames::add); + Jt1078MalformedRtpEventPublisher publisher = new Jt1078MalformedRtpEventPublisher( + new InMemoryVehicleIdentityService(), frames::add); Instant receivedAt = Instant.parse("2026-06-22T08:30:00Z"); byte[] raw = new byte[]{0x30, 0x31, 0x63, 0x64}; diff --git a/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/protocol/jt1078/Jt1078MalformedRtpHandlerTest.java b/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/protocol/jt1078/Jt1078MalformedRtpHandlerTest.java index f876e6e2..35f2bea0 100644 --- a/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/protocol/jt1078/Jt1078MalformedRtpHandlerTest.java +++ b/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/protocol/jt1078/Jt1078MalformedRtpHandlerTest.java @@ -9,13 +9,27 @@ import com.lingniu.ingest.protocol.jt1078.handler.Jt1078MalformedRtpHandler; import com.lingniu.ingest.protocol.jt1078.media.Jt1078MalformedRtpPacket; import org.junit.jupiter.api.Test; +import java.lang.reflect.Constructor; import java.time.Instant; +import java.util.Arrays; import java.util.List; +import java.util.function.Consumer; import static org.assertj.core.api.Assertions.assertThat; class Jt1078MalformedRtpHandlerTest { + @Test + void malformedRtpComponentsRequireInjectedIdentityResolver() { + assertThat(Arrays.stream(Jt1078MalformedRtpHandler.class.getConstructors()) + .map(Constructor::getParameterCount)) + .doesNotContain(0); + assertThat(Arrays.stream(com.lingniu.ingest.protocol.jt1078.media.Jt1078MalformedRtpEventPublisher.class + .getConstructors())) + .noneSatisfy(constructor -> assertThat(constructor.getParameterTypes()) + .containsExactly(Consumer.class)); + } + @Test void malformedRtpRawFramePayloadMapsToPassthroughWithDiagnostics() { byte[] raw = new byte[]{0x30, 0x31, 0x63, 0x64}; diff --git a/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/protocol/jt1078/config/Jt1078AutoConfigurationTest.java b/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/protocol/jt1078/config/Jt1078AutoConfigurationTest.java index 136f1bab..d93a981f 100644 --- a/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/protocol/jt1078/config/Jt1078AutoConfigurationTest.java +++ b/modules/protocols/protocol-jt1078/src/test/java/com/lingniu/ingest/protocol/jt1078/config/Jt1078AutoConfigurationTest.java @@ -15,6 +15,8 @@ import com.lingniu.ingest.protocol.jt1078.media.Jt1078MalformedRtpPacket; import com.lingniu.ingest.protocol.jt1078.media.Jt1078MediaArchiveService; import com.lingniu.ingest.protocol.jt808.config.Jt808AutoConfiguration; import com.lingniu.ingest.protocol.jt808.mapper.Jt808EventMapper; +import com.lingniu.ingest.session.DeviceSession; +import com.lingniu.ingest.session.SessionStore; import com.lingniu.ingest.session.config.SessionCoreAutoConfiguration; import com.lingniu.ingest.sink.archive.ArchiveStore; import org.junit.jupiter.api.Test; @@ -28,8 +30,10 @@ import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; import java.util.List; +import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CopyOnWriteArrayList; +import java.util.function.UnaryOperator; import static org.assertj.core.api.Assertions.assertThat; @@ -149,6 +153,11 @@ class Jt1078AutoConfigurationTest { }; } + @Bean + SessionStore sessionStore() { + return new EmptySessionStore(); + } + @Bean(destroyMethod = "close") DisruptorEventBus disruptorEventBus(CapturingSink sink) { return new DisruptorEventBus(1024, "blocking", List.of(sink)); @@ -160,6 +169,41 @@ class Jt1078AutoConfigurationTest { } } + private static final class EmptySessionStore implements SessionStore { + @Override + public void put(DeviceSession session) { + } + + @Override + public Optional findBySessionId(String sessionId) { + return Optional.empty(); + } + + @Override + public Optional findByVin(String vin) { + return Optional.empty(); + } + + @Override + public Optional findByPhone(String phone) { + return Optional.empty(); + } + + @Override + public Optional update(String sessionId, UnaryOperator updater) { + return Optional.empty(); + } + + @Override + public void remove(String sessionId) { + } + + @Override + public int size() { + return 0; + } + } + private static final class CapturingSink implements EventSink { private final CopyOnWriteArrayList events = new CopyOnWriteArrayList<>();