Skip to main content

allsource_core/infrastructure/persistence/
storage.rs

1use crate::{
2    domain::entities::Event,
3    error::{AllSourceError, Result},
4};
5use arrow::{
6    array::{
7        Array, ArrayRef, StringBuilder, TimestampMicrosecondArray, TimestampMicrosecondBuilder,
8        UInt64Builder,
9    },
10    datatypes::{DataType, Field, Schema, TimeUnit},
11    record_batch::RecordBatch,
12};
13use parquet::{arrow::ArrowWriter, file::properties::WriterProperties};
14use std::{
15    collections::HashMap,
16    fs::{self, File},
17    path::{Path, PathBuf},
18    sync::{
19        Arc, Mutex,
20        atomic::{AtomicU64, Ordering},
21    },
22    time::{Duration, Instant},
23};
24
25/// Default batch size for Parquet writes (10,000 events as per US-023)
26pub const DEFAULT_BATCH_SIZE: usize = 10_000;
27
28/// Default flush timeout in milliseconds
29pub const DEFAULT_FLUSH_TIMEOUT_MS: u64 = 5_000;
30
31/// Configuration for ParquetStorage batch processing
32#[derive(Debug, Clone)]
33pub struct ParquetStorageConfig {
34    /// Batch size before automatic flush (default: 10,000)
35    pub batch_size: usize,
36    /// Timeout before flushing partial batch (default: 5 seconds)
37    pub flush_timeout: Duration,
38    /// Compression codec for Parquet files
39    pub compression: parquet::basic::Compression,
40}
41
42impl Default for ParquetStorageConfig {
43    fn default() -> Self {
44        Self {
45            batch_size: DEFAULT_BATCH_SIZE,
46            flush_timeout: Duration::from_millis(DEFAULT_FLUSH_TIMEOUT_MS),
47            compression: parquet::basic::Compression::SNAPPY,
48        }
49    }
50}
51
52impl ParquetStorageConfig {
53    /// High-throughput configuration optimized for large batch writes
54    pub fn high_throughput() -> Self {
55        Self {
56            batch_size: 50_000,
57            flush_timeout: Duration::from_secs(10),
58            compression: parquet::basic::Compression::SNAPPY,
59        }
60    }
61
62    /// Low-latency configuration for smaller, more frequent writes
63    pub fn low_latency() -> Self {
64        Self {
65            batch_size: 1_000,
66            flush_timeout: Duration::from_secs(1),
67            compression: parquet::basic::Compression::SNAPPY,
68        }
69    }
70}
71
72/// Statistics for batch write operations
73#[derive(Debug, Clone, Default)]
74pub struct BatchWriteStats {
75    /// Total batches written
76    pub batches_written: u64,
77    /// Total events written
78    pub events_written: u64,
79    /// Total bytes written
80    pub bytes_written: u64,
81    /// Average batch size
82    pub avg_batch_size: f64,
83    /// Events per second (throughput)
84    pub events_per_sec: f64,
85    /// Total write time in nanoseconds
86    pub total_write_time_ns: u64,
87    /// Number of timeout-triggered flushes
88    pub timeout_flushes: u64,
89    /// Number of size-triggered flushes
90    pub size_flushes: u64,
91}
92
93/// Result of a batch write operation
94#[derive(Debug, Clone)]
95pub struct BatchWriteResult {
96    /// Number of events written
97    pub events_written: usize,
98    /// Number of batches flushed to disk
99    pub batches_flushed: usize,
100    /// Total duration of the write operation
101    pub duration: Duration,
102    /// Write throughput in events per second
103    pub events_per_sec: f64,
104}
105
106/// Parquet-based persistent storage for events with batch processing
107///
108/// Features:
109/// - Configurable batch size (default: 10,000 events per US-023)
110/// - Timeout-based flushing for partial batches
111/// - Thread-safe batch accumulation
112/// - SNAPPY compression for efficient storage
113/// - Automatic flush on shutdown via Drop
114pub struct ParquetStorage {
115    /// Base directory for storing parquet files
116    storage_dir: PathBuf,
117
118    /// Buffered events keyed by tenant_id. Each tenant accumulates its own
119    /// batch and flushes independently into its partition under
120    /// `storage_dir/<tenant_id>/<yyyy-mm>/`. Single outer mutex protects the
121    /// whole map: lookup is O(1), tenant cardinality is low (single digits
122    /// today; bounded by Step 3's cache budget later), so contention is
123    /// fine. We keep the mutex held only for the push, not for disk I/O —
124    /// flush takes ownership of a tenant's batch via remove() and writes
125    /// after the lock is released.
126    current_batches: Mutex<HashMap<String, Vec<Event>>>,
127
128    /// Configuration
129    config: ParquetStorageConfig,
130
131    /// Schema for Arrow/Parquet
132    schema: Arc<Schema>,
133
134    /// Last flush timestamp for timeout tracking
135    last_flush_time: Mutex<Instant>,
136
137    /// Statistics tracking
138    batches_written: AtomicU64,
139    events_written: AtomicU64,
140    bytes_written: AtomicU64,
141    total_write_time_ns: AtomicU64,
142    timeout_flushes: AtomicU64,
143    size_flushes: AtomicU64,
144}
145
146impl ParquetStorage {
147    /// Create a new ParquetStorage with default configuration (10,000 event batches)
148    pub fn new(storage_dir: impl AsRef<Path>) -> Result<Self> {
149        Self::with_config(storage_dir, ParquetStorageConfig::default())
150    }
151
152    /// Create a new ParquetStorage with custom configuration
153    pub fn with_config(
154        storage_dir: impl AsRef<Path>,
155        config: ParquetStorageConfig,
156    ) -> Result<Self> {
157        let storage_dir = storage_dir.as_ref().to_path_buf();
158
159        // Create storage directory if it doesn't exist
160        fs::create_dir_all(&storage_dir).map_err(|e| {
161            AllSourceError::StorageError(format!("Failed to create storage directory: {e}"))
162        })?;
163
164        // Define Arrow schema for events
165        let schema = Arc::new(Schema::new(vec![
166            Field::new("event_id", DataType::Utf8, false),
167            Field::new("event_type", DataType::Utf8, false),
168            Field::new("entity_id", DataType::Utf8, false),
169            Field::new("payload", DataType::Utf8, false),
170            Field::new(
171                "timestamp",
172                DataType::Timestamp(TimeUnit::Microsecond, None),
173                false,
174            ),
175            Field::new("metadata", DataType::Utf8, true),
176            Field::new("version", DataType::UInt64, false),
177        ]));
178
179        let storage = Self {
180            storage_dir,
181            current_batches: Mutex::new(HashMap::new()),
182            config,
183            schema,
184            last_flush_time: Mutex::new(Instant::now()),
185            batches_written: AtomicU64::new(0),
186            events_written: AtomicU64::new(0),
187            bytes_written: AtomicU64::new(0),
188            total_write_time_ns: AtomicU64::new(0),
189            timeout_flushes: AtomicU64::new(0),
190            size_flushes: AtomicU64::new(0),
191        };
192
193        // Boot-time crash recovery: heal *.parquet.tmp files from
194        // crashed atomic writes (snapshot/compaction or, post-#166-fix,
195        // checkpoint flushes) and quarantine any 0-byte *.parquet files
196        // left over from pre-fix checkpoint crashes (issue #166).
197        match storage.cleanup_partial_writes() {
198            Ok(0) => {}
199            Ok(n) => tracing::warn!(
200                "cleanup_partial_writes acted on {n} crash-detritus file(s) on boot — \
201                 see preceding logs for per-file detail"
202            ),
203            Err(e) => tracing::error!("cleanup_partial_writes failed on boot: {e}"),
204        }
205
206        Ok(storage)
207    }
208
209    /// Create storage with legacy batch size (1000) for backward compatibility
210    #[deprecated(note = "Use new() or with_config() instead - default batch size is now 10,000")]
211    pub fn with_legacy_batch_size(storage_dir: impl AsRef<Path>) -> Result<Self> {
212        Self::with_config(
213            storage_dir,
214            ParquetStorageConfig {
215                batch_size: 1000,
216                ..Default::default()
217            },
218        )
219    }
220
221    /// Add an event to the current batch
222    ///
223    /// Events are routed to a per-tenant batch keyed by `event.tenant_id_str()`.
224    /// A tenant's batch is buffered until any of:
225    /// - That tenant's batch hits the configured `batch_size` (default 10,000)
226    ///   — flushes only that tenant, not the whole world
227    /// - The flush timeout elapses — flushes every tenant with pending data
228    /// - `flush()` is called explicitly
229    /// - The process shuts down
230    #[cfg_attr(feature = "hotpath", hotpath::measure)]
231    pub fn append_event(&self, event: Event) -> Result<()> {
232        let tenant = event.tenant_id_str().to_string();
233        let should_flush_tenant = {
234            let mut batches = self.current_batches.lock().unwrap();
235            let entry = batches.entry(tenant.clone()).or_default();
236            entry.push(event);
237            entry.len() >= self.config.batch_size
238        };
239
240        if should_flush_tenant {
241            self.size_flushes.fetch_add(1, Ordering::Relaxed);
242            self.flush_tenant(&tenant)?;
243        }
244
245        Ok(())
246    }
247
248    /// Add multiple events to the batch (optimized batch insertion)
249    ///
250    /// Preferred entry point for high-throughput ingestion. Events are
251    /// grouped by tenant under a single mutex acquisition and any tenant
252    /// that crosses `batch_size` is flushed on the spot.
253    #[cfg_attr(feature = "hotpath", hotpath::measure)]
254    pub fn batch_write(&self, events: Vec<Event>) -> Result<BatchWriteResult> {
255        let start = Instant::now();
256        let event_count = events.len();
257
258        // Pre-group by tenant to keep the lock window short — one acquire,
259        // one extend per tenant, decide which tenants are over threshold.
260        let mut grouped: HashMap<String, Vec<Event>> = HashMap::new();
261        for event in events {
262            grouped
263                .entry(event.tenant_id_str().to_string())
264                .or_default()
265                .push(event);
266        }
267
268        let mut tenants_to_flush: Vec<String> = Vec::new();
269        {
270            let mut batches = self.current_batches.lock().unwrap();
271            for (tenant, mut new_events) in grouped {
272                let entry = batches.entry(tenant.clone()).or_default();
273                entry.append(&mut new_events);
274                if entry.len() >= self.config.batch_size {
275                    tenants_to_flush.push(tenant);
276                }
277            }
278        }
279
280        let mut batches_flushed = 0;
281        for tenant in tenants_to_flush {
282            self.size_flushes.fetch_add(1, Ordering::Relaxed);
283            self.flush_tenant(&tenant)?;
284            batches_flushed += 1;
285        }
286
287        let duration = start.elapsed();
288
289        Ok(BatchWriteResult {
290            events_written: event_count,
291            batches_flushed,
292            duration,
293            events_per_sec: event_count as f64 / duration.as_secs_f64(),
294        })
295    }
296
297    /// Check if a timeout-based flush is needed and perform it
298    ///
299    /// Call this periodically (e.g., from a background task) to ensure
300    /// partial batches are flushed within the configured timeout. When
301    /// triggered, every tenant with pending events flushes — the timer is
302    /// global, not per-tenant, so a slow-trickle tenant doesn't get
303    /// stranded waiting for its own batch to fill.
304    #[cfg_attr(feature = "hotpath", hotpath::measure)]
305    pub fn check_timeout_flush(&self) -> Result<bool> {
306        let should_flush = {
307            let last_flush = self.last_flush_time.lock().unwrap();
308            let batches = self.current_batches.lock().unwrap();
309            let any_pending = batches.values().any(|v| !v.is_empty());
310            any_pending && last_flush.elapsed() >= self.config.flush_timeout
311        };
312
313        if should_flush {
314            self.timeout_flushes.fetch_add(1, Ordering::Relaxed);
315            self.flush()?;
316            Ok(true)
317        } else {
318            Ok(false)
319        }
320    }
321
322    /// Flush every tenant's pending batch to its partition.
323    ///
324    /// Thread-safe: callable from any thread. A snapshot of which tenants
325    /// have pending data is taken under a short lock; each tenant is then
326    /// flushed individually with its own lock cycle, so disk I/O for one
327    /// tenant doesn't block writes against another.
328    #[cfg_attr(feature = "hotpath", hotpath::measure)]
329    pub fn flush(&self) -> Result<()> {
330        let tenants: Vec<String> = {
331            let batches = self.current_batches.lock().unwrap();
332            batches
333                .iter()
334                .filter(|(_, v)| !v.is_empty())
335                .map(|(k, _)| k.clone())
336                .collect()
337        };
338        if tenants.is_empty() {
339            return Ok(());
340        }
341        for tenant in tenants {
342            self.flush_tenant(&tenant)?;
343        }
344        Ok(())
345    }
346
347    /// Flush a single tenant's pending events into its partition file.
348    ///
349    /// File path: `storage_dir/<sanitized_tenant_id>/<yyyy-mm>/events-<ts>-<uuid>.parquet`.
350    /// `<yyyy-mm>` is taken from the wall-clock at flush time (matching the
351    /// pre-tenant filename's timestamp semantics) rather than from
352    /// individual event timestamps — keeps each flush to a single output
353    /// file even when buffered events span months.
354    ///
355    /// The file is written atomically via `write_record_batch_atomic`:
356    /// crash mid-write leaves a `.parquet.tmp` file (cleaned on next boot)
357    /// rather than a 0-byte `.parquet` file with the final name. Issue
358    /// #166 — pre-fix, a SIGKILL inside this function bricked subsequent
359    /// loads because `load_all_events` aborted on the first unreadable
360    /// file.
361    fn flush_tenant(&self, tenant_id: &str) -> Result<()> {
362        let events_to_write = {
363            let mut batches = self.current_batches.lock().unwrap();
364            match batches.get_mut(tenant_id) {
365                Some(v) if !v.is_empty() => std::mem::take(v),
366                _ => return Ok(()),
367            }
368        };
369
370        let batch_count = events_to_write.len();
371        let start = Instant::now();
372
373        let written = self.write_tenant_events(tenant_id, &events_to_write);
374        let (file_path, file_metadata) = match written {
375            Ok(written) => written,
376            Err(e) => {
377                // The WAL may be retired by the next successful checkpoint, so
378                // dropping these here would make the failure permanent (#287).
379                let mut batches = self.current_batches.lock().unwrap();
380                let pending = batches.entry(tenant_id.to_string()).or_default();
381                let arrived_during_write = std::mem::replace(pending, events_to_write);
382                pending.extend(arrived_during_write);
383                return Err(e);
384            }
385        };
386
387        let duration = start.elapsed();
388
389        self.batches_written.fetch_add(1, Ordering::Relaxed);
390        self.events_written
391            .fetch_add(batch_count as u64, Ordering::Relaxed);
392        if let Some(size) = file_metadata
393            .row_groups()
394            .first()
395            .map(parquet::file::metadata::RowGroupMetaData::total_byte_size)
396        {
397            self.bytes_written.fetch_add(size as u64, Ordering::Relaxed);
398        }
399        self.total_write_time_ns
400            .fetch_add(duration.as_nanos() as u64, Ordering::Relaxed);
401
402        {
403            let mut last_flush = self.last_flush_time.lock().unwrap();
404            *last_flush = Instant::now();
405        }
406
407        tracing::info!(
408            "Wrote {} events for tenant={} to {} in {:?}",
409            batch_count,
410            tenant_id,
411            file_path.display(),
412            duration
413        );
414
415        Ok(())
416    }
417
418    fn write_tenant_events(
419        &self,
420        tenant_id: &str,
421        events: &[Event],
422    ) -> Result<(PathBuf, parquet::file::metadata::ParquetMetaData)> {
423        let record_batch = self.events_to_record_batch(events)?;
424
425        let now = chrono::Utc::now();
426        let partition_dir = partition_path_for_tenant(&self.storage_dir, tenant_id, now)?;
427        fs::create_dir_all(&partition_dir).map_err(|e| {
428            AllSourceError::StorageError(format!(
429                "Failed to create tenant partition {}: {e}",
430                partition_dir.display()
431            ))
432        })?;
433        let file_stem = format!(
434            "events-{}-{}",
435            now.format("%Y%m%d-%H%M%S%3f"),
436            uuid::Uuid::new_v4().as_simple()
437        );
438
439        tracing::info!(
440            "Flushing {} events for tenant={} to {}/{}.parquet",
441            events.len(),
442            tenant_id,
443            partition_dir.display(),
444            file_stem
445        );
446
447        self.write_record_batch_atomic(&partition_dir, &file_stem, &record_batch)
448    }
449
450    /// Write a single record batch atomically to
451    /// `<partition_dir>/<file_stem>.parquet`, returning the final path
452    /// and the parquet writer's metadata (used for byte-size metrics by
453    /// the checkpoint flush path).
454    ///
455    /// Crash-safety contract — same four steps as `write_atomic_parquet`:
456    /// 1. Write to `<final_path>.tmp` first.
457    /// 2. fsync the .tmp file so the data is durably on disk.
458    /// 3. Rename `.tmp` → final name (atomic POSIX rename).
459    /// 4. fsync the parent directory so the rename survives crash.
460    ///
461    /// Caller is responsible for choosing `partition_dir` (and creating
462    /// it). On any failure mid-way, the `.tmp` file gets cleaned up by
463    /// `cleanup_partial_writes` on next boot. The final file appears
464    /// atomically — readers either see the old state (no file) or the
465    /// complete new file, never a half-written one.
466    fn write_record_batch_atomic(
467        &self,
468        partition_dir: &Path,
469        file_stem: &str,
470        record_batch: &RecordBatch,
471    ) -> Result<(PathBuf, parquet::file::metadata::ParquetMetaData)> {
472        let final_path = partition_dir.join(format!("{file_stem}.parquet"));
473        let tmp_path = partition_dir.join(format!("{file_stem}.parquet.tmp"));
474
475        // 1. Write to .tmp.
476        let metadata = {
477            let file = File::create(&tmp_path).map_err(|e| {
478                AllSourceError::StorageError(format!(
479                    "Failed to create parquet tmp file {}: {e}",
480                    tmp_path.display()
481                ))
482            })?;
483
484            let props = WriterProperties::builder()
485                .set_compression(self.config.compression)
486                .build();
487
488            let mut writer = ArrowWriter::try_new(file, self.schema.clone(), Some(props))?;
489            writer.write(record_batch)?;
490            // close() consumes writer and returns the file metadata; the
491            // underlying File is flushed and closed here.
492            writer.close()?
493        };
494
495        // 2. fsync the .tmp file.
496        let tmp_file = File::open(&tmp_path).map_err(|e| {
497            AllSourceError::StorageError(format!(
498                "Failed to reopen parquet tmp for fsync {}: {e}",
499                tmp_path.display()
500            ))
501        })?;
502        tmp_file.sync_all().map_err(|e| {
503            AllSourceError::StorageError(format!("fsync on parquet tmp failed: {e}"))
504        })?;
505        drop(tmp_file);
506
507        // 3. Atomic rename.
508        fs::rename(&tmp_path, &final_path).map_err(|e| {
509            AllSourceError::StorageError(format!(
510                "Failed to rename {} → {}: {e}",
511                tmp_path.display(),
512                final_path.display()
513            ))
514        })?;
515
516        // 4. fsync the parent directory so the rename survives crash.
517        // Linux fsync-on-dir is the canonical way to make a rename
518        // durable; macOS no-ops the dir fsync but doesn't error.
519        if let Ok(dir) = File::open(partition_dir) {
520            let _ = dir.sync_all();
521        }
522
523        Ok((final_path, metadata))
524    }
525
526    /// Atomically write `events` to a Parquet file under the tenant's
527    /// partition. Step 4 of the sustainable data strategy uses this to
528    /// emit per-tenant snapshot/compaction files (`snapshot.<tenant>.<from>-<to>`)
529    /// without risking partial files on crash.
530    ///
531    /// Crash-safety contract:
532    /// 1. Write to `<final_path>.tmp` first.
533    /// 2. fsync the .tmp file so the data is durably on disk.
534    /// 3. Rename `.tmp` → final name (atomic POSIX rename).
535    /// 4. fsync the parent directory so the rename is durable.
536    ///
537    /// On any failure mid-way, the .tmp file gets cleaned up by
538    /// `cleanup_partial_writes` on next boot. The final file appears
539    /// atomically — readers either see the old state (no file) or
540    /// the complete new file, never a half-written one.
541    ///
542    /// `file_stem` is the filename without extension or partition path
543    /// — caller-controlled so the snapshot naming convention
544    /// (`snapshot.<tenant>.<from>-<to>`) lives in the compaction layer,
545    /// not here. The `.parquet` extension is appended automatically.
546    ///
547    /// Returns the final (post-rename) path so the caller can
548    /// confirm the file landed where expected.
549    pub fn write_atomic_parquet(
550        &self,
551        tenant_id: &str,
552        file_stem: &str,
553        events: &[Event],
554    ) -> Result<PathBuf> {
555        if events.is_empty() {
556            return Err(AllSourceError::StorageError(
557                "write_atomic_parquet called with empty event slice".to_string(),
558            ));
559        }
560        // Anchor partition by the earliest event's month — keeps the
561        // file in the same yyyy-mm bucket as the data it represents.
562        // Callers that span multiple months can use the earliest
563        // event's month and trust the recursive walker to find it.
564        let anchor_ts = events
565            .iter()
566            .map(|e| e.timestamp)
567            .min()
568            .unwrap_or_else(chrono::Utc::now);
569        let partition_dir = partition_path_for_tenant(&self.storage_dir, tenant_id, anchor_ts)?;
570        fs::create_dir_all(&partition_dir).map_err(|e| {
571            AllSourceError::StorageError(format!(
572                "Failed to create tenant partition {}: {e}",
573                partition_dir.display()
574            ))
575        })?;
576
577        let record_batch = self.events_to_record_batch(events)?;
578        let (final_path, _meta) =
579            self.write_record_batch_atomic(&partition_dir, file_stem, &record_batch)?;
580
581        tracing::info!(
582            tenant_id = tenant_id,
583            file = %final_path.display(),
584            event_count = events.len(),
585            "wrote atomic snapshot file"
586        );
587
588        Ok(final_path)
589    }
590
591    /// Sweep the storage tree for crash detritus and heal it on boot.
592    /// Called on boot from `ParquetStorage::new` so every cold start
593    /// starts clean. Two kinds of leftovers are handled:
594    ///
595    /// 1. **`*.parquet.tmp` files** — a snapshot or flush write that
596    ///    crashed between fsync and rename. Safe to delete: the data is
597    ///    either still in the WAL (for checkpoint flushes) or in the
598    ///    constituent raw files (for snapshot/compaction, which only
599    ///    deletes the inputs AFTER the rename succeeds). Either way, the
600    ///    failed write will be retried.
601    ///
602    /// 2. **0-byte `*.parquet` files** — issue #166. Pre-fix releases
603    ///    (≤ 0.20.0 for the checkpoint path) wrote directly to the
604    ///    final filename, so SIGKILL between `File::create` and the
605    ///    first flush left a 0-byte `events-*.parquet` that bricked
606    ///    every subsequent load. Rather than delete, we rename to
607    ///    `<original>.parquet.corrupt-<unix-ts>` so an operator can
608    ///    inspect/forensicate before destroying the evidence — most
609    ///    will just `rm` after seeing the size.
610    ///
611    /// Returns the total number of files acted on, for boot-log
612    /// observability.
613    pub fn cleanup_partial_writes(&self) -> Result<usize> {
614        let mut acted = 0usize;
615        let mut stack: Vec<PathBuf> = vec![self.storage_dir.clone()];
616        while let Some(dir) = stack.pop() {
617            let Ok(entries) = fs::read_dir(&dir) else {
618                continue;
619            };
620            for entry in entries.flatten() {
621                let path = entry.path();
622                let Ok(ft) = entry.file_type() else { continue };
623                if ft.is_dir() {
624                    stack.push(path);
625                    continue;
626                }
627                if !ft.is_file() {
628                    continue;
629                }
630                let path_str = path.to_string_lossy();
631                if path_str.ends_with(".parquet.tmp") {
632                    match fs::remove_file(&path) {
633                        Ok(()) => {
634                            tracing::warn!(
635                                file = %path.display(),
636                                "cleaned up orphan snapshot tmp file (crash recovery)"
637                            );
638                            acted += 1;
639                        }
640                        Err(e) => {
641                            tracing::error!(
642                                file = %path.display(),
643                                "failed to remove orphan snapshot tmp file: {e}"
644                            );
645                        }
646                    }
647                } else if path_str.ends_with(".parquet")
648                    && fs::metadata(&path).is_ok_and(|m| m.len() == 0)
649                {
650                    // Issue #166: pre-fix checkpoint write left a 0-byte
651                    // .parquet file with the final name on crash. Quarantine
652                    // (rename) rather than delete — operator can inspect.
653                    let ts = chrono::Utc::now().timestamp();
654                    let quarantine_path = path.with_extension(format!("parquet.corrupt-{ts}"));
655                    match fs::rename(&path, &quarantine_path) {
656                        Ok(()) => {
657                            tracing::error!(
658                                from = %path.display(),
659                                to = %quarantine_path.display(),
660                                "quarantined 0-byte parquet file (issue #166 pre-fix crash). \
661                                 Operator: inspect and rm if not needed."
662                            );
663                            acted += 1;
664                        }
665                        Err(e) => {
666                            tracing::error!(
667                                file = %path.display(),
668                                "failed to quarantine 0-byte parquet file: {e}"
669                            );
670                        }
671                    }
672                }
673            }
674        }
675        Ok(acted)
676    }
677
678    /// Force flush any remaining events (for shutdown handling).
679    ///
680    /// Sums pending counts across every tenant's batch so the caller can
681    /// log "we flushed N events on shutdown" without caring about
682    /// per-tenant breakdown.
683    #[cfg_attr(feature = "hotpath", hotpath::measure)]
684    pub fn flush_on_shutdown(&self) -> Result<usize> {
685        let total_pending: usize = {
686            let batches = self.current_batches.lock().unwrap();
687            batches.values().map(Vec::len).sum()
688        };
689
690        if total_pending > 0 {
691            tracing::info!(
692                "Shutdown: flushing {} pending events across all tenants",
693                total_pending
694            );
695            self.flush()?;
696        }
697
698        Ok(total_pending)
699    }
700
701    /// Get batch write statistics
702    pub fn batch_stats(&self) -> BatchWriteStats {
703        let batches = self.batches_written.load(Ordering::Relaxed);
704        let events = self.events_written.load(Ordering::Relaxed);
705        let bytes = self.bytes_written.load(Ordering::Relaxed);
706        let time_ns = self.total_write_time_ns.load(Ordering::Relaxed);
707
708        let time_secs = time_ns as f64 / 1_000_000_000.0;
709
710        BatchWriteStats {
711            batches_written: batches,
712            events_written: events,
713            bytes_written: bytes,
714            avg_batch_size: if batches > 0 {
715                events as f64 / batches as f64
716            } else {
717                0.0
718            },
719            events_per_sec: if time_secs > 0.0 {
720                events as f64 / time_secs
721            } else {
722                0.0
723            },
724            total_write_time_ns: time_ns,
725            timeout_flushes: self.timeout_flushes.load(Ordering::Relaxed),
726            size_flushes: self.size_flushes.load(Ordering::Relaxed),
727        }
728    }
729
730    /// Total pending events across all tenant batches.
731    pub fn pending_count(&self) -> usize {
732        self.current_batches
733            .lock()
734            .unwrap()
735            .values()
736            .map(Vec::len)
737            .sum()
738    }
739
740    /// Get configured batch size
741    pub fn batch_size(&self) -> usize {
742        self.config.batch_size
743    }
744
745    /// Get configured flush timeout
746    pub fn flush_timeout(&self) -> Duration {
747        self.config.flush_timeout
748    }
749
750    /// Convert events to Arrow RecordBatch
751    #[cfg_attr(feature = "hotpath", hotpath::measure)]
752    fn events_to_record_batch(&self, events: &[Event]) -> Result<RecordBatch> {
753        let mut event_id_builder = StringBuilder::new();
754        let mut event_type_builder = StringBuilder::new();
755        let mut entity_id_builder = StringBuilder::new();
756        let mut payload_builder = StringBuilder::new();
757        let mut timestamp_builder = TimestampMicrosecondBuilder::new();
758        let mut metadata_builder = StringBuilder::new();
759        let mut version_builder = UInt64Builder::new();
760
761        for event in events {
762            event_id_builder.append_value(event.id.to_string());
763            event_type_builder.append_value(event.event_type_str());
764            entity_id_builder.append_value(event.entity_id_str());
765            payload_builder.append_value(serde_json::to_string(&event.payload)?);
766
767            // Convert timestamp to microseconds
768            let timestamp_micros = event.timestamp.timestamp_micros();
769            timestamp_builder.append_value(timestamp_micros);
770
771            if let Some(ref metadata) = event.metadata {
772                metadata_builder.append_value(serde_json::to_string(metadata)?);
773            } else {
774                metadata_builder.append_null();
775            }
776
777            version_builder.append_value(event.version as u64);
778        }
779
780        let arrays: Vec<ArrayRef> = vec![
781            Arc::new(event_id_builder.finish()),
782            Arc::new(event_type_builder.finish()),
783            Arc::new(entity_id_builder.finish()),
784            Arc::new(payload_builder.finish()),
785            Arc::new(timestamp_builder.finish()),
786            Arc::new(metadata_builder.finish()),
787            Arc::new(version_builder.finish()),
788        ];
789
790        let record_batch = RecordBatch::try_new(self.schema.clone(), arrays)?;
791
792        Ok(record_batch)
793    }
794
795    /// Load events from all Parquet files under the storage directory.
796    ///
797    /// Walks the tree recursively so both layouts work: legacy flat
798    /// (`storage_dir/events-*.parquet`) and the tenant-partitioned tree
799    /// introduced by the data-strategy work (`storage_dir/<tenant>/<yyyy-mm>/
800    /// events-*.parquet`). The two coexist on disk during migration.
801    #[cfg_attr(feature = "hotpath", hotpath::measure)]
802    pub fn load_all_events(&self) -> Result<Vec<Event>> {
803        let parquet_files = find_parquet_files_recursive(&self.storage_dir)?;
804
805        let mut all_events = Vec::with_capacity(parquet_files.len() * self.config.batch_size);
806        let mut skipped = 0usize;
807        for file_path in parquet_files {
808            tracing::info!("Loading events from {}", file_path.display());
809            let tenant_id = tenant_id_from_path(&self.storage_dir, &file_path);
810            // Issue #166: a single corrupt/0-byte parquet file (crash mid-write
811            // on pre-fix releases) would propagate up and abort the whole load,
812            // bricking access to every healthy file alongside it. Log and skip
813            // instead so one bad file can't take the rest down.
814            match self.load_events_from_file(&file_path, &tenant_id) {
815                Ok(file_events) => all_events.extend(file_events),
816                Err(e) => {
817                    tracing::error!(
818                        file = %file_path.display(),
819                        error = %e,
820                        "Skipping unreadable parquet file — other files will still load. \
821                         Likely a 0-byte or truncated file from an unclean shutdown; \
822                         inspect and remove manually after confirming."
823                    );
824                    skipped += 1;
825                }
826            }
827        }
828
829        if skipped > 0 {
830            tracing::warn!(
831                "Loaded {} events from storage; skipped {} unreadable file(s)",
832                all_events.len(),
833                skipped
834            );
835        } else {
836            tracing::info!("Loaded {} total events from storage", all_events.len());
837        }
838
839        Ok(all_events)
840    }
841
842    /// Load events from a single Parquet file. `tenant_id` is the value to
843    /// stamp onto each loaded event — derived from the file's location in
844    /// the tree by `load_all_events`. The Parquet schema doesn't include
845    /// tenant_id today (path is the source of truth), so this is how
846    /// per-tenant identity survives the round trip.
847    #[cfg_attr(feature = "hotpath", hotpath::measure)]
848    fn load_events_from_file(&self, file_path: &Path, tenant_id: &str) -> Result<Vec<Event>> {
849        use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
850
851        let file = File::open(file_path).map_err(|e| {
852            AllSourceError::StorageError(format!("Failed to open parquet file: {e}"))
853        })?;
854
855        let builder = ParquetRecordBatchReaderBuilder::try_new(file)?;
856        let mut reader = builder.build()?;
857
858        let mut events = Vec::new();
859
860        while let Some(Ok(batch)) = reader.next() {
861            let batch_events = self.record_batch_to_events(&batch, tenant_id)?;
862            events.extend(batch_events);
863        }
864
865        Ok(events)
866    }
867
868    /// Public wrapper around the internal single-file loader. Used
869    /// by the per-tenant compaction (Step 4) which needs to read a
870    /// specific candidate set rather than the whole tenant subtree.
871    /// Caller passes the tenant_id explicitly because the schema
872    /// doesn't carry it.
873    pub fn load_events_from_file_path(
874        &self,
875        file_path: &Path,
876        tenant_id: &str,
877    ) -> Result<Vec<Event>> {
878        self.load_events_from_file(file_path, tenant_id)
879    }
880
881    /// Convert Arrow RecordBatch back to events. `tenant_id` is stamped onto
882    /// each reconstructed event — the schema doesn't carry it today, so the
883    /// caller passes the value derived from the file path.
884    #[cfg_attr(feature = "hotpath", hotpath::measure)]
885    fn record_batch_to_events(&self, batch: &RecordBatch, tenant_id: &str) -> Result<Vec<Event>> {
886        let event_ids = batch
887            .column(0)
888            .as_any()
889            .downcast_ref::<arrow::array::StringArray>()
890            .ok_or_else(|| AllSourceError::StorageError("Invalid event_id column".to_string()))?;
891
892        let event_types = batch
893            .column(1)
894            .as_any()
895            .downcast_ref::<arrow::array::StringArray>()
896            .ok_or_else(|| AllSourceError::StorageError("Invalid event_type column".to_string()))?;
897
898        let entity_ids = batch
899            .column(2)
900            .as_any()
901            .downcast_ref::<arrow::array::StringArray>()
902            .ok_or_else(|| AllSourceError::StorageError("Invalid entity_id column".to_string()))?;
903
904        let payloads = batch
905            .column(3)
906            .as_any()
907            .downcast_ref::<arrow::array::StringArray>()
908            .ok_or_else(|| AllSourceError::StorageError("Invalid payload column".to_string()))?;
909
910        let timestamps = batch
911            .column(4)
912            .as_any()
913            .downcast_ref::<TimestampMicrosecondArray>()
914            .ok_or_else(|| AllSourceError::StorageError("Invalid timestamp column".to_string()))?;
915
916        let metadatas = batch
917            .column(5)
918            .as_any()
919            .downcast_ref::<arrow::array::StringArray>()
920            .ok_or_else(|| AllSourceError::StorageError("Invalid metadata column".to_string()))?;
921
922        let versions = batch
923            .column(6)
924            .as_any()
925            .downcast_ref::<arrow::array::UInt64Array>()
926            .ok_or_else(|| AllSourceError::StorageError("Invalid version column".to_string()))?;
927
928        let mut events = Vec::new();
929
930        for i in 0..batch.num_rows() {
931            let id = uuid::Uuid::parse_str(event_ids.value(i))
932                .map_err(|e| AllSourceError::StorageError(format!("Invalid UUID: {e}")))?;
933
934            let timestamp = chrono::DateTime::from_timestamp_micros(timestamps.value(i))
935                .ok_or_else(|| AllSourceError::StorageError("Invalid timestamp".to_string()))?;
936
937            let metadata = if metadatas.is_null(i) {
938                None
939            } else {
940                Some(serde_json::from_str(metadatas.value(i))?)
941            };
942
943            let event = Event::reconstruct_from_strings(
944                id,
945                event_types.value(i).to_string(),
946                entity_ids.value(i).to_string(),
947                tenant_id.to_string(),
948                serde_json::from_str(payloads.value(i))?,
949                timestamp,
950                metadata,
951                versions.value(i) as i64,
952            );
953
954            events.push(event);
955        }
956
957        Ok(events)
958    }
959
960    /// List all Parquet file paths under the storage directory, sorted by
961    /// the relative path so files in the same partition stay grouped.
962    ///
963    /// Used by the replication catch-up protocol to stream snapshot files
964    /// to followers that are too far behind for WAL-only catch-up.
965    pub fn list_parquet_files(&self) -> Result<Vec<PathBuf>> {
966        find_parquet_files_recursive(&self.storage_dir)
967    }
968
969    /// List Parquet files belonging to a single tenant — i.e. only files
970    /// under `<storage_dir>/<tenant>/...`. Legacy flat-layout files at the
971    /// root are intentionally excluded; the migration tool moves them under
972    /// `default/` so once it has run a `tenant=default` query sees them.
973    ///
974    /// Returns an empty vec when the tenant subtree doesn't exist (no data
975    /// for that tenant yet). Returns an error only if `tenant_id` fails the
976    /// path-safety whitelist.
977    ///
978    /// This is the building block for tenant-scoped reads: the caller knows
979    /// which files might contain the tenant's data without opening any of
980    /// the others.
981    pub fn list_parquet_files_for_tenant(&self, tenant_id: &str) -> Result<Vec<PathBuf>> {
982        let safe = sanitize_tenant_id_for_path(tenant_id)?;
983        let tenant_root = self.storage_dir.join(safe);
984        if !tenant_root.is_dir() {
985            return Ok(Vec::new());
986        }
987        find_parquet_files_recursive(&tenant_root)
988    }
989
990    /// Load only the events belonging to `tenant_id`, walking just that
991    /// tenant's subtree on disk. The full-storage loader
992    /// (`load_all_events`) opens every Parquet file regardless of tenant;
993    /// this one only opens files under `<storage_dir>/<tenant>/`.
994    ///
995    /// Returns an empty vec when the tenant has no on-disk data. Returns
996    /// an error if the tenant_id fails the path-safety whitelist or any
997    /// individual file fails to load.
998    ///
999    /// This is the read-side complement to per-tenant flushing. It's the
1000    /// foundation Step 2 (lazy per-tenant load on demand) needs: a way to
1001    /// hydrate one tenant without paying the cost of loading every other
1002    /// tenant's data into memory.
1003    ///
1004    /// Tenant identity for loaded events comes from the file path, the
1005    /// same as `load_all_events` — `record_batch_to_events` stamps the
1006    /// passed `tenant_id` onto every reconstructed event.
1007    #[cfg_attr(feature = "hotpath", hotpath::measure)]
1008    pub fn load_events_for_tenant(&self, tenant_id: &str) -> Result<Vec<Event>> {
1009        let parquet_files = self.list_parquet_files_for_tenant(tenant_id)?;
1010        tracing::info!(
1011            tenant_id = tenant_id,
1012            file_count = parquet_files.len(),
1013            "load_events_for_tenant: walking tenant subtree only"
1014        );
1015
1016        let mut events = Vec::with_capacity(parquet_files.len() * self.config.batch_size);
1017        let mut skipped = 0usize;
1018        for file_path in parquet_files {
1019            tracing::debug!(
1020                tenant_id = tenant_id,
1021                file = %file_path.display(),
1022                "load_events_for_tenant: opening file"
1023            );
1024            // Issue #166: same defensive log-and-skip as load_all_events.
1025            // The lazy-load path can hit a corrupt file just as easily as
1026            // boot can — failing the whole tenant load on one bad file
1027            // would brick every query for that tenant.
1028            match self.load_events_from_file(&file_path, tenant_id) {
1029                Ok(file_events) => events.extend(file_events),
1030                Err(e) => {
1031                    tracing::error!(
1032                        tenant_id = tenant_id,
1033                        file = %file_path.display(),
1034                        error = %e,
1035                        "Skipping unreadable parquet file in tenant subtree"
1036                    );
1037                    skipped += 1;
1038                }
1039            }
1040        }
1041
1042        tracing::info!(
1043            tenant_id = tenant_id,
1044            event_count = events.len(),
1045            skipped_files = skipped,
1046            "load_events_for_tenant: complete"
1047        );
1048        Ok(events)
1049    }
1050
1051    /// Get the storage directory path.
1052    pub fn storage_dir(&self) -> &Path {
1053        &self.storage_dir
1054    }
1055
1056    /// One-shot migration of flat-layout files into the tenant-partitioned
1057    /// tree. Run with Core stopped (no concurrent writes).
1058    ///
1059    /// Walks `storage_dir`'s top level (non-recursive) for the legacy
1060    /// `events-*.parquet` files. For each one it loads the events,
1061    /// regroups them by (tenant_id, yyyy-mm) — events from a flat file
1062    /// take the path-derived "default" tenant, since pre-partitioning
1063    /// data carried no tenant in its on-disk form — writes a fresh
1064    /// Parquet under the corresponding partition directory, and deletes
1065    /// the original flat file once the new file is closed.
1066    ///
1067    /// `dry_run = true` reports what would happen without touching disk.
1068    /// Run dry first; production data deserves the rehearsal.
1069    ///
1070    /// Crash safety: this writes the new partition file before deleting
1071    /// the flat file, so a crash between the two leaves both on disk.
1072    /// The recursive loader (`load_all_events`) will then return both,
1073    /// duplicating those events on next boot. Mitigation: stop Core
1074    /// before running, and re-run the migration after any crash so the
1075    /// flat file gets deleted. A future commit can add atomic rename +
1076    /// fsync semantics; for the one-time migration the stop-Core
1077    /// constraint is enough.
1078    pub fn migrate_flat_layout(&self, dry_run: bool) -> Result<MigrationReport> {
1079        let flat_files = list_flat_layout_files(&self.storage_dir)?;
1080        let mut report = MigrationReport {
1081            dry_run,
1082            ..Default::default()
1083        };
1084
1085        for flat_file in flat_files {
1086            // Pre-partition events used path-derived tenant. For flat-layout
1087            // files that's always "default" (`tenant_id_from_path` falls back
1088            // to "default" for single-component paths).
1089            let events = self.load_events_from_file(&flat_file, "default")?;
1090            report.flat_files_seen += 1;
1091
1092            if events.is_empty() {
1093                // Stale empty file (zero rows). Just remove it.
1094                if !dry_run {
1095                    fs::remove_file(&flat_file).map_err(|e| {
1096                        AllSourceError::StorageError(format!(
1097                            "Failed to remove empty flat file {}: {e}",
1098                            flat_file.display()
1099                        ))
1100                    })?;
1101                }
1102                report.flat_files_removed += 1;
1103                continue;
1104            }
1105
1106            // Group by (tenant, yyyy-mm-from-event-timestamp). The
1107            // partition month tracks the event's wall-clock time so that
1108            // post-migration the layout reflects when data happened, not
1109            // when migration ran. Step 4 (per-tenant snapshots) and Step
1110            // 5 (retention) will key on that.
1111            let mut groups: HashMap<(String, String), Vec<Event>> = HashMap::new();
1112            for event in events {
1113                let key = (
1114                    event.tenant_id_str().to_string(),
1115                    event.timestamp().format("%Y-%m").to_string(),
1116                );
1117                groups.entry(key).or_default().push(event);
1118            }
1119
1120            for ((tenant, yyyy_mm), group_events) in groups {
1121                let count = group_events.len();
1122                if !dry_run {
1123                    let safe_tenant = sanitize_tenant_id_for_path(&tenant)?;
1124                    let target_dir = self.storage_dir.join(safe_tenant).join(&yyyy_mm);
1125                    fs::create_dir_all(&target_dir).map_err(|e| {
1126                        AllSourceError::StorageError(format!(
1127                            "Failed to create partition {}: {e}",
1128                            target_dir.display()
1129                        ))
1130                    })?;
1131                    let filename = format!(
1132                        "events-{}-{}.parquet",
1133                        chrono::Utc::now().format("%Y%m%d-%H%M%S%3f"),
1134                        uuid::Uuid::new_v4().as_simple()
1135                    );
1136                    let target_path = target_dir.join(&filename);
1137                    let record_batch = self.events_to_record_batch(&group_events)?;
1138                    let file = File::create(&target_path).map_err(|e| {
1139                        AllSourceError::StorageError(format!(
1140                            "Failed to create migration target {}: {e}",
1141                            target_path.display()
1142                        ))
1143                    })?;
1144                    let props = WriterProperties::builder()
1145                        .set_compression(self.config.compression)
1146                        .build();
1147                    let mut writer = ArrowWriter::try_new(file, self.schema.clone(), Some(props))?;
1148                    writer.write(&record_batch)?;
1149                    writer.close()?;
1150                    report.partitions_written += 1;
1151                }
1152                report.events_migrated += count;
1153            }
1154
1155            if !dry_run {
1156                fs::remove_file(&flat_file).map_err(|e| {
1157                    AllSourceError::StorageError(format!(
1158                        "Failed to remove flat file {} after migration: {e}",
1159                        flat_file.display()
1160                    ))
1161                })?;
1162                report.flat_files_removed += 1;
1163            }
1164        }
1165
1166        Ok(report)
1167    }
1168
1169    /// Get storage statistics
1170    pub fn stats(&self) -> Result<StorageStats> {
1171        let parquet_files = find_parquet_files_recursive(&self.storage_dir)?;
1172        let mut total_size_bytes = 0u64;
1173        for path in &parquet_files {
1174            if let Ok(metadata) = fs::metadata(path) {
1175                total_size_bytes += metadata.len();
1176            }
1177        }
1178
1179        let current_batch_size: usize = self
1180            .current_batches
1181            .lock()
1182            .unwrap()
1183            .values()
1184            .map(Vec::len)
1185            .sum();
1186
1187        Ok(StorageStats {
1188            total_files: parquet_files.len(),
1189            total_size_bytes,
1190            storage_dir: self.storage_dir.clone(),
1191            current_batch_size,
1192        })
1193    }
1194}
1195
1196/// Validate a tenant ID for use as a filesystem path component.
1197///
1198/// Whitelist: ASCII letters, digits, `-`, `_`, `.` — covers UUIDs, the
1199/// hyphen-and-lowercase tenant strings the onboarding flow produces, and
1200/// the `system` tenant the heartbeat emitter uses. Rejects empty input,
1201/// any path separator (`/`, `\`), and any "..". The whitelist is the
1202/// primary defence against path traversal; the explicit ".." check is
1203/// belt-and-braces in case the whitelist ever loosens.
1204///
1205/// Length capped at 128 bytes — comfortably above the 36-byte UUID and
1206/// the longest onboarding tenant the system has produced, well below
1207/// every common filesystem's NAME_MAX (typically 255).
1208fn sanitize_tenant_id_for_path(tenant_id: &str) -> Result<&str> {
1209    if tenant_id.is_empty() {
1210        return Err(AllSourceError::StorageError(
1211            "tenant_id is empty (cannot derive partition path)".to_string(),
1212        ));
1213    }
1214    if tenant_id.len() > 128 {
1215        return Err(AllSourceError::StorageError(format!(
1216            "tenant_id is too long for partition path: {} bytes (max 128)",
1217            tenant_id.len()
1218        )));
1219    }
1220    if tenant_id == "." || tenant_id == ".." {
1221        return Err(AllSourceError::StorageError(format!(
1222            "tenant_id {tenant_id:?} is reserved"
1223        )));
1224    }
1225    for c in tenant_id.chars() {
1226        let ok = c.is_ascii_alphanumeric() || c == '-' || c == '_' || c == '.';
1227        if !ok {
1228            return Err(AllSourceError::StorageError(format!(
1229                "tenant_id {tenant_id:?} contains disallowed character {c:?} for partition path"
1230            )));
1231        }
1232    }
1233    Ok(tenant_id)
1234}
1235
1236/// Resolve the directory a flush should write into for `(tenant, when)`.
1237///
1238/// Returns `<root>/<tenant>/<yyyy-mm>/`. Caller is responsible for
1239/// `create_dir_all`-ing the result before opening files in it.
1240fn partition_path_for_tenant(
1241    root: &Path,
1242    tenant_id: &str,
1243    when: chrono::DateTime<chrono::Utc>,
1244) -> Result<PathBuf> {
1245    let safe = sanitize_tenant_id_for_path(tenant_id)?;
1246    Ok(root.join(safe).join(when.format("%Y-%m").to_string()))
1247}
1248
1249/// Reverse of `partition_path_for_tenant` — given a parquet file's full
1250/// path and the storage root, return the tenant_id stored in the path.
1251///
1252/// Tenant-partitioned shape: `<root>/<tenant>/<yyyy-mm>/events-*.parquet`
1253/// → first component after root is the tenant.
1254///
1255/// Legacy flat shape: `<root>/events-*.parquet` → no tenant in path, fall
1256/// back to `"default"` so events written before the partitioning change
1257/// keep loading with their original (and only ever) tenant identity.
1258fn tenant_id_from_path(root: &Path, file_path: &Path) -> String {
1259    let Ok(rel) = file_path.strip_prefix(root) else {
1260        return "default".to_string();
1261    };
1262    let mut comps = rel.components();
1263    let first = comps.next();
1264    let next = comps.next();
1265    match (first, next) {
1266        // Two or more components: <tenant>/<rest>... → tenant
1267        (Some(std::path::Component::Normal(tenant)), Some(_)) => {
1268            tenant.to_string_lossy().into_owned()
1269        }
1270        // Single component (the parquet file itself): legacy flat layout.
1271        _ => "default".to_string(),
1272    }
1273}
1274
1275/// List Parquet files at the top level of `root` only — i.e. the legacy
1276/// flat-layout files. Used by the one-shot migration tool to find data
1277/// that needs moving into the tenant-partitioned tree. The opposite of
1278/// `find_parquet_files_recursive`: stops at the first directory level so
1279/// already-partitioned data isn't included.
1280fn list_flat_layout_files(root: &Path) -> Result<Vec<PathBuf>> {
1281    let entries = fs::read_dir(root).map_err(|e| {
1282        AllSourceError::StorageError(format!("Failed to read storage directory: {e}"))
1283    })?;
1284    let mut out: Vec<PathBuf> = entries
1285        .filter_map(std::result::Result::ok)
1286        .filter_map(|entry| {
1287            let ft = entry.file_type().ok()?;
1288            if !ft.is_file() {
1289                return None;
1290            }
1291            let path = entry.path();
1292            if path.extension().and_then(|s| s.to_str()) == Some("parquet") {
1293                Some(path)
1294            } else {
1295                None
1296            }
1297        })
1298        .collect();
1299    out.sort();
1300    Ok(out)
1301}
1302
1303/// Recursively collect all `*.parquet` files under `root`, sorted by path so
1304/// callers see a deterministic, tenant-grouped order.
1305///
1306/// Existence rationale: the storage layout is moving from a flat
1307/// `storage_dir/events-*.parquet` pile to a tenant-partitioned tree of the
1308/// shape `storage_dir/<tenant>/<yyyy-mm>/events-*.parquet`. During the
1309/// migration both shapes coexist, so every code path that asks "what
1310/// parquet files do we have?" needs to walk subdirectories. Symlinks are
1311/// not followed — the storage tree is mounted from a single volume and
1312/// chasing symlinks invites cycles.
1313fn find_parquet_files_recursive(root: &Path) -> Result<Vec<PathBuf>> {
1314    let mut out = Vec::new();
1315    let mut stack: Vec<PathBuf> = vec![root.to_path_buf()];
1316
1317    while let Some(dir) = stack.pop() {
1318        let entries = match fs::read_dir(&dir) {
1319            Ok(e) => e,
1320            // Root must exist (we created it in `new`); subdirectories may
1321            // race a delete from compaction. Skip vanished subdirs rather
1322            // than failing the whole load.
1323            Err(e) if dir == root => {
1324                return Err(AllSourceError::StorageError(format!(
1325                    "Failed to read storage directory: {e}"
1326                )));
1327            }
1328            Err(_) => continue,
1329        };
1330
1331        for entry in entries.flatten() {
1332            let path = entry.path();
1333            // Use file_type() rather than metadata() so symlinks don't get
1334            // followed by accident (metadata() resolves symlinks, file_type()
1335            // doesn't).
1336            let Ok(ft) = entry.file_type() else {
1337                continue;
1338            };
1339            if ft.is_dir() {
1340                stack.push(path);
1341            } else if ft.is_file()
1342                && path
1343                    .extension()
1344                    .and_then(|ext| ext.to_str())
1345                    .is_some_and(|ext| ext == "parquet")
1346            {
1347                out.push(path);
1348            }
1349        }
1350    }
1351
1352    out.sort();
1353    Ok(out)
1354}
1355
1356impl Drop for ParquetStorage {
1357    fn drop(&mut self) {
1358        // Ensure any remaining events are flushed on shutdown
1359        if let Err(e) = self.flush_on_shutdown() {
1360            tracing::error!("Failed to flush events on drop: {}", e);
1361        }
1362    }
1363}
1364
1365/// Outcome of a `migrate_flat_layout` run.
1366#[derive(Debug, Default, Clone, serde::Serialize)]
1367pub struct MigrationReport {
1368    /// Whether the run was a rehearsal (no disk changes).
1369    pub dry_run: bool,
1370    /// Number of legacy flat-layout files discovered.
1371    pub flat_files_seen: usize,
1372    /// Number of legacy flat files deleted (always 0 when `dry_run`).
1373    pub flat_files_removed: usize,
1374    /// Number of new partition files written under the tenant tree.
1375    pub partitions_written: usize,
1376    /// Total events copied into the new tree (counted in dry-run too).
1377    pub events_migrated: usize,
1378}
1379
1380#[derive(Debug, serde::Serialize)]
1381pub struct StorageStats {
1382    pub total_files: usize,
1383    pub total_size_bytes: u64,
1384    pub storage_dir: PathBuf,
1385    pub current_batch_size: usize,
1386}
1387
1388#[cfg(test)]
1389mod tests {
1390    use super::*;
1391    use serde_json::json;
1392    use std::sync::Arc;
1393    use tempfile::TempDir;
1394
1395    fn create_test_event(entity_id: &str) -> Event {
1396        Event::reconstruct_from_strings(
1397            uuid::Uuid::new_v4(),
1398            "test.event".to_string(),
1399            entity_id.to_string(),
1400            "default".to_string(),
1401            json!({
1402                "test": "data",
1403                "value": 42
1404            }),
1405            chrono::Utc::now(),
1406            None,
1407            1,
1408        )
1409    }
1410
1411    #[test]
1412    fn test_parquet_storage_write_read() {
1413        let temp_dir = TempDir::new().unwrap();
1414        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1415
1416        // Add events
1417        for i in 0..10 {
1418            let event = create_test_event(&format!("entity-{i}"));
1419            storage.append_event(event).unwrap();
1420        }
1421
1422        // Flush to disk
1423        storage.flush().unwrap();
1424
1425        // Load back
1426        let loaded_events = storage.load_all_events().unwrap();
1427        assert_eq!(loaded_events.len(), 10);
1428    }
1429
1430    #[test]
1431    fn test_storage_stats() {
1432        let temp_dir = TempDir::new().unwrap();
1433        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1434
1435        // Add and flush events
1436        for i in 0..5 {
1437            storage
1438                .append_event(create_test_event(&format!("entity-{i}")))
1439                .unwrap();
1440        }
1441        storage.flush().unwrap();
1442
1443        let stats = storage.stats().unwrap();
1444        assert_eq!(stats.total_files, 1);
1445        assert!(stats.total_size_bytes > 0);
1446    }
1447
1448    #[test]
1449    fn test_default_batch_size() {
1450        let temp_dir = TempDir::new().unwrap();
1451        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1452
1453        // Default batch size should be 10,000 as per US-023
1454        assert_eq!(storage.batch_size(), DEFAULT_BATCH_SIZE);
1455        assert_eq!(storage.batch_size(), 10_000);
1456    }
1457
1458    #[test]
1459    fn test_custom_config() {
1460        let temp_dir = TempDir::new().unwrap();
1461        let config = ParquetStorageConfig {
1462            batch_size: 5_000,
1463            flush_timeout: Duration::from_secs(2),
1464            compression: parquet::basic::Compression::SNAPPY,
1465        };
1466        let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1467
1468        assert_eq!(storage.batch_size(), 5_000);
1469        assert_eq!(storage.flush_timeout(), Duration::from_secs(2));
1470    }
1471
1472    #[test]
1473    fn test_batch_write() {
1474        let temp_dir = TempDir::new().unwrap();
1475        let config = ParquetStorageConfig {
1476            batch_size: 100, // Small batch for testing
1477            ..Default::default()
1478        };
1479        let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1480
1481        // 250 events for a single tenant. With per-tenant flush, when the
1482        // tenant's pending batch crosses batch_size we drain the whole
1483        // tenant in one flush — not chunk at exactly batch_size like the
1484        // old global-batch path did. So 250 events triggers exactly one
1485        // size-flush (the appender pushes all 250 onto the tenant's batch
1486        // under one lock, sees length >= 100, schedules a flush which
1487        // drains everything). 0 left pending.
1488        let events: Vec<Event> = (0..250)
1489            .map(|i| create_test_event(&format!("entity-{i}")))
1490            .collect();
1491
1492        let result = storage.batch_write(events).unwrap();
1493        assert_eq!(result.events_written, 250);
1494        assert_eq!(result.batches_flushed, 1);
1495        assert_eq!(storage.pending_count(), 0);
1496
1497        // Manual flush is a no-op since nothing's pending.
1498        storage.flush().unwrap();
1499
1500        // All 250 events round-trip through the tenant-partitioned tree.
1501        let loaded = storage.load_all_events().unwrap();
1502        assert_eq!(loaded.len(), 250);
1503    }
1504
1505    #[test]
1506    fn test_auto_flush_on_batch_size() {
1507        let temp_dir = TempDir::new().unwrap();
1508        let config = ParquetStorageConfig {
1509            batch_size: 10, // Very small for testing
1510            ..Default::default()
1511        };
1512        let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1513
1514        // Add 15 events - should auto-flush at 10
1515        for i in 0..15 {
1516            storage
1517                .append_event(create_test_event(&format!("entity-{i}")))
1518                .unwrap();
1519        }
1520
1521        // Should have 5 pending, 10 written
1522        assert_eq!(storage.pending_count(), 5);
1523
1524        let stats = storage.batch_stats();
1525        assert_eq!(stats.events_written, 10);
1526        assert_eq!(stats.batches_written, 1);
1527        assert_eq!(stats.size_flushes, 1);
1528    }
1529
1530    #[test]
1531    fn test_flush_on_shutdown() {
1532        let temp_dir = TempDir::new().unwrap();
1533        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1534
1535        // Add some events without reaching batch size
1536        for i in 0..5 {
1537            storage
1538                .append_event(create_test_event(&format!("entity-{i}")))
1539                .unwrap();
1540        }
1541
1542        assert_eq!(storage.pending_count(), 5);
1543
1544        // Manually trigger shutdown flush
1545        let flushed = storage.flush_on_shutdown().unwrap();
1546        assert_eq!(flushed, 5);
1547        assert_eq!(storage.pending_count(), 0);
1548
1549        // Verify events are persisted
1550        let loaded = storage.load_all_events().unwrap();
1551        assert_eq!(loaded.len(), 5);
1552    }
1553
1554    #[test]
1555    fn test_thread_safe_writes() {
1556        let temp_dir = TempDir::new().unwrap();
1557        let config = ParquetStorageConfig {
1558            batch_size: 100,
1559            ..Default::default()
1560        };
1561        let storage = Arc::new(ParquetStorage::with_config(temp_dir.path(), config).unwrap());
1562
1563        let events_per_thread = 50;
1564        let thread_count = 4;
1565
1566        std::thread::scope(|s| {
1567            for t in 0..thread_count {
1568                let storage_ref = Arc::clone(&storage);
1569                s.spawn(move || {
1570                    for i in 0..events_per_thread {
1571                        let event = create_test_event(&format!("thread-{t}-entity-{i}"));
1572                        storage_ref.append_event(event).unwrap();
1573                    }
1574                });
1575            }
1576        });
1577
1578        // Flush remaining
1579        storage.flush().unwrap();
1580
1581        // All events should be written
1582        let loaded = storage.load_all_events().unwrap();
1583        assert_eq!(loaded.len(), events_per_thread * thread_count);
1584    }
1585
1586    #[test]
1587    fn test_batch_stats() {
1588        let temp_dir = TempDir::new().unwrap();
1589        let config = ParquetStorageConfig {
1590            batch_size: 50,
1591            ..Default::default()
1592        };
1593        let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1594
1595        // 100 events, single tenant, batch_size=50. Per-tenant flush
1596        // drains the whole tenant on the first size trigger, so this
1597        // produces exactly one size-flush and one batches_written event
1598        // (vs. the pre-tenant world's two).
1599        let events: Vec<Event> = (0..100)
1600            .map(|i| create_test_event(&format!("entity-{i}")))
1601            .collect();
1602
1603        storage.batch_write(events).unwrap();
1604
1605        let stats = storage.batch_stats();
1606        assert_eq!(stats.batches_written, 1);
1607        assert_eq!(stats.events_written, 100);
1608        assert!(stats.avg_batch_size > 0.0);
1609        assert!(stats.events_per_sec > 0.0);
1610        assert_eq!(stats.size_flushes, 1);
1611    }
1612
1613    #[test]
1614    fn test_config_presets() {
1615        let high_throughput = ParquetStorageConfig::high_throughput();
1616        assert_eq!(high_throughput.batch_size, 50_000);
1617        assert_eq!(high_throughput.flush_timeout, Duration::from_secs(10));
1618
1619        let low_latency = ParquetStorageConfig::low_latency();
1620        assert_eq!(low_latency.batch_size, 1_000);
1621        assert_eq!(low_latency.flush_timeout, Duration::from_secs(1));
1622
1623        let default = ParquetStorageConfig::default();
1624        assert_eq!(default.batch_size, DEFAULT_BATCH_SIZE);
1625        assert_eq!(default.batch_size, 10_000);
1626    }
1627
1628    /// Benchmark: Compare single-event writes vs batch writes
1629    /// Run with: cargo test --release -- --ignored test_batch_write_throughput
1630    #[test]
1631    #[ignore]
1632    fn test_batch_write_throughput() {
1633        let temp_dir = TempDir::new().unwrap();
1634        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1635
1636        let event_count = 50_000;
1637
1638        // Benchmark batch write
1639        let events: Vec<Event> = (0..event_count)
1640            .map(|i| create_test_event(&format!("entity-{i}")))
1641            .collect();
1642
1643        let start = std::time::Instant::now();
1644        let result = storage.batch_write(events).unwrap();
1645        storage.flush().unwrap(); // Flush any remaining
1646        let batch_duration = start.elapsed();
1647
1648        let batch_stats = storage.batch_stats();
1649
1650        println!("\n=== Parquet Batch Write Performance (BATCH_SIZE=10,000) ===");
1651        println!("Events: {event_count}");
1652        println!("Duration: {batch_duration:?}");
1653        println!("Events/sec: {:.0}", result.events_per_sec);
1654        println!("Batches written: {}", batch_stats.batches_written);
1655        println!("Avg batch size: {:.0}", batch_stats.avg_batch_size);
1656        println!("Bytes written: {} KB", batch_stats.bytes_written / 1024);
1657
1658        // Target: Batch writes should achieve at least 100K events/sec in release mode
1659        // This represents 40%+ improvement over single-event writes
1660        assert!(
1661            result.events_per_sec > 10_000.0,
1662            "Batch write throughput too low: {:.0} events/sec (expected >10K in debug, >100K in release)",
1663            result.events_per_sec
1664        );
1665    }
1666
1667    /// Benchmark: Single-event write baseline (for comparison)
1668    #[test]
1669    #[ignore]
1670    fn test_single_event_write_baseline() {
1671        let temp_dir = TempDir::new().unwrap();
1672        let config = ParquetStorageConfig {
1673            batch_size: 1, // Force flush after each event
1674            ..Default::default()
1675        };
1676        let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1677
1678        let event_count = 1_000; // Fewer events since this is slow
1679
1680        let start = std::time::Instant::now();
1681        for i in 0..event_count {
1682            let event = create_test_event(&format!("entity-{i}"));
1683            storage.append_event(event).unwrap();
1684        }
1685        let duration = start.elapsed();
1686
1687        let events_per_sec = f64::from(event_count) / duration.as_secs_f64();
1688
1689        println!("\n=== Single-Event Write Baseline ===");
1690        println!("Events: {event_count}");
1691        println!("Duration: {duration:?}");
1692        println!("Events/sec: {events_per_sec:.0}");
1693
1694        // This should be significantly slower than batch writes
1695        // Used as a baseline to demonstrate 40%+ improvement
1696    }
1697
1698    // -----------------------------------------------------------------
1699    // Tests for the recursive parquet walker (Step 1, commit #1: read-side
1700    // bidirectional layout support — see SUSTAINABLE_DATA_STRATEGY.md).
1701    // -----------------------------------------------------------------
1702
1703    /// Helper: write a tiny placeholder parquet file at an arbitrary path so
1704    /// the walker has something concrete to find. We only care that the
1705    /// walker discovers the path, not that the file is loadable here — the
1706    /// load path is exercised by the existing read tests.
1707    fn touch_parquet(path: &Path) {
1708        std::fs::create_dir_all(path.parent().unwrap()).unwrap();
1709        std::fs::write(path, b"").unwrap();
1710    }
1711
1712    #[test]
1713    fn test_walker_finds_files_in_flat_layout() {
1714        let temp_dir = TempDir::new().unwrap();
1715        let root = temp_dir.path();
1716        touch_parquet(&root.join("events-20260101-120000000-aaaa.parquet"));
1717        touch_parquet(&root.join("events-20260101-130000000-bbbb.parquet"));
1718
1719        let mut found = find_parquet_files_recursive(root).unwrap();
1720        found.sort();
1721        assert_eq!(found.len(), 2);
1722        assert!(
1723            found[0]
1724                .file_name()
1725                .unwrap()
1726                .to_str()
1727                .unwrap()
1728                .starts_with("events-"),
1729            "expected events-* file, got {found:?}"
1730        );
1731    }
1732
1733    #[test]
1734    fn test_walker_finds_files_in_tenant_partitioned_tree() {
1735        let temp_dir = TempDir::new().unwrap();
1736        let root = temp_dir.path();
1737        // Tenant-partitioned shape: storage_dir/<tenant>/<yyyy-mm>/events-*.parquet
1738        touch_parquet(&root.join("tenant-a/2026-01/events-20260101-120000000-aaaa.parquet"));
1739        touch_parquet(&root.join("tenant-a/2026-02/events-20260201-120000000-bbbb.parquet"));
1740        touch_parquet(&root.join("tenant-b/2026-01/events-20260103-120000000-cccc.parquet"));
1741
1742        let found = find_parquet_files_recursive(root).unwrap();
1743        assert_eq!(found.len(), 3);
1744        // Sort places tenant-a files before tenant-b — that's the
1745        // tenant-grouping the docs claim.
1746        assert!(found[0].to_str().unwrap().contains("tenant-a"));
1747        assert!(found[1].to_str().unwrap().contains("tenant-a"));
1748        assert!(found[2].to_str().unwrap().contains("tenant-b"));
1749    }
1750
1751    #[test]
1752    fn test_walker_handles_mixed_legacy_and_partitioned_layouts() {
1753        // The migration window: some tenants have been moved into the tree,
1754        // some flat files still sit at the root. The walker must surface
1755        // both so load_all_events sees every event regardless of where it
1756        // currently lives.
1757        let temp_dir = TempDir::new().unwrap();
1758        let root = temp_dir.path();
1759        touch_parquet(&root.join("events-legacy-aaaa.parquet"));
1760        touch_parquet(&root.join("tenant-a/2026-01/events-new-bbbb.parquet"));
1761
1762        let found = find_parquet_files_recursive(root).unwrap();
1763        assert_eq!(found.len(), 2);
1764    }
1765
1766    #[test]
1767    fn test_walker_ignores_non_parquet_files() {
1768        let temp_dir = TempDir::new().unwrap();
1769        let root = temp_dir.path();
1770        std::fs::write(root.join("README.md"), b"hello").unwrap();
1771        std::fs::write(root.join("events.json"), b"[]").unwrap();
1772        touch_parquet(&root.join("events-20260101-120000000-aaaa.parquet"));
1773        // Files that just happen to have "parquet" in the name but no
1774        // .parquet extension stay out — extension-only filter, no name match.
1775        std::fs::write(root.join("not-a-parquet-file.bin"), b"").unwrap();
1776
1777        let found = find_parquet_files_recursive(root).unwrap();
1778        assert_eq!(found.len(), 1);
1779        assert_eq!(
1780            found[0].extension().and_then(|s| s.to_str()),
1781            Some("parquet")
1782        );
1783    }
1784
1785    /// Build an event whose tenant_id and entity_id we control, so tests
1786    /// can verify per-tenant routing without depending on the helper that
1787    /// hardcodes "default".
1788    fn event_with_tenant(tenant: &str, entity_id: &str) -> Event {
1789        Event::reconstruct_from_strings(
1790            uuid::Uuid::new_v4(),
1791            "test.event".to_string(),
1792            entity_id.to_string(),
1793            tenant.to_string(),
1794            json!({"k": "v"}),
1795            chrono::Utc::now(),
1796            None,
1797            1,
1798        )
1799    }
1800
1801    #[test]
1802    fn test_flush_writes_into_per_tenant_partition() {
1803        // End-to-end check that the new write path produces
1804        // <root>/<tenant>/<yyyy-mm>/events-*.parquet — no flat file at the
1805        // root, no cross-tenant mixing.
1806        let temp_dir = TempDir::new().unwrap();
1807        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1808
1809        for i in 0..3 {
1810            storage
1811                .append_event(event_with_tenant("default", &format!("entity-{i}")))
1812                .unwrap();
1813        }
1814        storage.flush().unwrap();
1815
1816        let parquet_files = find_parquet_files_recursive(temp_dir.path()).unwrap();
1817        assert_eq!(parquet_files.len(), 1);
1818
1819        let rel = parquet_files[0]
1820            .strip_prefix(temp_dir.path())
1821            .unwrap()
1822            .to_string_lossy()
1823            .into_owned();
1824        // Path shape: default/<yyyy-mm>/events-*.parquet
1825        let parts: Vec<&str> = rel.split(std::path::MAIN_SEPARATOR).collect();
1826        assert_eq!(parts.len(), 3, "expected tenant/yyyy-mm/file, got {rel}");
1827        assert_eq!(parts[0], "default");
1828        // yyyy-mm is two digits dash four — loose check, exact month
1829        // varies with wall-clock at test runtime.
1830        assert!(
1831            parts[1].len() == 7 && parts[1].as_bytes()[4] == b'-',
1832            "expected yyyy-mm, got {}",
1833            parts[1]
1834        );
1835        assert!(parts[2].starts_with("events-") && parts[2].ends_with(".parquet"));
1836
1837        let loaded = storage.load_all_events().unwrap();
1838        assert_eq!(loaded.len(), 3);
1839    }
1840
1841    #[test]
1842    fn test_multiple_tenants_get_isolated_subtrees() {
1843        // Per-tenant flush must not mix tenants into the same Parquet file
1844        // and must put each tenant under its own subdirectory.
1845        let temp_dir = TempDir::new().unwrap();
1846        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1847
1848        for i in 0..2 {
1849            storage
1850                .append_event(event_with_tenant("alice", &format!("a-{i}")))
1851                .unwrap();
1852        }
1853        for i in 0..3 {
1854            storage
1855                .append_event(event_with_tenant("bob", &format!("b-{i}")))
1856                .unwrap();
1857        }
1858        storage.flush().unwrap();
1859
1860        let alice_subtree = temp_dir.path().join("alice");
1861        let bob_subtree = temp_dir.path().join("bob");
1862        assert!(alice_subtree.is_dir(), "alice should have its own subtree");
1863        assert!(bob_subtree.is_dir(), "bob should have its own subtree");
1864
1865        let alice_files = find_parquet_files_recursive(&alice_subtree).unwrap();
1866        let bob_files = find_parquet_files_recursive(&bob_subtree).unwrap();
1867        assert_eq!(alice_files.len(), 1);
1868        assert_eq!(bob_files.len(), 1);
1869
1870        // Loaded events keep their tenant_id — round-trip preserves which
1871        // tenant each event belonged to.
1872        let loaded = storage.load_all_events().unwrap();
1873        let (alice_count, bob_count) =
1874            loaded
1875                .iter()
1876                .fold((0, 0), |(a, b), e| match e.tenant_id_str() {
1877                    "alice" => (a + 1, b),
1878                    "bob" => (a, b + 1),
1879                    _ => (a, b),
1880                });
1881        assert_eq!(alice_count, 2);
1882        assert_eq!(bob_count, 3);
1883    }
1884
1885    #[test]
1886    fn test_size_flush_only_drains_full_tenant() {
1887        // When one tenant exactly hits batch_size, only that tenant
1888        // flushes; the other tenant keeps its events buffered. Prevents
1889        // one noisy tenant from causing fragmented writes for everyone.
1890        let temp_dir = TempDir::new().unwrap();
1891        let config = ParquetStorageConfig {
1892            batch_size: 5,
1893            ..Default::default()
1894        };
1895        let storage = ParquetStorage::with_config(temp_dir.path(), config).unwrap();
1896
1897        // Alice: 5 events → on the 5th, len == batch_size triggers flush
1898        // which drains all 5. Alice ends empty.
1899        for i in 0..5 {
1900            storage
1901                .append_event(event_with_tenant("alice", &format!("a-{i}")))
1902                .unwrap();
1903        }
1904        // Bob: 2 events → still under threshold, stays pending.
1905        for i in 0..2 {
1906            storage
1907                .append_event(event_with_tenant("bob", &format!("b-{i}")))
1908                .unwrap();
1909        }
1910
1911        assert_eq!(
1912            storage.pending_count(),
1913            2,
1914            "only bob's 2 events should be pending"
1915        );
1916
1917        let parquet_files = find_parquet_files_recursive(temp_dir.path()).unwrap();
1918        assert_eq!(parquet_files.len(), 1, "only alice should have flushed");
1919        assert!(
1920            parquet_files[0]
1921                .to_string_lossy()
1922                .contains(&format!("alice{}", std::path::MAIN_SEPARATOR)),
1923            "expected alice partition, got {}",
1924            parquet_files[0].display()
1925        );
1926    }
1927
1928    #[test]
1929    fn test_tenant_id_from_path_recovers_tenant_for_partitioned_files() {
1930        let root = Path::new("/data/storage");
1931        let f = Path::new("/data/storage/alice/2026-04/events-20260426-120000000-aaaa.parquet");
1932        assert_eq!(tenant_id_from_path(root, f), "alice");
1933    }
1934
1935    #[test]
1936    fn test_tenant_id_from_path_falls_back_to_default_for_legacy_flat_layout() {
1937        let root = Path::new("/data/storage");
1938        let f = Path::new("/data/storage/events-20260426-120000000-aaaa.parquet");
1939        // Legacy single-component path. Pre-tenant data was always
1940        // commingled with tenant=default, so default is the right fallback.
1941        assert_eq!(tenant_id_from_path(root, f), "default");
1942    }
1943
1944    #[test]
1945    fn test_sanitize_tenant_id_for_path_accepts_safe_inputs() {
1946        for ok in [
1947            "default",
1948            "system",
1949            "1e6b2d1c-2f64-4441-9cf9-42f2e451aa17",
1950            "onboard-diagnostic-160-at-example-com",
1951            "tenant_with_underscore",
1952            "v1.0",
1953        ] {
1954            assert!(
1955                sanitize_tenant_id_for_path(ok).is_ok(),
1956                "{ok:?} should be accepted"
1957            );
1958        }
1959    }
1960
1961    #[test]
1962    fn test_sanitize_tenant_id_for_path_rejects_unsafe_inputs() {
1963        for bad in [
1964            "",         // empty
1965            "..",       // parent traversal
1966            ".",        // current dir
1967            "foo/bar",  // path separator
1968            "foo\\bar", // windows-style separator
1969            "foo bar",  // whitespace
1970            "foo\nbar", // newline
1971            "foo\0bar", // null byte
1972            "tenant?",  // shell glob char
1973            "tenant*",  // shell glob char
1974        ] {
1975            assert!(
1976                sanitize_tenant_id_for_path(bad).is_err(),
1977                "{bad:?} should be rejected"
1978            );
1979        }
1980
1981        // Length cap.
1982        let too_long = "a".repeat(129);
1983        assert!(sanitize_tenant_id_for_path(&too_long).is_err());
1984    }
1985
1986    #[test]
1987    fn test_partition_path_for_tenant_shape() {
1988        let root = Path::new("/data");
1989        let when = chrono::DateTime::parse_from_rfc3339("2026-04-26T12:00:00Z")
1990            .unwrap()
1991            .with_timezone(&chrono::Utc);
1992        let path = partition_path_for_tenant(root, "alice", when).unwrap();
1993        assert_eq!(path, Path::new("/data/alice/2026-04"));
1994    }
1995
1996    #[test]
1997    fn test_append_event_rejects_unsafe_tenant_at_flush() {
1998        // Defence in depth: even if some upstream forgets to validate, the
1999        // sanitizer in flush_tenant catches it. Since append doesn't write
2000        // synchronously, the bad tenant is rejected on the first flush
2001        // attempt. We test that flush surfaces an error rather than
2002        // silently writing somewhere weird.
2003        let temp_dir = TempDir::new().unwrap();
2004        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2005
2006        // append accepts whatever tenant_id the event carries — domain
2007        // construction would normally reject this, but if it slipped
2008        // through, flush should refuse to derive a path from it.
2009        storage
2010            .append_event(event_with_tenant("../escape", "e-0"))
2011            .unwrap();
2012        let result = storage.flush();
2013        assert!(result.is_err(), "flush should reject unsafe tenant_id");
2014        let msg = format!("{}", result.unwrap_err());
2015        assert!(
2016            msg.contains("disallowed character") || msg.contains("reserved"),
2017            "expected sanitization error message, got: {msg}"
2018        );
2019    }
2020
2021    // -----------------------------------------------------------------
2022    // Tenant-pruned read tests (Step 1, commit #4).
2023    // -----------------------------------------------------------------
2024
2025    #[test]
2026    fn test_load_events_for_tenant_only_walks_target_subtree() {
2027        // Seed three tenants with distinct event counts. Loading one
2028        // tenant must return only that tenant's events — and the file
2029        // list helper must report only that tenant's files (the strong
2030        // form of "didn't open the others").
2031        let temp_dir = TempDir::new().unwrap();
2032        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2033
2034        for i in 0..2 {
2035            storage
2036                .append_event(event_with_tenant("alice", &format!("a-{i}")))
2037                .unwrap();
2038        }
2039        for i in 0..3 {
2040            storage
2041                .append_event(event_with_tenant("bob", &format!("b-{i}")))
2042                .unwrap();
2043        }
2044        for i in 0..1 {
2045            storage
2046                .append_event(event_with_tenant("carol", &format!("c-{i}")))
2047                .unwrap();
2048        }
2049        storage.flush().unwrap();
2050
2051        let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
2052        assert_eq!(alice_files.len(), 1);
2053        assert!(
2054            alice_files[0]
2055                .to_string_lossy()
2056                .contains(&format!("alice{}", std::path::MAIN_SEPARATOR)),
2057            "expected alice file, got {}",
2058            alice_files[0].display()
2059        );
2060        // The pruned listing must NOT include any bob/carol files — this
2061        // is the property Step 2 will rely on to avoid loading every
2062        // tenant's data on a single-tenant query.
2063        for f in &alice_files {
2064            let s = f.to_string_lossy();
2065            assert!(!s.contains("bob"), "alice listing leaked bob file: {s}");
2066            assert!(!s.contains("carol"), "alice listing leaked carol file: {s}");
2067        }
2068
2069        let alice_events = storage.load_events_for_tenant("alice").unwrap();
2070        assert_eq!(alice_events.len(), 2);
2071        for e in &alice_events {
2072            assert_eq!(e.tenant_id_str(), "alice");
2073        }
2074
2075        let bob_events = storage.load_events_for_tenant("bob").unwrap();
2076        assert_eq!(bob_events.len(), 3);
2077        for e in &bob_events {
2078            assert_eq!(e.tenant_id_str(), "bob");
2079        }
2080
2081        let carol_events = storage.load_events_for_tenant("carol").unwrap();
2082        assert_eq!(carol_events.len(), 1);
2083        assert_eq!(carol_events[0].tenant_id_str(), "carol");
2084    }
2085
2086    #[test]
2087    fn test_load_events_for_tenant_returns_empty_when_subtree_missing() {
2088        // Querying a tenant that has never written must not error — it's
2089        // a normal "no data" case, not a misconfiguration. Important for
2090        // first-query latency on a fresh tenant.
2091        let temp_dir = TempDir::new().unwrap();
2092        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2093
2094        // Seed only alice so the storage_dir isn't empty (rule out the
2095        // empty-dir trivial case).
2096        storage
2097            .append_event(event_with_tenant("alice", "a-0"))
2098            .unwrap();
2099        storage.flush().unwrap();
2100
2101        let files = storage
2102            .list_parquet_files_for_tenant("nobody-here")
2103            .unwrap();
2104        assert!(files.is_empty());
2105
2106        let events = storage.load_events_for_tenant("nobody-here").unwrap();
2107        assert!(events.is_empty());
2108    }
2109
2110    #[test]
2111    fn test_load_events_for_tenant_rejects_unsafe_tenant_id() {
2112        // Path traversal must fail at the API boundary, not after disk
2113        // reads. Same whitelist as the write path.
2114        let temp_dir = TempDir::new().unwrap();
2115        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2116
2117        for unsafe_tid in ["..", "a/b", "a\\b", "", "a..b/.."] {
2118            let result = storage.load_events_for_tenant(unsafe_tid);
2119            assert!(
2120                result.is_err(),
2121                "tenant_id {unsafe_tid:?} should have been rejected"
2122            );
2123        }
2124    }
2125
2126    #[test]
2127    fn test_load_events_for_tenant_ignores_legacy_flat_layout_files() {
2128        // Flat-layout files at the storage root predate partitioning. A
2129        // tenant-scoped load must not pick them up — the migration tool
2130        // is what relocates them under default/. Until it runs, those
2131        // files are invisible to per-tenant queries (correct behavior:
2132        // the system has no way to tell which tenant they belong to
2133        // beyond "default", and pretending otherwise would mis-attribute
2134        // them).
2135        let temp_dir = TempDir::new().unwrap();
2136        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2137
2138        // Seed a flat-layout file (relocates default/<yyyy-mm>/ → root).
2139        let _flat = seed_flat_layout_file(&storage, 4);
2140
2141        // Querying default returns nothing — the flat file at the root
2142        // isn't under default/.
2143        let default_events = storage.load_events_for_tenant("default").unwrap();
2144        assert!(
2145            default_events.is_empty(),
2146            "tenant-scoped load must not pick up flat-layout files; got {} events",
2147            default_events.len()
2148        );
2149
2150        // Sanity: the full loader still sees them via the recursive walk.
2151        let all_events = storage.load_all_events().unwrap();
2152        assert_eq!(all_events.len(), 4);
2153    }
2154
2155    // -----------------------------------------------------------------
2156    // Atomic snapshot write tests (Step 4, commit #1).
2157    // -----------------------------------------------------------------
2158
2159    #[test]
2160    fn test_write_atomic_parquet_emits_file_under_tenant_partition() {
2161        // Happy path: write 3 events for tenant alice, get back a
2162        // path under <root>/alice/<yyyy-mm>/. The .tmp file should
2163        // be gone (rename completed); the final file readable via
2164        // load_events_for_tenant.
2165        let temp_dir = TempDir::new().unwrap();
2166        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2167
2168        let events: Vec<Event> = (0..3)
2169            .map(|i| event_with_tenant("alice", &format!("a-{i}")))
2170            .collect();
2171
2172        let final_path = storage
2173            .write_atomic_parquet("alice", "snapshot.alice.range", &events)
2174            .unwrap();
2175
2176        // Path shape: <root>/alice/<yyyy-mm>/snapshot.alice.range.parquet
2177        let rel = final_path
2178            .strip_prefix(temp_dir.path())
2179            .unwrap()
2180            .to_string_lossy()
2181            .into_owned();
2182        let parts: Vec<&str> = rel.split(std::path::MAIN_SEPARATOR).collect();
2183        assert_eq!(parts.len(), 3, "expected tenant/yyyy-mm/file, got {rel}");
2184        assert_eq!(parts[0], "alice");
2185        assert_eq!(parts[2], "snapshot.alice.range.parquet");
2186
2187        // Final file must exist; .tmp must NOT.
2188        assert!(final_path.is_file());
2189        let tmp = final_path.with_extension("parquet.tmp");
2190        assert!(
2191            !tmp.exists(),
2192            "tmp should have been renamed away; still at {}",
2193            tmp.display()
2194        );
2195
2196        // Loadable via the tenant loader.
2197        let loaded = storage.load_events_for_tenant("alice").unwrap();
2198        assert_eq!(loaded.len(), 3);
2199    }
2200
2201    #[test]
2202    fn test_write_atomic_parquet_rejects_empty_events() {
2203        let temp_dir = TempDir::new().unwrap();
2204        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2205        let result = storage.write_atomic_parquet("alice", "snap", &[]);
2206        assert!(result.is_err());
2207    }
2208
2209    #[test]
2210    fn test_write_atomic_parquet_rejects_unsafe_tenant() {
2211        let temp_dir = TempDir::new().unwrap();
2212        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2213        let events = [event_with_tenant("alice", "e-0")];
2214        for unsafe_tid in ["..", "a/b", ""] {
2215            let result = storage.write_atomic_parquet(unsafe_tid, "snap", &events);
2216            assert!(
2217                result.is_err(),
2218                "unsafe tenant_id {unsafe_tid:?} should have been rejected"
2219            );
2220        }
2221    }
2222
2223    #[test]
2224    fn test_cleanup_partial_writes_removes_orphan_tmps() {
2225        // Simulate a crashed mid-snapshot: pretend we have a leftover
2226        // events-x.parquet.tmp file in a tenant partition. Cleanup
2227        // must delete it without touching real .parquet files.
2228        let temp_dir = TempDir::new().unwrap();
2229        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2230
2231        // Seed a real parquet via a normal flush, so we have one
2232        // legit file we don't want to delete.
2233        for i in 0..2 {
2234            storage
2235                .append_event(event_with_tenant("alice", &format!("a-{i}")))
2236                .unwrap();
2237        }
2238        storage.flush().unwrap();
2239        let real_files_before = find_parquet_files_recursive(temp_dir.path()).unwrap();
2240        assert_eq!(real_files_before.len(), 1);
2241
2242        // Manufacture an orphan .tmp file in alice's partition.
2243        let alice_subtree = temp_dir.path().join("alice");
2244        let orphan_dir = real_files_before[0].parent().unwrap();
2245        let orphan_path = orphan_dir.join("snapshot.alice.crashed.parquet.tmp");
2246        std::fs::write(&orphan_path, b"fake partial parquet").unwrap();
2247        assert!(orphan_path.is_file());
2248
2249        // And one nested deeper, just to confirm recursion.
2250        let nested_dir = alice_subtree.join("2099-01");
2251        std::fs::create_dir_all(&nested_dir).unwrap();
2252        let nested_orphan = nested_dir.join("events-x.parquet.tmp");
2253        std::fs::write(&nested_orphan, b"junk").unwrap();
2254
2255        let removed = storage.cleanup_partial_writes().unwrap();
2256        assert_eq!(removed, 2, "two orphan tmps should have been cleaned");
2257        assert!(!orphan_path.exists());
2258        assert!(!nested_orphan.exists());
2259
2260        // Real parquet untouched.
2261        let real_files_after = find_parquet_files_recursive(temp_dir.path()).unwrap();
2262        assert_eq!(real_files_after, real_files_before);
2263    }
2264
2265    #[test]
2266    fn test_cleanup_partial_writes_quarantines_zero_byte_parquet() {
2267        // Issue #166 regression: a pre-fix checkpoint crash left a 0-byte
2268        // events-*.parquet file with the FINAL extension (no .tmp). The
2269        // cleanup pass must rename it aside rather than delete, so an
2270        // operator can inspect.
2271        let temp_dir = TempDir::new().unwrap();
2272        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2273
2274        // Seed a real parquet alongside, to confirm we don't molest healthy data.
2275        for i in 0..2 {
2276            storage
2277                .append_event(event_with_tenant("alice", &format!("a-{i}")))
2278                .unwrap();
2279        }
2280        storage.flush().unwrap();
2281        let healthy_files_before = find_parquet_files_recursive(temp_dir.path()).unwrap();
2282        assert_eq!(healthy_files_before.len(), 1);
2283        let healthy_file = &healthy_files_before[0];
2284        let healthy_dir = healthy_file.parent().unwrap();
2285
2286        // Drop a 0-byte parquet next to the healthy one (the issue's failure mode).
2287        let bricked = healthy_dir.join("events-bricked-deadbeef.parquet");
2288        std::fs::write(&bricked, b"").unwrap();
2289        assert_eq!(std::fs::metadata(&bricked).unwrap().len(), 0);
2290
2291        let acted = storage.cleanup_partial_writes().unwrap();
2292        assert_eq!(acted, 1, "only the 0-byte file should have been acted on");
2293
2294        // 0-byte file is gone from its original name, but a quarantined
2295        // sibling exists.
2296        assert!(!bricked.exists(), "0-byte file should have been renamed");
2297        let quarantined: Vec<_> = std::fs::read_dir(healthy_dir)
2298            .unwrap()
2299            .flatten()
2300            .map(|e| e.path())
2301            .filter(|p| {
2302                p.file_name()
2303                    .and_then(|n| n.to_str())
2304                    .is_some_and(|n| n.starts_with("events-bricked-deadbeef.parquet.corrupt-"))
2305            })
2306            .collect();
2307        assert_eq!(
2308            quarantined.len(),
2309            1,
2310            "expected one .parquet.corrupt-<ts> sibling"
2311        );
2312
2313        // Healthy file untouched.
2314        assert!(
2315            healthy_file.exists(),
2316            "healthy parquet must not be molested"
2317        );
2318    }
2319
2320    #[test]
2321    fn test_load_all_events_skips_zero_byte_parquet() {
2322        // Issue #166 defensive layer: a single bad file must not abort
2323        // the whole load. With t-d8b1's match-and-continue, the healthy
2324        // events still come back and the bad file is logged-and-skipped.
2325        let temp_dir = TempDir::new().unwrap();
2326        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2327
2328        // Two healthy events, flushed.
2329        for i in 0..2 {
2330            storage
2331                .append_event(event_with_tenant("alice", &format!("a-{i}")))
2332                .unwrap();
2333        }
2334        storage.flush().unwrap();
2335
2336        // Now drop a 0-byte parquet alongside (sneaking it in AFTER
2337        // construction so cleanup_partial_writes doesn't rename it
2338        // before we get to test the load path).
2339        let healthy_files = find_parquet_files_recursive(temp_dir.path()).unwrap();
2340        let bricked = healthy_files[0]
2341            .parent()
2342            .unwrap()
2343            .join("events-bricked-cafef00d.parquet");
2344        std::fs::write(&bricked, b"").unwrap();
2345
2346        let loaded = storage.load_all_events().unwrap();
2347        assert_eq!(
2348            loaded.len(),
2349            2,
2350            "all healthy events should still load despite the 0-byte file"
2351        );
2352    }
2353
2354    #[test]
2355    fn test_load_events_for_tenant_skips_zero_byte_parquet() {
2356        // Same defensive layer for the lazy-load path. A bad file in
2357        // the tenant subtree must not poison subsequent queries.
2358        let temp_dir = TempDir::new().unwrap();
2359        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2360
2361        for i in 0..3 {
2362            storage
2363                .append_event(event_with_tenant("alice", &format!("a-{i}")))
2364                .unwrap();
2365        }
2366        storage.flush().unwrap();
2367
2368        let healthy = find_parquet_files_recursive(temp_dir.path()).unwrap();
2369        let bricked = healthy[0]
2370            .parent()
2371            .unwrap()
2372            .join("events-bricked-feedface.parquet");
2373        std::fs::write(&bricked, b"").unwrap();
2374
2375        let events = storage.load_events_for_tenant("alice").unwrap();
2376        assert_eq!(events.len(), 3, "lazy load must skip the 0-byte file");
2377    }
2378
2379    #[test]
2380    fn test_flush_tenant_leaves_no_zero_byte_parquet_after_normal_flush() {
2381        // Direct verification of t-c1d3: a successful flush_tenant
2382        // leaves exactly one healthy *.parquet — no .tmp survivors,
2383        // no 0-byte files.
2384        let temp_dir = TempDir::new().unwrap();
2385        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2386
2387        for i in 0..5 {
2388            storage
2389                .append_event(event_with_tenant("alice", &format!("a-{i}")))
2390                .unwrap();
2391        }
2392        storage.flush().unwrap();
2393
2394        let mut tmp_count = 0;
2395        let mut zero_byte_count = 0;
2396        let mut healthy_count = 0;
2397        let mut stack = vec![temp_dir.path().to_path_buf()];
2398        while let Some(d) = stack.pop() {
2399            for entry in std::fs::read_dir(&d).unwrap().flatten() {
2400                let p = entry.path();
2401                if p.is_dir() {
2402                    stack.push(p);
2403                    continue;
2404                }
2405                let name = p.file_name().unwrap().to_string_lossy().into_owned();
2406                if name.ends_with(".parquet.tmp") {
2407                    tmp_count += 1;
2408                } else if name.ends_with(".parquet") {
2409                    if std::fs::metadata(&p).unwrap().len() == 0 {
2410                        zero_byte_count += 1;
2411                    } else {
2412                        healthy_count += 1;
2413                    }
2414                }
2415            }
2416        }
2417        assert_eq!(tmp_count, 0, ".parquet.tmp survivors after flush");
2418        assert_eq!(zero_byte_count, 0, "0-byte .parquet survivors after flush");
2419        assert_eq!(healthy_count, 1, "expected exactly one healthy parquet");
2420    }
2421
2422    #[test]
2423    fn test_new_calls_cleanup_partial_writes_on_boot() {
2424        // Drop a stale .tmp into a directory, then construct a
2425        // fresh ParquetStorage on it. The constructor must clean
2426        // it up.
2427        let temp_dir = TempDir::new().unwrap();
2428        let stale = temp_dir.path().join("orphan.parquet.tmp");
2429        std::fs::write(&stale, b"crash detritus").unwrap();
2430        assert!(stale.is_file());
2431
2432        let _storage = ParquetStorage::new(temp_dir.path()).unwrap();
2433        assert!(
2434            !stale.exists(),
2435            "stale tmp should have been cleaned by ParquetStorage::new"
2436        );
2437    }
2438
2439    // -----------------------------------------------------------------
2440    // Migration tests (Step 1, commit #3: flat → tenant-tree migration).
2441    // -----------------------------------------------------------------
2442
2443    /// Helper: produce a flat-layout Parquet file at the storage root,
2444    /// matching what pre-#2 deploys wrote. Uses the existing flush path
2445    /// briefly and then relocates the resulting file from
2446    /// default/<yyyy-mm>/ back up to the root, simulating the legacy
2447    /// state.
2448    fn seed_flat_layout_file(storage: &ParquetStorage, count: usize) -> PathBuf {
2449        for i in 0..count {
2450            storage
2451                .append_event(create_test_event(&format!("entity-{i}")))
2452                .unwrap();
2453        }
2454        storage.flush().unwrap();
2455
2456        // create_test_event uses tenant="default", so the just-flushed file
2457        // landed under <root>/default/<yyyy-mm>/. Find the newest file in
2458        // that subtree to avoid picking up files from other tenants seeded
2459        // by the test before us.
2460        let default_subtree = storage.storage_dir().join("default");
2461        let candidates = find_parquet_files_recursive(&default_subtree).unwrap();
2462        assert!(
2463            !candidates.is_empty(),
2464            "seed expected at least one file under default/"
2465        );
2466        let src = candidates.into_iter().max().unwrap();
2467
2468        let dst = storage.storage_dir().join(src.file_name().unwrap());
2469        std::fs::rename(&src, &dst).unwrap();
2470        // Best-effort cleanup of the now-empty intermediate dirs so the
2471        // migration tool only sees the flat file. (`remove_dir` succeeds
2472        // only on empty dirs, which is exactly the safety we want here.)
2473        if let Some(month_dir) = src.parent() {
2474            let _ = std::fs::remove_dir(month_dir);
2475            if let Some(tenant_dir) = month_dir.parent() {
2476                let _ = std::fs::remove_dir(tenant_dir);
2477            }
2478        }
2479        dst
2480    }
2481
2482    #[test]
2483    fn test_migrate_flat_layout_dry_run_touches_nothing() {
2484        let temp_dir = TempDir::new().unwrap();
2485        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2486        let flat = seed_flat_layout_file(&storage, 7);
2487        assert!(flat.is_file(), "test setup: flat file should exist");
2488
2489        let report = storage.migrate_flat_layout(true).unwrap();
2490        assert!(report.dry_run);
2491        assert_eq!(report.flat_files_seen, 1);
2492        assert_eq!(report.events_migrated, 7);
2493        assert_eq!(report.flat_files_removed, 0);
2494        assert_eq!(report.partitions_written, 0);
2495        assert!(
2496            flat.is_file(),
2497            "flat file must still be present after dry run"
2498        );
2499    }
2500
2501    #[test]
2502    fn test_migrate_flat_layout_moves_events_into_default_tree_and_removes_flat() {
2503        let temp_dir = TempDir::new().unwrap();
2504        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2505        let flat = seed_flat_layout_file(&storage, 5);
2506
2507        let report = storage.migrate_flat_layout(false).unwrap();
2508        assert!(!report.dry_run);
2509        assert_eq!(report.flat_files_seen, 1);
2510        assert_eq!(report.flat_files_removed, 1);
2511        assert_eq!(report.events_migrated, 5);
2512        assert!(report.partitions_written >= 1);
2513        assert!(
2514            !flat.exists(),
2515            "flat file should be deleted after migration"
2516        );
2517
2518        let post = find_parquet_files_recursive(temp_dir.path()).unwrap();
2519        assert!(
2520            post.iter().all(|p| {
2521                let rel = p
2522                    .strip_prefix(temp_dir.path())
2523                    .unwrap()
2524                    .to_string_lossy()
2525                    .into_owned();
2526                rel.starts_with(&format!("default{}", std::path::MAIN_SEPARATOR))
2527            }),
2528            "all migrated files should be under default/"
2529        );
2530
2531        let loaded = storage.load_all_events().unwrap();
2532        assert_eq!(loaded.len(), 5);
2533    }
2534
2535    #[test]
2536    fn test_migrate_flat_layout_is_idempotent_when_re_run_after_completion() {
2537        let temp_dir = TempDir::new().unwrap();
2538        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2539        let _flat = seed_flat_layout_file(&storage, 4);
2540
2541        let first = storage.migrate_flat_layout(false).unwrap();
2542        assert_eq!(first.events_migrated, 4);
2543
2544        // Second run sees no flat files at the root, so it's a no-op —
2545        // events do not duplicate even if an operator runs the tool twice.
2546        let second = storage.migrate_flat_layout(false).unwrap();
2547        assert_eq!(second.flat_files_seen, 0);
2548        assert_eq!(second.events_migrated, 0);
2549        assert_eq!(second.flat_files_removed, 0);
2550
2551        let loaded = storage.load_all_events().unwrap();
2552        assert_eq!(loaded.len(), 4, "rerun must not duplicate or lose events");
2553    }
2554
2555    #[test]
2556    fn test_migrate_flat_layout_ignores_already_partitioned_data() {
2557        // Mixed state: a tenant tree already exists alongside one flat file.
2558        // Migration must touch only the flat file.
2559        let temp_dir = TempDir::new().unwrap();
2560        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2561
2562        for i in 0..3 {
2563            storage
2564                .append_event(event_with_tenant("alice", &format!("a-{i}")))
2565                .unwrap();
2566        }
2567        storage.flush().unwrap();
2568
2569        let _flat = seed_flat_layout_file(&storage, 2);
2570
2571        let report = storage.migrate_flat_layout(false).unwrap();
2572        assert_eq!(report.flat_files_seen, 1, "only the flat file is in scope");
2573        assert_eq!(report.events_migrated, 2);
2574
2575        let alice_files = find_parquet_files_recursive(&temp_dir.path().join("alice")).unwrap();
2576        assert_eq!(alice_files.len(), 1, "alice's tree must be untouched");
2577
2578        let loaded = storage.load_all_events().unwrap();
2579        assert_eq!(loaded.len(), 5);
2580        let alice_count = loaded
2581            .iter()
2582            .filter(|e| e.tenant_id_str() == "alice")
2583            .count();
2584        let default_count = loaded
2585            .iter()
2586            .filter(|e| e.tenant_id_str() == "default")
2587            .count();
2588        assert_eq!(alice_count, 3);
2589        assert_eq!(default_count, 2);
2590    }
2591
2592    #[test]
2593    fn test_migrate_flat_layout_with_no_flat_files_is_a_clean_noop() {
2594        let temp_dir = TempDir::new().unwrap();
2595        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
2596        let report = storage.migrate_flat_layout(false).unwrap();
2597        assert_eq!(report.flat_files_seen, 0);
2598        assert_eq!(report.events_migrated, 0);
2599    }
2600}