Skip to main content

allsource_core/infrastructure/persistence/
compaction.rs

1use crate::{
2    error::{AllSourceError, Result},
3    infrastructure::persistence::{cold_tier::ArchiveTarget, storage::ParquetStorage},
4};
5use chrono::{DateTime, Utc};
6use parking_lot::RwLock;
7use serde::{Deserialize, Serialize};
8use std::{collections::HashMap, fs, path::PathBuf, sync::Arc, time::Duration};
9
10/// Manages Parquet file compaction for optimal storage and query performance.
11///
12/// Step 4 of the sustainable data strategy moved compaction from a
13/// global-stream pass to a per-tenant pass. Each invocation iterates
14/// the tenants discovered under `<storage_dir>/<tenant>/...` and
15/// emits a `snapshot.<tenant>.<from>-<to>.parquet` per qualifying
16/// chunk. Snapshot files are written atomically (tmp + rename) and
17/// their constituent raw files are removed only after the rename
18/// succeeds — so a mid-compaction crash leaves data intact.
19pub struct CompactionManager {
20    /// Directory where Parquet files are stored
21    storage_dir: PathBuf,
22
23    /// Configuration
24    config: CompactionConfig,
25
26    /// Statistics
27    stats: Arc<RwLock<CompactionStats>>,
28
29    /// Last compaction time
30    last_compaction: Arc<RwLock<Option<DateTime<Utc>>>>,
31}
32
33/// Filename prefix that marks a file as already-compacted output.
34/// Excluded from the input set when picking compaction candidates
35/// (we don't re-compact a snapshot until the snapshot itself
36/// triggers the strategy criteria from a future commit).
37const SNAPSHOT_PREFIX: &str = "snapshot.";
38
39#[derive(Debug, Clone)]
40pub struct CompactionConfig {
41    /// Minimum number of files to trigger compaction
42    pub min_files_to_compact: usize,
43
44    /// Target size for compacted files (in bytes)
45    pub target_file_size: usize,
46
47    /// Maximum size for a single compacted file (in bytes)
48    pub max_file_size: usize,
49
50    /// Minimum file size to consider for compaction (small files)
51    pub small_file_threshold: usize,
52
53    /// Time interval between automatic compactions (in seconds)
54    pub compaction_interval_seconds: u64,
55
56    /// Enable automatic background compaction
57    pub auto_compact: bool,
58
59    /// Compaction strategy
60    pub strategy: CompactionStrategy,
61
62    /// Per-tenant retention TTLs (Step 5 of the sustainable data
63    /// strategy). Applied during the same compaction pass — events
64    /// older than `now - ttl` for that tenant are dropped from the
65    /// snapshot output and the originals are removed. Default
66    /// honors the bead: tenant `system` keeps 30 days; everyone
67    /// else keeps forever.
68    pub retention: RetentionConfig,
69
70    /// Optional cold-tier archive. When set, events that would be
71    /// dropped by retention are archived to this target BEFORE the
72    /// originals are deleted. A failed archive aborts the
73    /// compaction pass — originals stay on disk and the next run
74    /// retries. Default `None` preserves the pre-cold-tier behavior:
75    /// retention deletes outright. See
76    /// `infrastructure::persistence::cold_tier`.
77    pub archive: Option<Arc<dyn ArchiveTarget>>,
78}
79
80/// Per-tenant retention configuration. Look up a TTL via
81/// `ttl_for(tenant_id)`; `None` means "keep forever" for that
82/// tenant.
83///
84/// The default rule (from the bead): the CP heartbeat tenant
85/// (`system`) defaults to 30 days. The CP emits ~69k heartbeat
86/// events/day; without retention this grows unbounded for data
87/// that has no audit value past the dashboard window. Other
88/// tenants default to no TTL — user data stays put unless the
89/// owner opts in.
90///
91/// Per-tenant overrides win over `default_ttl`; "no entry" falls
92/// back to `default_ttl`.
93#[derive(Debug, Clone)]
94pub struct RetentionConfig {
95    /// Default TTL when no per-tenant override exists. `None` = keep forever.
96    pub default_ttl: Option<Duration>,
97    /// Per-tenant overrides. `Some(None)` would mean "explicitly no
98    /// TTL"; the API uses `Option<Duration>` directly so an entry
99    /// can record an explicit "keep forever" decision distinct
100    /// from "no entry".
101    pub per_tenant_ttl: HashMap<String, Option<Duration>>,
102}
103
104impl Default for RetentionConfig {
105    fn default() -> Self {
106        let mut per_tenant_ttl = HashMap::new();
107        per_tenant_ttl.insert("system".to_string(), Some(Duration::from_hours(30 * 24)));
108        Self {
109            default_ttl: None,
110            per_tenant_ttl,
111        }
112    }
113}
114
115impl RetentionConfig {
116    /// Effective TTL for `tenant_id`. Returns `None` if the tenant
117    /// has no TTL (keep forever).
118    ///
119    /// Lookup order:
120    /// 1. Per-tenant entry → that value (whether Some or explicit None).
121    /// 2. No entry → fall back to `default_ttl`.
122    pub fn ttl_for(&self, tenant_id: &str) -> Option<Duration> {
123        match self.per_tenant_ttl.get(tenant_id) {
124            Some(v) => *v,
125            None => self.default_ttl,
126        }
127    }
128
129    /// Override the TTL for a specific tenant. Use `None` to mean
130    /// "keep forever for this tenant".
131    pub fn set(&mut self, tenant_id: &str, ttl: Option<Duration>) {
132        self.per_tenant_ttl.insert(tenant_id.to_string(), ttl);
133    }
134}
135
136impl Default for CompactionConfig {
137    fn default() -> Self {
138        Self {
139            min_files_to_compact: 3,
140            target_file_size: 128 * 1024 * 1024,    // 128 MB
141            max_file_size: 256 * 1024 * 1024,       // 256 MB
142            small_file_threshold: 10 * 1024 * 1024, // 10 MB
143            compaction_interval_seconds: 3600,      // 1 hour
144            auto_compact: true,
145            strategy: CompactionStrategy::SizeBased,
146            retention: RetentionConfig::default(),
147            archive: None,
148        }
149    }
150}
151
152impl CompactionConfig {
153    /// Build a config from the relevant env vars:
154    /// - `ALLSOURCE_SNAPSHOT_INTERVAL_SECONDS`: per-pass cadence
155    ///   (default 3600).
156    /// - `ALLSOURCE_RETENTION_SYSTEM_DAYS`: TTL for the `system`
157    ///   tenant in days (default 30).
158    ///
159    /// Unparseable values log a warning and fall back to defaults
160    /// — boot doesn't fail.
161    pub fn from_env() -> Self {
162        Self::from_env_vars(
163            std::env::var("ALLSOURCE_SNAPSHOT_INTERVAL_SECONDS").ok(),
164            std::env::var("ALLSOURCE_RETENTION_SYSTEM_DAYS").ok(),
165        )
166    }
167
168    /// Testable variant of `from_env`. Production calls `from_env`;
169    /// tests pass explicit values.
170    pub fn from_env_vars(
171        interval_var: Option<String>,
172        system_retention_days_var: Option<String>,
173    ) -> Self {
174        let mut config = Self::default();
175        if let Some(s) = interval_var.filter(|s| !s.is_empty()) {
176            match s.parse::<u64>() {
177                Ok(v) => config.compaction_interval_seconds = v,
178                Err(e) => {
179                    tracing::warn!(
180                        "ALLSOURCE_SNAPSHOT_INTERVAL_SECONDS={s:?} could not be parsed as \
181                         u64: {e}; defaulting to {}s",
182                        config.compaction_interval_seconds
183                    );
184                }
185            }
186        }
187        if let Some(s) = system_retention_days_var.filter(|s| !s.is_empty()) {
188            match s.parse::<u64>() {
189                Ok(days) => {
190                    config
191                        .retention
192                        .set("system", Some(Duration::from_secs(days * 24 * 3600)));
193                }
194                Err(e) => {
195                    tracing::warn!(
196                        "ALLSOURCE_RETENTION_SYSTEM_DAYS={s:?} could not be parsed as u64: \
197                         {e}; defaulting to 30 days for tenant=system"
198                    );
199                }
200            }
201        }
202        config
203    }
204
205    /// Backwards-compatible single-arg variant for existing
206    /// callers that only set the snapshot interval.
207    pub fn from_env_var(interval_var: Option<String>) -> Self {
208        Self::from_env_vars(interval_var, None)
209    }
210}
211
212#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq)]
213#[serde(rename_all = "lowercase")]
214pub enum CompactionStrategy {
215    /// Compact based on file size (default)
216    SizeBased,
217    /// Compact based on file age
218    TimeBased,
219    /// Compact all files into one
220    FullCompaction,
221}
222
223#[derive(Debug, Clone, Default, Serialize)]
224pub struct CompactionStats {
225    pub total_compactions: u64,
226    pub total_files_compacted: u64,
227    pub total_bytes_before: u64,
228    pub total_bytes_after: u64,
229    pub total_events_compacted: u64,
230    pub last_compaction_duration_ms: u64,
231    pub space_saved_bytes: u64,
232}
233
234/// Information about a Parquet file candidate for compaction
235#[derive(Debug, Clone)]
236struct FileInfo {
237    path: PathBuf,
238    size: u64,
239    created: DateTime<Utc>,
240}
241
242impl CompactionManager {
243    /// Create a new compaction manager
244    pub fn new(storage_dir: impl Into<PathBuf>, config: CompactionConfig) -> Self {
245        let storage_dir = storage_dir.into();
246
247        tracing::info!(
248            "✅ Compaction manager initialized at: {}",
249            storage_dir.display()
250        );
251
252        Self {
253            storage_dir,
254            config,
255            stats: Arc::new(RwLock::new(CompactionStats::default())),
256            last_compaction: Arc::new(RwLock::new(None)),
257        }
258    }
259
260    /// List all Parquet files in the storage directory
261    fn list_parquet_files(&self) -> Result<Vec<FileInfo>> {
262        let entries = fs::read_dir(&self.storage_dir).map_err(|e| {
263            AllSourceError::StorageError(format!("Failed to read storage directory: {e}"))
264        })?;
265
266        let mut files = Vec::new();
267
268        for entry in entries {
269            let entry = entry.map_err(|e| {
270                AllSourceError::StorageError(format!("Failed to read directory entry: {e}"))
271            })?;
272
273            let path = entry.path();
274            if let Some(ext) = path.extension()
275                && ext == "parquet"
276            {
277                let metadata = entry.metadata().map_err(|e| {
278                    AllSourceError::StorageError(format!("Failed to read file metadata: {e}"))
279                })?;
280
281                let size = metadata.len();
282                let created = metadata
283                    .created()
284                    .ok()
285                    .and_then(|t| {
286                        t.duration_since(std::time::UNIX_EPOCH).ok().map(|d| {
287                            DateTime::from_timestamp(d.as_secs() as i64, 0).unwrap_or_else(Utc::now)
288                        })
289                    })
290                    .unwrap_or_else(Utc::now);
291
292                files.push(FileInfo {
293                    path,
294                    size,
295                    created,
296                });
297            }
298        }
299
300        // Sort by creation time (oldest first)
301        files.sort_by_key(|f| f.created);
302
303        Ok(files)
304    }
305
306    /// Identify files that should be compacted based on strategy
307    fn select_files_for_compaction(&self, files: &[FileInfo]) -> Vec<FileInfo> {
308        match self.config.strategy {
309            CompactionStrategy::SizeBased => self.select_small_files(files),
310            CompactionStrategy::TimeBased => self.select_old_files(files),
311            CompactionStrategy::FullCompaction => files.to_vec(),
312        }
313    }
314
315    /// Select small files for compaction
316    fn select_small_files(&self, files: &[FileInfo]) -> Vec<FileInfo> {
317        let small_files: Vec<FileInfo> = files
318            .iter()
319            .filter(|f| f.size < self.config.small_file_threshold as u64)
320            .cloned()
321            .collect();
322
323        // Only compact if we have enough small files
324        if small_files.len() >= self.config.min_files_to_compact {
325            small_files
326        } else {
327            Vec::new()
328        }
329    }
330
331    /// Select old files for time-based compaction
332    fn select_old_files(&self, files: &[FileInfo]) -> Vec<FileInfo> {
333        let now = Utc::now();
334        let age_threshold = chrono::Duration::hours(24); // Files older than 24 hours
335
336        let old_files: Vec<FileInfo> = files
337            .iter()
338            .filter(|f| now - f.created > age_threshold)
339            .cloned()
340            .collect();
341
342        if old_files.len() >= self.config.min_files_to_compact {
343            old_files
344        } else {
345            Vec::new()
346        }
347    }
348
349    /// Check if compaction should run
350    #[cfg_attr(feature = "hotpath", hotpath::measure)]
351    pub fn should_compact(&self) -> bool {
352        if !self.config.auto_compact {
353            return false;
354        }
355
356        let last = self.last_compaction.read();
357        match *last {
358            None => true, // Never compacted
359            Some(last_time) => {
360                let elapsed = (Utc::now() - last_time).num_seconds();
361                elapsed >= self.config.compaction_interval_seconds as i64
362            }
363        }
364    }
365
366    /// Perform compaction across every discovered tenant.
367    ///
368    /// Iterates the tenants under `<storage_dir>/<tenant>/...`, calls
369    /// `compact_tenant` for each, and aggregates the results. Step 4
370    /// of the sustainable data strategy: per-tenant compaction
371    /// instead of global, keyed off the per-tenant directory tree
372    /// Step 1 introduced.
373    ///
374    /// Errors compacting one tenant are logged but don't abort the
375    /// pass — other tenants still get compacted. The aggregate
376    /// result reflects what actually completed.
377    #[cfg_attr(feature = "hotpath", hotpath::measure)]
378    pub fn compact(&self) -> Result<CompactionResult> {
379        let start_time = std::time::Instant::now();
380        tracing::info!("🔄 Starting per-tenant compaction sweep...");
381
382        let tenants = self.discover_tenants()?;
383        if tenants.is_empty() {
384            tracing::debug!("No tenants found under {}", self.storage_dir.display());
385            return Ok(CompactionResult::default());
386        }
387
388        let mut aggregate = CompactionResult::default();
389        for tenant in &tenants {
390            match self.compact_tenant(tenant) {
391                Ok(r) => {
392                    aggregate.files_compacted += r.files_compacted;
393                    aggregate.bytes_before += r.bytes_before;
394                    aggregate.bytes_after += r.bytes_after;
395                    aggregate.events_compacted += r.events_compacted;
396                }
397                Err(e) => {
398                    tracing::error!(
399                        tenant_id = %tenant,
400                        "compact_tenant failed: {e}"
401                    );
402                }
403            }
404        }
405        aggregate.duration_ms = start_time.elapsed().as_millis() as u64;
406
407        if aggregate.files_compacted > 0 {
408            let mut stats = self.stats.write();
409            stats.total_compactions += 1;
410            stats.total_files_compacted += aggregate.files_compacted as u64;
411            stats.total_bytes_before += aggregate.bytes_before;
412            stats.total_bytes_after += aggregate.bytes_after;
413            stats.total_events_compacted += aggregate.events_compacted as u64;
414            stats.last_compaction_duration_ms = aggregate.duration_ms;
415            stats.space_saved_bytes += aggregate.bytes_before.saturating_sub(aggregate.bytes_after);
416        }
417        *self.last_compaction.write() = Some(Utc::now());
418
419        tracing::info!(
420            "✅ Compaction sweep complete: {} files → 1 snapshot per tenant, \
421             {:.2} MB → {:.2} MB, {} events, {} tenants in {}ms",
422            aggregate.files_compacted,
423            aggregate.bytes_before as f64 / (1024.0 * 1024.0),
424            aggregate.bytes_after as f64 / (1024.0 * 1024.0),
425            aggregate.events_compacted,
426            tenants.len(),
427            aggregate.duration_ms
428        );
429
430        Ok(aggregate)
431    }
432
433    /// Compact one tenant's raw event files into a single snapshot
434    /// file under that tenant's partition.
435    ///
436    /// Per-tenant pipeline:
437    /// 1. List `<storage>/<tenant>/...*.parquet` excluding existing
438    ///    `snapshot.*` files.
439    /// 2. Apply the configured strategy (size / time / full) to
440    ///    pick candidate raw files.
441    /// 3. If enough candidates: read events, sort by timestamp,
442    ///    atomically write `snapshot.<tenant>.<from>-<to>.parquet`
443    ///    via `ParquetStorage::write_atomic_parquet`.
444    /// 4. After the snapshot rename succeeds, delete the
445    ///    constituent raw files. Crash between snapshot and delete
446    ///    leaves both on disk; the dedupe in `append_loaded_event`
447    ///    keeps memory consistent on next load. A future commit can
448    ///    record the constituent file list in the snapshot's
449    ///    metadata and let a cleanup pass finish the deletion.
450    pub fn compact_tenant(&self, tenant_id: &str) -> Result<CompactionResult> {
451        let start_time = std::time::Instant::now();
452
453        // 1. List raw files for this tenant.
454        let storage = ParquetStorage::new(&self.storage_dir)?;
455        let all_files = storage.list_parquet_files_for_tenant(tenant_id)?;
456        let raw_files: Vec<FileInfo> = all_files
457            .into_iter()
458            .filter(|p| {
459                p.file_name()
460                    .and_then(|n| n.to_str())
461                    .is_none_or(|n| !n.starts_with(SNAPSHOT_PREFIX))
462            })
463            .filter_map(|p| {
464                let metadata = fs::metadata(&p).ok()?;
465                let size = metadata.len();
466                let created = metadata
467                    .created()
468                    .ok()
469                    .and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
470                    .and_then(|d| DateTime::from_timestamp(d.as_secs() as i64, 0))
471                    .unwrap_or_else(Utc::now);
472                Some(FileInfo {
473                    path: p,
474                    size,
475                    created,
476                })
477            })
478            .collect();
479
480        // 2. Strategy filter.
481        let candidates = self.select_files_for_compaction(&raw_files);
482        if candidates.is_empty() {
483            tracing::debug!(
484                tenant_id = tenant_id,
485                strategy = ?self.config.strategy,
486                "no files meet compaction criteria"
487            );
488            return Ok(CompactionResult::default());
489        }
490
491        let bytes_before: u64 = candidates.iter().map(|f| f.size).sum();
492        tracing::info!(
493            tenant_id = tenant_id,
494            files = candidates.len(),
495            mib = bytes_before as f64 / (1024.0 * 1024.0),
496            "compacting tenant"
497        );
498
499        // 3. Read events from candidates. We deliberately read each
500        // file individually (not via load_events_for_tenant) so we
501        // only pick up the candidate set — concurrent writes that
502        // produced new files since step 1 are skipped, deferred to
503        // the next interval (AC #6).
504        let mut events = Vec::new();
505        for fi in &candidates {
506            // Recover tenant from path so loaded events keep
507            // their identity (Step 1's path-as-tenant-source).
508            let event_tenant = match fi.path.strip_prefix(&self.storage_dir).ok() {
509                Some(rel) => rel
510                    .components()
511                    .next()
512                    .and_then(|c| match c {
513                        std::path::Component::Normal(t) => Some(t.to_string_lossy().into_owned()),
514                        _ => None,
515                    })
516                    .unwrap_or_else(|| "default".to_string()),
517                None => "default".to_string(),
518            };
519            match storage.load_events_from_file_path(&fi.path, &event_tenant) {
520                Ok(mut e) => events.append(&mut e),
521                Err(e) => {
522                    tracing::error!(
523                        file = %fi.path.display(),
524                        "failed to read parquet file for compaction: {e}"
525                    );
526                }
527            }
528        }
529
530        if events.is_empty() {
531            tracing::warn!(
532                tenant_id = tenant_id,
533                "candidate files had no readable events; skipping snapshot"
534            );
535            return Ok(CompactionResult::default());
536        }
537
538        // Apply retention (Step 5). Drop events older than the
539        // tenant's TTL before they get rewritten into the
540        // snapshot. The originals get deleted at the end of the
541        // happy path either way, so dropped events go away with
542        // the same crash-safe guarantee as the rest of the
543        // compaction pipeline (snapshot rename completes BEFORE
544        // any original is deleted; AC #6).
545        //
546        // Cold-tier (sustainability): when an `archive` target is
547        // configured, dropped events are written to it BEFORE this
548        // function deletes any original file. A failed archive
549        // returns `Err`, originals stay on disk, and the next
550        // compaction pass retries — same crash-safety contract as
551        // the snapshot path. Without an archive, retention behaves
552        // exactly as before (delete outright).
553        let dropped_by_retention = if let Some(ttl) = self.config.retention.ttl_for(tenant_id) {
554            let cutoff = Utc::now()
555                - chrono::Duration::from_std(ttl).unwrap_or_else(|_| chrono::Duration::zero());
556            let before = events.len();
557
558            // Partition: drained = events past TTL, kept = events to keep.
559            // Using drain_filter would be simpler but it's unstable;
560            // partition + reassign keeps events: Vec<Event> with the
561            // kept slice in original order.
562            let (drained, kept): (Vec<_>, Vec<_>) = std::mem::take(&mut events)
563                .into_iter()
564                .partition(|e| e.timestamp < cutoff);
565            events = kept;
566            let dropped = before - events.len();
567
568            if dropped > 0 {
569                tracing::info!(
570                    retention_tenant = tenant_id,
571                    dropped = dropped,
572                    kept = events.len(),
573                    cutoff = %cutoff.to_rfc3339(),
574                    ttl_secs = ttl.as_secs(),
575                    "retention: dropped events older than TTL"
576                );
577
578                if let Some(archive) = self.config.archive.as_ref() {
579                    let from = drained
580                        .iter()
581                        .map(|e| e.timestamp)
582                        .min()
583                        .expect("dropped > 0 guarantees non-empty drained");
584                    let to = drained
585                        .iter()
586                        .map(|e| e.timestamp)
587                        .max()
588                        .expect("dropped > 0 guarantees non-empty drained");
589                    archive.archive(tenant_id, from, to, &drained)?;
590                    tracing::info!(
591                        retention_tenant = tenant_id,
592                        archived_to = %archive.description(),
593                        archived = drained.len(),
594                        "retention: dropped events archived to cold tier"
595                    );
596                }
597            }
598            dropped
599        } else {
600            0
601        };
602
603        // Edge case: every event aged out. Skip the snapshot
604        // write and just delete originals — the data is gone by
605        // design. Delete-without-snapshot is safe here because
606        // every event we'd have written to the snapshot was
607        // already past its TTL, and the input files live under
608        // the tenant's partition with nothing else relying on
609        // them.
610        if events.is_empty() {
611            tracing::info!(
612                tenant_id = tenant_id,
613                files_dropped = candidates.len(),
614                events_dropped = dropped_by_retention,
615                "retention: every event aged out — deleting originals without snapshot"
616            );
617            for fi in &candidates {
618                if let Err(e) = fs::remove_file(&fi.path) {
619                    tracing::error!(
620                        file = %fi.path.display(),
621                        "failed to remove fully-aged raw file: {e}"
622                    );
623                }
624            }
625            return Ok(CompactionResult {
626                files_compacted: candidates.len(),
627                bytes_before,
628                bytes_after: 0,
629                events_compacted: 0,
630                duration_ms: start_time.elapsed().as_millis() as u64,
631            });
632        }
633
634        events.sort_by_key(|e| e.timestamp);
635        let from = events.first().expect("non-empty checked above").timestamp;
636        let to = events.last().expect("non-empty checked above").timestamp;
637
638        // Filename: snapshot.<tenant>.<from>-<to> with filesystem-safe
639        // ISO-basic timestamps (no colons).
640        let file_stem = format!(
641            "snapshot.{tenant_id}.{}-{}",
642            format_iso_basic(from),
643            format_iso_basic(to)
644        );
645        let snapshot_path = storage.write_atomic_parquet(tenant_id, &file_stem, &events)?;
646        let bytes_after = fs::metadata(&snapshot_path).map_or(0, |m| m.len());
647
648        // 4. Delete originals AFTER snapshot is durably renamed.
649        // AC #6: a snapshot-write failure short-circuits via the
650        // `?` above, so originals stay on disk and events remain
651        // queryable until the next successful pass.
652        for fi in &candidates {
653            if let Err(e) = fs::remove_file(&fi.path) {
654                tracing::error!(
655                    file = %fi.path.display(),
656                    "failed to remove pre-snapshot raw file: {e}"
657                );
658            }
659        }
660
661        let duration_ms = start_time.elapsed().as_millis() as u64;
662        tracing::info!(
663            tenant_id = tenant_id,
664            files_compacted = candidates.len(),
665            events = events.len(),
666            dropped_by_retention = dropped_by_retention,
667            mib_before = bytes_before as f64 / (1024.0 * 1024.0),
668            mib_after = bytes_after as f64 / (1024.0 * 1024.0),
669            duration_ms = duration_ms,
670            "tenant compaction complete"
671        );
672
673        Ok(CompactionResult {
674            files_compacted: candidates.len(),
675            bytes_before,
676            bytes_after,
677            events_compacted: events.len(),
678            duration_ms,
679        })
680    }
681
682    /// Discover tenant ids by scanning `<storage_dir>/<X>/` for
683    /// directories. The migration tool's flat-layout files at the
684    /// root are not picked up — Step 1 #3's migration moves them
685    /// under `default/`.
686    fn discover_tenants(&self) -> Result<Vec<String>> {
687        let Ok(entries) = fs::read_dir(&self.storage_dir) else {
688            return Ok(Vec::new());
689        };
690        let mut tenants: Vec<String> = entries
691            .filter_map(std::result::Result::ok)
692            .filter_map(|entry| {
693                let ft = entry.file_type().ok()?;
694                if !ft.is_dir() {
695                    return None;
696                }
697                let name = entry.file_name().to_string_lossy().into_owned();
698                // Skip the system metadata subtree (Core's own
699                // event-sourced repos) and any hidden folders.
700                if name.starts_with('.') || name == "__system" {
701                    return None;
702                }
703                Some(name)
704            })
705            .collect();
706        tenants.sort();
707        Ok(tenants)
708    }
709
710    /// Get compaction statistics
711    pub fn stats(&self) -> CompactionStats {
712        (*self.stats.read()).clone()
713    }
714
715    /// Get configuration
716    pub fn config(&self) -> &CompactionConfig {
717        &self.config
718    }
719
720    /// Trigger manual compaction
721    #[cfg_attr(feature = "hotpath", hotpath::measure)]
722    pub fn compact_now(&self) -> Result<CompactionResult> {
723        tracing::info!("Manual compaction triggered");
724        self.compact()
725    }
726}
727
728/// Result of a compaction operation
729#[derive(Debug, Clone, Default, Serialize)]
730pub struct CompactionResult {
731    pub files_compacted: usize,
732    pub bytes_before: u64,
733    pub bytes_after: u64,
734    pub events_compacted: usize,
735    pub duration_ms: u64,
736}
737
738/// Format a UTC timestamp as a filename-safe ISO-8601 basic-form
739/// string: `2026-04-27T134567Z` — no colons, no fractional second.
740/// Used for the `<from>-<to>` portion of snapshot filenames so
741/// they're portable across filesystems. `pub(super)` so the
742/// cold-tier archive can reuse the same naming convention.
743pub(super) fn format_iso_basic(t: DateTime<Utc>) -> String {
744    t.format("%Y-%m-%dT%H%M%SZ").to_string()
745}
746
747/// Background compaction task
748pub struct CompactionTask {
749    manager: Arc<CompactionManager>,
750    interval: Duration,
751}
752
753impl CompactionTask {
754    /// Create a new background compaction task
755    pub fn new(manager: Arc<CompactionManager>, interval_seconds: u64) -> Self {
756        Self {
757            manager,
758            interval: Duration::from_secs(interval_seconds),
759        }
760    }
761
762    /// Run the compaction task in a loop
763    #[cfg_attr(feature = "hotpath", hotpath::measure)]
764    pub async fn run(self) {
765        let mut interval = tokio::time::interval(self.interval);
766
767        loop {
768            interval.tick().await;
769
770            if self.manager.should_compact() {
771                tracing::debug!("Auto-compaction check triggered");
772
773                match self.manager.compact() {
774                    Ok(result) => {
775                        if result.files_compacted > 0 {
776                            tracing::info!(
777                                "Auto-compaction succeeded: {} files, {:.2} MB saved",
778                                result.files_compacted,
779                                (result.bytes_before - result.bytes_after) as f64
780                                    / (1024.0 * 1024.0)
781                            );
782                        }
783                    }
784                    Err(e) => {
785                        tracing::error!("Auto-compaction failed: {}", e);
786                    }
787                }
788            }
789        }
790    }
791}
792
793#[cfg(test)]
794mod tests {
795    use super::*;
796    use tempfile::TempDir;
797
798    #[test]
799    fn test_compaction_manager_creation() {
800        let temp_dir = TempDir::new().unwrap();
801        let config = CompactionConfig::default();
802        let manager = CompactionManager::new(temp_dir.path(), config);
803
804        assert_eq!(manager.stats().total_compactions, 0);
805    }
806
807    #[test]
808    fn test_should_compact() {
809        let temp_dir = TempDir::new().unwrap();
810        let config = CompactionConfig {
811            auto_compact: true,
812            compaction_interval_seconds: 1,
813            ..Default::default()
814        };
815        let manager = CompactionManager::new(temp_dir.path(), config);
816
817        // Should compact on first check (never compacted)
818        assert!(manager.should_compact());
819    }
820
821    #[test]
822    fn test_file_selection_size_based() {
823        let temp_dir = TempDir::new().unwrap();
824        let config = CompactionConfig {
825            small_file_threshold: 1024 * 1024, // 1 MB
826            min_files_to_compact: 2,
827            strategy: CompactionStrategy::SizeBased,
828            ..Default::default()
829        };
830        let manager = CompactionManager::new(temp_dir.path(), config);
831
832        let files = vec![
833            FileInfo {
834                path: PathBuf::from("small1.parquet"),
835                size: 500_000, // 500 KB
836                created: Utc::now(),
837            },
838            FileInfo {
839                path: PathBuf::from("small2.parquet"),
840                size: 600_000, // 600 KB
841                created: Utc::now(),
842            },
843            FileInfo {
844                path: PathBuf::from("large.parquet"),
845                size: 10_000_000, // 10 MB
846                created: Utc::now(),
847            },
848        ];
849
850        let selected = manager.select_files_for_compaction(&files);
851        assert_eq!(selected.len(), 2); // Only the 2 small files
852    }
853
854    #[test]
855    fn test_default_compaction_config() {
856        let config = CompactionConfig::default();
857        assert_eq!(config.min_files_to_compact, 3);
858        assert_eq!(config.target_file_size, 128 * 1024 * 1024);
859        assert_eq!(config.max_file_size, 256 * 1024 * 1024);
860        assert_eq!(config.small_file_threshold, 10 * 1024 * 1024);
861        assert_eq!(config.compaction_interval_seconds, 3600);
862        assert!(config.auto_compact);
863        assert_eq!(config.strategy, CompactionStrategy::SizeBased);
864    }
865
866    #[test]
867    fn test_should_compact_disabled() {
868        let temp_dir = TempDir::new().unwrap();
869        let config = CompactionConfig {
870            auto_compact: false,
871            ..Default::default()
872        };
873        let manager = CompactionManager::new(temp_dir.path(), config);
874
875        assert!(!manager.should_compact());
876    }
877
878    #[test]
879    fn test_compact_empty_directory() {
880        let temp_dir = TempDir::new().unwrap();
881        let config = CompactionConfig::default();
882        let manager = CompactionManager::new(temp_dir.path(), config);
883
884        let result = manager.compact().unwrap();
885        assert_eq!(result.files_compacted, 0);
886        assert_eq!(result.bytes_before, 0);
887        assert_eq!(result.bytes_after, 0);
888        assert_eq!(result.events_compacted, 0);
889    }
890
891    #[test]
892    fn test_compact_now() {
893        let temp_dir = TempDir::new().unwrap();
894        let config = CompactionConfig::default();
895        let manager = CompactionManager::new(temp_dir.path(), config);
896
897        let result = manager.compact_now().unwrap();
898        assert_eq!(result.files_compacted, 0);
899    }
900
901    #[test]
902    fn test_get_config() {
903        let temp_dir = TempDir::new().unwrap();
904        let config = CompactionConfig {
905            min_files_to_compact: 5,
906            ..Default::default()
907        };
908        let manager = CompactionManager::new(temp_dir.path(), config);
909
910        assert_eq!(manager.config().min_files_to_compact, 5);
911    }
912
913    #[test]
914    fn test_get_stats() {
915        let temp_dir = TempDir::new().unwrap();
916        let config = CompactionConfig::default();
917        let manager = CompactionManager::new(temp_dir.path(), config);
918
919        let stats = manager.stats();
920        assert_eq!(stats.total_compactions, 0);
921        assert_eq!(stats.total_files_compacted, 0);
922        assert_eq!(stats.total_bytes_before, 0);
923        assert_eq!(stats.total_bytes_after, 0);
924        assert_eq!(stats.total_events_compacted, 0);
925        assert_eq!(stats.last_compaction_duration_ms, 0);
926        assert_eq!(stats.space_saved_bytes, 0);
927    }
928
929    #[test]
930    fn test_file_selection_not_enough_small_files() {
931        let temp_dir = TempDir::new().unwrap();
932        let config = CompactionConfig {
933            small_file_threshold: 1024 * 1024,
934            min_files_to_compact: 3, // Need 3 files
935            strategy: CompactionStrategy::SizeBased,
936            ..Default::default()
937        };
938        let manager = CompactionManager::new(temp_dir.path(), config);
939
940        let files = vec![
941            FileInfo {
942                path: PathBuf::from("small1.parquet"),
943                size: 500_000,
944                created: Utc::now(),
945            },
946            FileInfo {
947                path: PathBuf::from("small2.parquet"),
948                size: 600_000,
949                created: Utc::now(),
950            },
951        ];
952
953        let selected = manager.select_files_for_compaction(&files);
954        assert_eq!(selected.len(), 0); // Not enough small files
955    }
956
957    #[test]
958    fn test_file_selection_time_based() {
959        let temp_dir = TempDir::new().unwrap();
960        let config = CompactionConfig {
961            min_files_to_compact: 2,
962            strategy: CompactionStrategy::TimeBased,
963            ..Default::default()
964        };
965        let manager = CompactionManager::new(temp_dir.path(), config);
966
967        let old_time = Utc::now() - chrono::Duration::hours(48);
968        let files = vec![
969            FileInfo {
970                path: PathBuf::from("old1.parquet"),
971                size: 1_000_000,
972                created: old_time,
973            },
974            FileInfo {
975                path: PathBuf::from("old2.parquet"),
976                size: 2_000_000,
977                created: old_time,
978            },
979            FileInfo {
980                path: PathBuf::from("new.parquet"),
981                size: 500_000,
982                created: Utc::now(),
983            },
984        ];
985
986        let selected = manager.select_files_for_compaction(&files);
987        assert_eq!(selected.len(), 2); // Only the 2 old files
988    }
989
990    #[test]
991    fn test_file_selection_time_based_not_enough() {
992        let temp_dir = TempDir::new().unwrap();
993        let config = CompactionConfig {
994            min_files_to_compact: 3,
995            strategy: CompactionStrategy::TimeBased,
996            ..Default::default()
997        };
998        let manager = CompactionManager::new(temp_dir.path(), config);
999
1000        let old_time = Utc::now() - chrono::Duration::hours(48);
1001        let files = vec![
1002            FileInfo {
1003                path: PathBuf::from("old1.parquet"),
1004                size: 1_000_000,
1005                created: old_time,
1006            },
1007            FileInfo {
1008                path: PathBuf::from("new.parquet"),
1009                size: 500_000,
1010                created: Utc::now(),
1011            },
1012        ];
1013
1014        let selected = manager.select_files_for_compaction(&files);
1015        assert_eq!(selected.len(), 0); // Not enough old files
1016    }
1017
1018    #[test]
1019    fn test_file_selection_full_compaction() {
1020        let temp_dir = TempDir::new().unwrap();
1021        let config = CompactionConfig {
1022            strategy: CompactionStrategy::FullCompaction,
1023            ..Default::default()
1024        };
1025        let manager = CompactionManager::new(temp_dir.path(), config);
1026
1027        let files = vec![
1028            FileInfo {
1029                path: PathBuf::from("file1.parquet"),
1030                size: 1_000_000,
1031                created: Utc::now(),
1032            },
1033            FileInfo {
1034                path: PathBuf::from("file2.parquet"),
1035                size: 2_000_000,
1036                created: Utc::now(),
1037            },
1038        ];
1039
1040        let selected = manager.select_files_for_compaction(&files);
1041        assert_eq!(selected.len(), 2); // All files selected
1042    }
1043
1044    #[test]
1045    fn test_compaction_strategy_serde() {
1046        let strategies = vec![
1047            CompactionStrategy::SizeBased,
1048            CompactionStrategy::TimeBased,
1049            CompactionStrategy::FullCompaction,
1050        ];
1051
1052        for strategy in strategies {
1053            let json = serde_json::to_string(&strategy).unwrap();
1054            let parsed: CompactionStrategy = serde_json::from_str(&json).unwrap();
1055            assert_eq!(parsed, strategy);
1056        }
1057    }
1058
1059    #[test]
1060    fn test_compaction_stats_default() {
1061        let stats = CompactionStats::default();
1062        assert_eq!(stats.total_compactions, 0);
1063        assert_eq!(stats.total_files_compacted, 0);
1064    }
1065
1066    #[test]
1067    fn test_compaction_stats_serde() {
1068        let stats = CompactionStats {
1069            total_compactions: 5,
1070            total_files_compacted: 20,
1071            total_bytes_before: 1000000,
1072            total_bytes_after: 500000,
1073            total_events_compacted: 10000,
1074            last_compaction_duration_ms: 500,
1075            space_saved_bytes: 500000,
1076        };
1077
1078        let json = serde_json::to_string(&stats).unwrap();
1079        assert!(json.contains("\"total_compactions\":5"));
1080        assert!(json.contains("\"space_saved_bytes\":500000"));
1081    }
1082
1083    #[test]
1084    fn test_compaction_result_serde() {
1085        let result = CompactionResult {
1086            files_compacted: 3,
1087            bytes_before: 1000000,
1088            bytes_after: 500000,
1089            events_compacted: 5000,
1090            duration_ms: 250,
1091        };
1092
1093        let json = serde_json::to_string(&result).unwrap();
1094        assert!(json.contains("\"files_compacted\":3"));
1095        assert!(json.contains("\"bytes_before\":1000000"));
1096    }
1097
1098    #[test]
1099    fn test_compaction_task_creation() {
1100        let temp_dir = TempDir::new().unwrap();
1101        let config = CompactionConfig::default();
1102        let manager = Arc::new(CompactionManager::new(temp_dir.path(), config));
1103
1104        let _task = CompactionTask::new(manager.clone(), 60);
1105        // Task created successfully
1106    }
1107
1108    #[test]
1109    fn test_list_parquet_files_empty() {
1110        let temp_dir = TempDir::new().unwrap();
1111        let config = CompactionConfig::default();
1112        let manager = CompactionManager::new(temp_dir.path(), config);
1113
1114        let files = manager.list_parquet_files().unwrap();
1115        assert!(files.is_empty());
1116    }
1117
1118    #[test]
1119    fn test_list_parquet_files_with_non_parquet() {
1120        let temp_dir = TempDir::new().unwrap();
1121        let config = CompactionConfig::default();
1122        let manager = CompactionManager::new(temp_dir.path(), config);
1123
1124        // Create non-parquet files
1125        std::fs::write(temp_dir.path().join("test.txt"), "test").unwrap();
1126        std::fs::write(temp_dir.path().join("data.json"), "{}").unwrap();
1127
1128        let files = manager.list_parquet_files().unwrap();
1129        assert!(files.is_empty()); // No parquet files
1130    }
1131
1132    // -----------------------------------------------------------------
1133    // Per-tenant compaction tests (Step 4, commit #2).
1134    // -----------------------------------------------------------------
1135
1136    fn ingest_and_flush_per_call(storage_dir: &std::path::Path, tenant: &str, count: usize) {
1137        // Each ingested batch produces one parquet file when its
1138        // ParquetStorage is dropped, giving us multiple raw files
1139        // per tenant for the strategy filter to pick up.
1140        for i in 0..count {
1141            let storage = ParquetStorage::with_config(
1142                storage_dir,
1143                crate::infrastructure::persistence::ParquetStorageConfig {
1144                    batch_size: 1,
1145                    ..Default::default()
1146                },
1147            )
1148            .unwrap();
1149            let event = crate::domain::entities::Event::from_strings(
1150                "test.event".to_string(),
1151                format!("{tenant}-{i}"),
1152                tenant.to_string(),
1153                serde_json::json!({"i": i}),
1154                None,
1155            )
1156            .unwrap();
1157            storage.append_event(event).unwrap();
1158            storage.flush().unwrap();
1159        }
1160    }
1161
1162    #[test]
1163    fn test_compact_tenant_emits_one_snapshot_and_removes_originals() {
1164        let temp_dir = TempDir::new().unwrap();
1165
1166        // Seed 4 raw files for alice.
1167        ingest_and_flush_per_call(temp_dir.path(), "alice", 4);
1168
1169        let config = CompactionConfig {
1170            min_files_to_compact: 2,
1171            small_file_threshold: 100 * 1024 * 1024,
1172            strategy: CompactionStrategy::SizeBased,
1173            ..Default::default()
1174        };
1175        let manager = CompactionManager::new(temp_dir.path(), config);
1176
1177        let result = manager.compact_tenant("alice").unwrap();
1178        assert_eq!(result.files_compacted, 4);
1179        assert_eq!(result.events_compacted, 4);
1180
1181        // After: zero raw files, exactly one snapshot.* file under
1182        // alice's subtree.
1183        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1184        let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1185        assert_eq!(
1186            alice_files.len(),
1187            1,
1188            "expected exactly one snapshot file for alice"
1189        );
1190
1191        let name = alice_files[0]
1192            .file_name()
1193            .and_then(|n| n.to_str())
1194            .unwrap()
1195            .to_string();
1196        assert!(
1197            name.starts_with("snapshot.alice."),
1198            "expected snapshot prefix, got {name}"
1199        );
1200        assert!(name.ends_with(".parquet"));
1201
1202        // No tmp files left behind.
1203        let tmps: Vec<_> = std::fs::read_dir(alice_files[0].parent().unwrap())
1204            .unwrap()
1205            .filter_map(std::result::Result::ok)
1206            .filter(|e| e.path().to_string_lossy().ends_with(".tmp"))
1207            .collect();
1208        assert!(tmps.is_empty());
1209
1210        // Loaded events round-trip correctly.
1211        let loaded = storage.load_events_for_tenant("alice").unwrap();
1212        assert_eq!(loaded.len(), 4);
1213        for e in &loaded {
1214            assert_eq!(e.tenant_id_str(), "alice");
1215        }
1216    }
1217
1218    #[test]
1219    fn test_compact_tenant_skips_existing_snapshot_files() {
1220        // Seed alice with 4 raw files, run compaction → 1 snapshot.
1221        // Run compaction AGAIN: the snapshot is excluded from the
1222        // candidate set, no qualifying raw files remain, the second
1223        // pass is a no-op.
1224        let temp_dir = TempDir::new().unwrap();
1225        ingest_and_flush_per_call(temp_dir.path(), "alice", 4);
1226
1227        let config = CompactionConfig {
1228            min_files_to_compact: 2,
1229            small_file_threshold: 100 * 1024 * 1024,
1230            ..Default::default()
1231        };
1232        let manager = CompactionManager::new(temp_dir.path(), config);
1233
1234        let r1 = manager.compact_tenant("alice").unwrap();
1235        assert_eq!(r1.files_compacted, 4);
1236
1237        let r2 = manager.compact_tenant("alice").unwrap();
1238        assert_eq!(r2.files_compacted, 0, "snapshot must not be re-compacted");
1239
1240        // Still exactly one file on disk for alice.
1241        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1242        let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1243        assert_eq!(alice_files.len(), 1);
1244    }
1245
1246    #[test]
1247    fn test_compact_tenant_below_threshold_is_a_noop() {
1248        // 1 file < min_files_to_compact (default 3) → nothing happens.
1249        let temp_dir = TempDir::new().unwrap();
1250        ingest_and_flush_per_call(temp_dir.path(), "alice", 1);
1251
1252        let manager = CompactionManager::new(temp_dir.path(), CompactionConfig::default());
1253        let result = manager.compact_tenant("alice").unwrap();
1254        assert_eq!(result.files_compacted, 0);
1255
1256        // Raw file untouched.
1257        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1258        let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1259        assert_eq!(alice_files.len(), 1);
1260        let name = alice_files[0]
1261            .file_name()
1262            .unwrap()
1263            .to_string_lossy()
1264            .into_owned();
1265        assert!(
1266            !name.starts_with("snapshot."),
1267            "raw file must not be renamed"
1268        );
1269    }
1270
1271    #[test]
1272    fn test_compact_iterates_every_tenant() {
1273        // Two tenants each with enough raw files. compact() must
1274        // produce one snapshot per tenant.
1275        let temp_dir = TempDir::new().unwrap();
1276        ingest_and_flush_per_call(temp_dir.path(), "alice", 3);
1277        ingest_and_flush_per_call(temp_dir.path(), "bob", 3);
1278
1279        let config = CompactionConfig {
1280            min_files_to_compact: 2,
1281            small_file_threshold: 100 * 1024 * 1024,
1282            strategy: CompactionStrategy::SizeBased,
1283            ..Default::default()
1284        };
1285        let manager = CompactionManager::new(temp_dir.path(), config);
1286
1287        let result = manager.compact().unwrap();
1288        assert_eq!(result.files_compacted, 6);
1289        assert_eq!(result.events_compacted, 6);
1290
1291        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1292        for tenant in ["alice", "bob"] {
1293            let files = storage.list_parquet_files_for_tenant(tenant).unwrap();
1294            assert_eq!(files.len(), 1, "{tenant} should have one snapshot");
1295            let name = files[0].file_name().unwrap().to_string_lossy().into_owned();
1296            assert!(name.starts_with(&format!("snapshot.{tenant}.")));
1297        }
1298    }
1299
1300    #[test]
1301    fn test_retention_drops_events_older_than_ttl() {
1302        // Bead's integration test: ingest 100 events spanning 60
1303        // days, run compaction with 30-day TTL, only the last
1304        // 30 days remain queryable.
1305        let temp_dir = TempDir::new().unwrap();
1306
1307        // Seed alice with 100 events, timestamps spread across
1308        // 60 days. We can't backdate via ingest (timestamp = now()
1309        // in domain), so write parquet directly via flush of
1310        // back-dated events.
1311        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1312        let now = Utc::now();
1313        for i in 0..100 {
1314            // Day 0 = 60 days ago; day 99 ≈ now. Spreads evenly.
1315            let day_offset = 60 - (i * 60 / 99);
1316            let ts = now - chrono::Duration::days(i64::from(day_offset));
1317            let event = crate::domain::entities::Event::reconstruct_from_strings(
1318                uuid::Uuid::new_v4(),
1319                "test.event".to_string(),
1320                format!("e-{i}"),
1321                "alice".to_string(),
1322                serde_json::json!({"i": i}),
1323                ts,
1324                None,
1325                1,
1326            );
1327            storage.append_event(event).unwrap();
1328            // Every 10 events get a fresh storage to produce a
1329            // separate file (so compaction has multiple files
1330            // to merge).
1331            if i % 10 == 9 {
1332                storage.flush().unwrap();
1333            }
1334        }
1335        storage.flush().unwrap();
1336
1337        // 30-day TTL for alice via per-tenant override.
1338        let mut retention = RetentionConfig::default();
1339        retention.set("alice", Some(Duration::from_hours(30 * 24)));
1340        let config = CompactionConfig {
1341            min_files_to_compact: 2,
1342            small_file_threshold: 100 * 1024 * 1024,
1343            strategy: CompactionStrategy::SizeBased,
1344            retention,
1345            ..Default::default()
1346        };
1347        let manager = CompactionManager::new(temp_dir.path(), config);
1348
1349        let result = manager.compact_tenant("alice").unwrap();
1350        assert!(result.events_compacted > 0);
1351        assert!(
1352            result.events_compacted < 100,
1353            "retention should have dropped some events; kept {} of 100",
1354            result.events_compacted
1355        );
1356
1357        // Re-load to confirm the dropped events are gone.
1358        let storage2 = ParquetStorage::new(temp_dir.path()).unwrap();
1359        let loaded = storage2.load_events_for_tenant("alice").unwrap();
1360        assert_eq!(loaded.len(), result.events_compacted);
1361
1362        // Every loaded event must be within the 30-day window
1363        // (with a generous fudge for test-clock drift).
1364        let cutoff = Utc::now() - chrono::Duration::days(30);
1365        for e in &loaded {
1366            assert!(
1367                e.timestamp >= cutoff - chrono::Duration::seconds(60),
1368                "event with ts {} survived retention but is older than cutoff {}",
1369                e.timestamp.to_rfc3339(),
1370                cutoff.to_rfc3339()
1371            );
1372        }
1373    }
1374
1375    #[test]
1376    fn test_retention_keeps_forever_by_default_for_non_system_tenants() {
1377        // alice has no override → falls through to default_ttl
1378        // which is None. All events kept regardless of age.
1379        let temp_dir = TempDir::new().unwrap();
1380        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1381        let now = Utc::now();
1382        for i in 0..6 {
1383            let ts = now - chrono::Duration::days(i * 365);
1384            let event = crate::domain::entities::Event::reconstruct_from_strings(
1385                uuid::Uuid::new_v4(),
1386                "test.event".to_string(),
1387                format!("e-{i}"),
1388                "alice".to_string(),
1389                serde_json::json!({"i": i}),
1390                ts,
1391                None,
1392                1,
1393            );
1394            storage.append_event(event).unwrap();
1395            if i % 2 == 1 {
1396                storage.flush().unwrap();
1397            }
1398        }
1399        storage.flush().unwrap();
1400
1401        // Default RetentionConfig: alice has no entry, default_ttl = None.
1402        let config = CompactionConfig {
1403            min_files_to_compact: 2,
1404            small_file_threshold: 100 * 1024 * 1024,
1405            strategy: CompactionStrategy::SizeBased,
1406            ..Default::default()
1407        };
1408        let manager = CompactionManager::new(temp_dir.path(), config);
1409        let result = manager.compact_tenant("alice").unwrap();
1410        assert_eq!(result.events_compacted, 6, "no events should be dropped");
1411    }
1412
1413    #[test]
1414    fn test_retention_system_tenant_default_is_30_days() {
1415        // The system tenant's 30-day TTL is the bead's headline
1416        // requirement. Default config without any overrides
1417        // should already enforce it.
1418        let cfg = RetentionConfig::default();
1419        let ttl = cfg.ttl_for("system").unwrap();
1420        assert_eq!(ttl.as_secs(), 30 * 24 * 3600);
1421        // No override for arbitrary tenants → keep forever.
1422        assert!(cfg.ttl_for("acme").is_none());
1423    }
1424
1425    #[test]
1426    fn test_retention_drops_all_events_deletes_originals_without_snapshot() {
1427        // Edge case: every event is past the TTL. We delete the
1428        // raw files and emit no snapshot — there's nothing to
1429        // write. Tenant ends with zero files on disk.
1430        let temp_dir = TempDir::new().unwrap();
1431        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1432        let very_old = Utc::now() - chrono::Duration::days(90);
1433        for i in 0..6 {
1434            let event = crate::domain::entities::Event::reconstruct_from_strings(
1435                uuid::Uuid::new_v4(),
1436                "test.event".to_string(),
1437                format!("e-{i}"),
1438                "alice".to_string(),
1439                serde_json::json!({"i": i}),
1440                very_old,
1441                None,
1442                1,
1443            );
1444            storage.append_event(event).unwrap();
1445            if i % 2 == 1 {
1446                storage.flush().unwrap();
1447            }
1448        }
1449        storage.flush().unwrap();
1450
1451        let mut retention = RetentionConfig::default();
1452        retention.set("alice", Some(Duration::from_hours(7 * 24)));
1453        let config = CompactionConfig {
1454            min_files_to_compact: 2,
1455            small_file_threshold: 100 * 1024 * 1024,
1456            strategy: CompactionStrategy::SizeBased,
1457            retention,
1458            ..Default::default()
1459        };
1460        let manager = CompactionManager::new(temp_dir.path(), config);
1461        let result = manager.compact_tenant("alice").unwrap();
1462        assert_eq!(result.events_compacted, 0);
1463        assert!(result.files_compacted >= 2); // originals were deleted
1464
1465        // Tenant subtree has zero parquet files now.
1466        let storage2 = ParquetStorage::new(temp_dir.path()).unwrap();
1467        let alice_files = storage2.list_parquet_files_for_tenant("alice").unwrap();
1468        assert!(alice_files.is_empty(), "all originals should be deleted");
1469    }
1470
1471    #[test]
1472    fn test_compaction_with_simulated_crash_leaves_data_recoverable() {
1473        // AC #7-#8: simulate "crash mid-snapshot" by manually
1474        // dropping a tmp file in the partition, then asserting
1475        // a fresh ParquetStorage::new boot cleans it up and the
1476        // raw files (which would still be present in a real
1477        // mid-rename crash) remain queryable.
1478        let temp_dir = TempDir::new().unwrap();
1479
1480        // Seed alice with raw files.
1481        ingest_and_flush_per_call(temp_dir.path(), "alice", 3);
1482
1483        // Locate alice's partition dir.
1484        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1485        let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1486        let partition = alice_files[0].parent().unwrap().to_path_buf();
1487
1488        // Simulate the crash state: write a partial .tmp file as
1489        // if write_atomic_parquet had crashed mid-rename.
1490        let crashed_tmp = partition.join("snapshot.alice.range.parquet.tmp");
1491        std::fs::write(&crashed_tmp, b"partial parquet bytes").unwrap();
1492        assert!(crashed_tmp.is_file());
1493
1494        // Reboot — ParquetStorage::new triggers cleanup_partial_writes.
1495        let storage2 = ParquetStorage::new(temp_dir.path()).unwrap();
1496        assert!(
1497            !crashed_tmp.exists(),
1498            "stale tmp file should have been cleaned by ParquetStorage::new"
1499        );
1500
1501        // Raw files survived; events queryable.
1502        let events = storage2.load_events_for_tenant("alice").unwrap();
1503        assert_eq!(events.len(), 3);
1504    }
1505
1506    #[test]
1507    fn test_cold_tier_archives_dropped_events_before_deletion() {
1508        // Cold-tier integration: when retention drops events AND an
1509        // archive target is configured, the dropped events end up
1510        // in the archive root before originals are removed. This is
1511        // the load-bearing property — without archive-before-delete
1512        // the cold tier would silently lose data on retention runs.
1513        use crate::infrastructure::persistence::cold_tier::LocalFsArchive;
1514
1515        let live_dir = TempDir::new().unwrap();
1516        let archive_dir = TempDir::new().unwrap();
1517
1518        // Seed alice with 50 events spread across 60 days.
1519        let storage = ParquetStorage::new(live_dir.path()).unwrap();
1520        let now = Utc::now();
1521        for i in 0..50 {
1522            let day_offset = 60 - (i * 60 / 49);
1523            let ts = now - chrono::Duration::days(i64::from(day_offset));
1524            let event = crate::domain::entities::Event::reconstruct_from_strings(
1525                uuid::Uuid::new_v4(),
1526                "test.event".to_string(),
1527                format!("e-{i}"),
1528                "alice".to_string(),
1529                serde_json::json!({"i": i}),
1530                ts,
1531                None,
1532                1,
1533            );
1534            storage.append_event(event).unwrap();
1535            if i % 5 == 4 {
1536                storage.flush().unwrap();
1537            }
1538        }
1539        storage.flush().unwrap();
1540
1541        // 30-day TTL + cold-tier archive.
1542        let mut retention = RetentionConfig::default();
1543        retention.set("alice", Some(Duration::from_hours(30 * 24)));
1544        let archive: Arc<dyn ArchiveTarget> =
1545            Arc::new(LocalFsArchive::new(archive_dir.path()).unwrap());
1546        let config = CompactionConfig {
1547            min_files_to_compact: 2,
1548            small_file_threshold: 100 * 1024 * 1024,
1549            strategy: CompactionStrategy::SizeBased,
1550            retention,
1551            archive: Some(archive),
1552            ..Default::default()
1553        };
1554        let manager = CompactionManager::new(live_dir.path(), config);
1555
1556        let result = manager.compact_tenant("alice").unwrap();
1557        assert!(result.events_compacted > 0, "some events kept");
1558        assert!(
1559            result.events_compacted < 50,
1560            "some events dropped to retention; kept {} of 50",
1561            result.events_compacted
1562        );
1563
1564        // Live storage: only events within the 30-day window.
1565        let live_after = ParquetStorage::new(live_dir.path())
1566            .unwrap()
1567            .load_events_for_tenant("alice")
1568            .unwrap();
1569        assert_eq!(live_after.len(), result.events_compacted);
1570
1571        // Archive: contains the dropped events. Walk the archive root
1572        // and load any archive.alice.* file we find.
1573        let mut archive_files = vec![];
1574        let mut stack = vec![archive_dir.path().to_path_buf()];
1575        while let Some(d) = stack.pop() {
1576            for entry in std::fs::read_dir(&d).unwrap().flatten() {
1577                let p = entry.path();
1578                if p.is_dir() {
1579                    stack.push(p);
1580                } else if p
1581                    .file_name()
1582                    .is_some_and(|n| n.to_string_lossy().starts_with("archive.alice."))
1583                {
1584                    archive_files.push(p);
1585                }
1586            }
1587        }
1588        assert!(
1589            !archive_files.is_empty(),
1590            "archive directory must contain at least one archive.alice.* file"
1591        );
1592
1593        // Sum events across archive files; total live + archived
1594        // must equal the original 50 (no events lost in the pipeline).
1595        let archive_storage = ParquetStorage::new(archive_dir.path()).unwrap();
1596        let archived = archive_storage.load_events_for_tenant("alice").unwrap();
1597        assert_eq!(
1598            live_after.len() + archived.len(),
1599            50,
1600            "live + archived must equal original event count (live={}, archived={})",
1601            live_after.len(),
1602            archived.len()
1603        );
1604    }
1605
1606    #[test]
1607    fn test_cold_tier_failure_keeps_originals_on_disk() {
1608        // Crash-safety contract: a failing archive must NOT delete
1609        // originals. We use a custom ArchiveTarget that always
1610        // returns Err to simulate an outage.
1611        let live_dir = TempDir::new().unwrap();
1612        let storage = ParquetStorage::new(live_dir.path()).unwrap();
1613        let now = Utc::now();
1614        for i in 0..20 {
1615            let ts = now - chrono::Duration::days(60 - i);
1616            let event = crate::domain::entities::Event::reconstruct_from_strings(
1617                uuid::Uuid::new_v4(),
1618                "test.event".to_string(),
1619                format!("e-{i}"),
1620                "alice".to_string(),
1621                serde_json::json!({"i": i}),
1622                ts,
1623                None,
1624                1,
1625            );
1626            storage.append_event(event).unwrap();
1627            if i % 5 == 4 {
1628                storage.flush().unwrap();
1629            }
1630        }
1631        storage.flush().unwrap();
1632
1633        // Count files before. The compaction pipeline should leave
1634        // them all on disk after the failed archive.
1635        let count_files = |dir: &std::path::Path| -> usize {
1636            let mut n = 0;
1637            let mut stack = vec![dir.to_path_buf()];
1638            while let Some(d) = stack.pop() {
1639                for entry in std::fs::read_dir(&d).unwrap().flatten() {
1640                    let p = entry.path();
1641                    if p.is_dir() {
1642                        stack.push(p);
1643                    } else if p.extension().is_some_and(|e| e == "parquet") {
1644                        n += 1;
1645                    }
1646                }
1647            }
1648            n
1649        };
1650        let before = count_files(live_dir.path());
1651        assert!(before > 0);
1652
1653        #[derive(Debug)]
1654        struct FailingArchive;
1655        impl ArchiveTarget for FailingArchive {
1656            fn archive(
1657                &self,
1658                _: &str,
1659                _: DateTime<Utc>,
1660                _: DateTime<Utc>,
1661                _: &[crate::domain::entities::Event],
1662            ) -> Result<()> {
1663                Err(AllSourceError::StorageError(
1664                    "simulated archive outage".to_string(),
1665                ))
1666            }
1667        }
1668
1669        let mut retention = RetentionConfig::default();
1670        retention.set("alice", Some(Duration::from_hours(30 * 24)));
1671        let config = CompactionConfig {
1672            min_files_to_compact: 2,
1673            small_file_threshold: 100 * 1024 * 1024,
1674            strategy: CompactionStrategy::SizeBased,
1675            retention,
1676            archive: Some(Arc::new(FailingArchive) as Arc<dyn ArchiveTarget>),
1677            ..Default::default()
1678        };
1679        let manager = CompactionManager::new(live_dir.path(), config);
1680
1681        let result = manager.compact_tenant("alice");
1682        assert!(result.is_err(), "compaction must fail when archive fails");
1683
1684        // Originals still on disk — every event still queryable.
1685        let after = count_files(live_dir.path());
1686        assert_eq!(
1687            before, after,
1688            "no files should be removed after archive failure"
1689        );
1690
1691        let storage2 = ParquetStorage::new(live_dir.path()).unwrap();
1692        let loaded = storage2.load_events_for_tenant("alice").unwrap();
1693        assert_eq!(
1694            loaded.len(),
1695            20,
1696            "all 20 events still present after failed archive"
1697        );
1698    }
1699
1700    #[test]
1701    fn test_cold_tier_not_invoked_when_no_events_dropped() {
1702        // If retention drops zero events (e.g. tenant has no TTL),
1703        // the archive target must NOT be called. We assert by using
1704        // an archive that panics on call.
1705        let live_dir = TempDir::new().unwrap();
1706        let storage = ParquetStorage::new(live_dir.path()).unwrap();
1707        let now = Utc::now();
1708        for i in 0..10 {
1709            let ts = now - chrono::Duration::hours(i);
1710            let event = crate::domain::entities::Event::reconstruct_from_strings(
1711                uuid::Uuid::new_v4(),
1712                "test.event".to_string(),
1713                format!("e-{i}"),
1714                "alice".to_string(),
1715                serde_json::json!({"i": i}),
1716                ts,
1717                None,
1718                1,
1719            );
1720            storage.append_event(event).unwrap();
1721            if i % 3 == 2 {
1722                storage.flush().unwrap();
1723            }
1724        }
1725        storage.flush().unwrap();
1726
1727        #[derive(Debug)]
1728        struct PanickingArchive;
1729        impl ArchiveTarget for PanickingArchive {
1730            fn archive(
1731                &self,
1732                _: &str,
1733                _: DateTime<Utc>,
1734                _: DateTime<Utc>,
1735                _: &[crate::domain::entities::Event],
1736            ) -> Result<()> {
1737                panic!("archive must not be called when no events are dropped");
1738            }
1739        }
1740
1741        // Default retention → alice has no TTL → no events dropped.
1742        let config = CompactionConfig {
1743            min_files_to_compact: 2,
1744            small_file_threshold: 100 * 1024 * 1024,
1745            strategy: CompactionStrategy::SizeBased,
1746            archive: Some(Arc::new(PanickingArchive) as Arc<dyn ArchiveTarget>),
1747            ..Default::default()
1748        };
1749        let manager = CompactionManager::new(live_dir.path(), config);
1750        let result = manager.compact_tenant("alice").unwrap();
1751        assert_eq!(result.events_compacted, 10);
1752    }
1753
1754    #[test]
1755    fn test_discover_tenants_skips_system_and_hidden() {
1756        // Ensure the tenant scan doesn't pick up __system or hidden
1757        // directories — they're internal, not real tenants.
1758        let temp_dir = TempDir::new().unwrap();
1759        std::fs::create_dir_all(temp_dir.path().join("alice")).unwrap();
1760        std::fs::create_dir_all(temp_dir.path().join("bob")).unwrap();
1761        std::fs::create_dir_all(temp_dir.path().join("__system")).unwrap();
1762        std::fs::create_dir_all(temp_dir.path().join(".hidden")).unwrap();
1763
1764        let manager = CompactionManager::new(temp_dir.path(), CompactionConfig::default());
1765        let tenants = manager.discover_tenants().unwrap();
1766        assert_eq!(tenants, vec!["alice".to_string(), "bob".to_string()]);
1767    }
1768}