refactor: route jt808 mileage through dedicated stream
This commit is contained in:
@@ -29,8 +29,10 @@ public final class VehicleStatEnvelopeIngestor implements EnvelopeIngestor {
|
|||||||
if (!envelope.hasTelemetrySnapshot()) {
|
if (!envelope.hasTelemetrySnapshot()) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
if (processJt808Mileage(envelope)) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
processor.process(envelope);
|
processor.process(envelope);
|
||||||
processJt808Mileage(envelope);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
@@ -42,8 +44,10 @@ public final class VehicleStatEnvelopeIngestor implements EnvelopeIngestor {
|
|||||||
return EnvelopeIngestResult.skipped(
|
return EnvelopeIngestResult.skipped(
|
||||||
envelope.getEventId(), envelope.getVin(), "envelope telemetry_snapshot is required");
|
envelope.getEventId(), envelope.getVin(), "envelope telemetry_snapshot is required");
|
||||||
}
|
}
|
||||||
|
if (processJt808Mileage(envelope)) {
|
||||||
|
return EnvelopeIngestResult.processed(envelope.getEventId(), envelope.getVin());
|
||||||
|
}
|
||||||
processor.process(envelope);
|
processor.process(envelope);
|
||||||
processJt808Mileage(envelope);
|
|
||||||
return EnvelopeIngestResult.processed(envelope.getEventId(), envelope.getVin());
|
return EnvelopeIngestResult.processed(envelope.getEventId(), envelope.getVin());
|
||||||
} catch (IllegalArgumentException ex) {
|
} catch (IllegalArgumentException ex) {
|
||||||
return envelope == null
|
return envelope == null
|
||||||
@@ -52,10 +56,12 @@ public final class VehicleStatEnvelopeIngestor implements EnvelopeIngestor {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private void processJt808Mileage(VehicleEnvelope envelope) {
|
private boolean processJt808Mileage(VehicleEnvelope envelope) {
|
||||||
if (jt808MileageProcessor != null) {
|
if (jt808MileageProcessor != null && "JT808".equalsIgnoreCase(envelope.getSource())) {
|
||||||
jt808MileageProcessor.process(envelope);
|
jt808MileageProcessor.process(envelope);
|
||||||
|
return true;
|
||||||
}
|
}
|
||||||
|
return false;
|
||||||
}
|
}
|
||||||
|
|
||||||
private static VehicleEnvelope parse(byte[] kafkaValue) {
|
private static VehicleEnvelope parse(byte[] kafkaValue) {
|
||||||
|
|||||||
@@ -17,6 +17,7 @@ import static org.assertj.core.api.Assertions.assertThat;
|
|||||||
import static org.assertj.core.api.Assertions.assertThatThrownBy;
|
import static org.assertj.core.api.Assertions.assertThatThrownBy;
|
||||||
import static org.mockito.Mockito.mock;
|
import static org.mockito.Mockito.mock;
|
||||||
import static org.mockito.Mockito.verify;
|
import static org.mockito.Mockito.verify;
|
||||||
|
import static org.mockito.Mockito.verifyNoInteractions;
|
||||||
|
|
||||||
class VehicleStatEnvelopeIngestorTest {
|
class VehicleStatEnvelopeIngestorTest {
|
||||||
|
|
||||||
@@ -35,10 +36,27 @@ class VehicleStatEnvelopeIngestorTest {
|
|||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void fansOutTelemetryEnvelopeToOptionalJt808MileageProcessor() {
|
void doesNotFanOutNonJt808TelemetryToJt808MileageProcessor() {
|
||||||
InMemoryVehicleStatRepository repository = new InMemoryVehicleStatRepository();
|
InMemoryVehicleStatRepository repository = new InMemoryVehicleStatRepository();
|
||||||
Jt808MileageStreamProcessor jt808Processor = mock(Jt808MileageStreamProcessor.class);
|
Jt808MileageStreamProcessor jt808Processor = mock(Jt808MileageStreamProcessor.class);
|
||||||
VehicleEnvelope envelope = envelope();
|
VehicleEnvelope envelope = envelope("GB32960");
|
||||||
|
VehicleStatEnvelopeIngestor ingestor = new VehicleStatEnvelopeIngestor(
|
||||||
|
new VehicleStatEventProcessor(repository),
|
||||||
|
jt808Processor);
|
||||||
|
|
||||||
|
ingestor.ingest(envelope.toByteArray());
|
||||||
|
|
||||||
|
verifyNoInteractions(jt808Processor);
|
||||||
|
assertThat(repository.mileagePoints("VIN001", LocalDate.of(2026, 6, 22)))
|
||||||
|
.extracting(MileagePoint::totalMileageKm)
|
||||||
|
.containsExactly(123.45);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void routesJt808TelemetryOnlyToJt808MileageProcessorWhenAvailable() {
|
||||||
|
InMemoryVehicleStatRepository repository = new InMemoryVehicleStatRepository();
|
||||||
|
Jt808MileageStreamProcessor jt808Processor = mock(Jt808MileageStreamProcessor.class);
|
||||||
|
VehicleEnvelope envelope = envelope("JT808");
|
||||||
VehicleStatEnvelopeIngestor ingestor = new VehicleStatEnvelopeIngestor(
|
VehicleStatEnvelopeIngestor ingestor = new VehicleStatEnvelopeIngestor(
|
||||||
new VehicleStatEventProcessor(repository),
|
new VehicleStatEventProcessor(repository),
|
||||||
jt808Processor);
|
jt808Processor);
|
||||||
@@ -46,6 +64,7 @@ class VehicleStatEnvelopeIngestorTest {
|
|||||||
ingestor.ingest(envelope.toByteArray());
|
ingestor.ingest(envelope.toByteArray());
|
||||||
|
|
||||||
verify(jt808Processor).process(envelope);
|
verify(jt808Processor).process(envelope);
|
||||||
|
assertThat(repository.mileagePoints("VIN001", LocalDate.of(2026, 6, 22))).isEmpty();
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
@@ -110,10 +129,14 @@ class VehicleStatEnvelopeIngestorTest {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private static VehicleEnvelope envelope() {
|
private static VehicleEnvelope envelope() {
|
||||||
|
return envelope("GB32960");
|
||||||
|
}
|
||||||
|
|
||||||
|
private static VehicleEnvelope envelope(String source) {
|
||||||
return VehicleEnvelope.newBuilder()
|
return VehicleEnvelope.newBuilder()
|
||||||
.setEventId("event-1")
|
.setEventId("event-1")
|
||||||
.setVin("VIN001")
|
.setVin("VIN001")
|
||||||
.setSource("GB32960")
|
.setSource(source)
|
||||||
.setEventTimeMs(1_782_112_400_000L)
|
.setEventTimeMs(1_782_112_400_000L)
|
||||||
.setIngestTimeMs(1_782_112_401_000L)
|
.setIngestTimeMs(1_782_112_401_000L)
|
||||||
.setTelemetrySnapshot(TelemetrySnapshot.newBuilder()
|
.setTelemetrySnapshot(TelemetrySnapshot.newBuilder()
|
||||||
|
|||||||
Reference in New Issue
Block a user