Files
lingniu-vehicle-ingest/docs/superpowers/plans/2026-06-29-vehicle-ingest-redesign-phase1.md
2026-06-29 12:50:41 +08:00

45 KiB

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:

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:

<module>modules/core/ingest-facts</module>

Add this dependency management entry after ingest-api:

<dependency>
    <groupId>com.lingniu.ingest</groupId>
    <artifactId>ingest-facts</artifactId>
    <version>${project.version}</version>
</dependency>
  • Step 3: Create module POM

Create modules/core/ingest-facts/pom.xml:

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0">
    <modelVersion>4.0.0</modelVersion>
    <parent>
        <groupId>com.lingniu.ingest</groupId>
        <artifactId>lingniu-vehicle-ingest</artifactId>
        <version>0.1.0-SNAPSHOT</version>
        <relativePath>../../../pom.xml</relativePath>
    </parent>
    <artifactId>ingest-facts</artifactId>
    <name>ingest-facts</name>
    <description>Protocol-neutral raw frame and decoded fact model for vehicle ingestion.</description>

    <dependencies>
        <dependency>
            <groupId>com.lingniu.ingest</groupId>
            <artifactId>ingest-api</artifactId>
        </dependency>
        <dependency>
            <groupId>org.junit.jupiter</groupId>
            <artifactId>junit-jupiter</artifactId>
            <scope>test</scope>
        </dependency>
        <dependency>
            <groupId>org.assertj</groupId>
            <artifactId>assertj-core</artifactId>
            <scope>test</scope>
        </dependency>
    </dependencies>
</project>
  • Step 4: Run Maven target check

Run:

mvn -pl modules/core/ingest-facts test

Expected: PASS with no tests.

  • Step 5: Commit
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:

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:

mvn -pl modules/core/ingest-facts -Dtest=VehicleKeyTest test

Expected: FAIL because VehicleKey does not exist.

  • Step 3: Implement VehicleKey

Create VehicleKey.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:

mvn -pl modules/core/ingest-facts -Dtest=VehicleKeyTest test

Expected: PASS.

  • Step 5: Commit
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:

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:

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:

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:

package com.lingniu.ingest.facts;

public enum ParseStatus {
    NOT_PARSED,
    SUCCEEDED,
    FAILED
}
  • Step 4: Implement FactIds

Create FactIds.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:

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<String, String> 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:

mvn -pl modules/core/ingest-facts -Dtest=FactIdsTest,RawFrameFactTest test

Expected: PASS.

  • Step 7: Commit
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:

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:

mvn -pl modules/core/ingest-facts -Dtest=DecodedFactTest test

Expected: FAIL because DecodedFact does not exist.

  • Step 3: Implement DecodedFact

Create DecodedFact.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<String, String> metadata();

    record Location(
            String factId,
            String frameId,
            ProtocolId protocol,
            String vehicleKey,
            String vin,
            String phone,
            Instant eventTime,
            Instant receivedAt,
            String rawUri,
            Map<String, String> 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<String, String> metadata,
            Map<String, String> 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<String, String> metadata,
            String alarmLevel,
            Long alarmCode,
            String alarmName,
            Map<String, String> 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<String, String> metadata,
            String sessionType,
            Integer resultCode,
            Map<String, String> 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<String, String> metadata,
            Map<String, String> 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<String, String> metadata,
            String extensionType,
            Map<String, String> 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<String, String> metadata(Map<String, String> 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:

mvn -pl modules/core/ingest-facts -Dtest=DecodedFactTest test

Expected: PASS.

  • Step 5: Commit
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:

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:

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:

<module>modules/sinks/raw-archive-store</module>

Add dependency management after sink-archive:

<dependency>
    <groupId>com.lingniu.ingest</groupId>
    <artifactId>raw-archive-store</artifactId>
    <version>${project.version}</version>
</dependency>

Create modules/sinks/raw-archive-store/pom.xml:

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0">
    <modelVersion>4.0.0</modelVersion>
    <parent>
        <groupId>com.lingniu.ingest</groupId>
        <artifactId>lingniu-vehicle-ingest</artifactId>
        <version>0.1.0-SNAPSHOT</version>
        <relativePath>../../../pom.xml</relativePath>
    </parent>
    <artifactId>raw-archive-store</artifactId>
    <name>raw-archive-store</name>
    <description>Raw frame archive writer and reader contracts.</description>

    <dependencies>
        <dependency>
            <groupId>com.lingniu.ingest</groupId>
            <artifactId>ingest-api</artifactId>
        </dependency>
        <dependency>
            <groupId>org.junit.jupiter</groupId>
            <artifactId>junit-jupiter</artifactId>
            <scope>test</scope>
        </dependency>
        <dependency>
            <groupId>org.assertj</groupId>
            <artifactId>assertj-core</artifactId>
            <scope>test</scope>
        </dependency>
    </dependencies>
</project>
  • Step 4: Implement contracts

Create RawArchiveWriteRequest.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:

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:

package com.lingniu.ingest.rawarchive;

import java.io.IOException;

public interface RawArchiveWriter {
    RawArchiveReceipt write(RawArchiveWriteRequest request) throws IOException;
}

Create RawArchiveReader.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:

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:

mvn -pl modules/sinks/raw-archive-store -Dtest=LocalRawArchiveStoreTest test

Expected: PASS.

  • Step 7: Commit
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:

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:

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:

<dependency>
    <groupId>com.lingniu.ingest</groupId>
    <artifactId>ingest-facts</artifactId>
</dependency>
  • Step 4: Extend proto schema

Modify vehicle_envelope.proto.

Add to VehicleEnvelope after TelemetrySnapshot telemetry_snapshot = 50;:

  RawFrameFactPayload raw_frame_fact = 60;
  DecodedFactPayload decoded_fact = 61;

Add these messages after TelemetryField:

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<string, string> 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<string, string> fields = 8;
  map<string, string> metadata = 9;
}
  • Step 5: Run test and verify GREEN

Run:

mvn -pl modules/sinks/sink-mq -Dtest=VehicleEnvelopeProtoCompatibilityTest test

Expected: PASS.

  • Step 6: Run existing sink-mq tests

Run:

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
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:

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:

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:

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:

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.