Skip to main content

allsource_core/infrastructure/persistence/
storage.rs

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
26/// Default batch size for Parquet writes (10,000 events as per US-023)
27pub const DEFAULT_BATCH_SIZE: usize = 10_000;
28
29/// Default flush timeout in milliseconds
30pub const DEFAULT_FLUSH_TIMEOUT_MS: u64 = 5_000;
31
32/// Configuration for ParquetStorage batch processing
33#[derive(Debug, Clone)]
34pub struct ParquetStorageConfig {
35    /// Batch size before automatic flush (default: 10,000)
36    pub batch_size: usize,
37    /// Timeout before flushing partial batch (default: 5 seconds)
38    pub flush_timeout: Duration,
39    /// Compression codec for Parquet files
40    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    /// High-throughput configuration optimized for large batch writes
55    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    /// Low-latency configuration for smaller, more frequent writes
64    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/// Statistics for batch write operations
74#[derive(Debug, Clone, Default)]
75pub struct BatchWriteStats {
76    /// Total batches written
77    pub batches_written: u64,
78    /// Total events written
79    pub events_written: u64,
80    /// Total bytes written
81    pub bytes_written: u64,
82    /// Average batch size
83    pub avg_batch_size: f64,
84    /// Events per second (throughput)
85    pub events_per_sec: f64,
86    /// Total write time in nanoseconds
87    pub total_write_time_ns: u64,
88    /// Number of timeout-triggered flushes
89    pub timeout_flushes: u64,
90    /// Number of size-triggered flushes
91    pub size_flushes: u64,
92}
93
94/// Result of a batch write operation
95#[derive(Debug, Clone)]
96pub struct BatchWriteResult {
97    /// Number of events written
98    pub events_written: usize,
99    /// Number of batches flushed to disk
100    pub batches_flushed: usize,
101    /// Total duration of the write operation
102    pub duration: Duration,
103    /// Write throughput in events per second
104    pub events_per_sec: f64,
105}
106
107/// Parquet-based persistent storage for events with batch processing
108///
109/// Features:
110/// - Configurable batch size (default: 10,000 events per US-023)
111/// - Timeout-based flushing for partial batches
112/// - Thread-safe batch accumulation
113/// - SNAPPY compression for efficient storage
114/// - Automatic flush on shutdown via Drop
115pub struct ParquetStorage {
116    /// Base directory for storing parquet files
117    storage_dir: PathBuf,
118
119    /// Buffered events keyed by tenant_id. Each tenant accumulates its own
120    /// batch and flushes independently into its partition under
121    /// `storage_dir/<tenant_id>/<yyyy-mm>/`. Single outer mutex protects the
122    /// whole map: lookup is O(1), tenant cardinality is low (single digits
123    /// today; bounded by Step 3's cache budget later), so contention is
124    /// fine. We keep the mutex held only for the push, not for disk I/O —
125    /// flush takes ownership of a tenant's batch via remove() and writes
126    /// after the lock is released.
127    current_batches: Mutex<HashMap<String, Vec<Event>>>,
128
129    /// Configuration
130    config: ParquetStorageConfig,
131
132    /// Schema for Arrow/Parquet
133    schema: Arc<Schema>,
134
135    /// Last flush timestamp for timeout tracking
136    last_flush_time: Mutex<Instant>,
137
138    /// Statistics tracking
139    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    /// Create a new ParquetStorage with default configuration (10,000 event batches)
149    pub fn new(storage_dir: impl AsRef<Path>) -> Result<Self> {
150        Self::with_config(storage_dir, ParquetStorageConfig::default())
151    }
152
153    /// Create a new ParquetStorage with custom configuration
154    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        // Create storage directory if it doesn't exist
161        fs::create_dir_all(&storage_dir).map_err(|e| {
162            AllSourceError::StorageError(format!("Failed to create storage directory: {e}"))
163        })?;
164
165        // Define Arrow schema for events
166        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        // Boot-time crash recovery: heal *.parquet.tmp files from
195        // crashed atomic writes (snapshot/compaction or, post-#166-fix,
196        // checkpoint flushes) and quarantine any 0-byte *.parquet files
197        // left over from pre-fix checkpoint crashes (issue #166).
198        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    /// Create storage with legacy batch size (1000) for backward compatibility
211    #[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    /// Add an event to the current batch
223    ///
224    /// Events are routed to a per-tenant batch keyed by `event.tenant_id_str()`.
225    /// A tenant's batch is buffered until any of:
226    /// - That tenant's batch hits the configured `batch_size` (default 10,000)
227    ///   — flushes only that tenant, not the whole world
228    /// - The flush timeout elapses — flushes every tenant with pending data
229    /// - `flush()` is called explicitly
230    /// - The process shuts down
231    #[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    /// Add multiple events to the batch (optimized batch insertion)
250    ///
251    /// Preferred entry point for high-throughput ingestion. Events are
252    /// grouped by tenant under a single mutex acquisition and any tenant
253    /// that crosses `batch_size` is flushed on the spot.
254    #[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        // Pre-group by tenant to keep the lock window short — one acquire,
260        // one extend per tenant, decide which tenants are over threshold.
261        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    /// Check if a timeout-based flush is needed and perform it
299    ///
300    /// Call this periodically (e.g., from a background task) to ensure
301    /// partial batches are flushed within the configured timeout. When
302    /// triggered, every tenant with pending events flushes — the timer is
303    /// global, not per-tenant, so a slow-trickle tenant doesn't get
304    /// stranded waiting for its own batch to fill.
305    #[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    /// Flush every tenant's pending batch to its partition.
324    ///
325    /// Thread-safe: callable from any thread. A snapshot of which tenants
326    /// have pending data is taken under a short lock; each tenant is then
327    /// flushed individually with its own lock cycle, so disk I/O for one
328    /// tenant doesn't block writes against another.
329    #[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    /// Call while holding exclusive access to the storage owner: a shared
349    /// access check alone cannot exclude a batch already drained by a flush.
350    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    /// Flush a single tenant's pending events into its partition file.
359    ///
360    /// File path: `storage_dir/<sanitized_tenant_id>/<yyyy-mm>/events-<ts>-<uuid>.parquet`.
361    /// `<yyyy-mm>` is taken from the wall-clock at flush time (matching the
362    /// pre-tenant filename's timestamp semantics) rather than from
363    /// individual event timestamps — keeps each flush to a single output
364    /// file even when buffered events span months.
365    ///
366    /// The file is written atomically via `write_record_batch_atomic`:
367    /// crash mid-write leaves a `.parquet.tmp` file (cleaned on next boot)
368    /// rather than a 0-byte `.parquet` file with the final name. Issue
369    /// #166 — pre-fix, a SIGKILL inside this function bricked subsequent
370    /// loads because `load_all_events` aborted on the first unreadable
371    /// file.
372    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                // The WAL may be retired by the next successful checkpoint, so
389                // dropping these here would make the failure permanent (#287).
390                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    /// Write a single record batch atomically to
462    /// `<partition_dir>/<file_stem>.parquet`, returning the final path
463    /// and the parquet writer's metadata (used for byte-size metrics by
464    /// the checkpoint flush path).
465    ///
466    /// Crash-safety contract — same four steps as `write_atomic_parquet`:
467    /// 1. Write to `<final_path>.tmp` first.
468    /// 2. fsync the .tmp file so the data is durably on disk.
469    /// 3. Rename `.tmp` → final name (atomic POSIX rename).
470    /// 4. fsync the parent directory so the rename survives crash.
471    ///
472    /// Caller is responsible for choosing `partition_dir` (and creating
473    /// it). On any failure mid-way, the `.tmp` file gets cleaned up by
474    /// `cleanup_partial_writes` on next boot. The final file appears
475    /// atomically — readers either see the old state (no file) or the
476    /// complete new file, never a half-written one.
477    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        // 1. Write to .tmp.
487        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            // close() consumes writer and returns the file metadata; the
502            // underlying File is flushed and closed here.
503            writer.close()?
504        };
505
506        // 2. fsync the .tmp file.
507        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        // 3. Atomic rename.
519        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        // 4. fsync the parent directory so the rename survives crash.
528        // Linux fsync-on-dir is the canonical way to make a rename
529        // durable; macOS no-ops the dir fsync but doesn't error.
530        if let Ok(dir) = File::open(partition_dir) {
531            let _ = dir.sync_all();
532        }
533
534        Ok((final_path, metadata))
535    }
536
537    /// Atomically write `events` to a Parquet file under the tenant's
538    /// partition. Step 4 of the sustainable data strategy uses this to
539    /// emit per-tenant snapshot/compaction files (`snapshot.<tenant>.<from>-<to>`)
540    /// without risking partial files on crash.
541    ///
542    /// Crash-safety contract:
543    /// 1. Write to `<final_path>.tmp` first.
544    /// 2. fsync the .tmp file so the data is durably on disk.
545    /// 3. Rename `.tmp` → final name (atomic POSIX rename).
546    /// 4. fsync the parent directory so the rename is durable.
547    ///
548    /// On any failure mid-way, the .tmp file gets cleaned up by
549    /// `cleanup_partial_writes` on next boot. The final file appears
550    /// atomically — readers either see the old state (no file) or
551    /// the complete new file, never a half-written one.
552    ///
553    /// `file_stem` is the filename without extension or partition path
554    /// — caller-controlled so the snapshot naming convention
555    /// (`snapshot.<tenant>.<from>-<to>`) lives in the compaction layer,
556    /// not here. The `.parquet` extension is appended automatically.
557    ///
558    /// Returns the final (post-rename) path so the caller can
559    /// confirm the file landed where expected.
560    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        // Anchor partition by the earliest event's month — keeps the
572        // file in the same yyyy-mm bucket as the data it represents.
573        // Callers that span multiple months can use the earliest
574        // event's month and trust the recursive walker to find it.
575        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    /// Sweep the storage tree for crash detritus and heal it on boot.
603    /// Called on boot from `ParquetStorage::new` so every cold start
604    /// starts clean. Two kinds of leftovers are handled:
605    ///
606    /// 1. **`*.parquet.tmp` files** — a snapshot or flush write that
607    ///    crashed between fsync and rename. Safe to delete: the data is
608    ///    either still in the WAL (for checkpoint flushes) or in the
609    ///    constituent raw files (for snapshot/compaction, which only
610    ///    deletes the inputs AFTER the rename succeeds). Either way, the
611    ///    failed write will be retried.
612    ///
613    /// 2. **0-byte `*.parquet` files** — issue #166. Pre-fix releases
614    ///    (≤ 0.20.0 for the checkpoint path) wrote directly to the
615    ///    final filename, so SIGKILL between `File::create` and the
616    ///    first flush left a 0-byte `events-*.parquet` that bricked
617    ///    every subsequent load. Rather than delete, we rename to
618    ///    `<original>.parquet.corrupt-<unix-ts>` so an operator can
619    ///    inspect/forensicate before destroying the evidence — most
620    ///    will just `rm` after seeing the size.
621    ///
622    /// Returns the total number of files acted on, for boot-log
623    /// observability.
624    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                    // Issue #166: pre-fix checkpoint write left a 0-byte
662                    // .parquet file with the final name on crash. Quarantine
663                    // (rename) rather than delete — operator can inspect.
664                    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    /// Force flush any remaining events (for shutdown handling).
690    ///
691    /// Sums pending counts across every tenant's batch so the caller can
692    /// log "we flushed N events on shutdown" without caring about
693    /// per-tenant breakdown.
694    #[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    /// Get batch write statistics
713    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    /// Total pending events across all tenant batches.
742    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    /// Get configured batch size
752    pub fn batch_size(&self) -> usize {
753        self.config.batch_size
754    }
755
756    /// Get configured flush timeout
757    pub fn flush_timeout(&self) -> Duration {
758        self.config.flush_timeout
759    }
760
761    /// Convert events to Arrow RecordBatch
762    #[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            // Convert timestamp to microseconds
779            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    /// Load events from all Parquet files under the storage directory.
807    ///
808    /// Walks the tree recursively so both layouts work: legacy flat
809    /// (`storage_dir/events-*.parquet`) and the tenant-partitioned tree
810    /// introduced by the data-strategy work (`storage_dir/<tenant>/<yyyy-mm>/
811    /// events-*.parquet`). The two coexist on disk during migration.
812    #[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            // Issue #166: a single corrupt/0-byte parquet file (crash mid-write
822            // on pre-fix releases) would propagate up and abort the whole load,
823            // bricking access to every healthy file alongside it. Log and skip
824            // instead so one bad file can't take the rest down.
825            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    /// Load events from a single Parquet file. `tenant_id` is the value to
854    /// stamp onto each loaded event — derived from the file's location in
855    /// the tree by `load_all_events`. The Parquet schema doesn't include
856    /// tenant_id today (path is the source of truth), so this is how
857    /// per-tenant identity survives the round trip.
858    #[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    /// Public wrapper around the internal single-file loader. Used
911    /// by the per-tenant compaction (Step 4) which needs to read a
912    /// specific candidate set rather than the whole tenant subtree.
913    /// Caller passes the tenant_id explicitly because the schema
914    /// doesn't carry it.
915    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    /// Convert Arrow RecordBatch back to events. `tenant_id` is stamped onto
924    /// each reconstructed event — the schema doesn't carry it today, so the
925    /// caller passes the value derived from the file path.
926    #[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    /// List all Parquet file paths under the storage directory, sorted by
1003    /// the relative path so files in the same partition stay grouped.
1004    ///
1005    /// Used by the replication catch-up protocol to stream snapshot files
1006    /// to followers that are too far behind for WAL-only catch-up.
1007    pub fn list_parquet_files(&self) -> Result<Vec<PathBuf>> {
1008        find_parquet_files_recursive(&self.storage_dir)
1009    }
1010
1011    /// The tenant a Parquet file belongs to, read from its place in the tree.
1012    pub fn tenant_id_for_file(&self, file_path: &Path) -> String {
1013        tenant_id_from_path(&self.storage_dir, file_path)
1014    }
1015
1016    /// List Parquet files belonging to a single tenant — i.e. only files
1017    /// under `<storage_dir>/<tenant>/...`. Legacy flat-layout files at the
1018    /// root are intentionally excluded; the migration tool moves them under
1019    /// `default/` so once it has run a `tenant=default` query sees them.
1020    ///
1021    /// Returns an empty vec when the tenant subtree doesn't exist (no data
1022    /// for that tenant yet). Returns an error only if `tenant_id` fails the
1023    /// path-safety whitelist.
1024    ///
1025    /// This is the building block for tenant-scoped reads: the caller knows
1026    /// which files might contain the tenant's data without opening any of
1027    /// the others.
1028    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    /// Load only the events belonging to `tenant_id`, walking just that
1058    /// tenant's subtree on disk. The full-storage loader
1059    /// (`load_all_events`) opens every Parquet file regardless of tenant;
1060    /// this one only opens files under `<storage_dir>/<tenant>/`.
1061    ///
1062    /// Returns an empty vec when the tenant has no on-disk data. Returns
1063    /// an error if the tenant_id fails the path-safety whitelist or listing
1064    /// fails. Individual unreadable files are logged and skipped for queries;
1065    /// this tolerant result cannot establish a complete version for OCC.
1066    ///
1067    /// This is the read-side complement to per-tenant flushing. It's the
1068    /// foundation Step 2 (lazy per-tenant load on demand) needs: a way to
1069    /// hydrate one tenant without paying the cost of loading every other
1070    /// tenant's data into memory.
1071    ///
1072    /// Tenant identity for loaded events comes from the file path, the
1073    /// same as `load_all_events` — `record_batch_to_events` stamps the
1074    /// passed `tenant_id` onto every reconstructed event.
1075    #[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            // Issue #166: same defensive log-and-skip as load_all_events.
1107            // The lazy-load path can hit a corrupt file just as easily as
1108            // boot can — failing the whole tenant load on one bad file
1109            // would brick every query for that tenant.
1110            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        // Tolerant enumeration may omit unreadable directories or entries.
1144        // Only strict enumeration can certify the retained archive for OCC.
1145        Ok((events, require_complete && skipped == 0))
1146    }
1147
1148    /// Get the storage directory path.
1149    pub fn storage_dir(&self) -> &Path {
1150        &self.storage_dir
1151    }
1152
1153    /// One-shot migration of flat-layout files into the tenant-partitioned
1154    /// tree. Run with Core stopped (no concurrent writes).
1155    ///
1156    /// Walks `storage_dir`'s top level (non-recursive) for the legacy
1157    /// `events-*.parquet` files. For each one it loads the events,
1158    /// regroups them by (tenant_id, yyyy-mm) — events from a flat file
1159    /// take the path-derived "default" tenant, since pre-partitioning
1160    /// data carried no tenant in its on-disk form — writes a fresh
1161    /// Parquet under the corresponding partition directory, and deletes
1162    /// the original flat file once the new file is closed.
1163    ///
1164    /// `dry_run = true` reports what would happen without touching disk.
1165    /// Run dry first; production data deserves the rehearsal.
1166    ///
1167    /// Crash safety: this writes the new partition file before deleting
1168    /// the flat file, so a crash between the two leaves both on disk.
1169    /// The recursive loader (`load_all_events`) will then return both,
1170    /// duplicating those events on next boot. Mitigation: stop Core
1171    /// before running, and re-run the migration after any crash so the
1172    /// flat file gets deleted. A future commit can add atomic rename +
1173    /// fsync semantics; for the one-time migration the stop-Core
1174    /// constraint is enough.
1175    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            // Pre-partition events used path-derived tenant. For flat-layout
1184            // files that's always "default" (`tenant_id_from_path` falls back
1185            // to "default" for single-component paths).
1186            let events = self.load_events_from_file(&flat_file, "default")?;
1187            report.flat_files_seen += 1;
1188
1189            if events.is_empty() {
1190                // Stale empty file (zero rows). Just remove it.
1191                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            // Group by (tenant, yyyy-mm-from-event-timestamp). The
1204            // partition month tracks the event's wall-clock time so that
1205            // post-migration the layout reflects when data happened, not
1206            // when migration ran. Step 4 (per-tenant snapshots) and Step
1207            // 5 (retention) will key on that.
1208            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    /// Get storage statistics
1267    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
1293/// Validate a tenant ID for use as a filesystem path component.
1294///
1295/// Whitelist: ASCII letters, digits, `-`, `_`, `.` — covers UUIDs, the
1296/// hyphen-and-lowercase tenant strings the onboarding flow produces, and
1297/// the `system` tenant the heartbeat emitter uses. Rejects empty input,
1298/// any path separator (`/`, `\`), and any "..". The whitelist is the
1299/// primary defence against path traversal; the explicit ".." check is
1300/// belt-and-braces in case the whitelist ever loosens.
1301///
1302/// Length capped at 128 bytes — comfortably above the 36-byte UUID and
1303/// the longest onboarding tenant the system has produced, well below
1304/// every common filesystem's NAME_MAX (typically 255).
1305fn 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
1333/// Resolve the directory a flush should write into for `(tenant, when)`.
1334///
1335/// Returns `<root>/<tenant>/<yyyy-mm>/`. Caller is responsible for
1336/// `create_dir_all`-ing the result before opening files in it.
1337fn 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
1346/// Reverse of `partition_path_for_tenant` — given a parquet file's full
1347/// path and the storage root, return the tenant_id stored in the path.
1348///
1349/// Tenant-partitioned shape: `<root>/<tenant>/<yyyy-mm>/events-*.parquet`
1350/// → first component after root is the tenant.
1351///
1352/// Legacy flat shape: `<root>/events-*.parquet` → no tenant in path, fall
1353/// back to `"default"` so events written before the partitioning change
1354/// keep loading with their original (and only ever) tenant identity.
1355fn 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        // Two or more components: <tenant>/<rest>... → tenant
1364        (Some(std::path::Component::Normal(tenant)), Some(_)) => {
1365            tenant.to_string_lossy().into_owned()
1366        }
1367        // Single component (the parquet file itself): legacy flat layout.
1368        _ => "default".to_string(),
1369    }
1370}
1371
1372/// List Parquet files at the top level of `root` only — i.e. the legacy
1373/// flat-layout files. Used by the one-shot migration tool to find data
1374/// that needs moving into the tenant-partitioned tree. The opposite of
1375/// `find_parquet_files_recursive`: stops at the first directory level so
1376/// already-partitioned data isn't included.
1377fn 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
1400/// Recursively collect all `*.parquet` files under `root`, sorted by path so
1401/// callers see a deterministic, tenant-grouped order.
1402///
1403/// Existence rationale: the storage layout is moving from a flat
1404/// `storage_dir/events-*.parquet` pile to a tenant-partitioned tree of the
1405/// shape `storage_dir/<tenant>/<yyyy-mm>/events-*.parquet`. During the
1406/// migration both shapes coexist, so every code path that asks "what
1407/// parquet files do we have?" needs to walk subdirectories. Symlinks are
1408/// not followed — the storage tree is mounted from a single volume and
1409/// chasing symlinks invites cycles.
1410fn 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            // Root must exist (we created it in `new`); subdirectories may
1429            // race a delete from compaction. Tolerant reads skip those;
1430            // strict reads cannot certify history after enumeration failed.
1431            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            // Use file_type() rather than metadata() so symlinks don't get
1454            // followed by accident (metadata() resolves symlinks, file_type()
1455            // doesn't).
1456            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        // Ensure any remaining events are flushed on shutdown
1491        if let Err(e) = self.flush_on_shutdown() {
1492            tracing::error!("Failed to flush events on drop: {}", e);
1493        }
1494    }
1495}
1496
1497/// Outcome of a `migrate_flat_layout` run.
1498#[derive(Debug, Default, Clone, serde::Serialize)]
1499pub struct MigrationReport {
1500    /// Whether the run was a rehearsal (no disk changes).
1501    pub dry_run: bool,
1502    /// Number of legacy flat-layout files discovered.
1503    pub flat_files_seen: usize,
1504    /// Number of legacy flat files deleted (always 0 when `dry_run`).
1505    pub flat_files_removed: usize,
1506    /// Number of new partition files written under the tenant tree.
1507    pub partitions_written: usize,
1508    /// Total events copied into the new tree (counted in dry-run too).
1509    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        // Add events
1549        for i in 0..10 {
1550            let event = create_test_event(&format!("entity-{i}"));
1551            storage.append_event(event).unwrap();
1552        }
1553
1554        // Flush to disk
1555        storage.flush().unwrap();
1556
1557        // Load back
1558        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        // Add and flush events
1568        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        // Default batch size should be 10,000 as per US-023
1586        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, // Small batch for testing
1609            ..Default::default()
1610        };
1611        let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1612
1613        // 250 events for a single tenant. With per-tenant flush, when the
1614        // tenant's pending batch crosses batch_size we drain the whole
1615        // tenant in one flush — not chunk at exactly batch_size like the
1616        // old global-batch path did. So 250 events triggers exactly one
1617        // size-flush (the appender pushes all 250 onto the tenant's batch
1618        // under one lock, sees length >= 100, schedules a flush which
1619        // drains everything). 0 left pending.
1620        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        // Manual flush is a no-op since nothing's pending.
1630        storage.flush().unwrap();
1631
1632        // All 250 events round-trip through the tenant-partitioned tree.
1633        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, // Very small for testing
1642            ..Default::default()
1643        };
1644        let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1645
1646        // Add 15 events - should auto-flush at 10
1647        for i in 0..15 {
1648            storage
1649                .append_event(create_test_event(&format!("entity-{i}")))
1650                .unwrap();
1651        }
1652
1653        // Should have 5 pending, 10 written
1654        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        // Add some events without reaching batch size
1668        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        // Manually trigger shutdown flush
1677        let flushed = storage.flush_on_shutdown().unwrap();
1678        assert_eq!(flushed, 5);
1679        assert_eq!(storage.pending_count(), 0);
1680
1681        // Verify events are persisted
1682        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        // Flush remaining
1711        storage.flush().unwrap();
1712
1713        // All events should be written
1714        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        // 100 events, single tenant, batch_size=50. Per-tenant flush
1728        // drains the whole tenant on the first size trigger, so this
1729        // produces exactly one size-flush and one batches_written event
1730        // (vs. the pre-tenant world's two).
1731        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    /// Benchmark: Compare single-event writes vs batch writes
1761    /// Run with: cargo test --release -- --ignored test_batch_write_throughput
1762    #[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        // Benchmark batch write
1771        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(); // Flush any remaining
1778        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        // Target: Batch writes should achieve at least 100K events/sec in release mode
1791        // This represents 40%+ improvement over single-event writes
1792        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    /// Benchmark: Single-event write baseline (for comparison)
1800    #[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, // Force flush after each event
1806            ..Default::default()
1807        };
1808        let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1809
1810        let event_count = 1_000; // Fewer events since this is slow
1811
1812        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        // This should be significantly slower than batch writes
1827        // Used as a baseline to demonstrate 40%+ improvement
1828    }
1829
1830    // -----------------------------------------------------------------
1831    // Tests for the recursive parquet walker (Step 1, commit #1: read-side
1832    // bidirectional layout support — see SUSTAINABLE_DATA_STRATEGY.md).
1833    // -----------------------------------------------------------------
1834
1835    /// Helper: write a tiny placeholder parquet file at an arbitrary path so
1836    /// the walker has something concrete to find. We only care that the
1837    /// walker discovers the path, not that the file is loadable here — the
1838    /// load path is exercised by the existing read tests.
1839    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        // Tenant-partitioned shape: storage_dir/<tenant>/<yyyy-mm>/events-*.parquet
1870        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        // Sort places tenant-a files before tenant-b — that's the
1877        // tenant-grouping the docs claim.
1878        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        // The migration window: some tenants have been moved into the tree,
1886        // some flat files still sit at the root. The walker must surface
1887        // both so load_all_events sees every event regardless of where it
1888        // currently lives.
1889        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        // Files that just happen to have "parquet" in the name but no
1906        // .parquet extension stay out — extension-only filter, no name match.
1907        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    /// Build an event whose tenant_id and entity_id we control, so tests
1918    /// can verify per-tenant routing without depending on the helper that
1919    /// hardcodes "default".
1920    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        // End-to-end check that the new write path produces
1936        // <root>/<tenant>/<yyyy-mm>/events-*.parquet — no flat file at the
1937        // root, no cross-tenant mixing.
1938        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        // Path shape: default/<yyyy-mm>/events-*.parquet
1957        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        // yyyy-mm is two digits dash four — loose check, exact month
1961        // varies with wall-clock at test runtime.
1962        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        // Per-tenant flush must not mix tenants into the same Parquet file
1976        // and must put each tenant under its own subdirectory.
1977        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        // Loaded events keep their tenant_id — round-trip preserves which
2003        // tenant each event belonged to.
2004        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        // When one tenant exactly hits batch_size, only that tenant
2020        // flushes; the other tenant keeps its events buffered. Prevents
2021        // one noisy tenant from causing fragmented writes for everyone.
2022        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        // Alice: 5 events → on the 5th, len == batch_size triggers flush
2030        // which drains all 5. Alice ends empty.
2031        for i in 0..5 {
2032            storage
2033                .append_event(event_with_tenant("alice", &format!("a-{i}")))
2034                .unwrap();
2035        }
2036        // Bob: 2 events → still under threshold, stays pending.
2037        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        // Legacy single-component path. Pre-tenant data was always
2072        // commingled with tenant=default, so default is the right fallback.
2073        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            "",         // empty
2097            "..",       // parent traversal
2098            ".",        // current dir
2099            "foo/bar",  // path separator
2100            "foo\\bar", // windows-style separator
2101            "foo bar",  // whitespace
2102            "foo\nbar", // newline
2103            "foo\0bar", // null byte
2104            "tenant?",  // shell glob char
2105            "tenant*",  // shell glob char
2106        ] {
2107            assert!(
2108                sanitize_tenant_id_for_path(bad).is_err(),
2109                "{bad:?} should be rejected"
2110            );
2111        }
2112
2113        // Length cap.
2114        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        // Defence in depth: even if some upstream forgets to validate, the
2131        // sanitizer in flush_tenant catches it. Since append doesn't write
2132        // synchronously, the bad tenant is rejected on the first flush
2133        // attempt. We test that flush surfaces an error rather than
2134        // silently writing somewhere weird.
2135        let temp_dir = TempDir::new().unwrap();
2136        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2137
2138        // append accepts whatever tenant_id the event carries — domain
2139        // construction would normally reject this, but if it slipped
2140        // through, flush should refuse to derive a path from it.
2141        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    // -----------------------------------------------------------------
2154    // Tenant-pruned read tests (Step 1, commit #4).
2155    // -----------------------------------------------------------------
2156
2157    #[test]
2158    fn test_load_events_for_tenant_only_walks_target_subtree() {
2159        // Seed three tenants with distinct event counts. Loading one
2160        // tenant must return only that tenant's events — and the file
2161        // list helper must report only that tenant's files (the strong
2162        // form of "didn't open the others").
2163        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        // The pruned listing must NOT include any bob/carol files — this
2193        // is the property Step 2 will rely on to avoid loading every
2194        // tenant's data on a single-tenant query.
2195        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        // Querying a tenant that has never written must not error — it's
2221        // a normal "no data" case, not a misconfiguration. Important for
2222        // first-query latency on a fresh tenant.
2223        let temp_dir = TempDir::new().unwrap();
2224        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2225
2226        // Seed only alice so the storage_dir isn't empty (rule out the
2227        // empty-dir trivial case).
2228        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        // Path traversal must fail at the API boundary, not after disk
2245        // reads. Same whitelist as the write path.
2246        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        // Flat-layout files at the storage root predate partitioning. A
2261        // tenant-scoped load must not pick them up — the migration tool
2262        // is what relocates them under default/. Until it runs, those
2263        // files are invisible to per-tenant queries (correct behavior:
2264        // the system has no way to tell which tenant they belong to
2265        // beyond "default", and pretending otherwise would mis-attribute
2266        // them).
2267        let temp_dir = TempDir::new().unwrap();
2268        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2269
2270        // Seed a flat-layout file (relocates default/<yyyy-mm>/ → root).
2271        let _flat = seed_flat_layout_file(&storage, 4);
2272
2273        // Querying default returns nothing — the flat file at the root
2274        // isn't under default/.
2275        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        // Sanity: the full loader still sees them via the recursive walk.
2283        let all_events = storage.load_all_events().unwrap();
2284        assert_eq!(all_events.len(), 4);
2285    }
2286
2287    // -----------------------------------------------------------------
2288    // Atomic snapshot write tests (Step 4, commit #1).
2289    // -----------------------------------------------------------------
2290
2291    #[test]
2292    fn test_write_atomic_parquet_emits_file_under_tenant_partition() {
2293        // Happy path: write 3 events for tenant alice, get back a
2294        // path under <root>/alice/<yyyy-mm>/. The .tmp file should
2295        // be gone (rename completed); the final file readable via
2296        // load_events_for_tenant.
2297        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        // Path shape: <root>/alice/<yyyy-mm>/snapshot.alice.range.parquet
2309        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        // Final file must exist; .tmp must NOT.
2320        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        // Loadable via the tenant loader.
2329        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        // Simulate a crashed mid-snapshot: pretend we have a leftover
2358        // events-x.parquet.tmp file in a tenant partition. Cleanup
2359        // must delete it without touching real .parquet files.
2360        let temp_dir = TempDir::new().unwrap();
2361        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2362
2363        // Seed a real parquet via a normal flush, so we have one
2364        // legit file we don't want to delete.
2365        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        // Manufacture an orphan .tmp file in alice's partition.
2375        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        // And one nested deeper, just to confirm recursion.
2382        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        // Real parquet untouched.
2393        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        // Issue #166 regression: a pre-fix checkpoint crash left a 0-byte
2400        // events-*.parquet file with the FINAL extension (no .tmp). The
2401        // cleanup pass must rename it aside rather than delete, so an
2402        // operator can inspect.
2403        let temp_dir = TempDir::new().unwrap();
2404        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2405
2406        // Seed a real parquet alongside, to confirm we don't molest healthy data.
2407        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        // Drop a 0-byte parquet next to the healthy one (the issue's failure mode).
2419        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        // 0-byte file is gone from its original name, but a quarantined
2427        // sibling exists.
2428        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        // Healthy file untouched.
2446        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        // Issue #166 defensive layer: a single bad file must not abort
2455        // the whole load. With t-d8b1's match-and-continue, the healthy
2456        // events still come back and the bad file is logged-and-skipped.
2457        let temp_dir = TempDir::new().unwrap();
2458        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2459
2460        // Two healthy events, flushed.
2461        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        // Now drop a 0-byte parquet alongside (sneaking it in AFTER
2469        // construction so cleanup_partial_writes doesn't rename it
2470        // before we get to test the load path).
2471        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        // Same defensive layer for the lazy-load path. A bad file in
2489        // the tenant subtree must not poison subsequent queries.
2490        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        // Direct verification of t-c1d3: a successful flush_tenant
2514        // leaves exactly one healthy *.parquet — no .tmp survivors,
2515        // no 0-byte files.
2516        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        // Drop a stale .tmp into a directory, then construct a
2557        // fresh ParquetStorage on it. The constructor must clean
2558        // it up.
2559        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    // -----------------------------------------------------------------
2572    // Migration tests (Step 1, commit #3: flat → tenant-tree migration).
2573    // -----------------------------------------------------------------
2574
2575    /// Helper: produce a flat-layout Parquet file at the storage root,
2576    /// matching what pre-#2 deploys wrote. Uses the existing flush path
2577    /// briefly and then relocates the resulting file from
2578    /// default/<yyyy-mm>/ back up to the root, simulating the legacy
2579    /// state.
2580    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        // create_test_event uses tenant="default", so the just-flushed file
2589        // landed under <root>/default/<yyyy-mm>/. Find the newest file in
2590        // that subtree to avoid picking up files from other tenants seeded
2591        // by the test before us.
2592        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        // Best-effort cleanup of the now-empty intermediate dirs so the
2603        // migration tool only sees the flat file. (`remove_dir` succeeds
2604        // only on empty dirs, which is exactly the safety we want here.)
2605        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        // Second run sees no flat files at the root, so it's a no-op —
2677        // events do not duplicate even if an operator runs the tool twice.
2678        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        // Mixed state: a tenant tree already exists alongside one flat file.
2690        // Migration must touch only the flat file.
2691        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}