From 310ef94a6e4b7dbd47d4f47fa65d69d6820934ab Mon Sep 17 00:00:00 2001 From: lingniu Date: Mon, 29 Jun 2026 12:50:41 +0800 Subject: [PATCH] docs: plan vehicle ingest redesign phase 1 --- ...26-06-29-vehicle-ingest-redesign-phase1.md | 1512 +++++++++++++++++ .../2026-06-29-vehicle-ingest-redesign.md | 8 + .../2026-06-29-vehicle-ingest-redesign.zh.md | 8 + 3 files changed, 1528 insertions(+) create mode 100644 docs/superpowers/plans/2026-06-29-vehicle-ingest-redesign-phase1.md diff --git a/docs/superpowers/plans/2026-06-29-vehicle-ingest-redesign-phase1.md b/docs/superpowers/plans/2026-06-29-vehicle-ingest-redesign-phase1.md new file mode 100644 index 00000000..a0fe0fe4 --- /dev/null +++ b/docs/superpowers/plans/2026-06-29-vehicle-ingest-redesign-phase1.md @@ -0,0 +1,1512 @@ +# Vehicle Ingest Redesign Phase 1 Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Build the foundation for the new GB32960/JT808 architecture: shared fact model, raw archive contract, and byte-light Kafka schema. + +**Architecture:** Phase 1 adds new modules beside existing runtime paths. It does not change current Netty handlers, ACK behavior, or TDengine writes yet. The output is a tested set of immutable facts and contracts that later phases will wire into protocol apps and the history writer. + +**Tech Stack:** Java 25, Maven, JUnit 5, AssertJ, Protobuf 4, existing `ProtocolId`. + +--- + +## Scope + +This plan implements only the first slice of the redesign: + +- `modules/core/ingest-facts`: protocol-neutral fact records and helpers. +- `modules/sinks/raw-archive-store`: raw archive write/read contract and local URI-safe implementation skeleton. +- `modules/sinks/sink-mq`: protobuf messages for raw frame facts and decoded facts. +- Parent Maven wiring and focused tests. + +It intentionally does not yet modify: + +- `gb32960-ingest-app` +- `jt808-ingest-app` +- existing `Dispatcher` +- existing `VehicleEvent` +- existing `vehicle-history-app` TDengine writer + +Those changes belong to later implementation plans. + +## File Structure + +Create: + +- `modules/core/ingest-facts/pom.xml`: Maven module for fact model. +- `modules/core/ingest-facts/src/main/java/com/lingniu/ingest/facts/ParseStatus.java`: parse status enum. +- `modules/core/ingest-facts/src/main/java/com/lingniu/ingest/facts/VehicleKey.java`: stable vehicle key derivation. +- `modules/core/ingest-facts/src/main/java/com/lingniu/ingest/facts/RawFrameFact.java`: raw frame fact record. +- `modules/core/ingest-facts/src/main/java/com/lingniu/ingest/facts/DecodedFact.java`: sealed decoded fact hierarchy. +- `modules/core/ingest-facts/src/main/java/com/lingniu/ingest/facts/FactIds.java`: deterministic id helpers. +- `modules/core/ingest-facts/src/test/java/com/lingniu/ingest/facts/VehicleKeyTest.java` +- `modules/core/ingest-facts/src/test/java/com/lingniu/ingest/facts/RawFrameFactTest.java` +- `modules/core/ingest-facts/src/test/java/com/lingniu/ingest/facts/FactIdsTest.java` +- `modules/sinks/raw-archive-store/pom.xml`: Maven module for raw archive contracts. +- `modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveWriteRequest.java` +- `modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveReceipt.java` +- `modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveWriter.java` +- `modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveReader.java` +- `modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/LocalRawArchiveStore.java` +- `modules/sinks/raw-archive-store/src/test/java/com/lingniu/ingest/rawarchive/LocalRawArchiveStoreTest.java` + +Modify: + +- `pom.xml`: add modules and dependency management entries. +- `modules/sinks/sink-mq/pom.xml`: add dependency on `ingest-facts`. +- `modules/sinks/sink-mq/src/main/proto/vehicle_envelope.proto`: add `RawFrameFactPayload`, `DecodedFactPayload`, and related enums/messages. +- `modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/VehicleEnvelopeProtoCompatibilityTest.java`: verify new proto messages can round-trip. + +--- + +### Task 1: Add `ingest-facts` Maven Module + +**Files:** + +- Modify: `pom.xml` +- Create: `modules/core/ingest-facts/pom.xml` +- Test later in Task 2 through Task 4. + +- [ ] **Step 1: Add failing Maven target check** + +Run: + +```bash +mvn -pl modules/core/ingest-facts test +``` + +Expected: FAIL because Maven cannot find the selected project. + +- [ ] **Step 2: Add module to root `pom.xml`** + +Add this module after `modules/core/ingest-api`: + +```xml +modules/core/ingest-facts +``` + +Add this dependency management entry after `ingest-api`: + +```xml + + com.lingniu.ingest + ingest-facts + ${project.version} + +``` + +- [ ] **Step 3: Create module POM** + +Create `modules/core/ingest-facts/pom.xml`: + +```xml + + + 4.0.0 + + com.lingniu.ingest + lingniu-vehicle-ingest + 0.1.0-SNAPSHOT + ../../../pom.xml + + ingest-facts + ingest-facts + Protocol-neutral raw frame and decoded fact model for vehicle ingestion. + + + + com.lingniu.ingest + ingest-api + + + org.junit.jupiter + junit-jupiter + test + + + org.assertj + assertj-core + test + + + +``` + +- [ ] **Step 4: Run Maven target check** + +Run: + +```bash +mvn -pl modules/core/ingest-facts test +``` + +Expected: PASS with no tests. + +- [ ] **Step 5: Commit** + +```bash +git add pom.xml modules/core/ingest-facts/pom.xml +git commit -m "build: add ingest facts module" +``` + +--- + +### Task 2: Implement Vehicle Key Derivation + +**Files:** + +- Create: `modules/core/ingest-facts/src/main/java/com/lingniu/ingest/facts/VehicleKey.java` +- Create: `modules/core/ingest-facts/src/test/java/com/lingniu/ingest/facts/VehicleKeyTest.java` + +- [ ] **Step 1: Write failing tests** + +Create `VehicleKeyTest.java`: + +```java +package com.lingniu.ingest.facts; + +import com.lingniu.ingest.api.ProtocolId; +import org.junit.jupiter.api.Test; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +class VehicleKeyTest { + + @Test + void usesVinForGb32960() { + assertThat(VehicleKey.derive(ProtocolId.GB32960, " LB9A32A20P0LS1257 ", "", "abc")) + .isEqualTo("LB9A32A20P0LS1257"); + } + + @Test + void usesResolvedVinForJt808WhenPresent() { + assertThat(VehicleKey.derive(ProtocolId.JT808, "VIN001", "013912345678", "abc")) + .isEqualTo("VIN001"); + } + + @Test + void usesPhoneForJt808WhenVinIsMissing() { + assertThat(VehicleKey.derive(ProtocolId.JT808, "", "013912345678", "abc")) + .isEqualTo("jt808:013912345678"); + } + + @Test + void usesStableUnknownHashWhenVehicleIdentifiersAreMissing() { + assertThat(VehicleKey.derive(ProtocolId.JT808, "", "", "abc")) + .isEqualTo(VehicleKey.derive(ProtocolId.JT808, null, null, "abc")) + .startsWith("unknown:JT808:"); + } + + @Test + void rejectsMissingProtocol() { + assertThatThrownBy(() -> VehicleKey.derive(null, "VIN001", "", "abc")) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("protocol"); + } +} +``` + +- [ ] **Step 2: Run test and verify RED** + +Run: + +```bash +mvn -pl modules/core/ingest-facts -Dtest=VehicleKeyTest test +``` + +Expected: FAIL because `VehicleKey` does not exist. + +- [ ] **Step 3: Implement `VehicleKey`** + +Create `VehicleKey.java`: + +```java +package com.lingniu.ingest.facts; + +import com.lingniu.ingest.api.ProtocolId; + +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.util.HexFormat; + +public final class VehicleKey { + + private VehicleKey() { + } + + public static String derive(ProtocolId protocol, String vin, String phone, String rawIdentitySeed) { + if (protocol == null) { + throw new IllegalArgumentException("protocol must not be null"); + } + String normalizedVin = clean(vin); + if (!normalizedVin.isBlank()) { + return normalizedVin; + } + String normalizedPhone = clean(phone); + if (protocol == ProtocolId.JT808 && !normalizedPhone.isBlank()) { + return "jt808:" + normalizedPhone; + } + String seed = clean(rawIdentitySeed); + if (seed.isBlank()) { + seed = protocol.name(); + } + return "unknown:" + protocol.name() + ":" + sha256(seed).substring(0, 16); + } + + private static String clean(String value) { + return value == null ? "" : value.trim(); + } + + private static String sha256(String value) { + try { + MessageDigest digest = MessageDigest.getInstance("SHA-256"); + return HexFormat.of().formatHex(digest.digest(value.getBytes(StandardCharsets.UTF_8))); + } catch (NoSuchAlgorithmException e) { + throw new IllegalStateException("SHA-256 is unavailable", e); + } + } +} +``` + +- [ ] **Step 4: Run test and verify GREEN** + +Run: + +```bash +mvn -pl modules/core/ingest-facts -Dtest=VehicleKeyTest test +``` + +Expected: PASS. + +- [ ] **Step 5: Commit** + +```bash +git add modules/core/ingest-facts/src/main/java/com/lingniu/ingest/facts/VehicleKey.java \ + modules/core/ingest-facts/src/test/java/com/lingniu/ingest/facts/VehicleKeyTest.java +git commit -m "feat: add vehicle key derivation" +``` + +--- + +### Task 3: Implement Raw Frame Fact and Id Helpers + +**Files:** + +- Create: `modules/core/ingest-facts/src/main/java/com/lingniu/ingest/facts/ParseStatus.java` +- Create: `modules/core/ingest-facts/src/main/java/com/lingniu/ingest/facts/FactIds.java` +- Create: `modules/core/ingest-facts/src/main/java/com/lingniu/ingest/facts/RawFrameFact.java` +- Create: `modules/core/ingest-facts/src/test/java/com/lingniu/ingest/facts/FactIdsTest.java` +- Create: `modules/core/ingest-facts/src/test/java/com/lingniu/ingest/facts/RawFrameFactTest.java` + +- [ ] **Step 1: Write failing tests** + +Create `FactIdsTest.java`: + +```java +package com.lingniu.ingest.facts; + +import com.lingniu.ingest.api.ProtocolId; +import org.junit.jupiter.api.Test; + +import java.time.Instant; + +import static org.assertj.core.api.Assertions.assertThat; + +class FactIdsTest { + + @Test + void rawFrameIdIsStableForSameInputs() { + String first = FactIds.rawFrameId( + ProtocolId.GB32960, + "VIN001", + 2, + 7, + Instant.parse("2026-06-29T00:00:00Z"), + "sha256:abc"); + String second = FactIds.rawFrameId( + ProtocolId.GB32960, + "VIN001", + 2, + 7, + Instant.parse("2026-06-29T00:00:00Z"), + "sha256:abc"); + + assertThat(first).isEqualTo(second).startsWith("rf_"); + } + + @Test + void decodedFactIdIncludesKind() { + String id = FactIds.decodedFactId("rf_123", "location", 0); + + assertThat(id).isEqualTo("df_rf_123_location_0"); + } +} +``` + +Create `RawFrameFactTest.java`: + +```java +package com.lingniu.ingest.facts; + +import com.lingniu.ingest.api.ProtocolId; +import org.junit.jupiter.api.Test; + +import java.time.Instant; +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +class RawFrameFactTest { + + @Test + void normalizesOptionalStringsAndMetadata() { + RawFrameFact fact = new RawFrameFact( + "frame-1", + ProtocolId.JT808, + "jt808:013912345678", + null, + "013912345678", + 0x0200, + 0, + null, + Instant.parse("2026-06-29T00:00:00Z"), + null, + "archive://2026/06/29/JT808/jt808-013912345678/frame-1.bin", + "sha256:abc", + 128, + ParseStatus.SUCCEEDED, + null, + Map.of("messageName", "LOCATION")); + + assertThat(fact.vin()).isEmpty(); + assertThat(fact.eventTime()).isEqualTo(fact.receivedAt()); + assertThat(fact.peer()).isEmpty(); + assertThat(fact.parseError()).isEmpty(); + assertThat(fact.metadata()).containsEntry("messageName", "LOCATION"); + } + + @Test + void rejectsMissingRequiredFields() { + assertThatThrownBy(() -> new RawFrameFact( + "", + ProtocolId.GB32960, + "VIN001", + "VIN001", + "", + 2, + 0, + Instant.parse("2026-06-29T00:00:00Z"), + Instant.parse("2026-06-29T00:00:00Z"), + "", + "archive://x", + "sha256:abc", + 1, + ParseStatus.SUCCEEDED, + "", + Map.of())) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("frameId"); + } +} +``` + +- [ ] **Step 2: Run tests and verify RED** + +Run: + +```bash +mvn -pl modules/core/ingest-facts -Dtest=FactIdsTest,RawFrameFactTest test +``` + +Expected: FAIL because `FactIds`, `RawFrameFact`, and `ParseStatus` do not exist. + +- [ ] **Step 3: Implement `ParseStatus`** + +Create `ParseStatus.java`: + +```java +package com.lingniu.ingest.facts; + +public enum ParseStatus { + NOT_PARSED, + SUCCEEDED, + FAILED +} +``` + +- [ ] **Step 4: Implement `FactIds`** + +Create `FactIds.java`: + +```java +package com.lingniu.ingest.facts; + +import com.lingniu.ingest.api.ProtocolId; + +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.time.Instant; +import java.util.HexFormat; + +public final class FactIds { + + private FactIds() { + } + + public static String rawFrameId(ProtocolId protocol, + String vehicleKey, + int messageId, + int subType, + Instant receivedAt, + String checksum) { + if (protocol == null) { + throw new IllegalArgumentException("protocol must not be null"); + } + String input = protocol.name() + "|" + + clean(vehicleKey) + "|" + + messageId + "|" + + subType + "|" + + (receivedAt == null ? "" : receivedAt.toString()) + "|" + + clean(checksum); + return "rf_" + sha256(input).substring(0, 32); + } + + public static String decodedFactId(String frameId, String kind, int ordinal) { + return "df_" + clean(frameId) + "_" + clean(kind) + "_" + ordinal; + } + + private static String clean(String value) { + return value == null ? "" : value.trim(); + } + + private static String sha256(String value) { + try { + MessageDigest digest = MessageDigest.getInstance("SHA-256"); + return HexFormat.of().formatHex(digest.digest(value.getBytes(StandardCharsets.UTF_8))); + } catch (NoSuchAlgorithmException e) { + throw new IllegalStateException("SHA-256 is unavailable", e); + } + } +} +``` + +- [ ] **Step 5: Implement `RawFrameFact`** + +Create `RawFrameFact.java`: + +```java +package com.lingniu.ingest.facts; + +import com.lingniu.ingest.api.ProtocolId; + +import java.time.Instant; +import java.util.Map; + +public record RawFrameFact( + String frameId, + ProtocolId protocol, + String vehicleKey, + String vin, + String phone, + int messageId, + int subType, + Instant eventTime, + Instant receivedAt, + String peer, + String rawUri, + String checksum, + long rawSizeBytes, + ParseStatus parseStatus, + String parseError, + Map metadata +) { + public RawFrameFact { + frameId = required(frameId, "frameId"); + if (protocol == null) { + throw new IllegalArgumentException("protocol must not be null"); + } + vehicleKey = required(vehicleKey, "vehicleKey"); + receivedAt = required(receivedAt, "receivedAt"); + eventTime = eventTime == null ? receivedAt : eventTime; + rawUri = required(rawUri, "rawUri"); + checksum = required(checksum, "checksum"); + if (rawSizeBytes < 0) { + throw new IllegalArgumentException("rawSizeBytes must not be negative"); + } + parseStatus = parseStatus == null ? ParseStatus.NOT_PARSED : parseStatus; + vin = clean(vin); + phone = clean(phone); + peer = clean(peer); + parseError = clean(parseError); + metadata = metadata == null ? Map.of() : Map.copyOf(metadata); + } + + private static String required(String value, String field) { + String cleaned = clean(value); + if (cleaned.isBlank()) { + throw new IllegalArgumentException(field + " must not be blank"); + } + return cleaned; + } + + private static Instant required(Instant value, String field) { + if (value == null) { + throw new IllegalArgumentException(field + " must not be null"); + } + return value; + } + + private static String clean(String value) { + return value == null ? "" : value.trim(); + } +} +``` + +- [ ] **Step 6: Run tests and verify GREEN** + +Run: + +```bash +mvn -pl modules/core/ingest-facts -Dtest=FactIdsTest,RawFrameFactTest test +``` + +Expected: PASS. + +- [ ] **Step 7: Commit** + +```bash +git add modules/core/ingest-facts/src/main/java/com/lingniu/ingest/facts \ + modules/core/ingest-facts/src/test/java/com/lingniu/ingest/facts +git commit -m "feat: add raw frame fact model" +``` + +--- + +### Task 4: Implement Decoded Fact Records + +**Files:** + +- Create: `modules/core/ingest-facts/src/main/java/com/lingniu/ingest/facts/DecodedFact.java` +- Create: `modules/core/ingest-facts/src/test/java/com/lingniu/ingest/facts/DecodedFactTest.java` + +- [ ] **Step 1: Write failing test** + +Create `DecodedFactTest.java`: + +```java +package com.lingniu.ingest.facts; + +import com.lingniu.ingest.api.ProtocolId; +import org.junit.jupiter.api.Test; + +import java.time.Instant; +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; + +class DecodedFactTest { + + @Test + void locationFactCarriesRawReferenceAndMetadata() { + DecodedFact.Location location = new DecodedFact.Location( + "fact-1", + "frame-1", + ProtocolId.JT808, + "jt808:013912345678", + "", + "013912345678", + Instant.parse("2026-06-29T00:00:00Z"), + Instant.parse("2026-06-29T00:00:01Z"), + "archive://jt808/frame-1.bin", + Map.of("messageId", "0x0200"), + 120.123456, + 30.123456, + 10.0, + 42.0, + 180.0, + 0L, + 2L, + 1000.5); + + assertThat(location.factId()).isEqualTo("fact-1"); + assertThat(location.rawUri()).isEqualTo("archive://jt808/frame-1.bin"); + assertThat(location.metadata()).containsEntry("messageId", "0x0200"); + } +} +``` + +- [ ] **Step 2: Run test and verify RED** + +Run: + +```bash +mvn -pl modules/core/ingest-facts -Dtest=DecodedFactTest test +``` + +Expected: FAIL because `DecodedFact` does not exist. + +- [ ] **Step 3: Implement `DecodedFact`** + +Create `DecodedFact.java`: + +```java +package com.lingniu.ingest.facts; + +import com.lingniu.ingest.api.ProtocolId; + +import java.time.Instant; +import java.util.Map; + +public sealed interface DecodedFact + permits DecodedFact.Location, + DecodedFact.Realtime, + DecodedFact.Alarm, + DecodedFact.Session, + DecodedFact.Register, + DecodedFact.Extension { + + String factId(); + String frameId(); + ProtocolId protocol(); + String vehicleKey(); + String vin(); + String phone(); + Instant eventTime(); + Instant receivedAt(); + String rawUri(); + Map metadata(); + + record Location( + String factId, + String frameId, + ProtocolId protocol, + String vehicleKey, + String vin, + String phone, + Instant eventTime, + Instant receivedAt, + String rawUri, + Map metadata, + Double longitude, + Double latitude, + Double altitudeM, + Double speedKmh, + Double directionDeg, + Long alarmFlag, + Long statusFlag, + Double totalMileageKm + ) implements DecodedFact { + public Location { + BaseFields.validate(factId, frameId, protocol, vehicleKey, eventTime, receivedAt, rawUri); + vin = BaseFields.clean(vin); + phone = BaseFields.clean(phone); + metadata = BaseFields.metadata(metadata); + } + } + + record Realtime( + String factId, + String frameId, + ProtocolId protocol, + String vehicleKey, + String vin, + String phone, + Instant eventTime, + Instant receivedAt, + String rawUri, + Map metadata, + Map fields + ) implements DecodedFact { + public Realtime { + BaseFields.validate(factId, frameId, protocol, vehicleKey, eventTime, receivedAt, rawUri); + vin = BaseFields.clean(vin); + phone = BaseFields.clean(phone); + metadata = BaseFields.metadata(metadata); + fields = fields == null ? Map.of() : Map.copyOf(fields); + } + } + + record Alarm( + String factId, + String frameId, + ProtocolId protocol, + String vehicleKey, + String vin, + String phone, + Instant eventTime, + Instant receivedAt, + String rawUri, + Map metadata, + String alarmLevel, + Long alarmCode, + String alarmName, + Map fields + ) implements DecodedFact { + public Alarm { + BaseFields.validate(factId, frameId, protocol, vehicleKey, eventTime, receivedAt, rawUri); + vin = BaseFields.clean(vin); + phone = BaseFields.clean(phone); + alarmLevel = BaseFields.clean(alarmLevel); + alarmName = BaseFields.clean(alarmName); + metadata = BaseFields.metadata(metadata); + fields = fields == null ? Map.of() : Map.copyOf(fields); + } + } + + record Session( + String factId, + String frameId, + ProtocolId protocol, + String vehicleKey, + String vin, + String phone, + Instant eventTime, + Instant receivedAt, + String rawUri, + Map metadata, + String sessionType, + Integer resultCode, + Map fields + ) implements DecodedFact { + public Session { + BaseFields.validate(factId, frameId, protocol, vehicleKey, eventTime, receivedAt, rawUri); + vin = BaseFields.clean(vin); + phone = BaseFields.clean(phone); + sessionType = BaseFields.clean(sessionType); + metadata = BaseFields.metadata(metadata); + fields = fields == null ? Map.of() : Map.copyOf(fields); + } + } + + record Register( + String factId, + String frameId, + ProtocolId protocol, + String vehicleKey, + String vin, + String phone, + Instant eventTime, + Instant receivedAt, + String rawUri, + Map metadata, + Map registrationFields + ) implements DecodedFact { + public Register { + BaseFields.validate(factId, frameId, protocol, vehicleKey, eventTime, receivedAt, rawUri); + vin = BaseFields.clean(vin); + phone = BaseFields.clean(phone); + metadata = BaseFields.metadata(metadata); + registrationFields = registrationFields == null ? Map.of() : Map.copyOf(registrationFields); + } + } + + record Extension( + String factId, + String frameId, + ProtocolId protocol, + String vehicleKey, + String vin, + String phone, + Instant eventTime, + Instant receivedAt, + String rawUri, + Map metadata, + String extensionType, + Map fields + ) implements DecodedFact { + public Extension { + BaseFields.validate(factId, frameId, protocol, vehicleKey, eventTime, receivedAt, rawUri); + vin = BaseFields.clean(vin); + phone = BaseFields.clean(phone); + extensionType = BaseFields.clean(extensionType); + metadata = BaseFields.metadata(metadata); + fields = fields == null ? Map.of() : Map.copyOf(fields); + } + } + + final class BaseFields { + private BaseFields() { + } + + static void validate(String factId, + String frameId, + ProtocolId protocol, + String vehicleKey, + Instant eventTime, + Instant receivedAt, + String rawUri) { + required(factId, "factId"); + required(frameId, "frameId"); + if (protocol == null) { + throw new IllegalArgumentException("protocol must not be null"); + } + required(vehicleKey, "vehicleKey"); + if (eventTime == null) { + throw new IllegalArgumentException("eventTime must not be null"); + } + if (receivedAt == null) { + throw new IllegalArgumentException("receivedAt must not be null"); + } + required(rawUri, "rawUri"); + } + + static String clean(String value) { + return value == null ? "" : value.trim(); + } + + static Map metadata(Map metadata) { + return metadata == null ? Map.of() : Map.copyOf(metadata); + } + + private static void required(String value, String field) { + if (clean(value).isBlank()) { + throw new IllegalArgumentException(field + " must not be blank"); + } + } + } +} +``` + +- [ ] **Step 4: Run test and verify GREEN** + +Run: + +```bash +mvn -pl modules/core/ingest-facts -Dtest=DecodedFactTest test +``` + +Expected: PASS. + +- [ ] **Step 5: Commit** + +```bash +git add modules/core/ingest-facts/src/main/java/com/lingniu/ingest/facts/DecodedFact.java \ + modules/core/ingest-facts/src/test/java/com/lingniu/ingest/facts/DecodedFactTest.java +git commit -m "feat: add decoded fact model" +``` + +--- + +### Task 5: Add Raw Archive Store Contract and Local Implementation + +**Files:** + +- Modify: `pom.xml` +- Create: `modules/sinks/raw-archive-store/pom.xml` +- Create: `modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveWriteRequest.java` +- Create: `modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveReceipt.java` +- Create: `modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveWriter.java` +- Create: `modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveReader.java` +- Create: `modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/LocalRawArchiveStore.java` +- Create: `modules/sinks/raw-archive-store/src/test/java/com/lingniu/ingest/rawarchive/LocalRawArchiveStoreTest.java` + +- [ ] **Step 1: Write failing local archive test** + +Create `LocalRawArchiveStoreTest.java`: + +```java +package com.lingniu.ingest.rawarchive; + +import com.lingniu.ingest.api.ProtocolId; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.nio.file.Path; +import java.time.Instant; +import java.util.HexFormat; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +class LocalRawArchiveStoreTest { + + @TempDir + Path tempDir; + + @Test + void writesReadableRawObjectWithArchiveUriAndChecksum() throws Exception { + LocalRawArchiveStore store = new LocalRawArchiveStore(tempDir, "archive://raw"); + byte[] bytes = HexFormat.of().parseHex("7e020000007e"); + + RawArchiveReceipt receipt = store.write(new RawArchiveWriteRequest( + ProtocolId.JT808, + "jt808:013912345678", + "frame-1", + Instant.parse("2026-06-29T00:00:00Z"), + bytes)); + + assertThat(receipt.rawUri()).isEqualTo("archive://raw/2026/06/29/JT808/jt808_013912345678/frame-1.bin"); + assertThat(receipt.checksum()).startsWith("sha256:"); + assertThat(receipt.sizeBytes()).isEqualTo(bytes.length); + assertThat(store.read(receipt.rawUri())).containsExactly(bytes); + } + + @Test + void rejectsBlankFrameId() { + LocalRawArchiveStore store = new LocalRawArchiveStore(tempDir, "archive://raw"); + + assertThatThrownBy(() -> store.write(new RawArchiveWriteRequest( + ProtocolId.GB32960, + "VIN001", + "", + Instant.parse("2026-06-29T00:00:00Z"), + new byte[]{1}))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("frameId"); + } +} +``` + +- [ ] **Step 2: Run test and verify RED** + +Run: + +```bash +mvn -pl modules/sinks/raw-archive-store test +``` + +Expected: FAIL because the module does not exist. + +- [ ] **Step 3: Add module and POM wiring** + +Add this root module after `modules/sinks/sink-archive`: + +```xml +modules/sinks/raw-archive-store +``` + +Add dependency management after `sink-archive`: + +```xml + + com.lingniu.ingest + raw-archive-store + ${project.version} + +``` + +Create `modules/sinks/raw-archive-store/pom.xml`: + +```xml + + + 4.0.0 + + com.lingniu.ingest + lingniu-vehicle-ingest + 0.1.0-SNAPSHOT + ../../../pom.xml + + raw-archive-store + raw-archive-store + Raw frame archive writer and reader contracts. + + + + com.lingniu.ingest + ingest-api + + + org.junit.jupiter + junit-jupiter + test + + + org.assertj + assertj-core + test + + + +``` + +- [ ] **Step 4: Implement contracts** + +Create `RawArchiveWriteRequest.java`: + +```java +package com.lingniu.ingest.rawarchive; + +import com.lingniu.ingest.api.ProtocolId; + +import java.time.Instant; + +public record RawArchiveWriteRequest( + ProtocolId protocol, + String vehicleKey, + String frameId, + Instant receivedAt, + byte[] rawBytes +) { + public RawArchiveWriteRequest { + if (protocol == null) { + throw new IllegalArgumentException("protocol must not be null"); + } + vehicleKey = required(vehicleKey, "vehicleKey"); + frameId = required(frameId, "frameId"); + if (receivedAt == null) { + throw new IllegalArgumentException("receivedAt must not be null"); + } + if (rawBytes == null || rawBytes.length == 0) { + throw new IllegalArgumentException("rawBytes must not be empty"); + } + rawBytes = rawBytes.clone(); + } + + @Override + public byte[] rawBytes() { + return rawBytes.clone(); + } + + private static String required(String value, String field) { + String cleaned = value == null ? "" : value.trim(); + if (cleaned.isBlank()) { + throw new IllegalArgumentException(field + " must not be blank"); + } + return cleaned; + } +} +``` + +Create `RawArchiveReceipt.java`: + +```java +package com.lingniu.ingest.rawarchive; + +public record RawArchiveReceipt(String rawUri, String checksum, long sizeBytes) { + public RawArchiveReceipt { + rawUri = required(rawUri, "rawUri"); + checksum = required(checksum, "checksum"); + if (sizeBytes < 0) { + throw new IllegalArgumentException("sizeBytes must not be negative"); + } + } + + private static String required(String value, String field) { + String cleaned = value == null ? "" : value.trim(); + if (cleaned.isBlank()) { + throw new IllegalArgumentException(field + " must not be blank"); + } + return cleaned; + } +} +``` + +Create `RawArchiveWriter.java`: + +```java +package com.lingniu.ingest.rawarchive; + +import java.io.IOException; + +public interface RawArchiveWriter { + RawArchiveReceipt write(RawArchiveWriteRequest request) throws IOException; +} +``` + +Create `RawArchiveReader.java`: + +```java +package com.lingniu.ingest.rawarchive; + +import java.io.IOException; + +public interface RawArchiveReader { + byte[] read(String rawUri) throws IOException; +} +``` + +- [ ] **Step 5: Implement local store** + +Create `LocalRawArchiveStore.java`: + +```java +package com.lingniu.ingest.rawarchive; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.LinkOption; +import java.nio.file.Path; +import java.nio.file.StandardOpenOption; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.time.ZoneId; +import java.util.HexFormat; + +public final class LocalRawArchiveStore implements RawArchiveWriter, RawArchiveReader { + + private static final ZoneId PARTITION_ZONE = ZoneId.of("Asia/Shanghai"); + + private final Path root; + private final String uriPrefix; + + public LocalRawArchiveStore(Path root, String uriPrefix) { + if (root == null) { + throw new IllegalArgumentException("root must not be null"); + } + this.root = root.toAbsolutePath().normalize(); + String prefix = uriPrefix == null || uriPrefix.isBlank() ? "archive://raw" : uriPrefix.trim(); + this.uriPrefix = prefix.endsWith("/") ? prefix.substring(0, prefix.length() - 1) : prefix; + } + + @Override + public RawArchiveReceipt write(RawArchiveWriteRequest request) throws IOException { + String key = key(request); + Path target = resolveKey(key); + Files.createDirectories(target.getParent()); + rejectSymlinkPath(target); + byte[] bytes = request.rawBytes(); + Files.write(target, bytes, StandardOpenOption.CREATE_NEW, StandardOpenOption.WRITE); + return new RawArchiveReceipt(uriPrefix + "/" + key, checksum(bytes), bytes.length); + } + + @Override + public byte[] read(String rawUri) throws IOException { + String key = keyFromUri(rawUri); + Path target = resolveKey(key); + rejectSymlinkPath(target); + return Files.readAllBytes(target); + } + + private String key(RawArchiveWriteRequest request) { + var date = request.receivedAt().atZone(PARTITION_ZONE).toLocalDate(); + return "%04d/%02d/%02d/%s/%s/%s.bin".formatted( + date.getYear(), + date.getMonthValue(), + date.getDayOfMonth(), + request.protocol().name(), + safeSegment(request.vehicleKey()), + safeSegment(request.frameId())); + } + + private String keyFromUri(String rawUri) { + if (rawUri == null || rawUri.isBlank()) { + throw new IllegalArgumentException("rawUri must not be blank"); + } + String normalized = rawUri.trim(); + String expected = uriPrefix + "/"; + if (!normalized.startsWith(expected)) { + throw new IllegalArgumentException("rawUri does not use expected prefix: " + rawUri); + } + return normalized.substring(expected.length()); + } + + private Path resolveKey(String key) { + Path target = root.resolve(key).normalize(); + if (!target.startsWith(root)) { + throw new IllegalArgumentException("raw archive key escapes root: " + key); + } + return target; + } + + private void rejectSymlinkPath(Path target) throws IOException { + Path current = root; + Path relative = root.relativize(target); + for (Path segment : relative) { + current = current.resolve(segment); + if (Files.isSymbolicLink(current)) { + throw new IOException("raw archive path traverses symlink: " + current); + } + if (Files.notExists(current, LinkOption.NOFOLLOW_LINKS)) { + return; + } + } + } + + private static String safeSegment(String value) { + String cleaned = value == null ? "" : value.trim(); + if (cleaned.isBlank()) { + return "unknown"; + } + return cleaned.replaceAll("[^A-Za-z0-9._-]", "_"); + } + + private static String checksum(byte[] bytes) { + try { + MessageDigest digest = MessageDigest.getInstance("SHA-256"); + return "sha256:" + HexFormat.of().formatHex(digest.digest(bytes)); + } catch (NoSuchAlgorithmException e) { + throw new IllegalStateException("SHA-256 is unavailable", e); + } + } +} +``` + +- [ ] **Step 6: Run tests and verify GREEN** + +Run: + +```bash +mvn -pl modules/sinks/raw-archive-store -Dtest=LocalRawArchiveStoreTest test +``` + +Expected: PASS. + +- [ ] **Step 7: Commit** + +```bash +git add pom.xml modules/sinks/raw-archive-store +git commit -m "feat: add raw archive store contract" +``` + +--- + +### Task 6: Add Kafka Proto Messages for Facts + +**Files:** + +- Modify: `modules/sinks/sink-mq/pom.xml` +- Modify: `modules/sinks/sink-mq/src/main/proto/vehicle_envelope.proto` +- Create: `modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/VehicleEnvelopeProtoCompatibilityTest.java` + +- [ ] **Step 1: Write failing proto compatibility test** + +Create `VehicleEnvelopeProtoCompatibilityTest.java`: + +```java +package com.lingniu.ingest.sink.mq; + +import com.lingniu.ingest.sink.mq.proto.DecodedFactPayload; +import com.lingniu.ingest.sink.mq.proto.ParseStatusProto; +import com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload; +import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import org.junit.jupiter.api.Test; + +import static org.assertj.core.api.Assertions.assertThat; + +class VehicleEnvelopeProtoCompatibilityTest { + + @Test + void rawFrameFactPayloadRoundTripsThroughEnvelope() throws Exception { + VehicleEnvelope envelope = VehicleEnvelope.newBuilder() + .setSchemaVersion("2.0") + .setEventId("frame-1") + .setVin("VIN001") + .setSource("GB32960") + .setEventTimeMs(1782662400000L) + .setIngestTimeMs(1782662400001L) + .setRawFrameFact(RawFrameFactPayload.newBuilder() + .setFrameId("frame-1") + .setVehicleKey("VIN001") + .setMessageId(2) + .setRawUri("archive://raw/2026/06/29/GB32960/VIN001/frame-1.bin") + .setChecksum("sha256:abc") + .setRawSizeBytes(128) + .setParseStatus(ParseStatusProto.PARSE_STATUS_SUCCEEDED) + .build()) + .build(); + + VehicleEnvelope parsed = VehicleEnvelope.parseFrom(envelope.toByteArray()); + + assertThat(parsed.hasRawFrameFact()).isTrue(); + assertThat(parsed.getRawFrameFact().getFrameId()).isEqualTo("frame-1"); + assertThat(parsed.getRawFrameFact().getParseStatus()).isEqualTo(ParseStatusProto.PARSE_STATUS_SUCCEEDED); + } + + @Test + void decodedFactPayloadRoundTripsThroughEnvelope() throws Exception { + VehicleEnvelope envelope = VehicleEnvelope.newBuilder() + .setSchemaVersion("2.0") + .setEventId("fact-1") + .setVin("VIN001") + .setSource("JT808") + .setEventTimeMs(1782662400000L) + .setIngestTimeMs(1782662400001L) + .setDecodedFact(DecodedFactPayload.newBuilder() + .setFactId("fact-1") + .setFrameId("frame-1") + .setFactType("location") + .setVehicleKey("jt808:013912345678") + .setRawUri("archive://raw/2026/06/29/JT808/jt808_013912345678/frame-1.bin") + .putFields("longitude", "120.1") + .putFields("latitude", "30.1") + .build()) + .build(); + + VehicleEnvelope parsed = VehicleEnvelope.parseFrom(envelope.toByteArray()); + + assertThat(parsed.hasDecodedFact()).isTrue(); + assertThat(parsed.getDecodedFact().getFieldsMap()).containsEntry("longitude", "120.1"); + } +} +``` + +- [ ] **Step 2: Run test and verify RED** + +Run: + +```bash +mvn -pl modules/sinks/sink-mq -Dtest=VehicleEnvelopeProtoCompatibilityTest test +``` + +Expected: FAIL because proto generated classes do not exist. + +- [ ] **Step 3: Add `ingest-facts` dependency to sink-mq** + +Add this dependency to `modules/sinks/sink-mq/pom.xml` after `ingest-api`: + +```xml + + com.lingniu.ingest + ingest-facts + +``` + +- [ ] **Step 4: Extend proto schema** + +Modify `vehicle_envelope.proto`. + +Add to `VehicleEnvelope` after `TelemetrySnapshot telemetry_snapshot = 50;`: + +```proto + RawFrameFactPayload raw_frame_fact = 60; + DecodedFactPayload decoded_fact = 61; +``` + +Add these messages after `TelemetryField`: + +```proto +enum ParseStatusProto { + PARSE_STATUS_UNSPECIFIED = 0; + PARSE_STATUS_NOT_PARSED = 1; + PARSE_STATUS_SUCCEEDED = 2; + PARSE_STATUS_FAILED = 3; +} + +message RawFrameFactPayload { + string frame_id = 1; + string vehicle_key = 2; + string vin = 3; + string phone = 4; + int32 message_id = 5; + int32 sub_type = 6; + string raw_uri = 7; + string checksum = 8; + int64 raw_size_bytes = 9; + ParseStatusProto parse_status = 10; + string parse_error = 11; + string peer = 12; + map metadata = 13; +} + +message DecodedFactPayload { + string fact_id = 1; + string frame_id = 2; + string fact_type = 3; + string vehicle_key = 4; + string vin = 5; + string phone = 6; + string raw_uri = 7; + map fields = 8; + map metadata = 9; +} +``` + +- [ ] **Step 5: Run test and verify GREEN** + +Run: + +```bash +mvn -pl modules/sinks/sink-mq -Dtest=VehicleEnvelopeProtoCompatibilityTest test +``` + +Expected: PASS. + +- [ ] **Step 6: Run existing sink-mq tests** + +Run: + +```bash +mvn -pl modules/sinks/sink-mq test +``` + +Expected: PASS. Existing envelope tests must continue to pass because new proto fields are additive. + +- [ ] **Step 7: Commit** + +```bash +git add modules/sinks/sink-mq/pom.xml \ + modules/sinks/sink-mq/src/main/proto/vehicle_envelope.proto \ + modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/VehicleEnvelopeProtoCompatibilityTest.java +git commit -m "feat: add fact payloads to kafka envelope" +``` + +--- + +### Task 7: Verify Phase 1 Build + +**Files:** + +- No new files. + +- [ ] **Step 1: Run focused module tests** + +Run: + +```bash +mvn -pl modules/core/ingest-facts,modules/sinks/raw-archive-store,modules/sinks/sink-mq -am test +``` + +Expected: PASS. + +- [ ] **Step 2: Run dependency tree smoke check** + +Run: + +```bash +mvn -pl modules/sinks/sink-mq -am -DskipTests package +``` + +Expected: PASS. Generated protobuf Java sources include `RawFrameFactPayload`, `DecodedFactPayload`, and `ParseStatusProto`. + +- [ ] **Step 3: Search for accidental runtime wiring** + +Run: + +```bash +rg -n "RawFrameFact|DecodedFact|RawArchiveWriter|vehicle.raw-frame.v1|vehicle.decoded-fact.v1" \ + modules/apps modules/protocols modules/services modules/core modules/sinks +``` + +Expected: Matches only in the new modules, sink-mq proto/test, and planned references. No existing handler behavior should be changed in Phase 1. + +- [ ] **Step 4: Commit final verification note if any docs changed** + +If the implementation updated this plan with checked boxes, commit it: + +```bash +git add docs/superpowers/plans/2026-06-29-vehicle-ingest-redesign-phase1.md +git commit -m "docs: update phase 1 implementation checklist" +``` + +If no docs changed, no commit is needed. + +--- + +## Self-Review + +Spec coverage: + +- Raw original frame preservation is covered by `RawFrameFact`, raw archive contracts, and local archive tests. +- Historical query storage is prepared by stable fact model and raw URI references; TDengine writer is intentionally deferred to Phase 3. +- Extension is covered by `DecodedFact.Extension` and byte-light proto maps. +- Performance foundations are covered by immutable facts, Kafka key design, and no handler wiring in Phase 1. + +Placeholder scan: + +- This plan contains no unfinished markers or unspecified implementation steps. + +Type consistency: + +- `RawFrameFact`, `ParseStatus`, `DecodedFact`, `RawArchiveWriteRequest`, `RawArchiveReceipt`, and proto payload names are used consistently across tasks. + +Next plans: + +- Phase 2: wire GB32960/JT808 durable raw boundary and fact publication. +- Phase 3: implement TDengine schema manager and history writer. +- Phase 4: implement history query APIs and raw replay endpoints. diff --git a/docs/superpowers/specs/2026-06-29-vehicle-ingest-redesign.md b/docs/superpowers/specs/2026-06-29-vehicle-ingest-redesign.md index 7a8fe1d7..801dab60 100644 --- a/docs/superpowers/specs/2026-06-29-vehicle-ingest-redesign.md +++ b/docs/superpowers/specs/2026-06-29-vehicle-ingest-redesign.md @@ -165,6 +165,14 @@ Rules: - `rawUri` must be resolvable by history API for single-frame replay. - Malformed frames are written with `parseStatus=FAILED` when bytes are available. +Deployment constraints: + +- Raw archive, TDengine, and the history service may run on different ECS instances, but should use private-network access in the same VPC and availability zone. +- High-frequency historical queries only read TDengine and must not read raw archive. +- Raw archive is only used by single-frame detail, troubleshooting, replay, and asynchronous export paths. +- `rawUri` must be globally resolvable by services and must not be a private local absolute path on one ECS instance. +- The first local raw archive implementation must stay behind a replaceable interface so it can later move to OSS, MinIO, or another object store. + ### Kafka Topics Use stable versioned topics: diff --git a/docs/superpowers/specs/2026-06-29-vehicle-ingest-redesign.zh.md b/docs/superpowers/specs/2026-06-29-vehicle-ingest-redesign.zh.md index 24598520..b3b457a7 100644 --- a/docs/superpowers/specs/2026-06-29-vehicle-ingest-redesign.zh.md +++ b/docs/superpowers/specs/2026-06-29-vehicle-ingest-redesign.zh.md @@ -165,6 +165,14 @@ public record RawArchiveReceipt( - `rawUri` 必须能被历史 API 解析,用于单帧回放。 - 异常帧只要有 bytes,也必须写入,并记录 `parseStatus=FAILED`。 +部署约束: + +- Raw archive、TDengine、history 服务可以不在同一台 ECS,但应在同一 VPC 和同一可用区内通过内网访问。 +- 高频历史查询只访问 TDengine,不访问 raw archive。 +- Raw archive 只进入单帧详情、问题排查、回放和异步导出路径。 +- `rawUri` 必须是全局可解析地址,不能是某台 ECS 私有的本地绝对路径。 +- 第一版本地 raw archive 需要通过可替换接口封装,后续可以切换到 OSS、MinIO 或其它对象存储。 + ### Kafka Topic 使用稳定的版本化 topic: