1use super::archive_budget::ArchiveReadBudget;
2use crate::{
3 domain::entities::Event,
4 error::{AllSourceError, Result},
5};
6use arrow::{
7 array::{
8 Array, ArrayRef, StringBuilder, TimestampMicrosecondArray, TimestampMicrosecondBuilder,
9 UInt64Builder,
10 },
11 datatypes::{DataType, Field, Schema, TimeUnit},
12 record_batch::RecordBatch,
13};
14use parquet::{arrow::ArrowWriter, file::properties::WriterProperties};
15use std::{
16 collections::HashMap,
17 fs::{self, File},
18 path::{Path, PathBuf},
19 sync::{
20 Arc, Mutex,
21 atomic::{AtomicU64, Ordering},
22 },
23 time::{Duration, Instant},
24};
25
26pub const DEFAULT_BATCH_SIZE: usize = 10_000;
28
29pub const DEFAULT_FLUSH_TIMEOUT_MS: u64 = 5_000;
31
32#[derive(Debug, Clone)]
34pub struct ParquetStorageConfig {
35 pub batch_size: usize,
37 pub flush_timeout: Duration,
39 pub compression: parquet::basic::Compression,
41}
42
43impl Default for ParquetStorageConfig {
44 fn default() -> Self {
45 Self {
46 batch_size: DEFAULT_BATCH_SIZE,
47 flush_timeout: Duration::from_millis(DEFAULT_FLUSH_TIMEOUT_MS),
48 compression: parquet::basic::Compression::SNAPPY,
49 }
50 }
51}
52
53impl ParquetStorageConfig {
54 pub fn high_throughput() -> Self {
56 Self {
57 batch_size: 50_000,
58 flush_timeout: Duration::from_secs(10),
59 compression: parquet::basic::Compression::SNAPPY,
60 }
61 }
62
63 pub fn low_latency() -> Self {
65 Self {
66 batch_size: 1_000,
67 flush_timeout: Duration::from_secs(1),
68 compression: parquet::basic::Compression::SNAPPY,
69 }
70 }
71}
72
73#[derive(Debug, Clone, Default)]
75pub struct BatchWriteStats {
76 pub batches_written: u64,
78 pub events_written: u64,
80 pub bytes_written: u64,
82 pub avg_batch_size: f64,
84 pub events_per_sec: f64,
86 pub total_write_time_ns: u64,
88 pub timeout_flushes: u64,
90 pub size_flushes: u64,
92}
93
94#[derive(Debug, Clone)]
96pub struct BatchWriteResult {
97 pub events_written: usize,
99 pub batches_flushed: usize,
101 pub duration: Duration,
103 pub events_per_sec: f64,
105}
106
107pub struct ParquetStorage {
116 storage_dir: PathBuf,
118
119 current_batches: Mutex<HashMap<String, Vec<Event>>>,
128
129 config: ParquetStorageConfig,
131
132 schema: Arc<Schema>,
134
135 last_flush_time: Mutex<Instant>,
137
138 batches_written: AtomicU64,
140 events_written: AtomicU64,
141 bytes_written: AtomicU64,
142 total_write_time_ns: AtomicU64,
143 timeout_flushes: AtomicU64,
144 size_flushes: AtomicU64,
145}
146
147impl ParquetStorage {
148 pub fn new(storage_dir: impl AsRef<Path>) -> Result<Self> {
150 Self::with_config(storage_dir, ParquetStorageConfig::default())
151 }
152
153 pub fn with_config(
155 storage_dir: impl AsRef<Path>,
156 config: ParquetStorageConfig,
157 ) -> Result<Self> {
158 let storage_dir = storage_dir.as_ref().to_path_buf();
159
160 fs::create_dir_all(&storage_dir).map_err(|e| {
162 AllSourceError::StorageError(format!("Failed to create storage directory: {e}"))
163 })?;
164
165 let schema = Arc::new(Schema::new(vec![
167 Field::new("event_id", DataType::Utf8, false),
168 Field::new("event_type", DataType::Utf8, false),
169 Field::new("entity_id", DataType::Utf8, false),
170 Field::new("payload", DataType::Utf8, false),
171 Field::new(
172 "timestamp",
173 DataType::Timestamp(TimeUnit::Microsecond, None),
174 false,
175 ),
176 Field::new("metadata", DataType::Utf8, true),
177 Field::new("version", DataType::UInt64, false),
178 ]));
179
180 let storage = Self {
181 storage_dir,
182 current_batches: Mutex::new(HashMap::new()),
183 config,
184 schema,
185 last_flush_time: Mutex::new(Instant::now()),
186 batches_written: AtomicU64::new(0),
187 events_written: AtomicU64::new(0),
188 bytes_written: AtomicU64::new(0),
189 total_write_time_ns: AtomicU64::new(0),
190 timeout_flushes: AtomicU64::new(0),
191 size_flushes: AtomicU64::new(0),
192 };
193
194 match storage.cleanup_partial_writes() {
199 Ok(0) => {}
200 Ok(n) => tracing::warn!(
201 "cleanup_partial_writes acted on {n} crash-detritus file(s) on boot — \
202 see preceding logs for per-file detail"
203 ),
204 Err(e) => tracing::error!("cleanup_partial_writes failed on boot: {e}"),
205 }
206
207 Ok(storage)
208 }
209
210 #[deprecated(note = "Use new() or with_config() instead - default batch size is now 10,000")]
212 pub fn with_legacy_batch_size(storage_dir: impl AsRef<Path>) -> Result<Self> {
213 Self::with_config(
214 storage_dir,
215 ParquetStorageConfig {
216 batch_size: 1000,
217 ..Default::default()
218 },
219 )
220 }
221
222 #[cfg_attr(feature = "hotpath", hotpath::measure)]
232 pub fn append_event(&self, event: Event) -> Result<()> {
233 let tenant = event.tenant_id_str().to_string();
234 let should_flush_tenant = {
235 let mut batches = self.current_batches.lock().unwrap();
236 let entry = batches.entry(tenant.clone()).or_default();
237 entry.push(event);
238 entry.len() >= self.config.batch_size
239 };
240
241 if should_flush_tenant {
242 self.size_flushes.fetch_add(1, Ordering::Relaxed);
243 self.flush_tenant(&tenant)?;
244 }
245
246 Ok(())
247 }
248
249 #[cfg_attr(feature = "hotpath", hotpath::measure)]
255 pub fn batch_write(&self, events: Vec<Event>) -> Result<BatchWriteResult> {
256 let start = Instant::now();
257 let event_count = events.len();
258
259 let mut grouped: HashMap<String, Vec<Event>> = HashMap::new();
262 for event in events {
263 grouped
264 .entry(event.tenant_id_str().to_string())
265 .or_default()
266 .push(event);
267 }
268
269 let mut tenants_to_flush: Vec<String> = Vec::new();
270 {
271 let mut batches = self.current_batches.lock().unwrap();
272 for (tenant, mut new_events) in grouped {
273 let entry = batches.entry(tenant.clone()).or_default();
274 entry.append(&mut new_events);
275 if entry.len() >= self.config.batch_size {
276 tenants_to_flush.push(tenant);
277 }
278 }
279 }
280
281 let mut batches_flushed = 0;
282 for tenant in tenants_to_flush {
283 self.size_flushes.fetch_add(1, Ordering::Relaxed);
284 self.flush_tenant(&tenant)?;
285 batches_flushed += 1;
286 }
287
288 let duration = start.elapsed();
289
290 Ok(BatchWriteResult {
291 events_written: event_count,
292 batches_flushed,
293 duration,
294 events_per_sec: event_count as f64 / duration.as_secs_f64(),
295 })
296 }
297
298 #[cfg_attr(feature = "hotpath", hotpath::measure)]
306 pub fn check_timeout_flush(&self) -> Result<bool> {
307 let should_flush = {
308 let last_flush = self.last_flush_time.lock().unwrap();
309 let batches = self.current_batches.lock().unwrap();
310 let any_pending = batches.values().any(|v| !v.is_empty());
311 any_pending && last_flush.elapsed() >= self.config.flush_timeout
312 };
313
314 if should_flush {
315 self.timeout_flushes.fetch_add(1, Ordering::Relaxed);
316 self.flush()?;
317 Ok(true)
318 } else {
319 Ok(false)
320 }
321 }
322
323 #[cfg_attr(feature = "hotpath", hotpath::measure)]
330 pub fn flush(&self) -> Result<()> {
331 let tenants: Vec<String> = {
332 let batches = self.current_batches.lock().unwrap();
333 batches
334 .iter()
335 .filter(|(_, v)| !v.is_empty())
336 .map(|(k, _)| k.clone())
337 .collect()
338 };
339 if tenants.is_empty() {
340 return Ok(());
341 }
342 for tenant in tenants {
343 self.flush_tenant(&tenant)?;
344 }
345 Ok(())
346 }
347
348 pub(crate) fn has_pending_tenant_events(&self, tenant_id: &str) -> bool {
351 self.current_batches
352 .lock()
353 .unwrap()
354 .get(tenant_id)
355 .is_some_and(|events| !events.is_empty())
356 }
357
358 fn flush_tenant(&self, tenant_id: &str) -> Result<()> {
373 let events_to_write = {
374 let mut batches = self.current_batches.lock().unwrap();
375 match batches.get_mut(tenant_id) {
376 Some(v) if !v.is_empty() => std::mem::take(v),
377 _ => return Ok(()),
378 }
379 };
380
381 let batch_count = events_to_write.len();
382 let start = Instant::now();
383
384 let written = self.write_tenant_events(tenant_id, &events_to_write);
385 let (file_path, file_metadata) = match written {
386 Ok(written) => written,
387 Err(e) => {
388 let mut batches = self.current_batches.lock().unwrap();
391 let pending = batches.entry(tenant_id.to_string()).or_default();
392 let arrived_during_write = std::mem::replace(pending, events_to_write);
393 pending.extend(arrived_during_write);
394 return Err(e);
395 }
396 };
397
398 let duration = start.elapsed();
399
400 self.batches_written.fetch_add(1, Ordering::Relaxed);
401 self.events_written
402 .fetch_add(batch_count as u64, Ordering::Relaxed);
403 if let Some(size) = file_metadata
404 .row_groups()
405 .first()
406 .map(parquet::file::metadata::RowGroupMetaData::total_byte_size)
407 {
408 self.bytes_written.fetch_add(size as u64, Ordering::Relaxed);
409 }
410 self.total_write_time_ns
411 .fetch_add(duration.as_nanos() as u64, Ordering::Relaxed);
412
413 {
414 let mut last_flush = self.last_flush_time.lock().unwrap();
415 *last_flush = Instant::now();
416 }
417
418 tracing::info!(
419 "Wrote {} events for tenant={} to {} in {:?}",
420 batch_count,
421 tenant_id,
422 file_path.display(),
423 duration
424 );
425
426 Ok(())
427 }
428
429 fn write_tenant_events(
430 &self,
431 tenant_id: &str,
432 events: &[Event],
433 ) -> Result<(PathBuf, parquet::file::metadata::ParquetMetaData)> {
434 let record_batch = self.events_to_record_batch(events)?;
435
436 let now = chrono::Utc::now();
437 let partition_dir = partition_path_for_tenant(&self.storage_dir, tenant_id, now)?;
438 fs::create_dir_all(&partition_dir).map_err(|e| {
439 AllSourceError::StorageError(format!(
440 "Failed to create tenant partition {}: {e}",
441 partition_dir.display()
442 ))
443 })?;
444 let file_stem = format!(
445 "events-{}-{}",
446 now.format("%Y%m%d-%H%M%S%3f"),
447 uuid::Uuid::new_v4().as_simple()
448 );
449
450 tracing::info!(
451 "Flushing {} events for tenant={} to {}/{}.parquet",
452 events.len(),
453 tenant_id,
454 partition_dir.display(),
455 file_stem
456 );
457
458 self.write_record_batch_atomic(&partition_dir, &file_stem, &record_batch)
459 }
460
461 fn write_record_batch_atomic(
478 &self,
479 partition_dir: &Path,
480 file_stem: &str,
481 record_batch: &RecordBatch,
482 ) -> Result<(PathBuf, parquet::file::metadata::ParquetMetaData)> {
483 let final_path = partition_dir.join(format!("{file_stem}.parquet"));
484 let tmp_path = partition_dir.join(format!("{file_stem}.parquet.tmp"));
485
486 let metadata = {
488 let file = File::create(&tmp_path).map_err(|e| {
489 AllSourceError::StorageError(format!(
490 "Failed to create parquet tmp file {}: {e}",
491 tmp_path.display()
492 ))
493 })?;
494
495 let props = WriterProperties::builder()
496 .set_compression(self.config.compression)
497 .build();
498
499 let mut writer = ArrowWriter::try_new(file, self.schema.clone(), Some(props))?;
500 writer.write(record_batch)?;
501 writer.close()?
504 };
505
506 let tmp_file = File::open(&tmp_path).map_err(|e| {
508 AllSourceError::StorageError(format!(
509 "Failed to reopen parquet tmp for fsync {}: {e}",
510 tmp_path.display()
511 ))
512 })?;
513 tmp_file.sync_all().map_err(|e| {
514 AllSourceError::StorageError(format!("fsync on parquet tmp failed: {e}"))
515 })?;
516 drop(tmp_file);
517
518 fs::rename(&tmp_path, &final_path).map_err(|e| {
520 AllSourceError::StorageError(format!(
521 "Failed to rename {} → {}: {e}",
522 tmp_path.display(),
523 final_path.display()
524 ))
525 })?;
526
527 if let Ok(dir) = File::open(partition_dir) {
531 let _ = dir.sync_all();
532 }
533
534 Ok((final_path, metadata))
535 }
536
537 pub fn write_atomic_parquet(
561 &self,
562 tenant_id: &str,
563 file_stem: &str,
564 events: &[Event],
565 ) -> Result<PathBuf> {
566 if events.is_empty() {
567 return Err(AllSourceError::StorageError(
568 "write_atomic_parquet called with empty event slice".to_string(),
569 ));
570 }
571 let anchor_ts = events
576 .iter()
577 .map(|e| e.timestamp)
578 .min()
579 .unwrap_or_else(chrono::Utc::now);
580 let partition_dir = partition_path_for_tenant(&self.storage_dir, tenant_id, anchor_ts)?;
581 fs::create_dir_all(&partition_dir).map_err(|e| {
582 AllSourceError::StorageError(format!(
583 "Failed to create tenant partition {}: {e}",
584 partition_dir.display()
585 ))
586 })?;
587
588 let record_batch = self.events_to_record_batch(events)?;
589 let (final_path, _meta) =
590 self.write_record_batch_atomic(&partition_dir, file_stem, &record_batch)?;
591
592 tracing::info!(
593 tenant_id = tenant_id,
594 file = %final_path.display(),
595 event_count = events.len(),
596 "wrote atomic snapshot file"
597 );
598
599 Ok(final_path)
600 }
601
602 pub fn cleanup_partial_writes(&self) -> Result<usize> {
625 let mut acted = 0usize;
626 let mut stack: Vec<PathBuf> = vec![self.storage_dir.clone()];
627 while let Some(dir) = stack.pop() {
628 let Ok(entries) = fs::read_dir(&dir) else {
629 continue;
630 };
631 for entry in entries.flatten() {
632 let path = entry.path();
633 let Ok(ft) = entry.file_type() else { continue };
634 if ft.is_dir() {
635 stack.push(path);
636 continue;
637 }
638 if !ft.is_file() {
639 continue;
640 }
641 let path_str = path.to_string_lossy();
642 if path_str.ends_with(".parquet.tmp") {
643 match fs::remove_file(&path) {
644 Ok(()) => {
645 tracing::warn!(
646 file = %path.display(),
647 "cleaned up orphan snapshot tmp file (crash recovery)"
648 );
649 acted += 1;
650 }
651 Err(e) => {
652 tracing::error!(
653 file = %path.display(),
654 "failed to remove orphan snapshot tmp file: {e}"
655 );
656 }
657 }
658 } else if path_str.ends_with(".parquet")
659 && fs::metadata(&path).is_ok_and(|m| m.len() == 0)
660 {
661 let ts = chrono::Utc::now().timestamp();
665 let quarantine_path = path.with_extension(format!("parquet.corrupt-{ts}"));
666 match fs::rename(&path, &quarantine_path) {
667 Ok(()) => {
668 tracing::error!(
669 from = %path.display(),
670 to = %quarantine_path.display(),
671 "quarantined 0-byte parquet file (issue #166 pre-fix crash). \
672 Operator: inspect and rm if not needed."
673 );
674 acted += 1;
675 }
676 Err(e) => {
677 tracing::error!(
678 file = %path.display(),
679 "failed to quarantine 0-byte parquet file: {e}"
680 );
681 }
682 }
683 }
684 }
685 }
686 Ok(acted)
687 }
688
689 #[cfg_attr(feature = "hotpath", hotpath::measure)]
695 pub fn flush_on_shutdown(&self) -> Result<usize> {
696 let total_pending: usize = {
697 let batches = self.current_batches.lock().unwrap();
698 batches.values().map(Vec::len).sum()
699 };
700
701 if total_pending > 0 {
702 tracing::info!(
703 "Shutdown: flushing {} pending events across all tenants",
704 total_pending
705 );
706 self.flush()?;
707 }
708
709 Ok(total_pending)
710 }
711
712 pub fn batch_stats(&self) -> BatchWriteStats {
714 let batches = self.batches_written.load(Ordering::Relaxed);
715 let events = self.events_written.load(Ordering::Relaxed);
716 let bytes = self.bytes_written.load(Ordering::Relaxed);
717 let time_ns = self.total_write_time_ns.load(Ordering::Relaxed);
718
719 let time_secs = time_ns as f64 / 1_000_000_000.0;
720
721 BatchWriteStats {
722 batches_written: batches,
723 events_written: events,
724 bytes_written: bytes,
725 avg_batch_size: if batches > 0 {
726 events as f64 / batches as f64
727 } else {
728 0.0
729 },
730 events_per_sec: if time_secs > 0.0 {
731 events as f64 / time_secs
732 } else {
733 0.0
734 },
735 total_write_time_ns: time_ns,
736 timeout_flushes: self.timeout_flushes.load(Ordering::Relaxed),
737 size_flushes: self.size_flushes.load(Ordering::Relaxed),
738 }
739 }
740
741 pub fn pending_count(&self) -> usize {
743 self.current_batches
744 .lock()
745 .unwrap()
746 .values()
747 .map(Vec::len)
748 .sum()
749 }
750
751 pub fn batch_size(&self) -> usize {
753 self.config.batch_size
754 }
755
756 pub fn flush_timeout(&self) -> Duration {
758 self.config.flush_timeout
759 }
760
761 #[cfg_attr(feature = "hotpath", hotpath::measure)]
763 fn events_to_record_batch(&self, events: &[Event]) -> Result<RecordBatch> {
764 let mut event_id_builder = StringBuilder::new();
765 let mut event_type_builder = StringBuilder::new();
766 let mut entity_id_builder = StringBuilder::new();
767 let mut payload_builder = StringBuilder::new();
768 let mut timestamp_builder = TimestampMicrosecondBuilder::new();
769 let mut metadata_builder = StringBuilder::new();
770 let mut version_builder = UInt64Builder::new();
771
772 for event in events {
773 event_id_builder.append_value(event.id.to_string());
774 event_type_builder.append_value(event.event_type_str());
775 entity_id_builder.append_value(event.entity_id_str());
776 payload_builder.append_value(serde_json::to_string(&event.payload)?);
777
778 let timestamp_micros = event.timestamp.timestamp_micros();
780 timestamp_builder.append_value(timestamp_micros);
781
782 if let Some(ref metadata) = event.metadata {
783 metadata_builder.append_value(serde_json::to_string(metadata)?);
784 } else {
785 metadata_builder.append_null();
786 }
787
788 version_builder.append_value(event.version as u64);
789 }
790
791 let arrays: Vec<ArrayRef> = vec![
792 Arc::new(event_id_builder.finish()),
793 Arc::new(event_type_builder.finish()),
794 Arc::new(entity_id_builder.finish()),
795 Arc::new(payload_builder.finish()),
796 Arc::new(timestamp_builder.finish()),
797 Arc::new(metadata_builder.finish()),
798 Arc::new(version_builder.finish()),
799 ];
800
801 let record_batch = RecordBatch::try_new(self.schema.clone(), arrays)?;
802
803 Ok(record_batch)
804 }
805
806 #[cfg_attr(feature = "hotpath", hotpath::measure)]
813 pub fn load_all_events(&self) -> Result<Vec<Event>> {
814 let parquet_files = find_parquet_files_recursive(&self.storage_dir)?;
815
816 let mut all_events = Vec::with_capacity(parquet_files.len() * self.config.batch_size);
817 let mut skipped = 0usize;
818 for file_path in parquet_files {
819 tracing::info!("Loading events from {}", file_path.display());
820 let tenant_id = tenant_id_from_path(&self.storage_dir, &file_path);
821 match self.load_events_from_file(&file_path, &tenant_id) {
826 Ok(file_events) => all_events.extend(file_events),
827 Err(e) => {
828 tracing::error!(
829 file = %file_path.display(),
830 error = %e,
831 "Skipping unreadable parquet file — other files will still load. \
832 Likely a 0-byte or truncated file from an unclean shutdown; \
833 inspect and remove manually after confirming."
834 );
835 skipped += 1;
836 }
837 }
838 }
839
840 if skipped > 0 {
841 tracing::warn!(
842 "Loaded {} events from storage; skipped {} unreadable file(s)",
843 all_events.len(),
844 skipped
845 );
846 } else {
847 tracing::info!("Loaded {} total events from storage", all_events.len());
848 }
849
850 Ok(all_events)
851 }
852
853 #[cfg_attr(feature = "hotpath", hotpath::measure)]
859 fn load_events_from_file(&self, file_path: &Path, tenant_id: &str) -> Result<Vec<Event>> {
860 self.load_events_from_file_with_budget(file_path, tenant_id, None)
861 }
862
863 fn load_events_from_file_with_budget(
864 &self,
865 file_path: &Path,
866 tenant_id: &str,
867 mut budget: Option<&mut ArchiveReadBudget>,
868 ) -> Result<Vec<Event>> {
869 use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
870
871 if let Some(budget) = budget.as_deref() {
872 budget.check()?;
873 }
874 let file = File::open(file_path).map_err(|e| {
875 AllSourceError::StorageError(format!("Failed to open parquet file: {e}"))
876 })?;
877
878 if let Some(budget) = budget.as_deref_mut() {
879 let metadata = file.metadata().map_err(|error| {
880 AllSourceError::StorageError(format!("Failed to inspect archive file: {error}"))
881 })?;
882 budget.compressed(metadata.len())?;
883 }
884 let mut builder = ParquetRecordBatchReaderBuilder::try_new(file)?;
885 if let Some(budget) = budget.as_deref_mut() {
886 for group in builder.metadata().row_groups() {
887 budget.decoded_metadata(group.num_rows(), group.total_byte_size())?;
888 }
889 builder = builder.with_batch_size(256);
890 }
891 let reader = builder.build()?;
892
893 let mut events = Vec::new();
894
895 for batch in reader {
896 if let Some(budget) = budget.as_deref() {
897 budget.check()?;
898 }
899 let batch_events = self.record_batch_to_events(&batch?, tenant_id)?;
900 events.extend(batch_events);
901 }
902
903 if let Some(budget) = budget {
904 budget.check()?;
905 }
906
907 Ok(events)
908 }
909
910 pub fn load_events_from_file_path(
916 &self,
917 file_path: &Path,
918 tenant_id: &str,
919 ) -> Result<Vec<Event>> {
920 self.load_events_from_file(file_path, tenant_id)
921 }
922
923 #[cfg_attr(feature = "hotpath", hotpath::measure)]
927 fn record_batch_to_events(&self, batch: &RecordBatch, tenant_id: &str) -> Result<Vec<Event>> {
928 let event_ids = batch
929 .column(0)
930 .as_any()
931 .downcast_ref::<arrow::array::StringArray>()
932 .ok_or_else(|| AllSourceError::StorageError("Invalid event_id column".to_string()))?;
933
934 let event_types = batch
935 .column(1)
936 .as_any()
937 .downcast_ref::<arrow::array::StringArray>()
938 .ok_or_else(|| AllSourceError::StorageError("Invalid event_type column".to_string()))?;
939
940 let entity_ids = batch
941 .column(2)
942 .as_any()
943 .downcast_ref::<arrow::array::StringArray>()
944 .ok_or_else(|| AllSourceError::StorageError("Invalid entity_id column".to_string()))?;
945
946 let payloads = batch
947 .column(3)
948 .as_any()
949 .downcast_ref::<arrow::array::StringArray>()
950 .ok_or_else(|| AllSourceError::StorageError("Invalid payload column".to_string()))?;
951
952 let timestamps = batch
953 .column(4)
954 .as_any()
955 .downcast_ref::<TimestampMicrosecondArray>()
956 .ok_or_else(|| AllSourceError::StorageError("Invalid timestamp column".to_string()))?;
957
958 let metadatas = batch
959 .column(5)
960 .as_any()
961 .downcast_ref::<arrow::array::StringArray>()
962 .ok_or_else(|| AllSourceError::StorageError("Invalid metadata column".to_string()))?;
963
964 let versions = batch
965 .column(6)
966 .as_any()
967 .downcast_ref::<arrow::array::UInt64Array>()
968 .ok_or_else(|| AllSourceError::StorageError("Invalid version column".to_string()))?;
969
970 let mut events = Vec::new();
971
972 for i in 0..batch.num_rows() {
973 let id = uuid::Uuid::parse_str(event_ids.value(i))
974 .map_err(|e| AllSourceError::StorageError(format!("Invalid UUID: {e}")))?;
975
976 let timestamp = chrono::DateTime::from_timestamp_micros(timestamps.value(i))
977 .ok_or_else(|| AllSourceError::StorageError("Invalid timestamp".to_string()))?;
978
979 let metadata = if metadatas.is_null(i) {
980 None
981 } else {
982 Some(serde_json::from_str(metadatas.value(i))?)
983 };
984
985 let event = Event::reconstruct_from_strings(
986 id,
987 event_types.value(i).to_string(),
988 entity_ids.value(i).to_string(),
989 tenant_id.to_string(),
990 serde_json::from_str(payloads.value(i))?,
991 timestamp,
992 metadata,
993 versions.value(i) as i64,
994 );
995
996 events.push(event);
997 }
998
999 Ok(events)
1000 }
1001
1002 pub fn list_parquet_files(&self) -> Result<Vec<PathBuf>> {
1008 find_parquet_files_recursive(&self.storage_dir)
1009 }
1010
1011 pub fn tenant_id_for_file(&self, file_path: &Path) -> String {
1013 tenant_id_from_path(&self.storage_dir, file_path)
1014 }
1015
1016 pub fn list_parquet_files_for_tenant(&self, tenant_id: &str) -> Result<Vec<PathBuf>> {
1029 let safe = sanitize_tenant_id_for_path(tenant_id)?;
1030 let tenant_root = self.storage_dir.join(safe);
1031 if !tenant_root.is_dir() {
1032 return Ok(Vec::new());
1033 }
1034 find_parquet_files_recursive(&tenant_root)
1035 }
1036
1037 fn list_complete_tenant_archive(
1038 &self,
1039 tenant_id: &str,
1040 budget: &mut ArchiveReadBudget,
1041 ) -> Result<Vec<PathBuf>> {
1042 budget.check()?;
1043 let safe = sanitize_tenant_id_for_path(tenant_id)?;
1044 let root = self.storage_dir.join(safe);
1045 match fs::symlink_metadata(&root) {
1046 Ok(metadata) if metadata.is_dir() => {
1047 find_parquet_files_recursive_with_integrity(&root, Some(budget))
1048 }
1049 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(Vec::new()),
1050 _ => Err(AllSourceError::StorageError(
1051 "Cannot verify conditional version: tenant archive is not an accessible directory"
1052 .into(),
1053 )),
1054 }
1055 }
1056
1057 #[cfg_attr(feature = "hotpath", hotpath::measure)]
1076 pub fn load_events_for_tenant(&self, tenant_id: &str) -> Result<Vec<Event>> {
1077 self.load_events_for_tenant_with_integrity(tenant_id, None)
1078 .map(|(events, _complete)| events)
1079 }
1080
1081 pub(crate) fn load_events_for_tenant_with_integrity(
1082 &self,
1083 tenant_id: &str,
1084 mut budget: Option<&mut ArchiveReadBudget>,
1085 ) -> Result<(Vec<Event>, bool)> {
1086 let require_complete = budget.is_some();
1087 let parquet_files = if let Some(budget) = budget.as_deref_mut() {
1088 self.list_complete_tenant_archive(tenant_id, budget)?
1089 } else {
1090 self.list_parquet_files_for_tenant(tenant_id)?
1091 };
1092 tracing::info!(
1093 tenant_id = tenant_id,
1094 file_count = parquet_files.len(),
1095 "load_events_for_tenant: walking tenant subtree only"
1096 );
1097
1098 let mut events = Vec::new();
1099 let mut skipped = 0usize;
1100 for file_path in parquet_files {
1101 tracing::debug!(
1102 tenant_id = tenant_id,
1103 file = %file_path.display(),
1104 "load_events_for_tenant: opening file"
1105 );
1106 match self.load_events_from_file_with_budget(
1111 &file_path,
1112 tenant_id,
1113 budget.as_deref_mut(),
1114 ) {
1115 Ok(file_events) => events.extend(file_events),
1116 Err(e) => {
1117 tracing::error!(
1118 tenant_id = tenant_id,
1119 file = %file_path.display(),
1120 error = %e,
1121 require_complete,
1122 "Unreadable parquet file in tenant subtree"
1123 );
1124 skipped += 1;
1125 if require_complete {
1126 if let Some(budget) = budget.as_deref() {
1127 budget.check()?;
1128 }
1129 return Err(AllSourceError::StorageError(format!(
1130 "Cannot verify conditional version from incomplete archive: {e}"
1131 )));
1132 }
1133 }
1134 }
1135 }
1136
1137 tracing::info!(
1138 tenant_id = tenant_id,
1139 event_count = events.len(),
1140 skipped_files = skipped,
1141 "load_events_for_tenant: complete"
1142 );
1143 Ok((events, require_complete && skipped == 0))
1146 }
1147
1148 pub fn storage_dir(&self) -> &Path {
1150 &self.storage_dir
1151 }
1152
1153 pub fn migrate_flat_layout(&self, dry_run: bool) -> Result<MigrationReport> {
1176 let flat_files = list_flat_layout_files(&self.storage_dir)?;
1177 let mut report = MigrationReport {
1178 dry_run,
1179 ..Default::default()
1180 };
1181
1182 for flat_file in flat_files {
1183 let events = self.load_events_from_file(&flat_file, "default")?;
1187 report.flat_files_seen += 1;
1188
1189 if events.is_empty() {
1190 if !dry_run {
1192 fs::remove_file(&flat_file).map_err(|e| {
1193 AllSourceError::StorageError(format!(
1194 "Failed to remove empty flat file {}: {e}",
1195 flat_file.display()
1196 ))
1197 })?;
1198 }
1199 report.flat_files_removed += 1;
1200 continue;
1201 }
1202
1203 let mut groups: HashMap<(String, String), Vec<Event>> = HashMap::new();
1209 for event in events {
1210 let key = (
1211 event.tenant_id_str().to_string(),
1212 event.timestamp().format("%Y-%m").to_string(),
1213 );
1214 groups.entry(key).or_default().push(event);
1215 }
1216
1217 for ((tenant, yyyy_mm), group_events) in groups {
1218 let count = group_events.len();
1219 if !dry_run {
1220 let safe_tenant = sanitize_tenant_id_for_path(&tenant)?;
1221 let target_dir = self.storage_dir.join(safe_tenant).join(&yyyy_mm);
1222 fs::create_dir_all(&target_dir).map_err(|e| {
1223 AllSourceError::StorageError(format!(
1224 "Failed to create partition {}: {e}",
1225 target_dir.display()
1226 ))
1227 })?;
1228 let filename = format!(
1229 "events-{}-{}.parquet",
1230 chrono::Utc::now().format("%Y%m%d-%H%M%S%3f"),
1231 uuid::Uuid::new_v4().as_simple()
1232 );
1233 let target_path = target_dir.join(&filename);
1234 let record_batch = self.events_to_record_batch(&group_events)?;
1235 let file = File::create(&target_path).map_err(|e| {
1236 AllSourceError::StorageError(format!(
1237 "Failed to create migration target {}: {e}",
1238 target_path.display()
1239 ))
1240 })?;
1241 let props = WriterProperties::builder()
1242 .set_compression(self.config.compression)
1243 .build();
1244 let mut writer = ArrowWriter::try_new(file, self.schema.clone(), Some(props))?;
1245 writer.write(&record_batch)?;
1246 writer.close()?;
1247 report.partitions_written += 1;
1248 }
1249 report.events_migrated += count;
1250 }
1251
1252 if !dry_run {
1253 fs::remove_file(&flat_file).map_err(|e| {
1254 AllSourceError::StorageError(format!(
1255 "Failed to remove flat file {} after migration: {e}",
1256 flat_file.display()
1257 ))
1258 })?;
1259 report.flat_files_removed += 1;
1260 }
1261 }
1262
1263 Ok(report)
1264 }
1265
1266 pub fn stats(&self) -> Result<StorageStats> {
1268 let parquet_files = find_parquet_files_recursive(&self.storage_dir)?;
1269 let mut total_size_bytes = 0u64;
1270 for path in &parquet_files {
1271 if let Ok(metadata) = fs::metadata(path) {
1272 total_size_bytes += metadata.len();
1273 }
1274 }
1275
1276 let current_batch_size: usize = self
1277 .current_batches
1278 .lock()
1279 .unwrap()
1280 .values()
1281 .map(Vec::len)
1282 .sum();
1283
1284 Ok(StorageStats {
1285 total_files: parquet_files.len(),
1286 total_size_bytes,
1287 storage_dir: self.storage_dir.clone(),
1288 current_batch_size,
1289 })
1290 }
1291}
1292
1293fn sanitize_tenant_id_for_path(tenant_id: &str) -> Result<&str> {
1306 if tenant_id.is_empty() {
1307 return Err(AllSourceError::StorageError(
1308 "tenant_id is empty (cannot derive partition path)".to_string(),
1309 ));
1310 }
1311 if tenant_id.len() > 128 {
1312 return Err(AllSourceError::StorageError(format!(
1313 "tenant_id is too long for partition path: {} bytes (max 128)",
1314 tenant_id.len()
1315 )));
1316 }
1317 if tenant_id == "." || tenant_id == ".." {
1318 return Err(AllSourceError::StorageError(format!(
1319 "tenant_id {tenant_id:?} is reserved"
1320 )));
1321 }
1322 for c in tenant_id.chars() {
1323 let ok = c.is_ascii_alphanumeric() || c == '-' || c == '_' || c == '.';
1324 if !ok {
1325 return Err(AllSourceError::StorageError(format!(
1326 "tenant_id {tenant_id:?} contains disallowed character {c:?} for partition path"
1327 )));
1328 }
1329 }
1330 Ok(tenant_id)
1331}
1332
1333fn partition_path_for_tenant(
1338 root: &Path,
1339 tenant_id: &str,
1340 when: chrono::DateTime<chrono::Utc>,
1341) -> Result<PathBuf> {
1342 let safe = sanitize_tenant_id_for_path(tenant_id)?;
1343 Ok(root.join(safe).join(when.format("%Y-%m").to_string()))
1344}
1345
1346fn tenant_id_from_path(root: &Path, file_path: &Path) -> String {
1356 let Ok(rel) = file_path.strip_prefix(root) else {
1357 return "default".to_string();
1358 };
1359 let mut comps = rel.components();
1360 let first = comps.next();
1361 let next = comps.next();
1362 match (first, next) {
1363 (Some(std::path::Component::Normal(tenant)), Some(_)) => {
1365 tenant.to_string_lossy().into_owned()
1366 }
1367 _ => "default".to_string(),
1369 }
1370}
1371
1372fn list_flat_layout_files(root: &Path) -> Result<Vec<PathBuf>> {
1378 let entries = fs::read_dir(root).map_err(|e| {
1379 AllSourceError::StorageError(format!("Failed to read storage directory: {e}"))
1380 })?;
1381 let mut out: Vec<PathBuf> = entries
1382 .filter_map(std::result::Result::ok)
1383 .filter_map(|entry| {
1384 let ft = entry.file_type().ok()?;
1385 if !ft.is_file() {
1386 return None;
1387 }
1388 let path = entry.path();
1389 if path.extension().and_then(|s| s.to_str()) == Some("parquet") {
1390 Some(path)
1391 } else {
1392 None
1393 }
1394 })
1395 .collect();
1396 out.sort();
1397 Ok(out)
1398}
1399
1400fn find_parquet_files_recursive(root: &Path) -> Result<Vec<PathBuf>> {
1411 find_parquet_files_recursive_with_integrity(root, None)
1412}
1413
1414fn find_parquet_files_recursive_with_integrity(
1415 root: &Path,
1416 mut budget: Option<&mut ArchiveReadBudget>,
1417) -> Result<Vec<PathBuf>> {
1418 let require_complete = budget.is_some();
1419 let mut out = Vec::new();
1420 let mut stack: Vec<PathBuf> = vec![root.to_path_buf()];
1421
1422 while let Some(dir) = stack.pop() {
1423 if let Some(budget) = budget.as_deref() {
1424 budget.check()?;
1425 }
1426 let entries = match fs::read_dir(&dir) {
1427 Ok(e) => e,
1428 Err(e) if dir == root || require_complete => {
1432 return Err(AllSourceError::StorageError(format!(
1433 "Failed to read storage directory: {e}"
1434 )));
1435 }
1436 Err(_) => continue,
1437 };
1438
1439 for entry in entries {
1440 if let Some(budget) = budget.as_deref_mut() {
1441 budget.entry()?;
1442 }
1443 let entry = match entry {
1444 Ok(entry) => entry,
1445 Err(error) if require_complete => {
1446 return Err(AllSourceError::StorageError(format!(
1447 "Failed to enumerate complete archive: {error}"
1448 )));
1449 }
1450 Err(_) => continue,
1451 };
1452 let path = entry.path();
1453 let ft = match entry.file_type() {
1457 Ok(file_type) => file_type,
1458 Err(error) if require_complete => {
1459 return Err(AllSourceError::StorageError(format!(
1460 "Failed to inspect complete archive entry: {error}"
1461 )));
1462 }
1463 Err(_) => continue,
1464 };
1465 if ft.is_dir() {
1466 stack.push(path);
1467 } else if ft.is_file()
1468 && path
1469 .extension()
1470 .and_then(|ext| ext.to_str())
1471 .is_some_and(|ext| ext == "parquet")
1472 {
1473 if let Some(budget) = budget.as_deref_mut() {
1474 budget.file()?;
1475 }
1476 out.push(path);
1477 }
1478 }
1479 }
1480
1481 out.sort();
1482 if let Some(budget) = budget {
1483 budget.check()?;
1484 }
1485 Ok(out)
1486}
1487
1488impl Drop for ParquetStorage {
1489 fn drop(&mut self) {
1490 if let Err(e) = self.flush_on_shutdown() {
1492 tracing::error!("Failed to flush events on drop: {}", e);
1493 }
1494 }
1495}
1496
1497#[derive(Debug, Default, Clone, serde::Serialize)]
1499pub struct MigrationReport {
1500 pub dry_run: bool,
1502 pub flat_files_seen: usize,
1504 pub flat_files_removed: usize,
1506 pub partitions_written: usize,
1508 pub events_migrated: usize,
1510}
1511
1512#[derive(Debug, serde::Serialize)]
1513pub struct StorageStats {
1514 pub total_files: usize,
1515 pub total_size_bytes: u64,
1516 pub storage_dir: PathBuf,
1517 pub current_batch_size: usize,
1518}
1519
1520#[cfg(test)]
1521mod tests {
1522 use super::*;
1523 use serde_json::json;
1524 use std::sync::Arc;
1525 use tempfile::TempDir;
1526
1527 fn create_test_event(entity_id: &str) -> Event {
1528 Event::reconstruct_from_strings(
1529 uuid::Uuid::new_v4(),
1530 "test.event".to_string(),
1531 entity_id.to_string(),
1532 "default".to_string(),
1533 json!({
1534 "test": "data",
1535 "value": 42
1536 }),
1537 chrono::Utc::now(),
1538 None,
1539 1,
1540 )
1541 }
1542
1543 #[test]
1544 fn test_parquet_storage_write_read() {
1545 let temp_dir = TempDir::new().unwrap();
1546 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1547
1548 for i in 0..10 {
1550 let event = create_test_event(&format!("entity-{i}"));
1551 storage.append_event(event).unwrap();
1552 }
1553
1554 storage.flush().unwrap();
1556
1557 let loaded_events = storage.load_all_events().unwrap();
1559 assert_eq!(loaded_events.len(), 10);
1560 }
1561
1562 #[test]
1563 fn test_storage_stats() {
1564 let temp_dir = TempDir::new().unwrap();
1565 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1566
1567 for i in 0..5 {
1569 storage
1570 .append_event(create_test_event(&format!("entity-{i}")))
1571 .unwrap();
1572 }
1573 storage.flush().unwrap();
1574
1575 let stats = storage.stats().unwrap();
1576 assert_eq!(stats.total_files, 1);
1577 assert!(stats.total_size_bytes > 0);
1578 }
1579
1580 #[test]
1581 fn test_default_batch_size() {
1582 let temp_dir = TempDir::new().unwrap();
1583 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1584
1585 assert_eq!(storage.batch_size(), DEFAULT_BATCH_SIZE);
1587 assert_eq!(storage.batch_size(), 10_000);
1588 }
1589
1590 #[test]
1591 fn test_custom_config() {
1592 let temp_dir = TempDir::new().unwrap();
1593 let config = ParquetStorageConfig {
1594 batch_size: 5_000,
1595 flush_timeout: Duration::from_secs(2),
1596 compression: parquet::basic::Compression::SNAPPY,
1597 };
1598 let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1599
1600 assert_eq!(storage.batch_size(), 5_000);
1601 assert_eq!(storage.flush_timeout(), Duration::from_secs(2));
1602 }
1603
1604 #[test]
1605 fn test_batch_write() {
1606 let temp_dir = TempDir::new().unwrap();
1607 let config = ParquetStorageConfig {
1608 batch_size: 100, ..Default::default()
1610 };
1611 let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1612
1613 let events: Vec<Event> = (0..250)
1621 .map(|i| create_test_event(&format!("entity-{i}")))
1622 .collect();
1623
1624 let result = storage.batch_write(events).unwrap();
1625 assert_eq!(result.events_written, 250);
1626 assert_eq!(result.batches_flushed, 1);
1627 assert_eq!(storage.pending_count(), 0);
1628
1629 storage.flush().unwrap();
1631
1632 let loaded = storage.load_all_events().unwrap();
1634 assert_eq!(loaded.len(), 250);
1635 }
1636
1637 #[test]
1638 fn test_auto_flush_on_batch_size() {
1639 let temp_dir = TempDir::new().unwrap();
1640 let config = ParquetStorageConfig {
1641 batch_size: 10, ..Default::default()
1643 };
1644 let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1645
1646 for i in 0..15 {
1648 storage
1649 .append_event(create_test_event(&format!("entity-{i}")))
1650 .unwrap();
1651 }
1652
1653 assert_eq!(storage.pending_count(), 5);
1655
1656 let stats = storage.batch_stats();
1657 assert_eq!(stats.events_written, 10);
1658 assert_eq!(stats.batches_written, 1);
1659 assert_eq!(stats.size_flushes, 1);
1660 }
1661
1662 #[test]
1663 fn test_flush_on_shutdown() {
1664 let temp_dir = TempDir::new().unwrap();
1665 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1666
1667 for i in 0..5 {
1669 storage
1670 .append_event(create_test_event(&format!("entity-{i}")))
1671 .unwrap();
1672 }
1673
1674 assert_eq!(storage.pending_count(), 5);
1675
1676 let flushed = storage.flush_on_shutdown().unwrap();
1678 assert_eq!(flushed, 5);
1679 assert_eq!(storage.pending_count(), 0);
1680
1681 let loaded = storage.load_all_events().unwrap();
1683 assert_eq!(loaded.len(), 5);
1684 }
1685
1686 #[test]
1687 fn test_thread_safe_writes() {
1688 let temp_dir = TempDir::new().unwrap();
1689 let config = ParquetStorageConfig {
1690 batch_size: 100,
1691 ..Default::default()
1692 };
1693 let storage = Arc::new(ParquetStorage::with_config(temp_dir.path(), config).unwrap());
1694
1695 let events_per_thread = 50;
1696 let thread_count = 4;
1697
1698 std::thread::scope(|s| {
1699 for t in 0..thread_count {
1700 let storage_ref = Arc::clone(&storage);
1701 s.spawn(move || {
1702 for i in 0..events_per_thread {
1703 let event = create_test_event(&format!("thread-{t}-entity-{i}"));
1704 storage_ref.append_event(event).unwrap();
1705 }
1706 });
1707 }
1708 });
1709
1710 storage.flush().unwrap();
1712
1713 let loaded = storage.load_all_events().unwrap();
1715 assert_eq!(loaded.len(), events_per_thread * thread_count);
1716 }
1717
1718 #[test]
1719 fn test_batch_stats() {
1720 let temp_dir = TempDir::new().unwrap();
1721 let config = ParquetStorageConfig {
1722 batch_size: 50,
1723 ..Default::default()
1724 };
1725 let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1726
1727 let events: Vec<Event> = (0..100)
1732 .map(|i| create_test_event(&format!("entity-{i}")))
1733 .collect();
1734
1735 storage.batch_write(events).unwrap();
1736
1737 let stats = storage.batch_stats();
1738 assert_eq!(stats.batches_written, 1);
1739 assert_eq!(stats.events_written, 100);
1740 assert!(stats.avg_batch_size > 0.0);
1741 assert!(stats.events_per_sec > 0.0);
1742 assert_eq!(stats.size_flushes, 1);
1743 }
1744
1745 #[test]
1746 fn test_config_presets() {
1747 let high_throughput = ParquetStorageConfig::high_throughput();
1748 assert_eq!(high_throughput.batch_size, 50_000);
1749 assert_eq!(high_throughput.flush_timeout, Duration::from_secs(10));
1750
1751 let low_latency = ParquetStorageConfig::low_latency();
1752 assert_eq!(low_latency.batch_size, 1_000);
1753 assert_eq!(low_latency.flush_timeout, Duration::from_secs(1));
1754
1755 let default = ParquetStorageConfig::default();
1756 assert_eq!(default.batch_size, DEFAULT_BATCH_SIZE);
1757 assert_eq!(default.batch_size, 10_000);
1758 }
1759
1760 #[test]
1763 #[ignore]
1764 fn test_batch_write_throughput() {
1765 let temp_dir = TempDir::new().unwrap();
1766 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1767
1768 let event_count = 50_000;
1769
1770 let events: Vec<Event> = (0..event_count)
1772 .map(|i| create_test_event(&format!("entity-{i}")))
1773 .collect();
1774
1775 let start = std::time::Instant::now();
1776 let result = storage.batch_write(events).unwrap();
1777 storage.flush().unwrap(); let batch_duration = start.elapsed();
1779
1780 let batch_stats = storage.batch_stats();
1781
1782 println!("\n=== Parquet Batch Write Performance (BATCH_SIZE=10,000) ===");
1783 println!("Events: {event_count}");
1784 println!("Duration: {batch_duration:?}");
1785 println!("Events/sec: {:.0}", result.events_per_sec);
1786 println!("Batches written: {}", batch_stats.batches_written);
1787 println!("Avg batch size: {:.0}", batch_stats.avg_batch_size);
1788 println!("Bytes written: {} KB", batch_stats.bytes_written / 1024);
1789
1790 assert!(
1793 result.events_per_sec > 10_000.0,
1794 "Batch write throughput too low: {:.0} events/sec (expected >10K in debug, >100K in release)",
1795 result.events_per_sec
1796 );
1797 }
1798
1799 #[test]
1801 #[ignore]
1802 fn test_single_event_write_baseline() {
1803 let temp_dir = TempDir::new().unwrap();
1804 let config = ParquetStorageConfig {
1805 batch_size: 1, ..Default::default()
1807 };
1808 let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1809
1810 let event_count = 1_000; let start = std::time::Instant::now();
1813 for i in 0..event_count {
1814 let event = create_test_event(&format!("entity-{i}"));
1815 storage.append_event(event).unwrap();
1816 }
1817 let duration = start.elapsed();
1818
1819 let events_per_sec = f64::from(event_count) / duration.as_secs_f64();
1820
1821 println!("\n=== Single-Event Write Baseline ===");
1822 println!("Events: {event_count}");
1823 println!("Duration: {duration:?}");
1824 println!("Events/sec: {events_per_sec:.0}");
1825
1826 }
1829
1830 fn touch_parquet(path: &Path) {
1840 std::fs::create_dir_all(path.parent().unwrap()).unwrap();
1841 std::fs::write(path, b"").unwrap();
1842 }
1843
1844 #[test]
1845 fn test_walker_finds_files_in_flat_layout() {
1846 let temp_dir = TempDir::new().unwrap();
1847 let root = temp_dir.path();
1848 touch_parquet(&root.join("events-20260101-120000000-aaaa.parquet"));
1849 touch_parquet(&root.join("events-20260101-130000000-bbbb.parquet"));
1850
1851 let mut found = find_parquet_files_recursive(root).unwrap();
1852 found.sort();
1853 assert_eq!(found.len(), 2);
1854 assert!(
1855 found[0]
1856 .file_name()
1857 .unwrap()
1858 .to_str()
1859 .unwrap()
1860 .starts_with("events-"),
1861 "expected events-* file, got {found:?}"
1862 );
1863 }
1864
1865 #[test]
1866 fn test_walker_finds_files_in_tenant_partitioned_tree() {
1867 let temp_dir = TempDir::new().unwrap();
1868 let root = temp_dir.path();
1869 touch_parquet(&root.join("tenant-a/2026-01/events-20260101-120000000-aaaa.parquet"));
1871 touch_parquet(&root.join("tenant-a/2026-02/events-20260201-120000000-bbbb.parquet"));
1872 touch_parquet(&root.join("tenant-b/2026-01/events-20260103-120000000-cccc.parquet"));
1873
1874 let found = find_parquet_files_recursive(root).unwrap();
1875 assert_eq!(found.len(), 3);
1876 assert!(found[0].to_str().unwrap().contains("tenant-a"));
1879 assert!(found[1].to_str().unwrap().contains("tenant-a"));
1880 assert!(found[2].to_str().unwrap().contains("tenant-b"));
1881 }
1882
1883 #[test]
1884 fn test_walker_handles_mixed_legacy_and_partitioned_layouts() {
1885 let temp_dir = TempDir::new().unwrap();
1890 let root = temp_dir.path();
1891 touch_parquet(&root.join("events-legacy-aaaa.parquet"));
1892 touch_parquet(&root.join("tenant-a/2026-01/events-new-bbbb.parquet"));
1893
1894 let found = find_parquet_files_recursive(root).unwrap();
1895 assert_eq!(found.len(), 2);
1896 }
1897
1898 #[test]
1899 fn test_walker_ignores_non_parquet_files() {
1900 let temp_dir = TempDir::new().unwrap();
1901 let root = temp_dir.path();
1902 std::fs::write(root.join("README.md"), b"hello").unwrap();
1903 std::fs::write(root.join("events.json"), b"[]").unwrap();
1904 touch_parquet(&root.join("events-20260101-120000000-aaaa.parquet"));
1905 std::fs::write(root.join("not-a-parquet-file.bin"), b"").unwrap();
1908
1909 let found = find_parquet_files_recursive(root).unwrap();
1910 assert_eq!(found.len(), 1);
1911 assert_eq!(
1912 found[0].extension().and_then(|s| s.to_str()),
1913 Some("parquet")
1914 );
1915 }
1916
1917 fn event_with_tenant(tenant: &str, entity_id: &str) -> Event {
1921 Event::reconstruct_from_strings(
1922 uuid::Uuid::new_v4(),
1923 "test.event".to_string(),
1924 entity_id.to_string(),
1925 tenant.to_string(),
1926 json!({"k": "v"}),
1927 chrono::Utc::now(),
1928 None,
1929 1,
1930 )
1931 }
1932
1933 #[test]
1934 fn test_flush_writes_into_per_tenant_partition() {
1935 let temp_dir = TempDir::new().unwrap();
1939 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1940
1941 for i in 0..3 {
1942 storage
1943 .append_event(event_with_tenant("default", &format!("entity-{i}")))
1944 .unwrap();
1945 }
1946 storage.flush().unwrap();
1947
1948 let parquet_files = find_parquet_files_recursive(temp_dir.path()).unwrap();
1949 assert_eq!(parquet_files.len(), 1);
1950
1951 let rel = parquet_files[0]
1952 .strip_prefix(temp_dir.path())
1953 .unwrap()
1954 .to_string_lossy()
1955 .into_owned();
1956 let parts: Vec<&str> = rel.split(std::path::MAIN_SEPARATOR).collect();
1958 assert_eq!(parts.len(), 3, "expected tenant/yyyy-mm/file, got {rel}");
1959 assert_eq!(parts[0], "default");
1960 assert!(
1963 parts[1].len() == 7 && parts[1].as_bytes()[4] == b'-',
1964 "expected yyyy-mm, got {}",
1965 parts[1]
1966 );
1967 assert!(parts[2].starts_with("events-") && parts[2].ends_with(".parquet"));
1968
1969 let loaded = storage.load_all_events().unwrap();
1970 assert_eq!(loaded.len(), 3);
1971 }
1972
1973 #[test]
1974 fn test_multiple_tenants_get_isolated_subtrees() {
1975 let temp_dir = TempDir::new().unwrap();
1978 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1979
1980 for i in 0..2 {
1981 storage
1982 .append_event(event_with_tenant("alice", &format!("a-{i}")))
1983 .unwrap();
1984 }
1985 for i in 0..3 {
1986 storage
1987 .append_event(event_with_tenant("bob", &format!("b-{i}")))
1988 .unwrap();
1989 }
1990 storage.flush().unwrap();
1991
1992 let alice_subtree = temp_dir.path().join("alice");
1993 let bob_subtree = temp_dir.path().join("bob");
1994 assert!(alice_subtree.is_dir(), "alice should have its own subtree");
1995 assert!(bob_subtree.is_dir(), "bob should have its own subtree");
1996
1997 let alice_files = find_parquet_files_recursive(&alice_subtree).unwrap();
1998 let bob_files = find_parquet_files_recursive(&bob_subtree).unwrap();
1999 assert_eq!(alice_files.len(), 1);
2000 assert_eq!(bob_files.len(), 1);
2001
2002 let loaded = storage.load_all_events().unwrap();
2005 let (alice_count, bob_count) =
2006 loaded
2007 .iter()
2008 .fold((0, 0), |(a, b), e| match e.tenant_id_str() {
2009 "alice" => (a + 1, b),
2010 "bob" => (a, b + 1),
2011 _ => (a, b),
2012 });
2013 assert_eq!(alice_count, 2);
2014 assert_eq!(bob_count, 3);
2015 }
2016
2017 #[test]
2018 fn test_size_flush_only_drains_full_tenant() {
2019 let temp_dir = TempDir::new().unwrap();
2023 let config = ParquetStorageConfig {
2024 batch_size: 5,
2025 ..Default::default()
2026 };
2027 let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
2028
2029 for i in 0..5 {
2032 storage
2033 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2034 .unwrap();
2035 }
2036 for i in 0..2 {
2038 storage
2039 .append_event(event_with_tenant("bob", &format!("b-{i}")))
2040 .unwrap();
2041 }
2042
2043 assert_eq!(
2044 storage.pending_count(),
2045 2,
2046 "only bob's 2 events should be pending"
2047 );
2048
2049 let parquet_files = find_parquet_files_recursive(temp_dir.path()).unwrap();
2050 assert_eq!(parquet_files.len(), 1, "only alice should have flushed");
2051 assert!(
2052 parquet_files[0]
2053 .to_string_lossy()
2054 .contains(&format!("alice{}", std::path::MAIN_SEPARATOR)),
2055 "expected alice partition, got {}",
2056 parquet_files[0].display()
2057 );
2058 }
2059
2060 #[test]
2061 fn test_tenant_id_from_path_recovers_tenant_for_partitioned_files() {
2062 let root = Path::new("/data/storage");
2063 let f = Path::new("/data/storage/alice/2026-04/events-20260426-120000000-aaaa.parquet");
2064 assert_eq!(tenant_id_from_path(root, f), "alice");
2065 }
2066
2067 #[test]
2068 fn test_tenant_id_from_path_falls_back_to_default_for_legacy_flat_layout() {
2069 let root = Path::new("/data/storage");
2070 let f = Path::new("/data/storage/events-20260426-120000000-aaaa.parquet");
2071 assert_eq!(tenant_id_from_path(root, f), "default");
2074 }
2075
2076 #[test]
2077 fn test_sanitize_tenant_id_for_path_accepts_safe_inputs() {
2078 for ok in [
2079 "default",
2080 "system",
2081 "1e6b2d1c-2f64-4441-9cf9-42f2e451aa17",
2082 "onboard-diagnostic-160-at-example-com",
2083 "tenant_with_underscore",
2084 "v1.0",
2085 ] {
2086 assert!(
2087 sanitize_tenant_id_for_path(ok).is_ok(),
2088 "{ok:?} should be accepted"
2089 );
2090 }
2091 }
2092
2093 #[test]
2094 fn test_sanitize_tenant_id_for_path_rejects_unsafe_inputs() {
2095 for bad in [
2096 "", "..", ".", "foo/bar", "foo\\bar", "foo bar", "foo\nbar", "foo\0bar", "tenant?", "tenant*", ] {
2107 assert!(
2108 sanitize_tenant_id_for_path(bad).is_err(),
2109 "{bad:?} should be rejected"
2110 );
2111 }
2112
2113 let too_long = "a".repeat(129);
2115 assert!(sanitize_tenant_id_for_path(&too_long).is_err());
2116 }
2117
2118 #[test]
2119 fn test_partition_path_for_tenant_shape() {
2120 let root = Path::new("/data");
2121 let when = chrono::DateTime::parse_from_rfc3339("2026-04-26T12:00:00Z")
2122 .unwrap()
2123 .with_timezone(&chrono::Utc);
2124 let path = partition_path_for_tenant(root, "alice", when).unwrap();
2125 assert_eq!(path, Path::new("/data/alice/2026-04"));
2126 }
2127
2128 #[test]
2129 fn test_append_event_rejects_unsafe_tenant_at_flush() {
2130 let temp_dir = TempDir::new().unwrap();
2136 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2137
2138 storage
2142 .append_event(event_with_tenant("../escape", "e-0"))
2143 .unwrap();
2144 let result = storage.flush();
2145 assert!(result.is_err(), "flush should reject unsafe tenant_id");
2146 let msg = format!("{}", result.unwrap_err());
2147 assert!(
2148 msg.contains("disallowed character") || msg.contains("reserved"),
2149 "expected sanitization error message, got: {msg}"
2150 );
2151 }
2152
2153 #[test]
2158 fn test_load_events_for_tenant_only_walks_target_subtree() {
2159 let temp_dir = TempDir::new().unwrap();
2164 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2165
2166 for i in 0..2 {
2167 storage
2168 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2169 .unwrap();
2170 }
2171 for i in 0..3 {
2172 storage
2173 .append_event(event_with_tenant("bob", &format!("b-{i}")))
2174 .unwrap();
2175 }
2176 for i in 0..1 {
2177 storage
2178 .append_event(event_with_tenant("carol", &format!("c-{i}")))
2179 .unwrap();
2180 }
2181 storage.flush().unwrap();
2182
2183 let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
2184 assert_eq!(alice_files.len(), 1);
2185 assert!(
2186 alice_files[0]
2187 .to_string_lossy()
2188 .contains(&format!("alice{}", std::path::MAIN_SEPARATOR)),
2189 "expected alice file, got {}",
2190 alice_files[0].display()
2191 );
2192 for f in &alice_files {
2196 let s = f.to_string_lossy();
2197 assert!(!s.contains("bob"), "alice listing leaked bob file: {s}");
2198 assert!(!s.contains("carol"), "alice listing leaked carol file: {s}");
2199 }
2200
2201 let alice_events = storage.load_events_for_tenant("alice").unwrap();
2202 assert_eq!(alice_events.len(), 2);
2203 for e in &alice_events {
2204 assert_eq!(e.tenant_id_str(), "alice");
2205 }
2206
2207 let bob_events = storage.load_events_for_tenant("bob").unwrap();
2208 assert_eq!(bob_events.len(), 3);
2209 for e in &bob_events {
2210 assert_eq!(e.tenant_id_str(), "bob");
2211 }
2212
2213 let carol_events = storage.load_events_for_tenant("carol").unwrap();
2214 assert_eq!(carol_events.len(), 1);
2215 assert_eq!(carol_events[0].tenant_id_str(), "carol");
2216 }
2217
2218 #[test]
2219 fn test_load_events_for_tenant_returns_empty_when_subtree_missing() {
2220 let temp_dir = TempDir::new().unwrap();
2224 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2225
2226 storage
2229 .append_event(event_with_tenant("alice", "a-0"))
2230 .unwrap();
2231 storage.flush().unwrap();
2232
2233 let files = storage
2234 .list_parquet_files_for_tenant("nobody-here")
2235 .unwrap();
2236 assert!(files.is_empty());
2237
2238 let events = storage.load_events_for_tenant("nobody-here").unwrap();
2239 assert!(events.is_empty());
2240 }
2241
2242 #[test]
2243 fn test_load_events_for_tenant_rejects_unsafe_tenant_id() {
2244 let temp_dir = TempDir::new().unwrap();
2247 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2248
2249 for unsafe_tid in ["..", "a/b", "a\\b", "", "a..b/.."] {
2250 let result = storage.load_events_for_tenant(unsafe_tid);
2251 assert!(
2252 result.is_err(),
2253 "tenant_id {unsafe_tid:?} should have been rejected"
2254 );
2255 }
2256 }
2257
2258 #[test]
2259 fn test_load_events_for_tenant_ignores_legacy_flat_layout_files() {
2260 let temp_dir = TempDir::new().unwrap();
2268 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2269
2270 let _flat = seed_flat_layout_file(&storage, 4);
2272
2273 let default_events = storage.load_events_for_tenant("default").unwrap();
2276 assert!(
2277 default_events.is_empty(),
2278 "tenant-scoped load must not pick up flat-layout files; got {} events",
2279 default_events.len()
2280 );
2281
2282 let all_events = storage.load_all_events().unwrap();
2284 assert_eq!(all_events.len(), 4);
2285 }
2286
2287 #[test]
2292 fn test_write_atomic_parquet_emits_file_under_tenant_partition() {
2293 let temp_dir = TempDir::new().unwrap();
2298 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2299
2300 let events: Vec<Event> = (0..3)
2301 .map(|i| event_with_tenant("alice", &format!("a-{i}")))
2302 .collect();
2303
2304 let final_path = storage
2305 .write_atomic_parquet("alice", "snapshot.alice.range", &events)
2306 .unwrap();
2307
2308 let rel = final_path
2310 .strip_prefix(temp_dir.path())
2311 .unwrap()
2312 .to_string_lossy()
2313 .into_owned();
2314 let parts: Vec<&str> = rel.split(std::path::MAIN_SEPARATOR).collect();
2315 assert_eq!(parts.len(), 3, "expected tenant/yyyy-mm/file, got {rel}");
2316 assert_eq!(parts[0], "alice");
2317 assert_eq!(parts[2], "snapshot.alice.range.parquet");
2318
2319 assert!(final_path.is_file());
2321 let tmp = final_path.with_extension("parquet.tmp");
2322 assert!(
2323 !tmp.exists(),
2324 "tmp should have been renamed away; still at {}",
2325 tmp.display()
2326 );
2327
2328 let loaded = storage.load_events_for_tenant("alice").unwrap();
2330 assert_eq!(loaded.len(), 3);
2331 }
2332
2333 #[test]
2334 fn test_write_atomic_parquet_rejects_empty_events() {
2335 let temp_dir = TempDir::new().unwrap();
2336 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2337 let result = storage.write_atomic_parquet("alice", "snap", &[]);
2338 assert!(result.is_err());
2339 }
2340
2341 #[test]
2342 fn test_write_atomic_parquet_rejects_unsafe_tenant() {
2343 let temp_dir = TempDir::new().unwrap();
2344 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2345 let events = [event_with_tenant("alice", "e-0")];
2346 for unsafe_tid in ["..", "a/b", ""] {
2347 let result = storage.write_atomic_parquet(unsafe_tid, "snap", &events);
2348 assert!(
2349 result.is_err(),
2350 "unsafe tenant_id {unsafe_tid:?} should have been rejected"
2351 );
2352 }
2353 }
2354
2355 #[test]
2356 fn test_cleanup_partial_writes_removes_orphan_tmps() {
2357 let temp_dir = TempDir::new().unwrap();
2361 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2362
2363 for i in 0..2 {
2366 storage
2367 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2368 .unwrap();
2369 }
2370 storage.flush().unwrap();
2371 let real_files_before = find_parquet_files_recursive(temp_dir.path()).unwrap();
2372 assert_eq!(real_files_before.len(), 1);
2373
2374 let alice_subtree = temp_dir.path().join("alice");
2376 let orphan_dir = real_files_before[0].parent().unwrap();
2377 let orphan_path = orphan_dir.join("snapshot.alice.crashed.parquet.tmp");
2378 std::fs::write(&orphan_path, b"fake partial parquet").unwrap();
2379 assert!(orphan_path.is_file());
2380
2381 let nested_dir = alice_subtree.join("2099-01");
2383 std::fs::create_dir_all(&nested_dir).unwrap();
2384 let nested_orphan = nested_dir.join("events-x.parquet.tmp");
2385 std::fs::write(&nested_orphan, b"junk").unwrap();
2386
2387 let removed = storage.cleanup_partial_writes().unwrap();
2388 assert_eq!(removed, 2, "two orphan tmps should have been cleaned");
2389 assert!(!orphan_path.exists());
2390 assert!(!nested_orphan.exists());
2391
2392 let real_files_after = find_parquet_files_recursive(temp_dir.path()).unwrap();
2394 assert_eq!(real_files_after, real_files_before);
2395 }
2396
2397 #[test]
2398 fn test_cleanup_partial_writes_quarantines_zero_byte_parquet() {
2399 let temp_dir = TempDir::new().unwrap();
2404 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2405
2406 for i in 0..2 {
2408 storage
2409 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2410 .unwrap();
2411 }
2412 storage.flush().unwrap();
2413 let healthy_files_before = find_parquet_files_recursive(temp_dir.path()).unwrap();
2414 assert_eq!(healthy_files_before.len(), 1);
2415 let healthy_file = &healthy_files_before[0];
2416 let healthy_dir = healthy_file.parent().unwrap();
2417
2418 let bricked = healthy_dir.join("events-bricked-deadbeef.parquet");
2420 std::fs::write(&bricked, b"").unwrap();
2421 assert_eq!(std::fs::metadata(&bricked).unwrap().len(), 0);
2422
2423 let acted = storage.cleanup_partial_writes().unwrap();
2424 assert_eq!(acted, 1, "only the 0-byte file should have been acted on");
2425
2426 assert!(!bricked.exists(), "0-byte file should have been renamed");
2429 let quarantined: Vec<_> = std::fs::read_dir(healthy_dir)
2430 .unwrap()
2431 .flatten()
2432 .map(|e| e.path())
2433 .filter(|p| {
2434 p.file_name()
2435 .and_then(|n| n.to_str())
2436 .is_some_and(|n| n.starts_with("events-bricked-deadbeef.parquet.corrupt-"))
2437 })
2438 .collect();
2439 assert_eq!(
2440 quarantined.len(),
2441 1,
2442 "expected one .parquet.corrupt-<ts> sibling"
2443 );
2444
2445 assert!(
2447 healthy_file.exists(),
2448 "healthy parquet must not be molested"
2449 );
2450 }
2451
2452 #[test]
2453 fn test_load_all_events_skips_zero_byte_parquet() {
2454 let temp_dir = TempDir::new().unwrap();
2458 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2459
2460 for i in 0..2 {
2462 storage
2463 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2464 .unwrap();
2465 }
2466 storage.flush().unwrap();
2467
2468 let healthy_files = find_parquet_files_recursive(temp_dir.path()).unwrap();
2472 let bricked = healthy_files[0]
2473 .parent()
2474 .unwrap()
2475 .join("events-bricked-cafef00d.parquet");
2476 std::fs::write(&bricked, b"").unwrap();
2477
2478 let loaded = storage.load_all_events().unwrap();
2479 assert_eq!(
2480 loaded.len(),
2481 2,
2482 "all healthy events should still load despite the 0-byte file"
2483 );
2484 }
2485
2486 #[test]
2487 fn test_load_events_for_tenant_skips_zero_byte_parquet() {
2488 let temp_dir = TempDir::new().unwrap();
2491 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2492
2493 for i in 0..3 {
2494 storage
2495 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2496 .unwrap();
2497 }
2498 storage.flush().unwrap();
2499
2500 let healthy = find_parquet_files_recursive(temp_dir.path()).unwrap();
2501 let bricked = healthy[0]
2502 .parent()
2503 .unwrap()
2504 .join("events-bricked-feedface.parquet");
2505 std::fs::write(&bricked, b"").unwrap();
2506
2507 let events = storage.load_events_for_tenant("alice").unwrap();
2508 assert_eq!(events.len(), 3, "lazy load must skip the 0-byte file");
2509 }
2510
2511 #[test]
2512 fn test_flush_tenant_leaves_no_zero_byte_parquet_after_normal_flush() {
2513 let temp_dir = TempDir::new().unwrap();
2517 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2518
2519 for i in 0..5 {
2520 storage
2521 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2522 .unwrap();
2523 }
2524 storage.flush().unwrap();
2525
2526 let mut tmp_count = 0;
2527 let mut zero_byte_count = 0;
2528 let mut healthy_count = 0;
2529 let mut stack = vec![temp_dir.path().to_path_buf()];
2530 while let Some(d) = stack.pop() {
2531 for entry in std::fs::read_dir(&d).unwrap().flatten() {
2532 let p = entry.path();
2533 if p.is_dir() {
2534 stack.push(p);
2535 continue;
2536 }
2537 let name = p.file_name().unwrap().to_string_lossy().into_owned();
2538 if name.ends_with(".parquet.tmp") {
2539 tmp_count += 1;
2540 } else if name.ends_with(".parquet") {
2541 if std::fs::metadata(&p).unwrap().len() == 0 {
2542 zero_byte_count += 1;
2543 } else {
2544 healthy_count += 1;
2545 }
2546 }
2547 }
2548 }
2549 assert_eq!(tmp_count, 0, ".parquet.tmp survivors after flush");
2550 assert_eq!(zero_byte_count, 0, "0-byte .parquet survivors after flush");
2551 assert_eq!(healthy_count, 1, "expected exactly one healthy parquet");
2552 }
2553
2554 #[test]
2555 fn test_new_calls_cleanup_partial_writes_on_boot() {
2556 let temp_dir = TempDir::new().unwrap();
2560 let stale = temp_dir.path().join("orphan.parquet.tmp");
2561 std::fs::write(&stale, b"crash detritus").unwrap();
2562 assert!(stale.is_file());
2563
2564 let _storage = ParquetStorage::new(temp_dir.path()).unwrap();
2565 assert!(
2566 !stale.exists(),
2567 "stale tmp should have been cleaned by ParquetStorage::new"
2568 );
2569 }
2570
2571 fn seed_flat_layout_file(storage: &ParquetStorage, count: usize) -> PathBuf {
2581 for i in 0..count {
2582 storage
2583 .append_event(create_test_event(&format!("entity-{i}")))
2584 .unwrap();
2585 }
2586 storage.flush().unwrap();
2587
2588 let default_subtree = storage.storage_dir().join("default");
2593 let candidates = find_parquet_files_recursive(&default_subtree).unwrap();
2594 assert!(
2595 !candidates.is_empty(),
2596 "seed expected at least one file under default/"
2597 );
2598 let src = candidates.into_iter().max().unwrap();
2599
2600 let dst = storage.storage_dir().join(src.file_name().unwrap());
2601 std::fs::rename(&src, &dst).unwrap();
2602 if let Some(month_dir) = src.parent() {
2606 let _ = std::fs::remove_dir(month_dir);
2607 if let Some(tenant_dir) = month_dir.parent() {
2608 let _ = std::fs::remove_dir(tenant_dir);
2609 }
2610 }
2611 dst
2612 }
2613
2614 #[test]
2615 fn test_migrate_flat_layout_dry_run_touches_nothing() {
2616 let temp_dir = TempDir::new().unwrap();
2617 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2618 let flat = seed_flat_layout_file(&storage, 7);
2619 assert!(flat.is_file(), "test setup: flat file should exist");
2620
2621 let report = storage.migrate_flat_layout(true).unwrap();
2622 assert!(report.dry_run);
2623 assert_eq!(report.flat_files_seen, 1);
2624 assert_eq!(report.events_migrated, 7);
2625 assert_eq!(report.flat_files_removed, 0);
2626 assert_eq!(report.partitions_written, 0);
2627 assert!(
2628 flat.is_file(),
2629 "flat file must still be present after dry run"
2630 );
2631 }
2632
2633 #[test]
2634 fn test_migrate_flat_layout_moves_events_into_default_tree_and_removes_flat() {
2635 let temp_dir = TempDir::new().unwrap();
2636 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2637 let flat = seed_flat_layout_file(&storage, 5);
2638
2639 let report = storage.migrate_flat_layout(false).unwrap();
2640 assert!(!report.dry_run);
2641 assert_eq!(report.flat_files_seen, 1);
2642 assert_eq!(report.flat_files_removed, 1);
2643 assert_eq!(report.events_migrated, 5);
2644 assert!(report.partitions_written >= 1);
2645 assert!(
2646 !flat.exists(),
2647 "flat file should be deleted after migration"
2648 );
2649
2650 let post = find_parquet_files_recursive(temp_dir.path()).unwrap();
2651 assert!(
2652 post.iter().all(|p| {
2653 let rel = p
2654 .strip_prefix(temp_dir.path())
2655 .unwrap()
2656 .to_string_lossy()
2657 .into_owned();
2658 rel.starts_with(&format!("default{}", std::path::MAIN_SEPARATOR))
2659 }),
2660 "all migrated files should be under default/"
2661 );
2662
2663 let loaded = storage.load_all_events().unwrap();
2664 assert_eq!(loaded.len(), 5);
2665 }
2666
2667 #[test]
2668 fn test_migrate_flat_layout_is_idempotent_when_re_run_after_completion() {
2669 let temp_dir = TempDir::new().unwrap();
2670 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2671 let _flat = seed_flat_layout_file(&storage, 4);
2672
2673 let first = storage.migrate_flat_layout(false).unwrap();
2674 assert_eq!(first.events_migrated, 4);
2675
2676 let second = storage.migrate_flat_layout(false).unwrap();
2679 assert_eq!(second.flat_files_seen, 0);
2680 assert_eq!(second.events_migrated, 0);
2681 assert_eq!(second.flat_files_removed, 0);
2682
2683 let loaded = storage.load_all_events().unwrap();
2684 assert_eq!(loaded.len(), 4, "rerun must not duplicate or lose events");
2685 }
2686
2687 #[test]
2688 fn test_migrate_flat_layout_ignores_already_partitioned_data() {
2689 let temp_dir = TempDir::new().unwrap();
2692 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2693
2694 for i in 0..3 {
2695 storage
2696 .append_event(event_with_tenant("alice", &format!("a-{i}")))
2697 .unwrap();
2698 }
2699 storage.flush().unwrap();
2700
2701 let _flat = seed_flat_layout_file(&storage, 2);
2702
2703 let report = storage.migrate_flat_layout(false).unwrap();
2704 assert_eq!(report.flat_files_seen, 1, "only the flat file is in scope");
2705 assert_eq!(report.events_migrated, 2);
2706
2707 let alice_files = find_parquet_files_recursive(&temp_dir.path().join("alice")).unwrap();
2708 assert_eq!(alice_files.len(), 1, "alice's tree must be untouched");
2709
2710 let loaded = storage.load_all_events().unwrap();
2711 assert_eq!(loaded.len(), 5);
2712 let alice_count = loaded
2713 .iter()
2714 .filter(|e| e.tenant_id_str() == "alice")
2715 .count();
2716 let default_count = loaded
2717 .iter()
2718 .filter(|e| e.tenant_id_str() == "default")
2719 .count();
2720 assert_eq!(alice_count, 3);
2721 assert_eq!(default_count, 2);
2722 }
2723
2724 #[test]
2725 fn test_migrate_flat_layout_with_no_flat_files_is_a_clean_noop() {
2726 let temp_dir = TempDir::new().unwrap();
2727 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2728 let report = storage.migrate_flat_layout(false).unwrap();
2729 assert_eq!(report.flat_files_seen, 0);
2730 assert_eq!(report.events_migrated, 0);
2731 }
2732}