refactor: require identity resolver for jt1078 malformed rtp

This commit is contained in:
lingniu
2026-07-01 09:07:12 +08:00
parent a321b6351d
commit 42e8742b10
9 changed files with 68 additions and 108 deletions

View File

@@ -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 入口坏包兜底")

View File

@@ -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<RawFrame> dispatcher;
public Jt1078MalformedRtpEventPublisher(Consumer<RawFrame> dispatcher) {
this(new StatelessVehicleIdentityResolver(), dispatcher);
}
public Jt1078MalformedRtpEventPublisher(VehicleIdentityResolver identityResolver,
Consumer<RawFrame> 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) {

View File

@@ -19,7 +19,8 @@ class Jt1078MalformedRtpEventPublisherTest {
@Test
void malformedRtpPacketProducesRawFrameForUnifiedArchiveAndPassthrough() {
ArrayList<RawFrame> 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};

View File

@@ -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};

View File

@@ -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<DeviceSession> findBySessionId(String sessionId) {
return Optional.empty();
}
@Override
public Optional<DeviceSession> findByVin(String vin) {
return Optional.empty();
}
@Override
public Optional<DeviceSession> findByPhone(String phone) {
return Optional.empty();
}
@Override
public Optional<DeviceSession> update(String sessionId, UnaryOperator<DeviceSession> 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<VehicleEvent> events = new CopyOnWriteArrayList<>();