diff --git a/ingest-core/src/main/java/com/lingniu/ingest/core/dispatcher/Dispatcher.java b/ingest-core/src/main/java/com/lingniu/ingest/core/dispatcher/Dispatcher.java index 720467bd..e4ac4dd9 100644 --- a/ingest-core/src/main/java/com/lingniu/ingest/core/dispatcher/Dispatcher.java +++ b/ingest-core/src/main/java/com/lingniu/ingest/core/dispatcher/Dispatcher.java @@ -9,7 +9,9 @@ import com.lingniu.ingest.core.pipeline.InterceptorChain; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.time.Instant; import java.util.List; +import java.util.Map; import java.util.UUID; /** @@ -46,6 +48,10 @@ public final class Dispatcher { public void dispatch(RawFrame frame) { IngestContext ctx = new IngestContext(UUID.randomUUID().toString()); try { + // 在 interceptor 之前发 RawArchive,保证原始字节被无条件落盘(dedup/rate-limit + // 不会过滤它),满足"原始可回放"目标。archive 的写盘靠下游 ArchiveEventSink 消费。 + emitRawArchive(frame, ctx); + if (!interceptors.before(frame, ctx)) { log.debug("frame aborted: {}", ctx.abortReason()); return; @@ -76,4 +82,33 @@ public final class Dispatcher { interceptors.onError(t, ctx); } } + + /** + * 从 {@link RawFrame} 构造一条 {@link VehicleEvent.RawArchive} 发到 EventBus。 + * 仅在 {@code rawBytes} 非空时发——有些入站适配器(未来)可能只传解析后的对象。 + * + *

VIN 取自 sourceMeta 里的 {@code vin} key(由各入站适配器负责填充);缺失时 + * 留空字符串,archive sink 会用 "unknown-vin" 占位保证 key 可解析。 + */ + private void emitRawArchive(RawFrame frame, IngestContext ctx) { + byte[] bytes = frame.rawBytes(); + if (bytes == null || bytes.length == 0) return; + + Map meta = frame.sourceMeta() == null ? Map.of() : frame.sourceMeta(); + String vin = meta.getOrDefault("vin", ""); + Instant ingestTime = frame.receivedAt() != null ? frame.receivedAt() : Instant.now(); + + VehicleEvent.RawArchive raw = new VehicleEvent.RawArchive( + UUID.randomUUID().toString(), + vin, + frame.protocolId(), + ingestTime, + ingestTime, + ctx.traceId(), + meta, + frame.command(), + frame.infoType(), + bytes); + eventBus.publish(raw); + } }