From 0ec29a732cc54174b4e7961bb53ca191e1fbf6d3 Mon Sep 17 00:00:00 2001 From: lingniu Date: Mon, 29 Jun 2026 12:56:50 +0800 Subject: [PATCH] feat: add raw archive store contract --- modules/sinks/raw-archive-store/pom.xml | 30 +++++ .../rawarchive/LocalRawArchiveStore.java | 110 ++++++++++++++++++ .../ingest/rawarchive/RawArchiveReader.java | 7 ++ .../ingest/rawarchive/RawArchiveReceipt.java | 19 +++ .../rawarchive/RawArchiveWriteRequest.java | 41 +++++++ .../ingest/rawarchive/RawArchiveWriter.java | 7 ++ .../rawarchive/LocalRawArchiveStoreTest.java | 50 ++++++++ pom.xml | 6 + 8 files changed, 270 insertions(+) create mode 100644 modules/sinks/raw-archive-store/pom.xml create mode 100644 modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/LocalRawArchiveStore.java create mode 100644 modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveReader.java create mode 100644 modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveReceipt.java create mode 100644 modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveWriteRequest.java create mode 100644 modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveWriter.java create mode 100644 modules/sinks/raw-archive-store/src/test/java/com/lingniu/ingest/rawarchive/LocalRawArchiveStoreTest.java diff --git a/modules/sinks/raw-archive-store/pom.xml b/modules/sinks/raw-archive-store/pom.xml new file mode 100644 index 00000000..e0179e52 --- /dev/null +++ b/modules/sinks/raw-archive-store/pom.xml @@ -0,0 +1,30 @@ + + + 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 + + + diff --git a/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/LocalRawArchiveStore.java b/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/LocalRawArchiveStore.java new file mode 100644 index 00000000..8ba89823 --- /dev/null +++ b/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/LocalRawArchiveStore.java @@ -0,0 +1,110 @@ +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); + } + } +} diff --git a/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveReader.java b/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveReader.java new file mode 100644 index 00000000..f5aad3ad --- /dev/null +++ b/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveReader.java @@ -0,0 +1,7 @@ +package com.lingniu.ingest.rawarchive; + +import java.io.IOException; + +public interface RawArchiveReader { + byte[] read(String rawUri) throws IOException; +} diff --git a/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveReceipt.java b/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveReceipt.java new file mode 100644 index 00000000..e65326c0 --- /dev/null +++ b/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveReceipt.java @@ -0,0 +1,19 @@ +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; + } +} diff --git a/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveWriteRequest.java b/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveWriteRequest.java new file mode 100644 index 00000000..bf42aba3 --- /dev/null +++ b/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveWriteRequest.java @@ -0,0 +1,41 @@ +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; + } +} diff --git a/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveWriter.java b/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveWriter.java new file mode 100644 index 00000000..655026fd --- /dev/null +++ b/modules/sinks/raw-archive-store/src/main/java/com/lingniu/ingest/rawarchive/RawArchiveWriter.java @@ -0,0 +1,7 @@ +package com.lingniu.ingest.rawarchive; + +import java.io.IOException; + +public interface RawArchiveWriter { + RawArchiveReceipt write(RawArchiveWriteRequest request) throws IOException; +} diff --git a/modules/sinks/raw-archive-store/src/test/java/com/lingniu/ingest/rawarchive/LocalRawArchiveStoreTest.java b/modules/sinks/raw-archive-store/src/test/java/com/lingniu/ingest/rawarchive/LocalRawArchiveStoreTest.java new file mode 100644 index 00000000..626a7aae --- /dev/null +++ b/modules/sinks/raw-archive-store/src/test/java/com/lingniu/ingest/rawarchive/LocalRawArchiveStoreTest.java @@ -0,0 +1,50 @@ +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"); + } +} diff --git a/pom.xml b/pom.xml index e5546c4a..6dccb62b 100644 --- a/pom.xml +++ b/pom.xml @@ -26,6 +26,7 @@ modules/core/observability modules/sinks/sink-mq modules/sinks/sink-archive + modules/sinks/raw-archive-store modules/sinks/event-file-store modules/services/event-history-service modules/services/vehicle-state-service @@ -184,6 +185,11 @@ sink-archive ${project.version} + + com.lingniu.ingest + raw-archive-store + ${project.version} + com.lingniu.ingest event-file-store