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
10pub struct CompactionManager {
20 storage_dir: PathBuf,
22
23 config: CompactionConfig,
25
26 stats: Arc<RwLock<CompactionStats>>,
28
29 last_compaction: Arc<RwLock<Option<DateTime<Utc>>>>,
31}
32
33const SNAPSHOT_PREFIX: &str = "snapshot.";
38
39#[derive(Debug, Clone)]
40pub struct CompactionConfig {
41 pub min_files_to_compact: usize,
43
44 pub target_file_size: usize,
46
47 pub max_file_size: usize,
49
50 pub small_file_threshold: usize,
52
53 pub compaction_interval_seconds: u64,
55
56 pub auto_compact: bool,
58
59 pub strategy: CompactionStrategy,
61
62 pub retention: RetentionConfig,
69
70 pub archive: Option<Arc<dyn ArchiveTarget>>,
78}
79
80#[derive(Debug, Clone)]
94pub struct RetentionConfig {
95 pub default_ttl: Option<Duration>,
97 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 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 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, max_file_size: 256 * 1024 * 1024, small_file_threshold: 10 * 1024 * 1024, compaction_interval_seconds: 3600, auto_compact: true,
145 strategy: CompactionStrategy::SizeBased,
146 retention: RetentionConfig::default(),
147 archive: None,
148 }
149 }
150}
151
152impl CompactionConfig {
153 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 #[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 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 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 SizeBased,
258 TimeBased,
260 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#[derive(Debug, Clone)]
277struct FileInfo {
278 path: PathBuf,
279 size: u64,
280 created: DateTime<Utc>,
281}
282
283impl CompactionManager {
284 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 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 files.sort_by_key(|f| f.created);
343
344 Ok(files)
345 }
346
347 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 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 if small_files.len() >= self.config.min_files_to_compact {
366 small_files
367 } else {
368 Vec::new()
369 }
370 }
371
372 fn select_old_files(&self, files: &[FileInfo]) -> Vec<FileInfo> {
374 let now = Utc::now();
375 let age_threshold = chrono::Duration::hours(24); 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 #[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, 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 #[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 pub fn compact_tenant(&self, tenant_id: &str) -> Result<CompactionResult> {
492 let start_time = std::time::Instant::now();
493
494 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 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 let mut events = Vec::new();
546 for fi in &candidates {
547 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 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 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 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 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 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 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 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 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 pub fn stats(&self) -> CompactionStats {
756 (*self.stats.read()).clone()
757 }
758
759 pub fn config(&self) -> &CompactionConfig {
761 &self.config
762 }
763
764 #[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#[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
782pub(super) fn format_iso_basic(t: DateTime<Utc>) -> String {
788 t.format("%Y-%m-%dT%H%M%SZ").to_string()
789}
790
791pub struct CompactionTask {
793 manager: Arc<CompactionManager>,
794 interval: Duration,
795}
796
797impl CompactionTask {
798 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 #[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 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, 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, created: Utc::now(),
881 },
882 FileInfo {
883 path: PathBuf::from("small2.parquet"),
884 size: 600_000, created: Utc::now(),
886 },
887 FileInfo {
888 path: PathBuf::from("large.parquet"),
889 size: 10_000_000, created: Utc::now(),
891 },
892 ];
893
894 let selected = manager.select_files_for_compaction(&files);
895 assert_eq!(selected.len(), 2); }
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, 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); }
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); }
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); }
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); }
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 }
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 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()); }
1175
1176 fn ingest_and_flush_per_call(storage_dir: &std::path::Path, tenant: &str, count: usize) {
1181 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 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 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 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 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 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 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 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 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 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 let temp_dir = TempDir::new().unwrap();
1350
1351 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1356 let now = Utc::now();
1357 for i in 0..100 {
1358 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 if i % 10 == 9 {
1376 storage.flush().unwrap();
1377 }
1378 }
1379 storage.flush().unwrap();
1380
1381 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 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 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 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 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 let cfg = RetentionConfig::default();
1463 let ttl = cfg.ttl_for("system").unwrap();
1464 assert_eq!(ttl.as_secs(), 30 * 24 * 3600);
1465 assert!(cfg.ttl_for("acme").is_none());
1467 }
1468
1469 #[test]
1470 fn test_retention_drops_all_events_deletes_originals_without_snapshot() {
1471 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); 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 let temp_dir = TempDir::new().unwrap();
1523
1524 ingest_and_flush_per_call(temp_dir.path(), "alice", 3);
1526
1527 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 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 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 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 use crate::infrastructure::persistence::cold_tier::LocalFsArchive;
1558
1559 let live_dir = TempDir::new().unwrap();
1560 let archive_dir = TempDir::new().unwrap();
1561
1562 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 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 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 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 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 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 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 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 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 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 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}