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