1use crate::{
2 domain::entities::Event,
3 error::{AllSourceError, Result},
4};
5use arrow::{
6 array::{
7 Array, ArrayRef, StringBuilder, TimestampMicrosecondArray, TimestampMicrosecondBuilder,
8 UInt64Builder,
9 },
10 datatypes::{DataType, Field, Schema, TimeUnit},
11 record_batch::RecordBatch,
12};
13use parquet::{arrow::ArrowWriter, file::properties::WriterProperties};
14use std::{
15 collections::HashMap,
16 fs::{self, File},
17 path::{Path, PathBuf},
18 sync::{
19 Arc, Mutex,
20 atomic::{AtomicU64, Ordering},
21 },
22 time::{Duration, Instant},
23};
24
25pub const DEFAULT_BATCH_SIZE: usize = 10_000;
27
28pub const DEFAULT_FLUSH_TIMEOUT_MS: u64 = 5_000;
30
31#[derive(Debug, Clone)]
33pub struct ParquetStorageConfig {
34 pub batch_size: usize,
36 pub flush_timeout: Duration,
38 pub compression: parquet::basic::Compression,
40}
41
42impl Default for ParquetStorageConfig {
43 fn default() -> Self {
44 Self {
45 batch_size: DEFAULT_BATCH_SIZE,
46 flush_timeout: Duration::from_millis(DEFAULT_FLUSH_TIMEOUT_MS),
47 compression: parquet::basic::Compression::SNAPPY,
48 }
49 }
50}
51
52impl ParquetStorageConfig {
53 pub fn high_throughput() -> Self {
55 Self {
56 batch_size: 50_000,
57 flush_timeout: Duration::from_secs(10),
58 compression: parquet::basic::Compression::SNAPPY,
59 }
60 }
61
62 pub fn low_latency() -> Self {
64 Self {
65 batch_size: 1_000,
66 flush_timeout: Duration::from_secs(1),
67 compression: parquet::basic::Compression::SNAPPY,
68 }
69 }
70}
71
72#[derive(Debug, Clone, Default)]
74pub struct BatchWriteStats {
75 pub batches_written: u64,
77 pub events_written: u64,
79 pub bytes_written: u64,
81 pub avg_batch_size: f64,
83 pub events_per_sec: f64,
85 pub total_write_time_ns: u64,
87 pub timeout_flushes: u64,
89 pub size_flushes: u64,
91}
92
93#[derive(Debug, Clone)]
95pub struct BatchWriteResult {
96 pub events_written: usize,
98 pub batches_flushed: usize,
100 pub duration: Duration,
102 pub events_per_sec: f64,
104}
105
106pub struct ParquetStorage {
115 storage_dir: PathBuf,
117
118 current_batches: Mutex<HashMap<String, Vec<Event>>>,
127
128 config: ParquetStorageConfig,
130
131 schema: Arc<Schema>,
133
134 last_flush_time: Mutex<Instant>,
136
137 batches_written: AtomicU64,
139 events_written: AtomicU64,
140 bytes_written: AtomicU64,
141 total_write_time_ns: AtomicU64,
142 timeout_flushes: AtomicU64,
143 size_flushes: AtomicU64,
144}
145
146impl ParquetStorage {
147 pub fn new(storage_dir: impl AsRef<Path>) -> Result<Self> {
149 Self::with_config(storage_dir, ParquetStorageConfig::default())
150 }
151
152 pub fn with_config(
154 storage_dir: impl AsRef<Path>,
155 config: ParquetStorageConfig,
156 ) -> Result<Self> {
157 let storage_dir = storage_dir.as_ref().to_path_buf();
158
159 fs::create_dir_all(&storage_dir).map_err(|e| {
161 AllSourceError::StorageError(format!("Failed to create storage directory: {e}"))
162 })?;
163
164 let schema = Arc::new(Schema::new(vec![
166 Field::new("event_id", DataType::Utf8, false),
167 Field::new("event_type", DataType::Utf8, false),
168 Field::new("entity_id", DataType::Utf8, false),
169 Field::new("payload", DataType::Utf8, false),
170 Field::new(
171 "timestamp",
172 DataType::Timestamp(TimeUnit::Microsecond, None),
173 false,
174 ),
175 Field::new("metadata", DataType::Utf8, true),
176 Field::new("version", DataType::UInt64, false),
177 ]));
178
179 let storage = Self {
180 storage_dir,
181 current_batches: Mutex::new(HashMap::new()),
182 config,
183 schema,
184 last_flush_time: Mutex::new(Instant::now()),
185 batches_written: AtomicU64::new(0),
186 events_written: AtomicU64::new(0),
187 bytes_written: AtomicU64::new(0),
188 total_write_time_ns: AtomicU64::new(0),
189 timeout_flushes: AtomicU64::new(0),
190 size_flushes: AtomicU64::new(0),
191 };
192
193 match storage.cleanup_partial_writes() {
198 Ok(0) => {}
199 Ok(n) => tracing::warn!(
200 "cleanup_partial_writes acted on {n} crash-detritus file(s) on boot — \
201 see preceding logs for per-file detail"
202 ),
203 Err(e) => tracing::error!("cleanup_partial_writes failed on boot: {e}"),
204 }
205
206 Ok(storage)
207 }
208
209 #[deprecated(note = "Use new() or with_config() instead - default batch size is now 10,000")]
211 pub fn with_legacy_batch_size(storage_dir: impl AsRef<Path>) -> Result<Self> {
212 Self::with_config(
213 storage_dir,
214 ParquetStorageConfig {
215 batch_size: 1000,
216 ..Default::default()
217 },
218 )
219 }
220
221 #[cfg_attr(feature = "hotpath", hotpath::measure)]
231 pub fn append_event(&self, event: Event) -> Result<()> {
232 let tenant = event.tenant_id_str().to_string();
233 let should_flush_tenant = {
234 let mut batches = self.current_batches.lock().unwrap();
235 let entry = batches.entry(tenant.clone()).or_default();
236 entry.push(event);
237 entry.len() >= self.config.batch_size
238 };
239
240 if should_flush_tenant {
241 self.size_flushes.fetch_add(1, Ordering::Relaxed);
242 self.flush_tenant(&tenant)?;
243 }
244
245 Ok(())
246 }
247
248 #[cfg_attr(feature = "hotpath", hotpath::measure)]
254 pub fn batch_write(&self, events: Vec<Event>) -> Result<BatchWriteResult> {
255 let start = Instant::now();
256 let event_count = events.len();
257
258 let mut grouped: HashMap<String, Vec<Event>> = HashMap::new();
261 for event in events {
262 grouped
263 .entry(event.tenant_id_str().to_string())
264 .or_default()
265 .push(event);
266 }
267
268 let mut tenants_to_flush: Vec<String> = Vec::new();
269 {
270 let mut batches = self.current_batches.lock().unwrap();
271 for (tenant, mut new_events) in grouped {
272 let entry = batches.entry(tenant.clone()).or_default();
273 entry.append(&mut new_events);
274 if entry.len() >= self.config.batch_size {
275 tenants_to_flush.push(tenant);
276 }
277 }
278 }
279
280 let mut batches_flushed = 0;
281 for tenant in tenants_to_flush {
282 self.size_flushes.fetch_add(1, Ordering::Relaxed);
283 self.flush_tenant(&tenant)?;
284 batches_flushed += 1;
285 }
286
287 let duration = start.elapsed();
288
289 Ok(BatchWriteResult {
290 events_written: event_count,
291 batches_flushed,
292 duration,
293 events_per_sec: event_count as f64 / duration.as_secs_f64(),
294 })
295 }
296
297 #[cfg_attr(feature = "hotpath", hotpath::measure)]
305 pub fn check_timeout_flush(&self) -> Result<bool> {
306 let should_flush = {
307 let last_flush = self.last_flush_time.lock().unwrap();
308 let batches = self.current_batches.lock().unwrap();
309 let any_pending = batches.values().any(|v| !v.is_empty());
310 any_pending && last_flush.elapsed() >= self.config.flush_timeout
311 };
312
313 if should_flush {
314 self.timeout_flushes.fetch_add(1, Ordering::Relaxed);
315 self.flush()?;
316 Ok(true)
317 } else {
318 Ok(false)
319 }
320 }
321
322 #[cfg_attr(feature = "hotpath", hotpath::measure)]
329 pub fn flush(&self) -> Result<()> {
330 let tenants: Vec<String> = {
331 let batches = self.current_batches.lock().unwrap();
332 batches
333 .iter()
334 .filter(|(_, v)| !v.is_empty())
335 .map(|(k, _)| k.clone())
336 .collect()
337 };
338 if tenants.is_empty() {
339 return Ok(());
340 }
341 for tenant in tenants {
342 self.flush_tenant(&tenant)?;
343 }
344 Ok(())
345 }
346
347 fn flush_tenant(&self, tenant_id: &str) -> Result<()> {
362 let events_to_write = {
363 let mut batches = self.current_batches.lock().unwrap();
364 match batches.get_mut(tenant_id) {
365 Some(v) if !v.is_empty() => std::mem::take(v),
366 _ => return Ok(()),
367 }
368 };
369
370 let batch_count = events_to_write.len();
371 let start = Instant::now();
372
373 let written = self.write_tenant_events(tenant_id, &events_to_write);
374 let (file_path, file_metadata) = match written {
375 Ok(written) => written,
376 Err(e) => {
377 let mut batches = self.current_batches.lock().unwrap();
380 let pending = batches.entry(tenant_id.to_string()).or_default();
381 let arrived_during_write = std::mem::replace(pending, events_to_write);
382 pending.extend(arrived_during_write);
383 return Err(e);
384 }
385 };
386
387 let duration = start.elapsed();
388
389 self.batches_written.fetch_add(1, Ordering::Relaxed);
390 self.events_written
391 .fetch_add(batch_count as u64, Ordering::Relaxed);
392 if let Some(size) = file_metadata
393 .row_groups()
394 .first()
395 .map(parquet::file::metadata::RowGroupMetaData::total_byte_size)
396 {
397 self.bytes_written.fetch_add(size as u64, Ordering::Relaxed);
398 }
399 self.total_write_time_ns
400 .fetch_add(duration.as_nanos() as u64, Ordering::Relaxed);
401
402 {
403 let mut last_flush = self.last_flush_time.lock().unwrap();
404 *last_flush = Instant::now();
405 }
406
407 tracing::info!(
408 "Wrote {} events for tenant={} to {} in {:?}",
409 batch_count,
410 tenant_id,
411 file_path.display(),
412 duration
413 );
414
415 Ok(())
416 }
417
418 fn write_tenant_events(
419 &self,
420 tenant_id: &str,
421 events: &[Event],
422 ) -> Result<(PathBuf, parquet::file::metadata::ParquetMetaData)> {
423 let record_batch = self.events_to_record_batch(events)?;
424
425 let now = chrono::Utc::now();
426 let partition_dir = partition_path_for_tenant(&self.storage_dir, tenant_id, now)?;
427 fs::create_dir_all(&partition_dir).map_err(|e| {
428 AllSourceError::StorageError(format!(
429 "Failed to create tenant partition {}: {e}",
430 partition_dir.display()
431 ))
432 })?;
433 let file_stem = format!(
434 "events-{}-{}",
435 now.format("%Y%m%d-%H%M%S%3f"),
436 uuid::Uuid::new_v4().as_simple()
437 );
438
439 tracing::info!(
440 "Flushing {} events for tenant={} to {}/{}.parquet",
441 events.len(),
442 tenant_id,
443 partition_dir.display(),
444 file_stem
445 );
446
447 self.write_record_batch_atomic(&partition_dir, &file_stem, &record_batch)
448 }
449
450 fn write_record_batch_atomic(
467 &self,
468 partition_dir: &Path,
469 file_stem: &str,
470 record_batch: &RecordBatch,
471 ) -> Result<(PathBuf, parquet::file::metadata::ParquetMetaData)> {
472 let final_path = partition_dir.join(format!("{file_stem}.parquet"));
473 let tmp_path = partition_dir.join(format!("{file_stem}.parquet.tmp"));
474
475 let metadata = {
477 let file = File::create(&tmp_path).map_err(|e| {
478 AllSourceError::StorageError(format!(
479 "Failed to create parquet tmp file {}: {e}",
480 tmp_path.display()
481 ))
482 })?;
483
484 let props = WriterProperties::builder()
485 .set_compression(self.config.compression)
486 .build();
487
488 let mut writer = ArrowWriter::try_new(file, self.schema.clone(), Some(props))?;
489 writer.write(record_batch)?;
490 writer.close()?
493 };
494
495 let tmp_file = File::open(&tmp_path).map_err(|e| {
497 AllSourceError::StorageError(format!(
498 "Failed to reopen parquet tmp for fsync {}: {e}",
499 tmp_path.display()
500 ))
501 })?;
502 tmp_file.sync_all().map_err(|e| {
503 AllSourceError::StorageError(format!("fsync on parquet tmp failed: {e}"))
504 })?;
505 drop(tmp_file);
506
507 fs::rename(&tmp_path, &final_path).map_err(|e| {
509 AllSourceError::StorageError(format!(
510 "Failed to rename {} → {}: {e}",
511 tmp_path.display(),
512 final_path.display()
513 ))
514 })?;
515
516 if let Ok(dir) = File::open(partition_dir) {
520 let _ = dir.sync_all();
521 }
522
523 Ok((final_path, metadata))
524 }
525
526 pub fn write_atomic_parquet(
550 &self,
551 tenant_id: &str,
552 file_stem: &str,
553 events: &[Event],
554 ) -> Result<PathBuf> {
555 if events.is_empty() {
556 return Err(AllSourceError::StorageError(
557 "write_atomic_parquet called with empty event slice".to_string(),
558 ));
559 }
560 let anchor_ts = events
565 .iter()
566 .map(|e| e.timestamp)
567 .min()
568 .unwrap_or_else(chrono::Utc::now);
569 let partition_dir = partition_path_for_tenant(&self.storage_dir, tenant_id, anchor_ts)?;
570 fs::create_dir_all(&partition_dir).map_err(|e| {
571 AllSourceError::StorageError(format!(
572 "Failed to create tenant partition {}: {e}",
573 partition_dir.display()
574 ))
575 })?;
576
577 let record_batch = self.events_to_record_batch(events)?;
578 let (final_path, _meta) =
579 self.write_record_batch_atomic(&partition_dir, file_stem, &record_batch)?;
580
581 tracing::info!(
582 tenant_id = tenant_id,
583 file = %final_path.display(),
584 event_count = events.len(),
585 "wrote atomic snapshot file"
586 );
587
588 Ok(final_path)
589 }
590
591 pub fn cleanup_partial_writes(&self) -> Result<usize> {
614 let mut acted = 0usize;
615 let mut stack: Vec<PathBuf> = vec![self.storage_dir.clone()];
616 while let Some(dir) = stack.pop() {
617 let Ok(entries) = fs::read_dir(&dir) else {
618 continue;
619 };
620 for entry in entries.flatten() {
621 let path = entry.path();
622 let Ok(ft) = entry.file_type() else { continue };
623 if ft.is_dir() {
624 stack.push(path);
625 continue;
626 }
627 if !ft.is_file() {
628 continue;
629 }
630 let path_str = path.to_string_lossy();
631 if path_str.ends_with(".parquet.tmp") {
632 match fs::remove_file(&path) {
633 Ok(()) => {
634 tracing::warn!(
635 file = %path.display(),
636 "cleaned up orphan snapshot tmp file (crash recovery)"
637 );
638 acted += 1;
639 }
640 Err(e) => {
641 tracing::error!(
642 file = %path.display(),
643 "failed to remove orphan snapshot tmp file: {e}"
644 );
645 }
646 }
647 } else if path_str.ends_with(".parquet")
648 && fs::metadata(&path).is_ok_and(|m| m.len() == 0)
649 {
650 let ts = chrono::Utc::now().timestamp();
654 let quarantine_path = path.with_extension(format!("parquet.corrupt-{ts}"));
655 match fs::rename(&path, &quarantine_path) {
656 Ok(()) => {
657 tracing::error!(
658 from = %path.display(),
659 to = %quarantine_path.display(),
660 "quarantined 0-byte parquet file (issue #166 pre-fix crash). \
661 Operator: inspect and rm if not needed."
662 );
663 acted += 1;
664 }
665 Err(e) => {
666 tracing::error!(
667 file = %path.display(),
668 "failed to quarantine 0-byte parquet file: {e}"
669 );
670 }
671 }
672 }
673 }
674 }
675 Ok(acted)
676 }
677
678 #[cfg_attr(feature = "hotpath", hotpath::measure)]
684 pub fn flush_on_shutdown(&self) -> Result<usize> {
685 let total_pending: usize = {
686 let batches = self.current_batches.lock().unwrap();
687 batches.values().map(Vec::len).sum()
688 };
689
690 if total_pending > 0 {
691 tracing::info!(
692 "Shutdown: flushing {} pending events across all tenants",
693 total_pending
694 );
695 self.flush()?;
696 }
697
698 Ok(total_pending)
699 }
700
701 pub fn batch_stats(&self) -> BatchWriteStats {
703 let batches = self.batches_written.load(Ordering::Relaxed);
704 let events = self.events_written.load(Ordering::Relaxed);
705 let bytes = self.bytes_written.load(Ordering::Relaxed);
706 let time_ns = self.total_write_time_ns.load(Ordering::Relaxed);
707
708 let time_secs = time_ns as f64 / 1_000_000_000.0;
709
710 BatchWriteStats {
711 batches_written: batches,
712 events_written: events,
713 bytes_written: bytes,
714 avg_batch_size: if batches > 0 {
715 events as f64 / batches as f64
716 } else {
717 0.0
718 },
719 events_per_sec: if time_secs > 0.0 {
720 events as f64 / time_secs
721 } else {
722 0.0
723 },
724 total_write_time_ns: time_ns,
725 timeout_flushes: self.timeout_flushes.load(Ordering::Relaxed),
726 size_flushes: self.size_flushes.load(Ordering::Relaxed),
727 }
728 }
729
730 pub fn pending_count(&self) -> usize {
732 self.current_batches
733 .lock()
734 .unwrap()
735 .values()
736 .map(Vec::len)
737 .sum()
738 }
739
740 pub fn batch_size(&self) -> usize {
742 self.config.batch_size
743 }
744
745 pub fn flush_timeout(&self) -> Duration {
747 self.config.flush_timeout
748 }
749
750 #[cfg_attr(feature = "hotpath", hotpath::measure)]
752 fn events_to_record_batch(&self, events: &[Event]) -> Result<RecordBatch> {
753 let mut event_id_builder = StringBuilder::new();
754 let mut event_type_builder = StringBuilder::new();
755 let mut entity_id_builder = StringBuilder::new();
756 let mut payload_builder = StringBuilder::new();
757 let mut timestamp_builder = TimestampMicrosecondBuilder::new();
758 let mut metadata_builder = StringBuilder::new();
759 let mut version_builder = UInt64Builder::new();
760
761 for event in events {
762 event_id_builder.append_value(event.id.to_string());
763 event_type_builder.append_value(event.event_type_str());
764 entity_id_builder.append_value(event.entity_id_str());
765 payload_builder.append_value(serde_json::to_string(&event.payload)?);
766
767 let timestamp_micros = event.timestamp.timestamp_micros();
769 timestamp_builder.append_value(timestamp_micros);
770
771 if let Some(ref metadata) = event.metadata {
772 metadata_builder.append_value(serde_json::to_string(metadata)?);
773 } else {
774 metadata_builder.append_null();
775 }
776
777 version_builder.append_value(event.version as u64);
778 }
779
780 let arrays: Vec<ArrayRef> = vec![
781 Arc::new(event_id_builder.finish()),
782 Arc::new(event_type_builder.finish()),
783 Arc::new(entity_id_builder.finish()),
784 Arc::new(payload_builder.finish()),
785 Arc::new(timestamp_builder.finish()),
786 Arc::new(metadata_builder.finish()),
787 Arc::new(version_builder.finish()),
788 ];
789
790 let record_batch = RecordBatch::try_new(self.schema.clone(), arrays)?;
791
792 Ok(record_batch)
793 }
794
795 #[cfg_attr(feature = "hotpath", hotpath::measure)]
802 pub fn load_all_events(&self) -> Result<Vec<Event>> {
803 let parquet_files = find_parquet_files_recursive(&self.storage_dir)?;
804
805 let mut all_events = Vec::with_capacity(parquet_files.len() * self.config.batch_size);
806 let mut skipped = 0usize;
807 for file_path in parquet_files {
808 tracing::info!("Loading events from {}", file_path.display());
809 let tenant_id = tenant_id_from_path(&self.storage_dir, &file_path);
810 match self.load_events_from_file(&file_path, &tenant_id) {
815 Ok(file_events) => all_events.extend(file_events),
816 Err(e) => {
817 tracing::error!(
818 file = %file_path.display(),
819 error = %e,
820 "Skipping unreadable parquet file — other files will still load. \
821 Likely a 0-byte or truncated file from an unclean shutdown; \
822 inspect and remove manually after confirming."
823 );
824 skipped += 1;
825 }
826 }
827 }
828
829 if skipped > 0 {
830 tracing::warn!(
831 "Loaded {} events from storage; skipped {} unreadable file(s)",
832 all_events.len(),
833 skipped
834 );
835 } else {
836 tracing::info!("Loaded {} total events from storage", all_events.len());
837 }
838
839 Ok(all_events)
840 }
841
842 #[cfg_attr(feature = "hotpath", hotpath::measure)]
848 fn load_events_from_file(&self, file_path: &Path, tenant_id: &str) -> Result<Vec<Event>> {
849 use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
850
851 let file = File::open(file_path).map_err(|e| {
852 AllSourceError::StorageError(format!("Failed to open parquet file: {e}"))
853 })?;
854
855 let builder = ParquetRecordBatchReaderBuilder::try_new(file)?;
856 let mut reader = builder.build()?;
857
858 let mut events = Vec::new();
859
860 while let Some(Ok(batch)) = reader.next() {
861 let batch_events = self.record_batch_to_events(&batch, tenant_id)?;
862 events.extend(batch_events);
863 }
864
865 Ok(events)
866 }
867
868 pub fn load_events_from_file_path(
874 &self,
875 file_path: &Path,
876 tenant_id: &str,
877 ) -> Result<Vec<Event>> {
878 self.load_events_from_file(file_path, tenant_id)
879 }
880
881 #[cfg_attr(feature = "hotpath", hotpath::measure)]
885 fn record_batch_to_events(&self, batch: &RecordBatch, tenant_id: &str) -> Result<Vec<Event>> {
886 let event_ids = batch
887 .column(0)
888 .as_any()
889 .downcast_ref::<arrow::array::StringArray>()
890 .ok_or_else(|| AllSourceError::StorageError("Invalid event_id column".to_string()))?;
891
892 let event_types = batch
893 .column(1)
894 .as_any()
895 .downcast_ref::<arrow::array::StringArray>()
896 .ok_or_else(|| AllSourceError::StorageError("Invalid event_type column".to_string()))?;
897
898 let entity_ids = batch
899 .column(2)
900 .as_any()
901 .downcast_ref::<arrow::array::StringArray>()
902 .ok_or_else(|| AllSourceError::StorageError("Invalid entity_id column".to_string()))?;
903
904 let payloads = batch
905 .column(3)
906 .as_any()
907 .downcast_ref::<arrow::array::StringArray>()
908 .ok_or_else(|| AllSourceError::StorageError("Invalid payload column".to_string()))?;
909
910 let timestamps = batch
911 .column(4)
912 .as_any()
913 .downcast_ref::<TimestampMicrosecondArray>()
914 .ok_or_else(|| AllSourceError::StorageError("Invalid timestamp column".to_string()))?;
915
916 let metadatas = batch
917 .column(5)
918 .as_any()
919 .downcast_ref::<arrow::array::StringArray>()
920 .ok_or_else(|| AllSourceError::StorageError("Invalid metadata column".to_string()))?;
921
922 let versions = batch
923 .column(6)
924 .as_any()
925 .downcast_ref::<arrow::array::UInt64Array>()
926 .ok_or_else(|| AllSourceError::StorageError("Invalid version column".to_string()))?;
927
928 let mut events = Vec::new();
929
930 for i in 0..batch.num_rows() {
931 let id = uuid::Uuid::parse_str(event_ids.value(i))
932 .map_err(|e| AllSourceError::StorageError(format!("Invalid UUID: {e}")))?;
933
934 let timestamp = chrono::DateTime::from_timestamp_micros(timestamps.value(i))
935 .ok_or_else(|| AllSourceError::StorageError("Invalid timestamp".to_string()))?;
936
937 let metadata = if metadatas.is_null(i) {
938 None
939 } else {
940 Some(serde_json::from_str(metadatas.value(i))?)
941 };
942
943 let event = Event::reconstruct_from_strings(
944 id,
945 event_types.value(i).to_string(),
946 entity_ids.value(i).to_string(),
947 tenant_id.to_string(),
948 serde_json::from_str(payloads.value(i))?,
949 timestamp,
950 metadata,
951 versions.value(i) as i64,
952 );
953
954 events.push(event);
955 }
956
957 Ok(events)
958 }
959
960 pub fn list_parquet_files(&self) -> Result<Vec<PathBuf>> {
966 find_parquet_files_recursive(&self.storage_dir)
967 }
968
969 pub fn list_parquet_files_for_tenant(&self, tenant_id: &str) -> Result<Vec<PathBuf>> {
982 let safe = sanitize_tenant_id_for_path(tenant_id)?;
983 let tenant_root = self.storage_dir.join(safe);
984 if !tenant_root.is_dir() {
985 return Ok(Vec::new());
986 }
987 find_parquet_files_recursive(&tenant_root)
988 }
989
990 #[cfg_attr(feature = "hotpath", hotpath::measure)]
1008 pub fn load_events_for_tenant(&self, tenant_id: &str) -> Result<Vec<Event>> {
1009 let parquet_files = self.list_parquet_files_for_tenant(tenant_id)?;
1010 tracing::info!(
1011 tenant_id = tenant_id,
1012 file_count = parquet_files.len(),
1013 "load_events_for_tenant: walking tenant subtree only"
1014 );
1015
1016 let mut events = Vec::with_capacity(parquet_files.len() * self.config.batch_size);
1017 let mut skipped = 0usize;
1018 for file_path in parquet_files {
1019 tracing::debug!(
1020 tenant_id = tenant_id,
1021 file = %file_path.display(),
1022 "load_events_for_tenant: opening file"
1023 );
1024 match self.load_events_from_file(&file_path, tenant_id) {
1029 Ok(file_events) => events.extend(file_events),
1030 Err(e) => {
1031 tracing::error!(
1032 tenant_id = tenant_id,
1033 file = %file_path.display(),
1034 error = %e,
1035 "Skipping unreadable parquet file in tenant subtree"
1036 );
1037 skipped += 1;
1038 }
1039 }
1040 }
1041
1042 tracing::info!(
1043 tenant_id = tenant_id,
1044 event_count = events.len(),
1045 skipped_files = skipped,
1046 "load_events_for_tenant: complete"
1047 );
1048 Ok(events)
1049 }
1050
1051 pub fn storage_dir(&self) -> &Path {
1053 &self.storage_dir
1054 }
1055
1056 pub fn migrate_flat_layout(&self, dry_run: bool) -> Result<MigrationReport> {
1079 let flat_files = list_flat_layout_files(&self.storage_dir)?;
1080 let mut report = MigrationReport {
1081 dry_run,
1082 ..Default::default()
1083 };
1084
1085 for flat_file in flat_files {
1086 let events = self.load_events_from_file(&flat_file, "default")?;
1090 report.flat_files_seen += 1;
1091
1092 if events.is_empty() {
1093 if !dry_run {
1095 fs::remove_file(&flat_file).map_err(|e| {
1096 AllSourceError::StorageError(format!(
1097 "Failed to remove empty flat file {}: {e}",
1098 flat_file.display()
1099 ))
1100 })?;
1101 }
1102 report.flat_files_removed += 1;
1103 continue;
1104 }
1105
1106 let mut groups: HashMap<(String, String), Vec<Event>> = HashMap::new();
1112 for event in events {
1113 let key = (
1114 event.tenant_id_str().to_string(),
1115 event.timestamp().format("%Y-%m").to_string(),
1116 );
1117 groups.entry(key).or_default().push(event);
1118 }
1119
1120 for ((tenant, yyyy_mm), group_events) in groups {
1121 let count = group_events.len();
1122 if !dry_run {
1123 let safe_tenant = sanitize_tenant_id_for_path(&tenant)?;
1124 let target_dir = self.storage_dir.join(safe_tenant).join(&yyyy_mm);
1125 fs::create_dir_all(&target_dir).map_err(|e| {
1126 AllSourceError::StorageError(format!(
1127 "Failed to create partition {}: {e}",
1128 target_dir.display()
1129 ))
1130 })?;
1131 let filename = format!(
1132 "events-{}-{}.parquet",
1133 chrono::Utc::now().format("%Y%m%d-%H%M%S%3f"),
1134 uuid::Uuid::new_v4().as_simple()
1135 );
1136 let target_path = target_dir.join(&filename);
1137 let record_batch = self.events_to_record_batch(&group_events)?;
1138 let file = File::create(&target_path).map_err(|e| {
1139 AllSourceError::StorageError(format!(
1140 "Failed to create migration target {}: {e}",
1141 target_path.display()
1142 ))
1143 })?;
1144 let props = WriterProperties::builder()
1145 .set_compression(self.config.compression)
1146 .build();
1147 let mut writer = ArrowWriter::try_new(file, self.schema.clone(), Some(props))?;
1148 writer.write(&record_batch)?;
1149 writer.close()?;
1150 report.partitions_written += 1;
1151 }
1152 report.events_migrated += count;
1153 }
1154
1155 if !dry_run {
1156 fs::remove_file(&flat_file).map_err(|e| {
1157 AllSourceError::StorageError(format!(
1158 "Failed to remove flat file {} after migration: {e}",
1159 flat_file.display()
1160 ))
1161 })?;
1162 report.flat_files_removed += 1;
1163 }
1164 }
1165
1166 Ok(report)
1167 }
1168
1169 pub fn stats(&self) -> Result<StorageStats> {
1171 let parquet_files = find_parquet_files_recursive(&self.storage_dir)?;
1172 let mut total_size_bytes = 0u64;
1173 for path in &parquet_files {
1174 if let Ok(metadata) = fs::metadata(path) {
1175 total_size_bytes += metadata.len();
1176 }
1177 }
1178
1179 let current_batch_size: usize = self
1180 .current_batches
1181 .lock()
1182 .unwrap()
1183 .values()
1184 .map(Vec::len)
1185 .sum();
1186
1187 Ok(StorageStats {
1188 total_files: parquet_files.len(),
1189 total_size_bytes,
1190 storage_dir: self.storage_dir.clone(),
1191 current_batch_size,
1192 })
1193 }
1194}
1195
1196fn sanitize_tenant_id_for_path(tenant_id: &str) -> Result<&str> {
1209 if tenant_id.is_empty() {
1210 return Err(AllSourceError::StorageError(
1211 "tenant_id is empty (cannot derive partition path)".to_string(),
1212 ));
1213 }
1214 if tenant_id.len() > 128 {
1215 return Err(AllSourceError::StorageError(format!(
1216 "tenant_id is too long for partition path: {} bytes (max 128)",
1217 tenant_id.len()
1218 )));
1219 }
1220 if tenant_id == "." || tenant_id == ".." {
1221 return Err(AllSourceError::StorageError(format!(
1222 "tenant_id {tenant_id:?} is reserved"
1223 )));
1224 }
1225 for c in tenant_id.chars() {
1226 let ok = c.is_ascii_alphanumeric() || c == '-' || c == '_' || c == '.';
1227 if !ok {
1228 return Err(AllSourceError::StorageError(format!(
1229 "tenant_id {tenant_id:?} contains disallowed character {c:?} for partition path"
1230 )));
1231 }
1232 }
1233 Ok(tenant_id)
1234}
1235
1236fn partition_path_for_tenant(
1241 root: &Path,
1242 tenant_id: &str,
1243 when: chrono::DateTime<chrono::Utc>,
1244) -> Result<PathBuf> {
1245 let safe = sanitize_tenant_id_for_path(tenant_id)?;
1246 Ok(root.join(safe).join(when.format("%Y-%m").to_string()))
1247}
1248
1249fn tenant_id_from_path(root: &Path, file_path: &Path) -> String {
1259 let Ok(rel) = file_path.strip_prefix(root) else {
1260 return "default".to_string();
1261 };
1262 let mut comps = rel.components();
1263 let first = comps.next();
1264 let next = comps.next();
1265 match (first, next) {
1266 (Some(std::path::Component::Normal(tenant)), Some(_)) => {
1268 tenant.to_string_lossy().into_owned()
1269 }
1270 _ => "default".to_string(),
1272 }
1273}
1274
1275fn list_flat_layout_files(root: &Path) -> Result<Vec<PathBuf>> {
1281 let entries = fs::read_dir(root).map_err(|e| {
1282 AllSourceError::StorageError(format!("Failed to read storage directory: {e}"))
1283 })?;
1284 let mut out: Vec<PathBuf> = entries
1285 .filter_map(std::result::Result::ok)
1286 .filter_map(|entry| {
1287 let ft = entry.file_type().ok()?;
1288 if !ft.is_file() {
1289 return None;
1290 }
1291 let path = entry.path();
1292 if path.extension().and_then(|s| s.to_str()) == Some("parquet") {
1293 Some(path)
1294 } else {
1295 None
1296 }
1297 })
1298 .collect();
1299 out.sort();
1300 Ok(out)
1301}
1302
1303fn find_parquet_files_recursive(root: &Path) -> Result<Vec<PathBuf>> {
1314 let mut out = Vec::new();
1315 let mut stack: Vec<PathBuf> = vec![root.to_path_buf()];
1316
1317 while let Some(dir) = stack.pop() {
1318 let entries = match fs::read_dir(&dir) {
1319 Ok(e) => e,
1320 Err(e) if dir == root => {
1324 return Err(AllSourceError::StorageError(format!(
1325 "Failed to read storage directory: {e}"
1326 )));
1327 }
1328 Err(_) => continue,
1329 };
1330
1331 for entry in entries.flatten() {
1332 let path = entry.path();
1333 let Ok(ft) = entry.file_type() else {
1337 continue;
1338 };
1339 if ft.is_dir() {
1340 stack.push(path);
1341 } else if ft.is_file()
1342 && path
1343 .extension()
1344 .and_then(|ext| ext.to_str())
1345 .is_some_and(|ext| ext == "parquet")
1346 {
1347 out.push(path);
1348 }
1349 }
1350 }
1351
1352 out.sort();
1353 Ok(out)
1354}
1355
1356impl Drop for ParquetStorage {
1357 fn drop(&mut self) {
1358 if let Err(e) = self.flush_on_shutdown() {
1360 tracing::error!("Failed to flush events on drop: {}", e);
1361 }
1362 }
1363}
1364
1365#[derive(Debug, Default, Clone, serde::Serialize)]
1367pub struct MigrationReport {
1368 pub dry_run: bool,
1370 pub flat_files_seen: usize,
1372 pub flat_files_removed: usize,
1374 pub partitions_written: usize,
1376 pub events_migrated: usize,
1378}
1379
1380#[derive(Debug, serde::Serialize)]
1381pub struct StorageStats {
1382 pub total_files: usize,
1383 pub total_size_bytes: u64,
1384 pub storage_dir: PathBuf,
1385 pub current_batch_size: usize,
1386}
1387
1388#[cfg(test)]
1389mod tests {
1390 use super::*;
1391 use serde_json::json;
1392 use std::sync::Arc;
1393 use tempfile::TempDir;
1394
1395 fn create_test_event(entity_id: &str) -> Event {
1396 Event::reconstruct_from_strings(
1397 uuid::Uuid::new_v4(),
1398 "test.event".to_string(),
1399 entity_id.to_string(),
1400 "default".to_string(),
1401 json!({
1402 "test": "data",
1403 "value": 42
1404 }),
1405 chrono::Utc::now(),
1406 None,
1407 1,
1408 )
1409 }
1410
1411 #[test]
1412 fn test_parquet_storage_write_read() {
1413 let temp_dir = TempDir::new().unwrap();
1414 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1415
1416 for i in 0..10 {
1418 let event = create_test_event(&format!("entity-{i}"));
1419 storage.append_event(event).unwrap();
1420 }
1421
1422 storage.flush().unwrap();
1424
1425 let loaded_events = storage.load_all_events().unwrap();
1427 assert_eq!(loaded_events.len(), 10);
1428 }
1429
1430 #[test]
1431 fn test_storage_stats() {
1432 let temp_dir = TempDir::new().unwrap();
1433 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1434
1435 for i in 0..5 {
1437 storage
1438 .append_event(create_test_event(&format!("entity-{i}")))
1439 .unwrap();
1440 }
1441 storage.flush().unwrap();
1442
1443 let stats = storage.stats().unwrap();
1444 assert_eq!(stats.total_files, 1);
1445 assert!(stats.total_size_bytes > 0);
1446 }
1447
1448 #[test]
1449 fn test_default_batch_size() {
1450 let temp_dir = TempDir::new().unwrap();
1451 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1452
1453 assert_eq!(storage.batch_size(), DEFAULT_BATCH_SIZE);
1455 assert_eq!(storage.batch_size(), 10_000);
1456 }
1457
1458 #[test]
1459 fn test_custom_config() {
1460 let temp_dir = TempDir::new().unwrap();
1461 let config = ParquetStorageConfig {
1462 batch_size: 5_000,
1463 flush_timeout: Duration::from_secs(2),
1464 compression: parquet::basic::Compression::SNAPPY,
1465 };
1466 let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1467
1468 assert_eq!(storage.batch_size(), 5_000);
1469 assert_eq!(storage.flush_timeout(), Duration::from_secs(2));
1470 }
1471
1472 #[test]
1473 fn test_batch_write() {
1474 let temp_dir = TempDir::new().unwrap();
1475 let config = ParquetStorageConfig {
1476 batch_size: 100, ..Default::default()
1478 };
1479 let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1480
1481 let events: Vec<Event> = (0..250)
1489 .map(|i| create_test_event(&format!("entity-{i}")))
1490 .collect();
1491
1492 let result = storage.batch_write(events).unwrap();
1493 assert_eq!(result.events_written, 250);
1494 assert_eq!(result.batches_flushed, 1);
1495 assert_eq!(storage.pending_count(), 0);
1496
1497 storage.flush().unwrap();
1499
1500 let loaded = storage.load_all_events().unwrap();
1502 assert_eq!(loaded.len(), 250);
1503 }
1504
1505 #[test]
1506 fn test_auto_flush_on_batch_size() {
1507 let temp_dir = TempDir::new().unwrap();
1508 let config = ParquetStorageConfig {
1509 batch_size: 10, ..Default::default()
1511 };
1512 let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1513
1514 for i in 0..15 {
1516 storage
1517 .append_event(create_test_event(&format!("entity-{i}")))
1518 .unwrap();
1519 }
1520
1521 assert_eq!(storage.pending_count(), 5);
1523
1524 let stats = storage.batch_stats();
1525 assert_eq!(stats.events_written, 10);
1526 assert_eq!(stats.batches_written, 1);
1527 assert_eq!(stats.size_flushes, 1);
1528 }
1529
1530 #[test]
1531 fn test_flush_on_shutdown() {
1532 let temp_dir = TempDir::new().unwrap();
1533 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1534
1535 for i in 0..5 {
1537 storage
1538 .append_event(create_test_event(&format!("entity-{i}")))
1539 .unwrap();
1540 }
1541
1542 assert_eq!(storage.pending_count(), 5);
1543
1544 let flushed = storage.flush_on_shutdown().unwrap();
1546 assert_eq!(flushed, 5);
1547 assert_eq!(storage.pending_count(), 0);
1548
1549 let loaded = storage.load_all_events().unwrap();
1551 assert_eq!(loaded.len(), 5);
1552 }
1553
1554 #[test]
1555 fn test_thread_safe_writes() {
1556 let temp_dir = TempDir::new().unwrap();
1557 let config = ParquetStorageConfig {
1558 batch_size: 100,
1559 ..Default::default()
1560 };
1561 let storage = Arc::new(ParquetStorage::with_config(temp_dir.path(), config).unwrap());
1562
1563 let events_per_thread = 50;
1564 let thread_count = 4;
1565
1566 std::thread::scope(|s| {
1567 for t in 0..thread_count {
1568 let storage_ref = Arc::clone(&storage);
1569 s.spawn(move || {
1570 for i in 0..events_per_thread {
1571 let event = create_test_event(&format!("thread-{t}-entity-{i}"));
1572 storage_ref.append_event(event).unwrap();
1573 }
1574 });
1575 }
1576 });
1577
1578 storage.flush().unwrap();
1580
1581 let loaded = storage.load_all_events().unwrap();
1583 assert_eq!(loaded.len(), events_per_thread * thread_count);
1584 }
1585
1586 #[test]
1587 fn test_batch_stats() {
1588 let temp_dir = TempDir::new().unwrap();
1589 let config = ParquetStorageConfig {
1590 batch_size: 50,
1591 ..Default::default()
1592 };
1593 let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1594
1595 let events: Vec<Event> = (0..100)
1600 .map(|i| create_test_event(&format!("entity-{i}")))
1601 .collect();
1602
1603 storage.batch_write(events).unwrap();
1604
1605 let stats = storage.batch_stats();
1606 assert_eq!(stats.batches_written, 1);
1607 assert_eq!(stats.events_written, 100);
1608 assert!(stats.avg_batch_size > 0.0);
1609 assert!(stats.events_per_sec > 0.0);
1610 assert_eq!(stats.size_flushes, 1);
1611 }
1612
1613 #[test]
1614 fn test_config_presets() {
1615 let high_throughput = ParquetStorageConfig::high_throughput();
1616 assert_eq!(high_throughput.batch_size, 50_000);
1617 assert_eq!(high_throughput.flush_timeout, Duration::from_secs(10));
1618
1619 let low_latency = ParquetStorageConfig::low_latency();
1620 assert_eq!(low_latency.batch_size, 1_000);
1621 assert_eq!(low_latency.flush_timeout, Duration::from_secs(1));
1622
1623 let default = ParquetStorageConfig::default();
1624 assert_eq!(default.batch_size, DEFAULT_BATCH_SIZE);
1625 assert_eq!(default.batch_size, 10_000);
1626 }
1627
1628 #[test]
1631 #[ignore]
1632 fn test_batch_write_throughput() {
1633 let temp_dir = TempDir::new().unwrap();
1634 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1635
1636 let event_count = 50_000;
1637
1638 let events: Vec<Event> = (0..event_count)
1640 .map(|i| create_test_event(&format!("entity-{i}")))
1641 .collect();
1642
1643 let start = std::time::Instant::now();
1644 let result = storage.batch_write(events).unwrap();
1645 storage.flush().unwrap(); let batch_duration = start.elapsed();
1647
1648 let batch_stats = storage.batch_stats();
1649
1650 println!("\n=== Parquet Batch Write Performance (BATCH_SIZE=10,000) ===");
1651 println!("Events: {event_count}");
1652 println!("Duration: {batch_duration:?}");
1653 println!("Events/sec: {:.0}", result.events_per_sec);
1654 println!("Batches written: {}", batch_stats.batches_written);
1655 println!("Avg batch size: {:.0}", batch_stats.avg_batch_size);
1656 println!("Bytes written: {} KB", batch_stats.bytes_written / 1024);
1657
1658 assert!(
1661 result.events_per_sec > 10_000.0,
1662 "Batch write throughput too low: {:.0} events/sec (expected >10K in debug, >100K in release)",
1663 result.events_per_sec
1664 );
1665 }
1666
1667 #[test]
1669 #[ignore]
1670 fn test_single_event_write_baseline() {
1671 let temp_dir = TempDir::new().unwrap();
1672 let config = ParquetStorageConfig {
1673 batch_size: 1, ..Default::default()
1675 };
1676 let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1677
1678 let event_count = 1_000; let start = std::time::Instant::now();
1681 for i in 0..event_count {
1682 let event = create_test_event(&format!("entity-{i}"));
1683 storage.append_event(event).unwrap();
1684 }
1685 let duration = start.elapsed();
1686
1687 let events_per_sec = f64::from(event_count) / duration.as_secs_f64();
1688
1689 println!("\n=== Single-Event Write Baseline ===");
1690 println!("Events: {event_count}");
1691 println!("Duration: {duration:?}");
1692 println!("Events/sec: {events_per_sec:.0}");
1693
1694 }
1697
1698 fn touch_parquet(path: &Path) {
1708 std::fs::create_dir_all(path.parent().unwrap()).unwrap();
1709 std::fs::write(path, b"").unwrap();
1710 }
1711
1712 #[test]
1713 fn test_walker_finds_files_in_flat_layout() {
1714 let temp_dir = TempDir::new().unwrap();
1715 let root = temp_dir.path();
1716 touch_parquet(&root.join("events-20260101-120000000-aaaa.parquet"));
1717 touch_parquet(&root.join("events-20260101-130000000-bbbb.parquet"));
1718
1719 let mut found = find_parquet_files_recursive(root).unwrap();
1720 found.sort();
1721 assert_eq!(found.len(), 2);
1722 assert!(
1723 found[0]
1724 .file_name()
1725 .unwrap()
1726 .to_str()
1727 .unwrap()
1728 .starts_with("events-"),
1729 "expected events-* file, got {found:?}"
1730 );
1731 }
1732
1733 #[test]
1734 fn test_walker_finds_files_in_tenant_partitioned_tree() {
1735 let temp_dir = TempDir::new().unwrap();
1736 let root = temp_dir.path();
1737 touch_parquet(&root.join("tenant-a/2026-01/events-20260101-120000000-aaaa.parquet"));
1739 touch_parquet(&root.join("tenant-a/2026-02/events-20260201-120000000-bbbb.parquet"));
1740 touch_parquet(&root.join("tenant-b/2026-01/events-20260103-120000000-cccc.parquet"));
1741
1742 let found = find_parquet_files_recursive(root).unwrap();
1743 assert_eq!(found.len(), 3);
1744 assert!(found[0].to_str().unwrap().contains("tenant-a"));
1747 assert!(found[1].to_str().unwrap().contains("tenant-a"));
1748 assert!(found[2].to_str().unwrap().contains("tenant-b"));
1749 }
1750
1751 #[test]
1752 fn test_walker_handles_mixed_legacy_and_partitioned_layouts() {
1753 let temp_dir = TempDir::new().unwrap();
1758 let root = temp_dir.path();
1759 touch_parquet(&root.join("events-legacy-aaaa.parquet"));
1760 touch_parquet(&root.join("tenant-a/2026-01/events-new-bbbb.parquet"));
1761
1762 let found = find_parquet_files_recursive(root).unwrap();
1763 assert_eq!(found.len(), 2);
1764 }
1765
1766 #[test]
1767 fn test_walker_ignores_non_parquet_files() {
1768 let temp_dir = TempDir::new().unwrap();
1769 let root = temp_dir.path();
1770 std::fs::write(root.join("README.md"), b"hello").unwrap();
1771 std::fs::write(root.join("events.json"), b"[]").unwrap();
1772 touch_parquet(&root.join("events-20260101-120000000-aaaa.parquet"));
1773 std::fs::write(root.join("not-a-parquet-file.bin"), b"").unwrap();
1776
1777 let found = find_parquet_files_recursive(root).unwrap();
1778 assert_eq!(found.len(), 1);
1779 assert_eq!(
1780 found[0].extension().and_then(|s| s.to_str()),
1781 Some("parquet")
1782 );
1783 }
1784
1785 fn event_with_tenant(tenant: &str, entity_id: &str) -> Event {
1789 Event::reconstruct_from_strings(
1790 uuid::Uuid::new_v4(),
1791 "test.event".to_string(),
1792 entity_id.to_string(),
1793 tenant.to_string(),
1794 json!({"k": "v"}),
1795 chrono::Utc::now(),
1796 None,
1797 1,
1798 )
1799 }
1800
1801 #[test]
1802 fn test_flush_writes_into_per_tenant_partition() {
1803 let temp_dir = TempDir::new().unwrap();
1807 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1808
1809 for i in 0..3 {
1810 storage
1811 .append_event(event_with_tenant("default", &format!("entity-{i}")))
1812 .unwrap();
1813 }
1814 storage.flush().unwrap();
1815
1816 let parquet_files = find_parquet_files_recursive(temp_dir.path()).unwrap();
1817 assert_eq!(parquet_files.len(), 1);
1818
1819 let rel = parquet_files[0]
1820 .strip_prefix(temp_dir.path())
1821 .unwrap()
1822 .to_string_lossy()
1823 .into_owned();
1824 let parts: Vec<&str> = rel.split(std::path::MAIN_SEPARATOR).collect();
1826 assert_eq!(parts.len(), 3, "expected tenant/yyyy-mm/file, got {rel}");
1827 assert_eq!(parts[0], "default");
1828 assert!(
1831 parts[1].len() == 7 && parts[1].as_bytes()[4] == b'-',
1832 "expected yyyy-mm, got {}",
1833 parts[1]
1834 );
1835 assert!(parts[2].starts_with("events-") && parts[2].ends_with(".parquet"));
1836
1837 let loaded = storage.load_all_events().unwrap();
1838 assert_eq!(loaded.len(), 3);
1839 }
1840
1841 #[test]
1842 fn test_multiple_tenants_get_isolated_subtrees() {
1843 let temp_dir = TempDir::new().unwrap();
1846 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1847
1848 for i in 0..2 {
1849 storage
1850 .append_event(event_with_tenant("alice", &format!("a-{i}")))
1851 .unwrap();
1852 }
1853 for i in 0..3 {
1854 storage
1855 .append_event(event_with_tenant("bob", &format!("b-{i}")))
1856 .unwrap();
1857 }
1858 storage.flush().unwrap();
1859
1860 let alice_subtree = temp_dir.path().join("alice");
1861 let bob_subtree = temp_dir.path().join("bob");
1862 assert!(alice_subtree.is_dir(), "alice should have its own subtree");
1863 assert!(bob_subtree.is_dir(), "bob should have its own subtree");
1864
1865 let alice_files = find_parquet_files_recursive(&alice_subtree).unwrap();
1866 let bob_files = find_parquet_files_recursive(&bob_subtree).unwrap();
1867 assert_eq!(alice_files.len(), 1);
1868 assert_eq!(bob_files.len(), 1);
1869
1870 let loaded = storage.load_all_events().unwrap();
1873 let (alice_count, bob_count) =
1874 loaded
1875 .iter()
1876 .fold((0, 0), |(a, b), e| match e.tenant_id_str() {
1877 "alice" => (a + 1, b),
1878 "bob" => (a, b + 1),
1879 _ => (a, b),
1880 });
1881 assert_eq!(alice_count, 2);
1882 assert_eq!(bob_count, 3);
1883 }
1884
1885 #[test]
1886 fn test_size_flush_only_drains_full_tenant() {
1887 let temp_dir = TempDir::new().unwrap();
1891 let config = ParquetStorageConfig {
1892 batch_size: 5,
1893 ..Default::default()
1894 };
1895 let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1896
1897 for i in 0..5 {
1900 storage
1901 .append_event(event_with_tenant("alice", &format!("a-{i}")))
1902 .unwrap();
1903 }
1904 for i in 0..2 {
1906 storage
1907 .append_event(event_with_tenant("bob", &format!("b-{i}")))
1908 .unwrap();
1909 }
1910
1911 assert_eq!(
1912 storage.pending_count(),
1913 2,
1914 "only bob's 2 events should be pending"
1915 );
1916
1917 let parquet_files = find_parquet_files_recursive(temp_dir.path()).unwrap();
1918 assert_eq!(parquet_files.len(), 1, "only alice should have flushed");
1919 assert!(
1920 parquet_files[0]
1921 .to_string_lossy()
1922 .contains(&format!("alice{}", std::path::MAIN_SEPARATOR)),
1923 "expected alice partition, got {}",
1924 parquet_files[0].display()
1925 );
1926 }
1927
1928 #[test]
1929 fn test_tenant_id_from_path_recovers_tenant_for_partitioned_files() {
1930 let root = Path::new("/data/storage");
1931 let f = Path::new("/data/storage/alice/2026-04/events-20260426-120000000-aaaa.parquet");
1932 assert_eq!(tenant_id_from_path(root, f), "alice");
1933 }
1934
1935 #[test]
1936 fn test_tenant_id_from_path_falls_back_to_default_for_legacy_flat_layout() {
1937 let root = Path::new("/data/storage");
1938 let f = Path::new("/data/storage/events-20260426-120000000-aaaa.parquet");
1939 assert_eq!(tenant_id_from_path(root, f), "default");
1942 }
1943
1944 #[test]
1945 fn test_sanitize_tenant_id_for_path_accepts_safe_inputs() {
1946 for ok in [
1947 "default",
1948 "system",
1949 "1e6b2d1c-2f64-4441-9cf9-42f2e451aa17",
1950 "onboard-diagnostic-160-at-example-com",
1951 "tenant_with_underscore",
1952 "v1.0",
1953 ] {
1954 assert!(
1955 sanitize_tenant_id_for_path(ok).is_ok(),
1956 "{ok:?} should be accepted"
1957 );
1958 }
1959 }
1960
1961 #[test]
1962 fn test_sanitize_tenant_id_for_path_rejects_unsafe_inputs() {
1963 for bad in [
1964 "", "..", ".", "foo/bar", "foo\\bar", "foo bar", "foo\nbar", "foo\0bar", "tenant?", "tenant*", ] {
1975 assert!(
1976 sanitize_tenant_id_for_path(bad).is_err(),
1977 "{bad:?} should be rejected"
1978 );
1979 }
1980
1981 let too_long = "a".repeat(129);
1983 assert!(sanitize_tenant_id_for_path(&too_long).is_err());
1984 }
1985
1986 #[test]
1987 fn test_partition_path_for_tenant_shape() {
1988 let root = Path::new("/data");
1989 let when = chrono::DateTime::parse_from_rfc3339("2026-04-26T12:00:00Z")
1990 .unwrap()
1991 .with_timezone(&chrono::Utc);
1992 let path = partition_path_for_tenant(root, "alice", when).unwrap();
1993 assert_eq!(path, Path::new("/data/alice/2026-04"));
1994 }
1995
1996 #[test]
1997 fn test_append_event_rejects_unsafe_tenant_at_flush() {
1998 let temp_dir = TempDir::new().unwrap();
2004 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2005
2006 storage
2010 .append_event(event_with_tenant("../escape", "e-0"))
2011 .unwrap();
2012 let result = storage.flush();
2013 assert!(result.is_err(), "flush should reject unsafe tenant_id");
2014 let msg = format!("{}", result.unwrap_err());
2015 assert!(
2016 msg.contains("disallowed character") || msg.contains("reserved"),
2017 "expected sanitization error message, got: {msg}"
2018 );
2019 }
2020
2021 #[test]
2026 fn test_load_events_for_tenant_only_walks_target_subtree() {
2027 let temp_dir = TempDir::new().unwrap();
2032 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2033
2034 for i in 0..2 {
2035 storage
2036 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2037 .unwrap();
2038 }
2039 for i in 0..3 {
2040 storage
2041 .append_event(event_with_tenant("bob", &format!("b-{i}")))
2042 .unwrap();
2043 }
2044 for i in 0..1 {
2045 storage
2046 .append_event(event_with_tenant("carol", &format!("c-{i}")))
2047 .unwrap();
2048 }
2049 storage.flush().unwrap();
2050
2051 let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
2052 assert_eq!(alice_files.len(), 1);
2053 assert!(
2054 alice_files[0]
2055 .to_string_lossy()
2056 .contains(&format!("alice{}", std::path::MAIN_SEPARATOR)),
2057 "expected alice file, got {}",
2058 alice_files[0].display()
2059 );
2060 for f in &alice_files {
2064 let s = f.to_string_lossy();
2065 assert!(!s.contains("bob"), "alice listing leaked bob file: {s}");
2066 assert!(!s.contains("carol"), "alice listing leaked carol file: {s}");
2067 }
2068
2069 let alice_events = storage.load_events_for_tenant("alice").unwrap();
2070 assert_eq!(alice_events.len(), 2);
2071 for e in &alice_events {
2072 assert_eq!(e.tenant_id_str(), "alice");
2073 }
2074
2075 let bob_events = storage.load_events_for_tenant("bob").unwrap();
2076 assert_eq!(bob_events.len(), 3);
2077 for e in &bob_events {
2078 assert_eq!(e.tenant_id_str(), "bob");
2079 }
2080
2081 let carol_events = storage.load_events_for_tenant("carol").unwrap();
2082 assert_eq!(carol_events.len(), 1);
2083 assert_eq!(carol_events[0].tenant_id_str(), "carol");
2084 }
2085
2086 #[test]
2087 fn test_load_events_for_tenant_returns_empty_when_subtree_missing() {
2088 let temp_dir = TempDir::new().unwrap();
2092 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2093
2094 storage
2097 .append_event(event_with_tenant("alice", "a-0"))
2098 .unwrap();
2099 storage.flush().unwrap();
2100
2101 let files = storage
2102 .list_parquet_files_for_tenant("nobody-here")
2103 .unwrap();
2104 assert!(files.is_empty());
2105
2106 let events = storage.load_events_for_tenant("nobody-here").unwrap();
2107 assert!(events.is_empty());
2108 }
2109
2110 #[test]
2111 fn test_load_events_for_tenant_rejects_unsafe_tenant_id() {
2112 let temp_dir = TempDir::new().unwrap();
2115 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2116
2117 for unsafe_tid in ["..", "a/b", "a\\b", "", "a..b/.."] {
2118 let result = storage.load_events_for_tenant(unsafe_tid);
2119 assert!(
2120 result.is_err(),
2121 "tenant_id {unsafe_tid:?} should have been rejected"
2122 );
2123 }
2124 }
2125
2126 #[test]
2127 fn test_load_events_for_tenant_ignores_legacy_flat_layout_files() {
2128 let temp_dir = TempDir::new().unwrap();
2136 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2137
2138 let _flat = seed_flat_layout_file(&storage, 4);
2140
2141 let default_events = storage.load_events_for_tenant("default").unwrap();
2144 assert!(
2145 default_events.is_empty(),
2146 "tenant-scoped load must not pick up flat-layout files; got {} events",
2147 default_events.len()
2148 );
2149
2150 let all_events = storage.load_all_events().unwrap();
2152 assert_eq!(all_events.len(), 4);
2153 }
2154
2155 #[test]
2160 fn test_write_atomic_parquet_emits_file_under_tenant_partition() {
2161 let temp_dir = TempDir::new().unwrap();
2166 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2167
2168 let events: Vec<Event> = (0..3)
2169 .map(|i| event_with_tenant("alice", &format!("a-{i}")))
2170 .collect();
2171
2172 let final_path = storage
2173 .write_atomic_parquet("alice", "snapshot.alice.range", &events)
2174 .unwrap();
2175
2176 let rel = final_path
2178 .strip_prefix(temp_dir.path())
2179 .unwrap()
2180 .to_string_lossy()
2181 .into_owned();
2182 let parts: Vec<&str> = rel.split(std::path::MAIN_SEPARATOR).collect();
2183 assert_eq!(parts.len(), 3, "expected tenant/yyyy-mm/file, got {rel}");
2184 assert_eq!(parts[0], "alice");
2185 assert_eq!(parts[2], "snapshot.alice.range.parquet");
2186
2187 assert!(final_path.is_file());
2189 let tmp = final_path.with_extension("parquet.tmp");
2190 assert!(
2191 !tmp.exists(),
2192 "tmp should have been renamed away; still at {}",
2193 tmp.display()
2194 );
2195
2196 let loaded = storage.load_events_for_tenant("alice").unwrap();
2198 assert_eq!(loaded.len(), 3);
2199 }
2200
2201 #[test]
2202 fn test_write_atomic_parquet_rejects_empty_events() {
2203 let temp_dir = TempDir::new().unwrap();
2204 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2205 let result = storage.write_atomic_parquet("alice", "snap", &[]);
2206 assert!(result.is_err());
2207 }
2208
2209 #[test]
2210 fn test_write_atomic_parquet_rejects_unsafe_tenant() {
2211 let temp_dir = TempDir::new().unwrap();
2212 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2213 let events = [event_with_tenant("alice", "e-0")];
2214 for unsafe_tid in ["..", "a/b", ""] {
2215 let result = storage.write_atomic_parquet(unsafe_tid, "snap", &events);
2216 assert!(
2217 result.is_err(),
2218 "unsafe tenant_id {unsafe_tid:?} should have been rejected"
2219 );
2220 }
2221 }
2222
2223 #[test]
2224 fn test_cleanup_partial_writes_removes_orphan_tmps() {
2225 let temp_dir = TempDir::new().unwrap();
2229 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2230
2231 for i in 0..2 {
2234 storage
2235 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2236 .unwrap();
2237 }
2238 storage.flush().unwrap();
2239 let real_files_before = find_parquet_files_recursive(temp_dir.path()).unwrap();
2240 assert_eq!(real_files_before.len(), 1);
2241
2242 let alice_subtree = temp_dir.path().join("alice");
2244 let orphan_dir = real_files_before[0].parent().unwrap();
2245 let orphan_path = orphan_dir.join("snapshot.alice.crashed.parquet.tmp");
2246 std::fs::write(&orphan_path, b"fake partial parquet").unwrap();
2247 assert!(orphan_path.is_file());
2248
2249 let nested_dir = alice_subtree.join("2099-01");
2251 std::fs::create_dir_all(&nested_dir).unwrap();
2252 let nested_orphan = nested_dir.join("events-x.parquet.tmp");
2253 std::fs::write(&nested_orphan, b"junk").unwrap();
2254
2255 let removed = storage.cleanup_partial_writes().unwrap();
2256 assert_eq!(removed, 2, "two orphan tmps should have been cleaned");
2257 assert!(!orphan_path.exists());
2258 assert!(!nested_orphan.exists());
2259
2260 let real_files_after = find_parquet_files_recursive(temp_dir.path()).unwrap();
2262 assert_eq!(real_files_after, real_files_before);
2263 }
2264
2265 #[test]
2266 fn test_cleanup_partial_writes_quarantines_zero_byte_parquet() {
2267 let temp_dir = TempDir::new().unwrap();
2272 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2273
2274 for i in 0..2 {
2276 storage
2277 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2278 .unwrap();
2279 }
2280 storage.flush().unwrap();
2281 let healthy_files_before = find_parquet_files_recursive(temp_dir.path()).unwrap();
2282 assert_eq!(healthy_files_before.len(), 1);
2283 let healthy_file = &healthy_files_before[0];
2284 let healthy_dir = healthy_file.parent().unwrap();
2285
2286 let bricked = healthy_dir.join("events-bricked-deadbeef.parquet");
2288 std::fs::write(&bricked, b"").unwrap();
2289 assert_eq!(std::fs::metadata(&bricked).unwrap().len(), 0);
2290
2291 let acted = storage.cleanup_partial_writes().unwrap();
2292 assert_eq!(acted, 1, "only the 0-byte file should have been acted on");
2293
2294 assert!(!bricked.exists(), "0-byte file should have been renamed");
2297 let quarantined: Vec<_> = std::fs::read_dir(healthy_dir)
2298 .unwrap()
2299 .flatten()
2300 .map(|e| e.path())
2301 .filter(|p| {
2302 p.file_name()
2303 .and_then(|n| n.to_str())
2304 .is_some_and(|n| n.starts_with("events-bricked-deadbeef.parquet.corrupt-"))
2305 })
2306 .collect();
2307 assert_eq!(
2308 quarantined.len(),
2309 1,
2310 "expected one .parquet.corrupt-<ts> sibling"
2311 );
2312
2313 assert!(
2315 healthy_file.exists(),
2316 "healthy parquet must not be molested"
2317 );
2318 }
2319
2320 #[test]
2321 fn test_load_all_events_skips_zero_byte_parquet() {
2322 let temp_dir = TempDir::new().unwrap();
2326 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2327
2328 for i in 0..2 {
2330 storage
2331 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2332 .unwrap();
2333 }
2334 storage.flush().unwrap();
2335
2336 let healthy_files = find_parquet_files_recursive(temp_dir.path()).unwrap();
2340 let bricked = healthy_files[0]
2341 .parent()
2342 .unwrap()
2343 .join("events-bricked-cafef00d.parquet");
2344 std::fs::write(&bricked, b"").unwrap();
2345
2346 let loaded = storage.load_all_events().unwrap();
2347 assert_eq!(
2348 loaded.len(),
2349 2,
2350 "all healthy events should still load despite the 0-byte file"
2351 );
2352 }
2353
2354 #[test]
2355 fn test_load_events_for_tenant_skips_zero_byte_parquet() {
2356 let temp_dir = TempDir::new().unwrap();
2359 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2360
2361 for i in 0..3 {
2362 storage
2363 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2364 .unwrap();
2365 }
2366 storage.flush().unwrap();
2367
2368 let healthy = find_parquet_files_recursive(temp_dir.path()).unwrap();
2369 let bricked = healthy[0]
2370 .parent()
2371 .unwrap()
2372 .join("events-bricked-feedface.parquet");
2373 std::fs::write(&bricked, b"").unwrap();
2374
2375 let events = storage.load_events_for_tenant("alice").unwrap();
2376 assert_eq!(events.len(), 3, "lazy load must skip the 0-byte file");
2377 }
2378
2379 #[test]
2380 fn test_flush_tenant_leaves_no_zero_byte_parquet_after_normal_flush() {
2381 let temp_dir = TempDir::new().unwrap();
2385 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2386
2387 for i in 0..5 {
2388 storage
2389 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2390 .unwrap();
2391 }
2392 storage.flush().unwrap();
2393
2394 let mut tmp_count = 0;
2395 let mut zero_byte_count = 0;
2396 let mut healthy_count = 0;
2397 let mut stack = vec![temp_dir.path().to_path_buf()];
2398 while let Some(d) = stack.pop() {
2399 for entry in std::fs::read_dir(&d).unwrap().flatten() {
2400 let p = entry.path();
2401 if p.is_dir() {
2402 stack.push(p);
2403 continue;
2404 }
2405 let name = p.file_name().unwrap().to_string_lossy().into_owned();
2406 if name.ends_with(".parquet.tmp") {
2407 tmp_count += 1;
2408 } else if name.ends_with(".parquet") {
2409 if std::fs::metadata(&p).unwrap().len() == 0 {
2410 zero_byte_count += 1;
2411 } else {
2412 healthy_count += 1;
2413 }
2414 }
2415 }
2416 }
2417 assert_eq!(tmp_count, 0, ".parquet.tmp survivors after flush");
2418 assert_eq!(zero_byte_count, 0, "0-byte .parquet survivors after flush");
2419 assert_eq!(healthy_count, 1, "expected exactly one healthy parquet");
2420 }
2421
2422 #[test]
2423 fn test_new_calls_cleanup_partial_writes_on_boot() {
2424 let temp_dir = TempDir::new().unwrap();
2428 let stale = temp_dir.path().join("orphan.parquet.tmp");
2429 std::fs::write(&stale, b"crash detritus").unwrap();
2430 assert!(stale.is_file());
2431
2432 let _storage = ParquetStorage::new(temp_dir.path()).unwrap();
2433 assert!(
2434 !stale.exists(),
2435 "stale tmp should have been cleaned by ParquetStorage::new"
2436 );
2437 }
2438
2439 fn seed_flat_layout_file(storage: &ParquetStorage, count: usize) -> PathBuf {
2449 for i in 0..count {
2450 storage
2451 .append_event(create_test_event(&format!("entity-{i}")))
2452 .unwrap();
2453 }
2454 storage.flush().unwrap();
2455
2456 let default_subtree = storage.storage_dir().join("default");
2461 let candidates = find_parquet_files_recursive(&default_subtree).unwrap();
2462 assert!(
2463 !candidates.is_empty(),
2464 "seed expected at least one file under default/"
2465 );
2466 let src = candidates.into_iter().max().unwrap();
2467
2468 let dst = storage.storage_dir().join(src.file_name().unwrap());
2469 std::fs::rename(&src, &dst).unwrap();
2470 if let Some(month_dir) = src.parent() {
2474 let _ = std::fs::remove_dir(month_dir);
2475 if let Some(tenant_dir) = month_dir.parent() {
2476 let _ = std::fs::remove_dir(tenant_dir);
2477 }
2478 }
2479 dst
2480 }
2481
2482 #[test]
2483 fn test_migrate_flat_layout_dry_run_touches_nothing() {
2484 let temp_dir = TempDir::new().unwrap();
2485 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2486 let flat = seed_flat_layout_file(&storage, 7);
2487 assert!(flat.is_file(), "test setup: flat file should exist");
2488
2489 let report = storage.migrate_flat_layout(true).unwrap();
2490 assert!(report.dry_run);
2491 assert_eq!(report.flat_files_seen, 1);
2492 assert_eq!(report.events_migrated, 7);
2493 assert_eq!(report.flat_files_removed, 0);
2494 assert_eq!(report.partitions_written, 0);
2495 assert!(
2496 flat.is_file(),
2497 "flat file must still be present after dry run"
2498 );
2499 }
2500
2501 #[test]
2502 fn test_migrate_flat_layout_moves_events_into_default_tree_and_removes_flat() {
2503 let temp_dir = TempDir::new().unwrap();
2504 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2505 let flat = seed_flat_layout_file(&storage, 5);
2506
2507 let report = storage.migrate_flat_layout(false).unwrap();
2508 assert!(!report.dry_run);
2509 assert_eq!(report.flat_files_seen, 1);
2510 assert_eq!(report.flat_files_removed, 1);
2511 assert_eq!(report.events_migrated, 5);
2512 assert!(report.partitions_written >= 1);
2513 assert!(
2514 !flat.exists(),
2515 "flat file should be deleted after migration"
2516 );
2517
2518 let post = find_parquet_files_recursive(temp_dir.path()).unwrap();
2519 assert!(
2520 post.iter().all(|p| {
2521 let rel = p
2522 .strip_prefix(temp_dir.path())
2523 .unwrap()
2524 .to_string_lossy()
2525 .into_owned();
2526 rel.starts_with(&format!("default{}", std::path::MAIN_SEPARATOR))
2527 }),
2528 "all migrated files should be under default/"
2529 );
2530
2531 let loaded = storage.load_all_events().unwrap();
2532 assert_eq!(loaded.len(), 5);
2533 }
2534
2535 #[test]
2536 fn test_migrate_flat_layout_is_idempotent_when_re_run_after_completion() {
2537 let temp_dir = TempDir::new().unwrap();
2538 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2539 let _flat = seed_flat_layout_file(&storage, 4);
2540
2541 let first = storage.migrate_flat_layout(false).unwrap();
2542 assert_eq!(first.events_migrated, 4);
2543
2544 let second = storage.migrate_flat_layout(false).unwrap();
2547 assert_eq!(second.flat_files_seen, 0);
2548 assert_eq!(second.events_migrated, 0);
2549 assert_eq!(second.flat_files_removed, 0);
2550
2551 let loaded = storage.load_all_events().unwrap();
2552 assert_eq!(loaded.len(), 4, "rerun must not duplicate or lose events");
2553 }
2554
2555 #[test]
2556 fn test_migrate_flat_layout_ignores_already_partitioned_data() {
2557 let temp_dir = TempDir::new().unwrap();
2560 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2561
2562 for i in 0..3 {
2563 storage
2564 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2565 .unwrap();
2566 }
2567 storage.flush().unwrap();
2568
2569 let _flat = seed_flat_layout_file(&storage, 2);
2570
2571 let report = storage.migrate_flat_layout(false).unwrap();
2572 assert_eq!(report.flat_files_seen, 1, "only the flat file is in scope");
2573 assert_eq!(report.events_migrated, 2);
2574
2575 let alice_files = find_parquet_files_recursive(&temp_dir.path().join("alice")).unwrap();
2576 assert_eq!(alice_files.len(), 1, "alice's tree must be untouched");
2577
2578 let loaded = storage.load_all_events().unwrap();
2579 assert_eq!(loaded.len(), 5);
2580 let alice_count = loaded
2581 .iter()
2582 .filter(|e| e.tenant_id_str() == "alice")
2583 .count();
2584 let default_count = loaded
2585 .iter()
2586 .filter(|e| e.tenant_id_str() == "default")
2587 .count();
2588 assert_eq!(alice_count, 3);
2589 assert_eq!(default_count, 2);
2590 }
2591
2592 #[test]
2593 fn test_migrate_flat_layout_with_no_flat_files_is_a_clean_noop() {
2594 let temp_dir = TempDir::new().unwrap();
2595 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2596 let report = storage.migrate_flat_layout(false).unwrap();
2597 assert_eq!(report.flat_files_seen, 0);
2598 assert_eq!(report.events_migrated, 0);
2599 }
2600}