diff --git a/modules/sinks/sink-archive/src/main/java/com/lingniu/ingest/sink/archive/LocalArchiveStore.java b/modules/sinks/sink-archive/src/main/java/com/lingniu/ingest/sink/archive/LocalArchiveStore.java index e220920b..89e6f04e 100644 --- a/modules/sinks/sink-archive/src/main/java/com/lingniu/ingest/sink/archive/LocalArchiveStore.java +++ b/modules/sinks/sink-archive/src/main/java/com/lingniu/ingest/sink/archive/LocalArchiveStore.java @@ -5,12 +5,16 @@ import java.io.InputStream; import java.io.OutputStream; import java.net.URI; import java.nio.file.Files; +import java.nio.file.LinkOption; import java.nio.file.Path; import java.nio.file.StandardOpenOption; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; public final class LocalArchiveStore implements ArchiveStore { private final Path root; + private final ConcurrentMap appendLocks = new ConcurrentHashMap<>(); public LocalArchiveStore(String root) { this(rootPath(root)); @@ -26,12 +30,14 @@ public final class LocalArchiveStore implements ArchiveStore { throw new IllegalArgumentException("data must not be null"); } Path target = resolve(key); - Files.createDirectories(target.getParent()); + createDirectoriesInsideRoot(target.getParent()); + rejectSymlinkPath(target); try (OutputStream out = Files.newOutputStream( target, StandardOpenOption.CREATE, StandardOpenOption.TRUNCATE_EXISTING, - StandardOpenOption.WRITE)) { + StandardOpenOption.WRITE, + LinkOption.NOFOLLOW_LINKS)) { data.transferTo(out); } return target.toUri().toString(); @@ -40,29 +46,38 @@ public final class LocalArchiveStore implements ArchiveStore { @Override public String append(String key, byte[] chunk) throws IOException { Path target = resolve(key); - Files.createDirectories(target.getParent()); - Files.write( - target, - chunk == null ? new byte[0] : chunk, - StandardOpenOption.CREATE, - StandardOpenOption.APPEND, - StandardOpenOption.WRITE); + synchronized (appendLocks.computeIfAbsent(target, ignored -> new Object())) { + createDirectoriesInsideRoot(target.getParent()); + rejectSymlinkPath(target); + Files.write( + target, + chunk == null ? new byte[0] : chunk, + StandardOpenOption.CREATE, + StandardOpenOption.APPEND, + StandardOpenOption.WRITE, + LinkOption.NOFOLLOW_LINKS); + } return target.toUri().toString(); } @Override public InputStream get(String key) throws IOException { - return Files.newInputStream(resolve(key), StandardOpenOption.READ); + Path target = resolveExisting(key); + return Files.newInputStream(target, StandardOpenOption.READ, LinkOption.NOFOLLOW_LINKS); } @Override public boolean exists(String key) { - return Files.exists(resolve(key)); + try { + return Files.exists(resolveExisting(key), LinkOption.NOFOLLOW_LINKS); + } catch (IOException | IllegalArgumentException e) { + return false; + } } @Override public long size(String key) throws IOException { - return Files.size(resolve(key)); + return Files.size(resolveExisting(key)); } private Path resolve(String key) { @@ -76,6 +91,53 @@ public final class LocalArchiveStore implements ArchiveStore { return target; } + private Path resolveExisting(String key) throws IOException { + Path target = resolve(key); + rejectSymlinkPath(target); + return target; + } + + private void createDirectoriesInsideRoot(Path directory) throws IOException { + rejectEscapedPath(directory); + if (Files.notExists(root, LinkOption.NOFOLLOW_LINKS)) { + Files.createDirectories(root); + } + Path current = root; + Path relative = root.relativize(directory); + for (Path segment : relative) { + current = current.resolve(segment); + if (Files.isSymbolicLink(current)) { + throw new IOException("archive path traverses symlink: " + current); + } + if (Files.notExists(current, LinkOption.NOFOLLOW_LINKS)) { + Files.createDirectory(current); + } else if (!Files.isDirectory(current, LinkOption.NOFOLLOW_LINKS)) { + throw new IOException("archive path segment is not a directory: " + current); + } + } + } + + private void rejectSymlinkPath(Path target) throws IOException { + rejectEscapedPath(target); + Path current = root; + Path relative = root.relativize(target); + for (Path segment : relative) { + current = current.resolve(segment); + if (Files.isSymbolicLink(current)) { + throw new IOException("archive path traverses symlink: " + current); + } + if (Files.notExists(current, LinkOption.NOFOLLOW_LINKS)) { + return; + } + } + } + + private void rejectEscapedPath(Path path) { + if (!path.normalize().startsWith(root)) { + throw new IllegalArgumentException("archive path escapes root: " + path); + } + } + private static Path rootPath(String value) { if (value == null || value.isBlank()) { return Path.of(System.getProperty("java.io.tmpdir"), "lingniu-archive"); diff --git a/modules/sinks/sink-archive/src/main/java/com/lingniu/ingest/sink/archive/RawArchiveEventSink.java b/modules/sinks/sink-archive/src/main/java/com/lingniu/ingest/sink/archive/RawArchiveEventSink.java index 06af1305..e148fd3b 100644 --- a/modules/sinks/sink-archive/src/main/java/com/lingniu/ingest/sink/archive/RawArchiveEventSink.java +++ b/modules/sinks/sink-archive/src/main/java/com/lingniu/ingest/sink/archive/RawArchiveEventSink.java @@ -25,8 +25,11 @@ public final class RawArchiveEventSink implements EventSink { if (!(event instanceof VehicleEvent.RawArchive raw)) { return CompletableFuture.completedFuture(null); } - byte[] bytes = raw.rawBytes() == null ? new byte[0] : raw.rawBytes(); try { + byte[] bytes = raw.rawBytes(); + if (bytes == null) { + throw new IllegalArgumentException("raw archive bytes must not be null"); + } store.put(RawArchiveKeys.key(raw), new ByteArrayInputStream(bytes), bytes.length); return CompletableFuture.completedFuture(null); } catch (Exception e) { diff --git a/modules/sinks/sink-archive/src/main/java/com/lingniu/ingest/sink/archive/config/SinkArchiveAutoConfiguration.java b/modules/sinks/sink-archive/src/main/java/com/lingniu/ingest/sink/archive/config/SinkArchiveAutoConfiguration.java index 73efc1d9..f64c7c86 100644 --- a/modules/sinks/sink-archive/src/main/java/com/lingniu/ingest/sink/archive/config/SinkArchiveAutoConfiguration.java +++ b/modules/sinks/sink-archive/src/main/java/com/lingniu/ingest/sink/archive/config/SinkArchiveAutoConfiguration.java @@ -5,6 +5,7 @@ import com.lingniu.ingest.sink.archive.LocalArchiveStore; import com.lingniu.ingest.sink.archive.RawArchiveEventSink; import org.springframework.boot.autoconfigure.AutoConfiguration; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.context.annotation.Bean; @@ -27,6 +28,7 @@ public class SinkArchiveAutoConfiguration { @Bean @ConditionalOnMissingBean + @ConditionalOnBean(ArchiveStore.class) public RawArchiveEventSink rawArchiveEventSink(ArchiveStore store) { return new RawArchiveEventSink(store); } diff --git a/modules/sinks/sink-archive/src/test/java/com/lingniu/ingest/sink/archive/LocalArchiveStoreTest.java b/modules/sinks/sink-archive/src/test/java/com/lingniu/ingest/sink/archive/LocalArchiveStoreTest.java new file mode 100644 index 00000000..f23e90c2 --- /dev/null +++ b/modules/sinks/sink-archive/src/test/java/com/lingniu/ingest/sink/archive/LocalArchiveStoreTest.java @@ -0,0 +1,74 @@ +package com.lingniu.ingest.sink.archive; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.io.ByteArrayInputStream; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.Executors; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +class LocalArchiveStoreTest { + + @TempDir + Path tempDir; + + @Test + void rejectsSymlinkTraversalOutsideRootForWriteReadAndSize() throws Exception { + Path root = tempDir.resolve("archive"); + Path outside = tempDir.resolve("outside"); + Files.createDirectories(root); + Files.createDirectories(outside); + Files.writeString(outside.resolve("existing.bin"), "secret", StandardCharsets.UTF_8); + Files.createSymbolicLink(root.resolve("link"), outside); + + LocalArchiveStore store = new LocalArchiveStore(root); + + assertThatThrownBy(() -> store.put( + "link/new.bin", + new ByteArrayInputStream("payload".getBytes(StandardCharsets.UTF_8)), + 7)) + .isInstanceOfAny(IllegalArgumentException.class, java.io.IOException.class); + assertThat(outside.resolve("new.bin")).doesNotExist(); + + assertThatThrownBy(() -> store.get("link/existing.bin")) + .isInstanceOfAny(IllegalArgumentException.class, java.io.IOException.class); + assertThatThrownBy(() -> store.size("link/existing.bin")) + .isInstanceOfAny(IllegalArgumentException.class, java.io.IOException.class); + } + + @Test + void appendsConcurrentChunksWithoutLosingData() throws Exception { + LocalArchiveStore store = new LocalArchiveStore(tempDir.resolve("archive")); + int chunks = 96; + + try (var executor = Executors.newFixedThreadPool(12)) { + List> futures = new ArrayList<>(); + for (int i = 0; i < chunks; i++) { + String chunk = "%03d\n".formatted(i); + futures.add(CompletableFuture.runAsync(() -> { + try { + store.append("same/key.bin", chunk.getBytes(StandardCharsets.UTF_8)); + } catch (Exception e) { + throw new RuntimeException(e); + } + }, executor)); + } + CompletableFuture.allOf(futures.toArray(CompletableFuture[]::new)).join(); + } + + String content = Files.readString(tempDir.resolve("archive/same/key.bin"), StandardCharsets.UTF_8); + assertThat(content).hasSize(chunks * 4); + assertThat(content.lines()).containsExactlyInAnyOrderElementsOf( + java.util.stream.IntStream.range(0, chunks) + .mapToObj("%03d"::formatted) + .toList()); + } +} diff --git a/modules/sinks/sink-archive/src/test/java/com/lingniu/ingest/sink/archive/RawArchiveEventSinkTest.java b/modules/sinks/sink-archive/src/test/java/com/lingniu/ingest/sink/archive/RawArchiveEventSinkTest.java new file mode 100644 index 00000000..7ffe5385 --- /dev/null +++ b/modules/sinks/sink-archive/src/test/java/com/lingniu/ingest/sink/archive/RawArchiveEventSinkTest.java @@ -0,0 +1,63 @@ +package com.lingniu.ingest.sink.archive; + +import com.lingniu.ingest.api.ProtocolId; +import com.lingniu.ingest.api.event.VehicleEvent; +import org.junit.jupiter.api.Test; + +import java.io.InputStream; +import java.time.Instant; +import java.util.Map; +import java.util.concurrent.CompletionException; + +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +class RawArchiveEventSinkTest { + + @Test + void failsFastWhenRawBytesAreNull() { + RawArchiveEventSink sink = new RawArchiveEventSink(new CapturingArchiveStore()); + VehicleEvent.RawArchive raw = new VehicleEvent.RawArchive( + "event-1", + "VIN123", + ProtocolId.GB32960, + Instant.parse("2026-06-23T00:00:00Z"), + Instant.parse("2026-06-23T00:00:01Z"), + "trace-1", + Map.of(), + 2, + 1, + null); + + assertThatThrownBy(() -> sink.publish(raw).join()) + .isInstanceOf(CompletionException.class) + .hasCauseInstanceOf(IllegalArgumentException.class); + } + + private static final class CapturingArchiveStore implements ArchiveStore { + + @Override + public String put(String key, InputStream data, long length) { + return "archive://" + key; + } + + @Override + public String append(String key, byte[] chunk) { + return "archive://" + key; + } + + @Override + public InputStream get(String key) { + throw new UnsupportedOperationException(); + } + + @Override + public boolean exists(String key) { + return false; + } + + @Override + public long size(String key) { + return 0; + } + } +} diff --git a/modules/sinks/sink-archive/src/test/java/com/lingniu/ingest/sink/archive/config/SinkArchiveAutoConfigurationTest.java b/modules/sinks/sink-archive/src/test/java/com/lingniu/ingest/sink/archive/config/SinkArchiveAutoConfigurationTest.java new file mode 100644 index 00000000..2156cb9e --- /dev/null +++ b/modules/sinks/sink-archive/src/test/java/com/lingniu/ingest/sink/archive/config/SinkArchiveAutoConfigurationTest.java @@ -0,0 +1,40 @@ +package com.lingniu.ingest.sink.archive.config; + +import com.lingniu.ingest.sink.archive.ArchiveStore; +import com.lingniu.ingest.sink.archive.RawArchiveEventSink; +import org.junit.jupiter.api.Test; +import org.springframework.context.annotation.AnnotationConfigApplicationContext; +import org.springframework.core.env.MapPropertySource; + +import static org.assertj.core.api.Assertions.assertThat; + +class SinkArchiveAutoConfigurationTest { + + @Test + void localArchiveTypeCreatesStoreAndSink() { + try (AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext()) { + context.getEnvironment().getPropertySources().addFirst(new MapPropertySource( + "test", + java.util.Map.of("lingniu.ingest.sink.archive.path", "target/test-archive"))); + context.register(SinkArchiveAutoConfiguration.class); + context.refresh(); + + assertThat(context.getBeansOfType(ArchiveStore.class)).hasSize(1); + assertThat(context.getBeansOfType(RawArchiveEventSink.class)).hasSize(1); + } + } + + @Test + void nonLocalArchiveTypeDoesNotCreateSinkWithoutStore() { + try (AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext()) { + context.getEnvironment().getPropertySources().addFirst(new MapPropertySource( + "test", + java.util.Map.of("lingniu.ingest.sink.archive.type", "oss"))); + context.register(SinkArchiveAutoConfiguration.class); + context.refresh(); + + assertThat(context.getBeansOfType(ArchiveStore.class)).isEmpty(); + assertThat(context.getBeansOfType(RawArchiveEventSink.class)).isEmpty(); + } + } +}