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            // Every selected input must be readable before any retention,
561            // archive, snapshot or deletion step can run. Skipping a failed
562            // input would later delete it as if its contents were preserved.
563            let mut loaded = storage
564                .load_events_from_file_path(&fi.path, &event_tenant)
565                .map_err(|error| {
566                    AllSourceError::StorageError(format!(
567                        "Compaction refused unreadable candidate {}: {error}",
568                        fi.path.display()
569                    ))
570                })?;
571            events.append(&mut loaded);
572        }
573
574        if events.is_empty() {
575            tracing::warn!(
576                tenant_id = tenant_id,
577                "candidate files had no readable events; skipping snapshot"
578            );
579            return Ok(CompactionResult::default());
580        }
581
582        // Apply retention (Step 5). Drop events older than the
583        // tenant's TTL before they get rewritten into the
584        // snapshot. The originals get deleted at the end of the
585        // happy path either way, so dropped events go away with
586        // the same crash-safe guarantee as the rest of the
587        // compaction pipeline (snapshot rename completes BEFORE
588        // any original is deleted; AC #6).
589        //
590        // Cold-tier (sustainability): when an `archive` target is
591        // configured, dropped events are written to it BEFORE this
592        // function deletes any original file. A failed archive
593        // returns `Err`, originals stay on disk, and the next
594        // compaction pass retries — same crash-safety contract as
595        // the snapshot path. Without an archive, retention behaves
596        // exactly as before (delete outright).
597        let dropped_by_retention = if let Some(ttl) = self.config.retention.ttl_for(tenant_id) {
598            let cutoff = Utc::now()
599                - chrono::Duration::from_std(ttl).unwrap_or_else(|_| chrono::Duration::zero());
600            let before = events.len();
601
602            // Partition: drained = events past TTL, kept = events to keep.
603            // Using drain_filter would be simpler but it's unstable;
604            // partition + reassign keeps events: Vec<Event> with the
605            // kept slice in original order.
606            let (drained, kept): (Vec<_>, Vec<_>) = std::mem::take(&mut events)
607                .into_iter()
608                .partition(|e| e.timestamp < cutoff);
609            events = kept;
610            let dropped = before - events.len();
611
612            if dropped > 0 {
613                tracing::info!(
614                    retention_tenant = tenant_id,
615                    dropped = dropped,
616                    kept = events.len(),
617                    cutoff = %cutoff.to_rfc3339(),
618                    ttl_secs = ttl.as_secs(),
619                    "retention: dropped events older than TTL"
620                );
621
622                if let Some(archive) = self.config.archive.as_ref() {
623                    let from = drained
624                        .iter()
625                        .map(|e| e.timestamp)
626                        .min()
627                        .expect("dropped > 0 guarantees non-empty drained");
628                    let to = drained
629                        .iter()
630                        .map(|e| e.timestamp)
631                        .max()
632                        .expect("dropped > 0 guarantees non-empty drained");
633                    archive.archive(tenant_id, from, to, &drained)?;
634                    tracing::info!(
635                        retention_tenant = tenant_id,
636                        archived_to = %archive.description(),
637                        archived = drained.len(),
638                        "retention: dropped events archived to cold tier"
639                    );
640                }
641            }
642            dropped
643        } else {
644            0
645        };
646
647        // Edge case: every event aged out. Skip the snapshot
648        // write and just delete originals — the data is gone by
649        // design. Delete-without-snapshot is safe here because
650        // every event we'd have written to the snapshot was
651        // already past its TTL, and the input files live under
652        // the tenant's partition with nothing else relying on
653        // them.
654        if events.is_empty() {
655            tracing::info!(
656                tenant_id = tenant_id,
657                files_dropped = candidates.len(),
658                events_dropped = dropped_by_retention,
659                "retention: every event aged out — deleting originals without snapshot"
660            );
661            for fi in &candidates {
662                if let Err(e) = fs::remove_file(&fi.path) {
663                    tracing::error!(
664                        file = %fi.path.display(),
665                        "failed to remove fully-aged raw file: {e}"
666                    );
667                }
668            }
669            return Ok(CompactionResult {
670                files_compacted: candidates.len(),
671                bytes_before,
672                bytes_after: 0,
673                events_compacted: 0,
674                duration_ms: start_time.elapsed().as_millis() as u64,
675            });
676        }
677
678        events.sort_by_key(|e| e.timestamp);
679        let from = events.first().expect("non-empty checked above").timestamp;
680        let to = events.last().expect("non-empty checked above").timestamp;
681
682        // Filename: snapshot.<tenant>.<from>-<to> with filesystem-safe
683        // ISO-basic timestamps (no colons).
684        let file_stem = format!(
685            "snapshot.{tenant_id}.{}-{}",
686            format_iso_basic(from),
687            format_iso_basic(to)
688        );
689        let snapshot_path = storage.write_atomic_parquet(tenant_id, &file_stem, &events)?;
690        let bytes_after = fs::metadata(&snapshot_path).map_or(0, |m| m.len());
691
692        // 4. Delete originals AFTER snapshot is durably renamed.
693        // AC #6: a snapshot-write failure short-circuits via the
694        // `?` above, so originals stay on disk and events remain
695        // queryable until the next successful pass.
696        for fi in &candidates {
697            if let Err(e) = fs::remove_file(&fi.path) {
698                tracing::error!(
699                    file = %fi.path.display(),
700                    "failed to remove pre-snapshot raw file: {e}"
701                );
702            }
703        }
704
705        let duration_ms = start_time.elapsed().as_millis() as u64;
706        tracing::info!(
707            tenant_id = tenant_id,
708            files_compacted = candidates.len(),
709            events = events.len(),
710            dropped_by_retention = dropped_by_retention,
711            mib_before = bytes_before as f64 / (1024.0 * 1024.0),
712            mib_after = bytes_after as f64 / (1024.0 * 1024.0),
713            duration_ms = duration_ms,
714            "tenant compaction complete"
715        );
716
717        Ok(CompactionResult {
718            files_compacted: candidates.len(),
719            bytes_before,
720            bytes_after,
721            events_compacted: events.len(),
722            duration_ms,
723        })
724    }
725
726    /// Discover tenant ids by scanning `<storage_dir>/<X>/` for
727    /// directories. The migration tool's flat-layout files at the
728    /// root are not picked up — Step 1 #3's migration moves them
729    /// under `default/`.
730    fn discover_tenants(&self) -> Result<Vec<String>> {
731        let Ok(entries) = fs::read_dir(&self.storage_dir) else {
732            return Ok(Vec::new());
733        };
734        let mut tenants: Vec<String> = entries
735            .filter_map(std::result::Result::ok)
736            .filter_map(|entry| {
737                let ft = entry.file_type().ok()?;
738                if !ft.is_dir() {
739                    return None;
740                }
741                let name = entry.file_name().to_string_lossy().into_owned();
742                // Skip the system metadata subtree (Core's own
743                // event-sourced repos) and any hidden folders.
744                if name.starts_with('.') || name == "__system" {
745                    return None;
746                }
747                Some(name)
748            })
749            .collect();
750        tenants.sort();
751        Ok(tenants)
752    }
753
754    /// Get compaction statistics
755    pub fn stats(&self) -> CompactionStats {
756        (*self.stats.read()).clone()
757    }
758
759    /// Get configuration
760    pub fn config(&self) -> &CompactionConfig {
761        &self.config
762    }
763
764    /// Trigger manual compaction
765    #[cfg_attr(feature = "hotpath", hotpath::measure)]
766    pub fn compact_now(&self) -> Result<CompactionResult> {
767        tracing::info!("Manual compaction triggered");
768        self.compact()
769    }
770}
771
772/// Result of a compaction operation
773#[derive(Debug, Clone, Default, Serialize)]
774pub struct CompactionResult {
775    pub files_compacted: usize,
776    pub bytes_before: u64,
777    pub bytes_after: u64,
778    pub events_compacted: usize,
779    pub duration_ms: u64,
780}
781
782/// Format a UTC timestamp as a filename-safe ISO-8601 basic-form
783/// string: `2026-04-27T134567Z` — no colons, no fractional second.
784/// Used for the `<from>-<to>` portion of snapshot filenames so
785/// they're portable across filesystems. `pub(super)` so the
786/// cold-tier archive can reuse the same naming convention.
787pub(super) fn format_iso_basic(t: DateTime<Utc>) -> String {
788    t.format("%Y-%m-%dT%H%M%SZ").to_string()
789}
790
791/// Background compaction task
792pub struct CompactionTask {
793    manager: Arc<CompactionManager>,
794    interval: Duration,
795}
796
797impl CompactionTask {
798    /// Create a new background compaction task
799    pub fn new(manager: Arc<CompactionManager>, interval_seconds: u64) -> Self {
800        Self {
801            manager,
802            interval: Duration::from_secs(interval_seconds),
803        }
804    }
805
806    /// Run the compaction task in a loop
807    #[cfg_attr(feature = "hotpath", hotpath::measure)]
808    pub async fn run(self) {
809        let mut interval = tokio::time::interval(self.interval);
810
811        loop {
812            interval.tick().await;
813
814            if self.manager.should_compact() {
815                tracing::debug!("Auto-compaction check triggered");
816
817                match self.manager.compact() {
818                    Ok(result) => {
819                        if result.files_compacted > 0 {
820                            tracing::info!(
821                                "Auto-compaction succeeded: {} files, {:.2} MB saved",
822                                result.files_compacted,
823                                (result.bytes_before - result.bytes_after) as f64
824                                    / (1024.0 * 1024.0)
825                            );
826                        }
827                    }
828                    Err(e) => {
829                        tracing::error!("Auto-compaction failed: {}", e);
830                    }
831                }
832            }
833        }
834    }
835}
836
837#[cfg(test)]
838mod tests {
839    use super::*;
840    use tempfile::TempDir;
841
842    #[test]
843    fn test_compaction_manager_creation() {
844        let temp_dir = TempDir::new().unwrap();
845        let config = CompactionConfig::default();
846        let manager = CompactionManager::new(temp_dir.path(), config);
847
848        assert_eq!(manager.stats().total_compactions, 0);
849    }
850
851    #[test]
852    fn test_should_compact() {
853        let temp_dir = TempDir::new().unwrap();
854        let config = CompactionConfig {
855            auto_compact: true,
856            compaction_interval_seconds: 1,
857            ..Default::default()
858        };
859        let manager = CompactionManager::new(temp_dir.path(), config);
860
861        // Should compact on first check (never compacted)
862        assert!(manager.should_compact());
863    }
864
865    #[test]
866    fn test_file_selection_size_based() {
867        let temp_dir = TempDir::new().unwrap();
868        let config = CompactionConfig {
869            small_file_threshold: 1024 * 1024, // 1 MB
870            min_files_to_compact: 2,
871            strategy: CompactionStrategy::SizeBased,
872            ..Default::default()
873        };
874        let manager = CompactionManager::new(temp_dir.path(), config);
875
876        let files = vec![
877            FileInfo {
878                path: PathBuf::from("small1.parquet"),
879                size: 500_000, // 500 KB
880                created: Utc::now(),
881            },
882            FileInfo {
883                path: PathBuf::from("small2.parquet"),
884                size: 600_000, // 600 KB
885                created: Utc::now(),
886            },
887            FileInfo {
888                path: PathBuf::from("large.parquet"),
889                size: 10_000_000, // 10 MB
890                created: Utc::now(),
891            },
892        ];
893
894        let selected = manager.select_files_for_compaction(&files);
895        assert_eq!(selected.len(), 2); // Only the 2 small files
896    }
897
898    #[test]
899    fn test_default_compaction_config() {
900        let config = CompactionConfig::default();
901        assert_eq!(config.min_files_to_compact, 3);
902        assert_eq!(config.target_file_size, 128 * 1024 * 1024);
903        assert_eq!(config.max_file_size, 256 * 1024 * 1024);
904        assert_eq!(config.small_file_threshold, 10 * 1024 * 1024);
905        assert_eq!(config.compaction_interval_seconds, 3600);
906        assert!(config.auto_compact);
907        assert_eq!(config.strategy, CompactionStrategy::SizeBased);
908    }
909
910    #[test]
911    fn test_should_compact_disabled() {
912        let temp_dir = TempDir::new().unwrap();
913        let config = CompactionConfig {
914            auto_compact: false,
915            ..Default::default()
916        };
917        let manager = CompactionManager::new(temp_dir.path(), config);
918
919        assert!(!manager.should_compact());
920    }
921
922    #[test]
923    fn test_compact_empty_directory() {
924        let temp_dir = TempDir::new().unwrap();
925        let config = CompactionConfig::default();
926        let manager = CompactionManager::new(temp_dir.path(), config);
927
928        let result = manager.compact().unwrap();
929        assert_eq!(result.files_compacted, 0);
930        assert_eq!(result.bytes_before, 0);
931        assert_eq!(result.bytes_after, 0);
932        assert_eq!(result.events_compacted, 0);
933    }
934
935    #[test]
936    fn test_compact_now() {
937        let temp_dir = TempDir::new().unwrap();
938        let config = CompactionConfig::default();
939        let manager = CompactionManager::new(temp_dir.path(), config);
940
941        let result = manager.compact_now().unwrap();
942        assert_eq!(result.files_compacted, 0);
943    }
944
945    #[test]
946    fn test_get_config() {
947        let temp_dir = TempDir::new().unwrap();
948        let config = CompactionConfig {
949            min_files_to_compact: 5,
950            ..Default::default()
951        };
952        let manager = CompactionManager::new(temp_dir.path(), config);
953
954        assert_eq!(manager.config().min_files_to_compact, 5);
955    }
956
957    #[test]
958    fn test_get_stats() {
959        let temp_dir = TempDir::new().unwrap();
960        let config = CompactionConfig::default();
961        let manager = CompactionManager::new(temp_dir.path(), config);
962
963        let stats = manager.stats();
964        assert_eq!(stats.total_compactions, 0);
965        assert_eq!(stats.total_files_compacted, 0);
966        assert_eq!(stats.total_bytes_before, 0);
967        assert_eq!(stats.total_bytes_after, 0);
968        assert_eq!(stats.total_events_compacted, 0);
969        assert_eq!(stats.last_compaction_duration_ms, 0);
970        assert_eq!(stats.space_saved_bytes, 0);
971    }
972
973    #[test]
974    fn test_file_selection_not_enough_small_files() {
975        let temp_dir = TempDir::new().unwrap();
976        let config = CompactionConfig {
977            small_file_threshold: 1024 * 1024,
978            min_files_to_compact: 3, // Need 3 files
979            strategy: CompactionStrategy::SizeBased,
980            ..Default::default()
981        };
982        let manager = CompactionManager::new(temp_dir.path(), config);
983
984        let files = vec![
985            FileInfo {
986                path: PathBuf::from("small1.parquet"),
987                size: 500_000,
988                created: Utc::now(),
989            },
990            FileInfo {
991                path: PathBuf::from("small2.parquet"),
992                size: 600_000,
993                created: Utc::now(),
994            },
995        ];
996
997        let selected = manager.select_files_for_compaction(&files);
998        assert_eq!(selected.len(), 0); // Not enough small files
999    }
1000
1001    #[test]
1002    fn test_file_selection_time_based() {
1003        let temp_dir = TempDir::new().unwrap();
1004        let config = CompactionConfig {
1005            min_files_to_compact: 2,
1006            strategy: CompactionStrategy::TimeBased,
1007            ..Default::default()
1008        };
1009        let manager = CompactionManager::new(temp_dir.path(), config);
1010
1011        let old_time = Utc::now() - chrono::Duration::hours(48);
1012        let files = vec![
1013            FileInfo {
1014                path: PathBuf::from("old1.parquet"),
1015                size: 1_000_000,
1016                created: old_time,
1017            },
1018            FileInfo {
1019                path: PathBuf::from("old2.parquet"),
1020                size: 2_000_000,
1021                created: old_time,
1022            },
1023            FileInfo {
1024                path: PathBuf::from("new.parquet"),
1025                size: 500_000,
1026                created: Utc::now(),
1027            },
1028        ];
1029
1030        let selected = manager.select_files_for_compaction(&files);
1031        assert_eq!(selected.len(), 2); // Only the 2 old files
1032    }
1033
1034    #[test]
1035    fn test_file_selection_time_based_not_enough() {
1036        let temp_dir = TempDir::new().unwrap();
1037        let config = CompactionConfig {
1038            min_files_to_compact: 3,
1039            strategy: CompactionStrategy::TimeBased,
1040            ..Default::default()
1041        };
1042        let manager = CompactionManager::new(temp_dir.path(), config);
1043
1044        let old_time = Utc::now() - chrono::Duration::hours(48);
1045        let files = vec![
1046            FileInfo {
1047                path: PathBuf::from("old1.parquet"),
1048                size: 1_000_000,
1049                created: old_time,
1050            },
1051            FileInfo {
1052                path: PathBuf::from("new.parquet"),
1053                size: 500_000,
1054                created: Utc::now(),
1055            },
1056        ];
1057
1058        let selected = manager.select_files_for_compaction(&files);
1059        assert_eq!(selected.len(), 0); // Not enough old files
1060    }
1061
1062    #[test]
1063    fn test_file_selection_full_compaction() {
1064        let temp_dir = TempDir::new().unwrap();
1065        let config = CompactionConfig {
1066            strategy: CompactionStrategy::FullCompaction,
1067            ..Default::default()
1068        };
1069        let manager = CompactionManager::new(temp_dir.path(), config);
1070
1071        let files = vec![
1072            FileInfo {
1073                path: PathBuf::from("file1.parquet"),
1074                size: 1_000_000,
1075                created: Utc::now(),
1076            },
1077            FileInfo {
1078                path: PathBuf::from("file2.parquet"),
1079                size: 2_000_000,
1080                created: Utc::now(),
1081            },
1082        ];
1083
1084        let selected = manager.select_files_for_compaction(&files);
1085        assert_eq!(selected.len(), 2); // All files selected
1086    }
1087
1088    #[test]
1089    fn test_compaction_strategy_serde() {
1090        let strategies = vec![
1091            CompactionStrategy::SizeBased,
1092            CompactionStrategy::TimeBased,
1093            CompactionStrategy::FullCompaction,
1094        ];
1095
1096        for strategy in strategies {
1097            let json = serde_json::to_string(&strategy).unwrap();
1098            let parsed: CompactionStrategy = serde_json::from_str(&json).unwrap();
1099            assert_eq!(parsed, strategy);
1100        }
1101    }
1102
1103    #[test]
1104    fn test_compaction_stats_default() {
1105        let stats = CompactionStats::default();
1106        assert_eq!(stats.total_compactions, 0);
1107        assert_eq!(stats.total_files_compacted, 0);
1108    }
1109
1110    #[test]
1111    fn test_compaction_stats_serde() {
1112        let stats = CompactionStats {
1113            total_compactions: 5,
1114            total_files_compacted: 20,
1115            total_bytes_before: 1000000,
1116            total_bytes_after: 500000,
1117            total_events_compacted: 10000,
1118            last_compaction_duration_ms: 500,
1119            space_saved_bytes: 500000,
1120        };
1121
1122        let json = serde_json::to_string(&stats).unwrap();
1123        assert!(json.contains("\"total_compactions\":5"));
1124        assert!(json.contains("\"space_saved_bytes\":500000"));
1125    }
1126
1127    #[test]
1128    fn test_compaction_result_serde() {
1129        let result = CompactionResult {
1130            files_compacted: 3,
1131            bytes_before: 1000000,
1132            bytes_after: 500000,
1133            events_compacted: 5000,
1134            duration_ms: 250,
1135        };
1136
1137        let json = serde_json::to_string(&result).unwrap();
1138        assert!(json.contains("\"files_compacted\":3"));
1139        assert!(json.contains("\"bytes_before\":1000000"));
1140    }
1141
1142    #[test]
1143    fn test_compaction_task_creation() {
1144        let temp_dir = TempDir::new().unwrap();
1145        let config = CompactionConfig::default();
1146        let manager = Arc::new(CompactionManager::new(temp_dir.path(), config));
1147
1148        let _task = CompactionTask::new(manager.clone(), 60);
1149        // Task created successfully
1150    }
1151
1152    #[test]
1153    fn test_list_parquet_files_empty() {
1154        let temp_dir = TempDir::new().unwrap();
1155        let config = CompactionConfig::default();
1156        let manager = CompactionManager::new(temp_dir.path(), config);
1157
1158        let files = manager.list_parquet_files().unwrap();
1159        assert!(files.is_empty());
1160    }
1161
1162    #[test]
1163    fn test_list_parquet_files_with_non_parquet() {
1164        let temp_dir = TempDir::new().unwrap();
1165        let config = CompactionConfig::default();
1166        let manager = CompactionManager::new(temp_dir.path(), config);
1167
1168        // Create non-parquet files
1169        std::fs::write(temp_dir.path().join("test.txt"), "test").unwrap();
1170        std::fs::write(temp_dir.path().join("data.json"), "{}").unwrap();
1171
1172        let files = manager.list_parquet_files().unwrap();
1173        assert!(files.is_empty()); // No parquet files
1174    }
1175
1176    // -----------------------------------------------------------------
1177    // Per-tenant compaction tests (Step 4, commit #2).
1178    // -----------------------------------------------------------------
1179
1180    fn ingest_and_flush_per_call(storage_dir: &std::path::Path, tenant: &str, count: usize) {
1181        // Each ingested batch produces one parquet file when its
1182        // ParquetStorage is dropped, giving us multiple raw files
1183        // per tenant for the strategy filter to pick up.
1184        for i in 0..count {
1185            let storage = ParquetStorage::with_config(
1186                storage_dir,
1187                crate::infrastructure::persistence::ParquetStorageConfig {
1188                    batch_size: 1,
1189                    ..Default::default()
1190                },
1191            )
1192            .unwrap();
1193            let event = crate::domain::entities::Event::from_strings(
1194                "test.event".to_string(),
1195                format!("{tenant}-{i}"),
1196                tenant.to_string(),
1197                serde_json::json!({"i": i}),
1198                None,
1199            )
1200            .unwrap();
1201            storage.append_event(event).unwrap();
1202            storage.flush().unwrap();
1203        }
1204    }
1205
1206    #[test]
1207    fn test_compact_tenant_emits_one_snapshot_and_removes_originals() {
1208        let temp_dir = TempDir::new().unwrap();
1209
1210        // Seed 4 raw files for alice.
1211        ingest_and_flush_per_call(temp_dir.path(), "alice", 4);
1212
1213        let config = CompactionConfig {
1214            min_files_to_compact: 2,
1215            small_file_threshold: 100 * 1024 * 1024,
1216            strategy: CompactionStrategy::SizeBased,
1217            ..Default::default()
1218        };
1219        let manager = CompactionManager::new(temp_dir.path(), config);
1220
1221        let result = manager.compact_tenant("alice").unwrap();
1222        assert_eq!(result.files_compacted, 4);
1223        assert_eq!(result.events_compacted, 4);
1224
1225        // After: zero raw files, exactly one snapshot.* file under
1226        // alice's subtree.
1227        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1228        let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1229        assert_eq!(
1230            alice_files.len(),
1231            1,
1232            "expected exactly one snapshot file for alice"
1233        );
1234
1235        let name = alice_files[0]
1236            .file_name()
1237            .and_then(|n| n.to_str())
1238            .unwrap()
1239            .to_string();
1240        assert!(
1241            name.starts_with("snapshot.alice."),
1242            "expected snapshot prefix, got {name}"
1243        );
1244        assert!(name.ends_with(".parquet"));
1245
1246        // No tmp files left behind.
1247        let tmps: Vec<_> = std::fs::read_dir(alice_files[0].parent().unwrap())
1248            .unwrap()
1249            .filter_map(std::result::Result::ok)
1250            .filter(|e| e.path().to_string_lossy().ends_with(".tmp"))
1251            .collect();
1252        assert!(tmps.is_empty());
1253
1254        // Loaded events round-trip correctly.
1255        let loaded = storage.load_events_for_tenant("alice").unwrap();
1256        assert_eq!(loaded.len(), 4);
1257        for e in &loaded {
1258            assert_eq!(e.tenant_id_str(), "alice");
1259        }
1260    }
1261
1262    #[test]
1263    fn test_compact_tenant_skips_existing_snapshot_files() {
1264        // Seed alice with 4 raw files, run compaction → 1 snapshot.
1265        // Run compaction AGAIN: the snapshot is excluded from the
1266        // candidate set, no qualifying raw files remain, the second
1267        // pass is a no-op.
1268        let temp_dir = TempDir::new().unwrap();
1269        ingest_and_flush_per_call(temp_dir.path(), "alice", 4);
1270
1271        let config = CompactionConfig {
1272            min_files_to_compact: 2,
1273            small_file_threshold: 100 * 1024 * 1024,
1274            ..Default::default()
1275        };
1276        let manager = CompactionManager::new(temp_dir.path(), config);
1277
1278        let r1 = manager.compact_tenant("alice").unwrap();
1279        assert_eq!(r1.files_compacted, 4);
1280
1281        let r2 = manager.compact_tenant("alice").unwrap();
1282        assert_eq!(r2.files_compacted, 0, "snapshot must not be re-compacted");
1283
1284        // Still exactly one file on disk for alice.
1285        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1286        let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1287        assert_eq!(alice_files.len(), 1);
1288    }
1289
1290    #[test]
1291    fn test_compact_tenant_below_threshold_is_a_noop() {
1292        // 1 file < min_files_to_compact (default 3) → nothing happens.
1293        let temp_dir = TempDir::new().unwrap();
1294        ingest_and_flush_per_call(temp_dir.path(), "alice", 1);
1295
1296        let manager = CompactionManager::new(temp_dir.path(), CompactionConfig::default());
1297        let result = manager.compact_tenant("alice").unwrap();
1298        assert_eq!(result.files_compacted, 0);
1299
1300        // Raw file untouched.
1301        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1302        let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1303        assert_eq!(alice_files.len(), 1);
1304        let name = alice_files[0]
1305            .file_name()
1306            .unwrap()
1307            .to_string_lossy()
1308            .into_owned();
1309        assert!(
1310            !name.starts_with("snapshot."),
1311            "raw file must not be renamed"
1312        );
1313    }
1314
1315    #[test]
1316    fn test_compact_iterates_every_tenant() {
1317        // Two tenants each with enough raw files. compact() must
1318        // produce one snapshot per tenant.
1319        let temp_dir = TempDir::new().unwrap();
1320        ingest_and_flush_per_call(temp_dir.path(), "alice", 3);
1321        ingest_and_flush_per_call(temp_dir.path(), "bob", 3);
1322
1323        let config = CompactionConfig {
1324            min_files_to_compact: 2,
1325            small_file_threshold: 100 * 1024 * 1024,
1326            strategy: CompactionStrategy::SizeBased,
1327            ..Default::default()
1328        };
1329        let manager = CompactionManager::new(temp_dir.path(), config);
1330
1331        let result = manager.compact().unwrap();
1332        assert_eq!(result.files_compacted, 6);
1333        assert_eq!(result.events_compacted, 6);
1334
1335        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1336        for tenant in ["alice", "bob"] {
1337            let files = storage.list_parquet_files_for_tenant(tenant).unwrap();
1338            assert_eq!(files.len(), 1, "{tenant} should have one snapshot");
1339            let name = files[0].file_name().unwrap().to_string_lossy().into_owned();
1340            assert!(name.starts_with(&format!("snapshot.{tenant}.")));
1341        }
1342    }
1343
1344    #[test]
1345    fn test_retention_drops_events_older_than_ttl() {
1346        // Bead's integration test: ingest 100 events spanning 60
1347        // days, run compaction with 30-day TTL, only the last
1348        // 30 days remain queryable.
1349        let temp_dir = TempDir::new().unwrap();
1350
1351        // Seed alice with 100 events, timestamps spread across
1352        // 60 days. We can't backdate via ingest (timestamp = now()
1353        // in domain), so write parquet directly via flush of
1354        // back-dated events.
1355        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1356        let now = Utc::now();
1357        for i in 0..100 {
1358            // Day 0 = 60 days ago; day 99 ≈ now. Spreads evenly.
1359            let day_offset = 60 - (i * 60 / 99);
1360            let ts = now - chrono::Duration::days(i64::from(day_offset));
1361            let event = crate::domain::entities::Event::reconstruct_from_strings(
1362                uuid::Uuid::new_v4(),
1363                "test.event".to_string(),
1364                format!("e-{i}"),
1365                "alice".to_string(),
1366                serde_json::json!({"i": i}),
1367                ts,
1368                None,
1369                1,
1370            );
1371            storage.append_event(event).unwrap();
1372            // Every 10 events get a fresh storage to produce a
1373            // separate file (so compaction has multiple files
1374            // to merge).
1375            if i % 10 == 9 {
1376                storage.flush().unwrap();
1377            }
1378        }
1379        storage.flush().unwrap();
1380
1381        // 30-day TTL for alice via per-tenant override.
1382        let mut retention = RetentionConfig::default();
1383        retention.set("alice", Some(Duration::from_hours(30 * 24)));
1384        let config = CompactionConfig {
1385            min_files_to_compact: 2,
1386            small_file_threshold: 100 * 1024 * 1024,
1387            strategy: CompactionStrategy::SizeBased,
1388            retention,
1389            ..Default::default()
1390        };
1391        let manager = CompactionManager::new(temp_dir.path(), config);
1392
1393        let result = manager.compact_tenant("alice").unwrap();
1394        assert!(result.events_compacted > 0);
1395        assert!(
1396            result.events_compacted < 100,
1397            "retention should have dropped some events; kept {} of 100",
1398            result.events_compacted
1399        );
1400
1401        // Re-load to confirm the dropped events are gone.
1402        let storage2 = ParquetStorage::new(temp_dir.path()).unwrap();
1403        let loaded = storage2.load_events_for_tenant("alice").unwrap();
1404        assert_eq!(loaded.len(), result.events_compacted);
1405
1406        // Every loaded event must be within the 30-day window
1407        // (with a generous fudge for test-clock drift).
1408        let cutoff = Utc::now() - chrono::Duration::days(30);
1409        for e in &loaded {
1410            assert!(
1411                e.timestamp >= cutoff - chrono::Duration::seconds(60),
1412                "event with ts {} survived retention but is older than cutoff {}",
1413                e.timestamp.to_rfc3339(),
1414                cutoff.to_rfc3339()
1415            );
1416        }
1417    }
1418
1419    #[test]
1420    fn test_retention_keeps_forever_by_default_for_non_system_tenants() {
1421        // alice has no override → falls through to default_ttl
1422        // which is None. All events kept regardless of age.
1423        let temp_dir = TempDir::new().unwrap();
1424        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1425        let now = Utc::now();
1426        for i in 0..6 {
1427            let ts = now - chrono::Duration::days(i * 365);
1428            let event = crate::domain::entities::Event::reconstruct_from_strings(
1429                uuid::Uuid::new_v4(),
1430                "test.event".to_string(),
1431                format!("e-{i}"),
1432                "alice".to_string(),
1433                serde_json::json!({"i": i}),
1434                ts,
1435                None,
1436                1,
1437            );
1438            storage.append_event(event).unwrap();
1439            if i % 2 == 1 {
1440                storage.flush().unwrap();
1441            }
1442        }
1443        storage.flush().unwrap();
1444
1445        // Default RetentionConfig: alice has no entry, default_ttl = None.
1446        let config = CompactionConfig {
1447            min_files_to_compact: 2,
1448            small_file_threshold: 100 * 1024 * 1024,
1449            strategy: CompactionStrategy::SizeBased,
1450            ..Default::default()
1451        };
1452        let manager = CompactionManager::new(temp_dir.path(), config);
1453        let result = manager.compact_tenant("alice").unwrap();
1454        assert_eq!(result.events_compacted, 6, "no events should be dropped");
1455    }
1456
1457    #[test]
1458    fn test_retention_system_tenant_default_is_30_days() {
1459        // The system tenant's 30-day TTL is the bead's headline
1460        // requirement. Default config without any overrides
1461        // should already enforce it.
1462        let cfg = RetentionConfig::default();
1463        let ttl = cfg.ttl_for("system").unwrap();
1464        assert_eq!(ttl.as_secs(), 30 * 24 * 3600);
1465        // No override for arbitrary tenants → keep forever.
1466        assert!(cfg.ttl_for("acme").is_none());
1467    }
1468
1469    #[test]
1470    fn test_retention_drops_all_events_deletes_originals_without_snapshot() {
1471        // Edge case: every event is past the TTL. We delete the
1472        // raw files and emit no snapshot — there's nothing to
1473        // write. Tenant ends with zero files on disk.
1474        let temp_dir = TempDir::new().unwrap();
1475        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1476        let very_old = Utc::now() - chrono::Duration::days(90);
1477        for i in 0..6 {
1478            let event = crate::domain::entities::Event::reconstruct_from_strings(
1479                uuid::Uuid::new_v4(),
1480                "test.event".to_string(),
1481                format!("e-{i}"),
1482                "alice".to_string(),
1483                serde_json::json!({"i": i}),
1484                very_old,
1485                None,
1486                1,
1487            );
1488            storage.append_event(event).unwrap();
1489            if i % 2 == 1 {
1490                storage.flush().unwrap();
1491            }
1492        }
1493        storage.flush().unwrap();
1494
1495        let mut retention = RetentionConfig::default();
1496        retention.set("alice", Some(Duration::from_hours(7 * 24)));
1497        let config = CompactionConfig {
1498            min_files_to_compact: 2,
1499            small_file_threshold: 100 * 1024 * 1024,
1500            strategy: CompactionStrategy::SizeBased,
1501            retention,
1502            ..Default::default()
1503        };
1504        let manager = CompactionManager::new(temp_dir.path(), config);
1505        let result = manager.compact_tenant("alice").unwrap();
1506        assert_eq!(result.events_compacted, 0);
1507        assert!(result.files_compacted >= 2); // originals were deleted
1508
1509        // Tenant subtree has zero parquet files now.
1510        let storage2 = ParquetStorage::new(temp_dir.path()).unwrap();
1511        let alice_files = storage2.list_parquet_files_for_tenant("alice").unwrap();
1512        assert!(alice_files.is_empty(), "all originals should be deleted");
1513    }
1514
1515    #[test]
1516    fn test_compaction_with_simulated_crash_leaves_data_recoverable() {
1517        // AC #7-#8: simulate "crash mid-snapshot" by manually
1518        // dropping a tmp file in the partition, then asserting
1519        // a fresh ParquetStorage::new boot cleans it up and the
1520        // raw files (which would still be present in a real
1521        // mid-rename crash) remain queryable.
1522        let temp_dir = TempDir::new().unwrap();
1523
1524        // Seed alice with raw files.
1525        ingest_and_flush_per_call(temp_dir.path(), "alice", 3);
1526
1527        // Locate alice's partition dir.
1528        let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1529        let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1530        let partition = alice_files[0].parent().unwrap().to_path_buf();
1531
1532        // Simulate the crash state: write a partial .tmp file as
1533        // if write_atomic_parquet had crashed mid-rename.
1534        let crashed_tmp = partition.join("snapshot.alice.range.parquet.tmp");
1535        std::fs::write(&crashed_tmp, b"partial parquet bytes").unwrap();
1536        assert!(crashed_tmp.is_file());
1537
1538        // Reboot — ParquetStorage::new triggers cleanup_partial_writes.
1539        let storage2 = ParquetStorage::new(temp_dir.path()).unwrap();
1540        assert!(
1541            !crashed_tmp.exists(),
1542            "stale tmp file should have been cleaned by ParquetStorage::new"
1543        );
1544
1545        // Raw files survived; events queryable.
1546        let events = storage2.load_events_for_tenant("alice").unwrap();
1547        assert_eq!(events.len(), 3);
1548    }
1549
1550    #[test]
1551    fn test_cold_tier_archives_dropped_events_before_deletion() {
1552        // Cold-tier integration: when retention drops events AND an
1553        // archive target is configured, the dropped events end up
1554        // in the archive root before originals are removed. This is
1555        // the load-bearing property — without archive-before-delete
1556        // the cold tier would silently lose data on retention runs.
1557        use crate::infrastructure::persistence::cold_tier::LocalFsArchive;
1558
1559        let live_dir = TempDir::new().unwrap();
1560        let archive_dir = TempDir::new().unwrap();
1561
1562        // Seed alice with 50 events spread across 60 days.
1563        let storage = ParquetStorage::new(live_dir.path()).unwrap();
1564        let now = Utc::now();
1565        for i in 0..50 {
1566            let day_offset = 60 - (i * 60 / 49);
1567            let ts = now - chrono::Duration::days(i64::from(day_offset));
1568            let event = crate::domain::entities::Event::reconstruct_from_strings(
1569                uuid::Uuid::new_v4(),
1570                "test.event".to_string(),
1571                format!("e-{i}"),
1572                "alice".to_string(),
1573                serde_json::json!({"i": i}),
1574                ts,
1575                None,
1576                1,
1577            );
1578            storage.append_event(event).unwrap();
1579            if i % 5 == 4 {
1580                storage.flush().unwrap();
1581            }
1582        }
1583        storage.flush().unwrap();
1584
1585        // 30-day TTL + cold-tier archive.
1586        let mut retention = RetentionConfig::default();
1587        retention.set("alice", Some(Duration::from_hours(30 * 24)));
1588        let archive: Arc<dyn ArchiveTarget> =
1589            Arc::new(LocalFsArchive::new(archive_dir.path()).unwrap());
1590        let config = CompactionConfig {
1591            min_files_to_compact: 2,
1592            small_file_threshold: 100 * 1024 * 1024,
1593            strategy: CompactionStrategy::SizeBased,
1594            retention,
1595            archive: Some(archive),
1596            ..Default::default()
1597        };
1598        let manager = CompactionManager::new(live_dir.path(), config);
1599
1600        let result = manager.compact_tenant("alice").unwrap();
1601        assert!(result.events_compacted > 0, "some events kept");
1602        assert!(
1603            result.events_compacted < 50,
1604            "some events dropped to retention; kept {} of 50",
1605            result.events_compacted
1606        );
1607
1608        // Live storage: only events within the 30-day window.
1609        let live_after = ParquetStorage::new(live_dir.path())
1610            .unwrap()
1611            .load_events_for_tenant("alice")
1612            .unwrap();
1613        assert_eq!(live_after.len(), result.events_compacted);
1614
1615        // Archive: contains the dropped events. Walk the archive root
1616        // and load any archive.alice.* file we find.
1617        let mut archive_files = vec![];
1618        let mut stack = vec![archive_dir.path().to_path_buf()];
1619        while let Some(d) = stack.pop() {
1620            for entry in std::fs::read_dir(&d).unwrap().flatten() {
1621                let p = entry.path();
1622                if p.is_dir() {
1623                    stack.push(p);
1624                } else if p
1625                    .file_name()
1626                    .is_some_and(|n| n.to_string_lossy().starts_with("archive.alice."))
1627                {
1628                    archive_files.push(p);
1629                }
1630            }
1631        }
1632        assert!(
1633            !archive_files.is_empty(),
1634            "archive directory must contain at least one archive.alice.* file"
1635        );
1636
1637        // Sum events across archive files; total live + archived
1638        // must equal the original 50 (no events lost in the pipeline).
1639        let archive_storage = ParquetStorage::new(archive_dir.path()).unwrap();
1640        let archived = archive_storage.load_events_for_tenant("alice").unwrap();
1641        assert_eq!(
1642            live_after.len() + archived.len(),
1643            50,
1644            "live + archived must equal original event count (live={}, archived={})",
1645            live_after.len(),
1646            archived.len()
1647        );
1648    }
1649
1650    #[test]
1651    fn test_cold_tier_failure_keeps_originals_on_disk() {
1652        // Crash-safety contract: a failing archive must NOT delete
1653        // originals. We use a custom ArchiveTarget that always
1654        // returns Err to simulate an outage.
1655        let live_dir = TempDir::new().unwrap();
1656        let storage = ParquetStorage::new(live_dir.path()).unwrap();
1657        let now = Utc::now();
1658        for i in 0..20 {
1659            let ts = now - chrono::Duration::days(60 - i);
1660            let event = crate::domain::entities::Event::reconstruct_from_strings(
1661                uuid::Uuid::new_v4(),
1662                "test.event".to_string(),
1663                format!("e-{i}"),
1664                "alice".to_string(),
1665                serde_json::json!({"i": i}),
1666                ts,
1667                None,
1668                1,
1669            );
1670            storage.append_event(event).unwrap();
1671            if i % 5 == 4 {
1672                storage.flush().unwrap();
1673            }
1674        }
1675        storage.flush().unwrap();
1676
1677        // Count files before. The compaction pipeline should leave
1678        // them all on disk after the failed archive.
1679        let count_files = |dir: &std::path::Path| -> usize {
1680            let mut n = 0;
1681            let mut stack = vec![dir.to_path_buf()];
1682            while let Some(d) = stack.pop() {
1683                for entry in std::fs::read_dir(&d).unwrap().flatten() {
1684                    let p = entry.path();
1685                    if p.is_dir() {
1686                        stack.push(p);
1687                    } else if p.extension().is_some_and(|e| e == "parquet") {
1688                        n += 1;
1689                    }
1690                }
1691            }
1692            n
1693        };
1694        let before = count_files(live_dir.path());
1695        assert!(before > 0);
1696
1697        #[derive(Debug)]
1698        struct FailingArchive;
1699        impl ArchiveTarget for FailingArchive {
1700            fn archive(
1701                &self,
1702                _: &str,
1703                _: DateTime<Utc>,
1704                _: DateTime<Utc>,
1705                _: &[crate::domain::entities::Event],
1706            ) -> Result<()> {
1707                Err(AllSourceError::StorageError(
1708                    "simulated archive outage".to_string(),
1709                ))
1710            }
1711        }
1712
1713        let mut retention = RetentionConfig::default();
1714        retention.set("alice", Some(Duration::from_hours(30 * 24)));
1715        let config = CompactionConfig {
1716            min_files_to_compact: 2,
1717            small_file_threshold: 100 * 1024 * 1024,
1718            strategy: CompactionStrategy::SizeBased,
1719            retention,
1720            archive: Some(Arc::new(FailingArchive) as Arc<dyn ArchiveTarget>),
1721            ..Default::default()
1722        };
1723        let manager = CompactionManager::new(live_dir.path(), config);
1724
1725        let result = manager.compact_tenant("alice");
1726        assert!(result.is_err(), "compaction must fail when archive fails");
1727
1728        // Originals still on disk — every event still queryable.
1729        let after = count_files(live_dir.path());
1730        assert_eq!(
1731            before, after,
1732            "no files should be removed after archive failure"
1733        );
1734
1735        let storage2 = ParquetStorage::new(live_dir.path()).unwrap();
1736        let loaded = storage2.load_events_for_tenant("alice").unwrap();
1737        assert_eq!(
1738            loaded.len(),
1739            20,
1740            "all 20 events still present after failed archive"
1741        );
1742    }
1743
1744    #[test]
1745    fn test_cold_tier_not_invoked_when_no_events_dropped() {
1746        // If retention drops zero events (e.g. tenant has no TTL),
1747        // the archive target must NOT be called. We assert by using
1748        // an archive that panics on call.
1749        let live_dir = TempDir::new().unwrap();
1750        let storage = ParquetStorage::new(live_dir.path()).unwrap();
1751        let now = Utc::now();
1752        for i in 0..10 {
1753            let ts = now - chrono::Duration::hours(i);
1754            let event = crate::domain::entities::Event::reconstruct_from_strings(
1755                uuid::Uuid::new_v4(),
1756                "test.event".to_string(),
1757                format!("e-{i}"),
1758                "alice".to_string(),
1759                serde_json::json!({"i": i}),
1760                ts,
1761                None,
1762                1,
1763            );
1764            storage.append_event(event).unwrap();
1765            if i % 3 == 2 {
1766                storage.flush().unwrap();
1767            }
1768        }
1769        storage.flush().unwrap();
1770
1771        #[derive(Debug)]
1772        struct PanickingArchive;
1773        impl ArchiveTarget for PanickingArchive {
1774            fn archive(
1775                &self,
1776                _: &str,
1777                _: DateTime<Utc>,
1778                _: DateTime<Utc>,
1779                _: &[crate::domain::entities::Event],
1780            ) -> Result<()> {
1781                panic!("archive must not be called when no events are dropped");
1782            }
1783        }
1784
1785        // Default retention → alice has no TTL → no events dropped.
1786        let config = CompactionConfig {
1787            min_files_to_compact: 2,
1788            small_file_threshold: 100 * 1024 * 1024,
1789            strategy: CompactionStrategy::SizeBased,
1790            archive: Some(Arc::new(PanickingArchive) as Arc<dyn ArchiveTarget>),
1791            ..Default::default()
1792        };
1793        let manager = CompactionManager::new(live_dir.path(), config);
1794        let result = manager.compact_tenant("alice").unwrap();
1795        assert_eq!(result.events_compacted, 10);
1796    }
1797
1798    #[test]
1799    fn test_discover_tenants_skips_system_and_hidden() {
1800        // Ensure the tenant scan doesn't pick up __system or hidden
1801        // directories — they're internal, not real tenants.
1802        let temp_dir = TempDir::new().unwrap();
1803        std::fs::create_dir_all(temp_dir.path().join("alice")).unwrap();
1804        std::fs::create_dir_all(temp_dir.path().join("bob")).unwrap();
1805        std::fs::create_dir_all(temp_dir.path().join("__system")).unwrap();
1806        std::fs::create_dir_all(temp_dir.path().join(".hidden")).unwrap();
1807
1808        let manager = CompactionManager::new(temp_dir.path(), CompactionConfig::default());
1809        let tenants = manager.discover_tenants().unwrap();
1810        assert_eq!(tenants, vec!["alice".to_string(), "bob".to_string()]);
1811    }
1812}