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 Self::from_env_vars(
163 std::env::var("ALLSOURCE_SNAPSHOT_INTERVAL_SECONDS").ok(),
164 std::env::var("ALLSOURCE_RETENTION_SYSTEM_DAYS").ok(),
165 )
166 }
167
168 pub fn from_env_vars(
171 interval_var: Option<String>,
172 system_retention_days_var: Option<String>,
173 ) -> Self {
174 let mut config = Self::default();
175 if let Some(s) = interval_var.filter(|s| !s.is_empty()) {
176 match s.parse::<u64>() {
177 Ok(v) => config.compaction_interval_seconds = v,
178 Err(e) => {
179 tracing::warn!(
180 "ALLSOURCE_SNAPSHOT_INTERVAL_SECONDS={s:?} could not be parsed as \
181 u64: {e}; defaulting to {}s",
182 config.compaction_interval_seconds
183 );
184 }
185 }
186 }
187 if let Some(s) = system_retention_days_var.filter(|s| !s.is_empty()) {
188 match s.parse::<u64>() {
189 Ok(days) => {
190 config
191 .retention
192 .set("system", Some(Duration::from_secs(days * 24 * 3600)));
193 }
194 Err(e) => {
195 tracing::warn!(
196 "ALLSOURCE_RETENTION_SYSTEM_DAYS={s:?} could not be parsed as u64: \
197 {e}; defaulting to 30 days for tenant=system"
198 );
199 }
200 }
201 }
202 config
203 }
204
205 pub fn from_env_var(interval_var: Option<String>) -> Self {
208 Self::from_env_vars(interval_var, None)
209 }
210}
211
212#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq)]
213#[serde(rename_all = "lowercase")]
214pub enum CompactionStrategy {
215 SizeBased,
217 TimeBased,
219 FullCompaction,
221}
222
223#[derive(Debug, Clone, Default, Serialize)]
224pub struct CompactionStats {
225 pub total_compactions: u64,
226 pub total_files_compacted: u64,
227 pub total_bytes_before: u64,
228 pub total_bytes_after: u64,
229 pub total_events_compacted: u64,
230 pub last_compaction_duration_ms: u64,
231 pub space_saved_bytes: u64,
232}
233
234#[derive(Debug, Clone)]
236struct FileInfo {
237 path: PathBuf,
238 size: u64,
239 created: DateTime<Utc>,
240}
241
242impl CompactionManager {
243 pub fn new(storage_dir: impl Into<PathBuf>, config: CompactionConfig) -> Self {
245 let storage_dir = storage_dir.into();
246
247 tracing::info!(
248 "✅ Compaction manager initialized at: {}",
249 storage_dir.display()
250 );
251
252 Self {
253 storage_dir,
254 config,
255 stats: Arc::new(RwLock::new(CompactionStats::default())),
256 last_compaction: Arc::new(RwLock::new(None)),
257 }
258 }
259
260 fn list_parquet_files(&self) -> Result<Vec<FileInfo>> {
262 let entries = fs::read_dir(&self.storage_dir).map_err(|e| {
263 AllSourceError::StorageError(format!("Failed to read storage directory: {e}"))
264 })?;
265
266 let mut files = Vec::new();
267
268 for entry in entries {
269 let entry = entry.map_err(|e| {
270 AllSourceError::StorageError(format!("Failed to read directory entry: {e}"))
271 })?;
272
273 let path = entry.path();
274 if let Some(ext) = path.extension()
275 && ext == "parquet"
276 {
277 let metadata = entry.metadata().map_err(|e| {
278 AllSourceError::StorageError(format!("Failed to read file metadata: {e}"))
279 })?;
280
281 let size = metadata.len();
282 let created = metadata
283 .created()
284 .ok()
285 .and_then(|t| {
286 t.duration_since(std::time::UNIX_EPOCH).ok().map(|d| {
287 DateTime::from_timestamp(d.as_secs() as i64, 0).unwrap_or_else(Utc::now)
288 })
289 })
290 .unwrap_or_else(Utc::now);
291
292 files.push(FileInfo {
293 path,
294 size,
295 created,
296 });
297 }
298 }
299
300 files.sort_by_key(|f| f.created);
302
303 Ok(files)
304 }
305
306 fn select_files_for_compaction(&self, files: &[FileInfo]) -> Vec<FileInfo> {
308 match self.config.strategy {
309 CompactionStrategy::SizeBased => self.select_small_files(files),
310 CompactionStrategy::TimeBased => self.select_old_files(files),
311 CompactionStrategy::FullCompaction => files.to_vec(),
312 }
313 }
314
315 fn select_small_files(&self, files: &[FileInfo]) -> Vec<FileInfo> {
317 let small_files: Vec<FileInfo> = files
318 .iter()
319 .filter(|f| f.size < self.config.small_file_threshold as u64)
320 .cloned()
321 .collect();
322
323 if small_files.len() >= self.config.min_files_to_compact {
325 small_files
326 } else {
327 Vec::new()
328 }
329 }
330
331 fn select_old_files(&self, files: &[FileInfo]) -> Vec<FileInfo> {
333 let now = Utc::now();
334 let age_threshold = chrono::Duration::hours(24); let old_files: Vec<FileInfo> = files
337 .iter()
338 .filter(|f| now - f.created > age_threshold)
339 .cloned()
340 .collect();
341
342 if old_files.len() >= self.config.min_files_to_compact {
343 old_files
344 } else {
345 Vec::new()
346 }
347 }
348
349 #[cfg_attr(feature = "hotpath", hotpath::measure)]
351 pub fn should_compact(&self) -> bool {
352 if !self.config.auto_compact {
353 return false;
354 }
355
356 let last = self.last_compaction.read();
357 match *last {
358 None => true, Some(last_time) => {
360 let elapsed = (Utc::now() - last_time).num_seconds();
361 elapsed >= self.config.compaction_interval_seconds as i64
362 }
363 }
364 }
365
366 #[cfg_attr(feature = "hotpath", hotpath::measure)]
378 pub fn compact(&self) -> Result<CompactionResult> {
379 let start_time = std::time::Instant::now();
380 tracing::info!("🔄 Starting per-tenant compaction sweep...");
381
382 let tenants = self.discover_tenants()?;
383 if tenants.is_empty() {
384 tracing::debug!("No tenants found under {}", self.storage_dir.display());
385 return Ok(CompactionResult::default());
386 }
387
388 let mut aggregate = CompactionResult::default();
389 for tenant in &tenants {
390 match self.compact_tenant(tenant) {
391 Ok(r) => {
392 aggregate.files_compacted += r.files_compacted;
393 aggregate.bytes_before += r.bytes_before;
394 aggregate.bytes_after += r.bytes_after;
395 aggregate.events_compacted += r.events_compacted;
396 }
397 Err(e) => {
398 tracing::error!(
399 tenant_id = %tenant,
400 "compact_tenant failed: {e}"
401 );
402 }
403 }
404 }
405 aggregate.duration_ms = start_time.elapsed().as_millis() as u64;
406
407 if aggregate.files_compacted > 0 {
408 let mut stats = self.stats.write();
409 stats.total_compactions += 1;
410 stats.total_files_compacted += aggregate.files_compacted as u64;
411 stats.total_bytes_before += aggregate.bytes_before;
412 stats.total_bytes_after += aggregate.bytes_after;
413 stats.total_events_compacted += aggregate.events_compacted as u64;
414 stats.last_compaction_duration_ms = aggregate.duration_ms;
415 stats.space_saved_bytes += aggregate.bytes_before.saturating_sub(aggregate.bytes_after);
416 }
417 *self.last_compaction.write() = Some(Utc::now());
418
419 tracing::info!(
420 "✅ Compaction sweep complete: {} files → 1 snapshot per tenant, \
421 {:.2} MB → {:.2} MB, {} events, {} tenants in {}ms",
422 aggregate.files_compacted,
423 aggregate.bytes_before as f64 / (1024.0 * 1024.0),
424 aggregate.bytes_after as f64 / (1024.0 * 1024.0),
425 aggregate.events_compacted,
426 tenants.len(),
427 aggregate.duration_ms
428 );
429
430 Ok(aggregate)
431 }
432
433 pub fn compact_tenant(&self, tenant_id: &str) -> Result<CompactionResult> {
451 let start_time = std::time::Instant::now();
452
453 let storage = ParquetStorage::new(&self.storage_dir)?;
455 let all_files = storage.list_parquet_files_for_tenant(tenant_id)?;
456 let raw_files: Vec<FileInfo> = all_files
457 .into_iter()
458 .filter(|p| {
459 p.file_name()
460 .and_then(|n| n.to_str())
461 .is_none_or(|n| !n.starts_with(SNAPSHOT_PREFIX))
462 })
463 .filter_map(|p| {
464 let metadata = fs::metadata(&p).ok()?;
465 let size = metadata.len();
466 let created = metadata
467 .created()
468 .ok()
469 .and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
470 .and_then(|d| DateTime::from_timestamp(d.as_secs() as i64, 0))
471 .unwrap_or_else(Utc::now);
472 Some(FileInfo {
473 path: p,
474 size,
475 created,
476 })
477 })
478 .collect();
479
480 let candidates = self.select_files_for_compaction(&raw_files);
482 if candidates.is_empty() {
483 tracing::debug!(
484 tenant_id = tenant_id,
485 strategy = ?self.config.strategy,
486 "no files meet compaction criteria"
487 );
488 return Ok(CompactionResult::default());
489 }
490
491 let bytes_before: u64 = candidates.iter().map(|f| f.size).sum();
492 tracing::info!(
493 tenant_id = tenant_id,
494 files = candidates.len(),
495 mib = bytes_before as f64 / (1024.0 * 1024.0),
496 "compacting tenant"
497 );
498
499 let mut events = Vec::new();
505 for fi in &candidates {
506 let event_tenant = match fi.path.strip_prefix(&self.storage_dir).ok() {
509 Some(rel) => rel
510 .components()
511 .next()
512 .and_then(|c| match c {
513 std::path::Component::Normal(t) => Some(t.to_string_lossy().into_owned()),
514 _ => None,
515 })
516 .unwrap_or_else(|| "default".to_string()),
517 None => "default".to_string(),
518 };
519 match storage.load_events_from_file_path(&fi.path, &event_tenant) {
520 Ok(mut e) => events.append(&mut e),
521 Err(e) => {
522 tracing::error!(
523 file = %fi.path.display(),
524 "failed to read parquet file for compaction: {e}"
525 );
526 }
527 }
528 }
529
530 if events.is_empty() {
531 tracing::warn!(
532 tenant_id = tenant_id,
533 "candidate files had no readable events; skipping snapshot"
534 );
535 return Ok(CompactionResult::default());
536 }
537
538 let dropped_by_retention = if let Some(ttl) = self.config.retention.ttl_for(tenant_id) {
554 let cutoff = Utc::now()
555 - chrono::Duration::from_std(ttl).unwrap_or_else(|_| chrono::Duration::zero());
556 let before = events.len();
557
558 let (drained, kept): (Vec<_>, Vec<_>) = std::mem::take(&mut events)
563 .into_iter()
564 .partition(|e| e.timestamp < cutoff);
565 events = kept;
566 let dropped = before - events.len();
567
568 if dropped > 0 {
569 tracing::info!(
570 retention_tenant = tenant_id,
571 dropped = dropped,
572 kept = events.len(),
573 cutoff = %cutoff.to_rfc3339(),
574 ttl_secs = ttl.as_secs(),
575 "retention: dropped events older than TTL"
576 );
577
578 if let Some(archive) = self.config.archive.as_ref() {
579 let from = drained
580 .iter()
581 .map(|e| e.timestamp)
582 .min()
583 .expect("dropped > 0 guarantees non-empty drained");
584 let to = drained
585 .iter()
586 .map(|e| e.timestamp)
587 .max()
588 .expect("dropped > 0 guarantees non-empty drained");
589 archive.archive(tenant_id, from, to, &drained)?;
590 tracing::info!(
591 retention_tenant = tenant_id,
592 archived_to = %archive.description(),
593 archived = drained.len(),
594 "retention: dropped events archived to cold tier"
595 );
596 }
597 }
598 dropped
599 } else {
600 0
601 };
602
603 if events.is_empty() {
611 tracing::info!(
612 tenant_id = tenant_id,
613 files_dropped = candidates.len(),
614 events_dropped = dropped_by_retention,
615 "retention: every event aged out — deleting originals without snapshot"
616 );
617 for fi in &candidates {
618 if let Err(e) = fs::remove_file(&fi.path) {
619 tracing::error!(
620 file = %fi.path.display(),
621 "failed to remove fully-aged raw file: {e}"
622 );
623 }
624 }
625 return Ok(CompactionResult {
626 files_compacted: candidates.len(),
627 bytes_before,
628 bytes_after: 0,
629 events_compacted: 0,
630 duration_ms: start_time.elapsed().as_millis() as u64,
631 });
632 }
633
634 events.sort_by_key(|e| e.timestamp);
635 let from = events.first().expect("non-empty checked above").timestamp;
636 let to = events.last().expect("non-empty checked above").timestamp;
637
638 let file_stem = format!(
641 "snapshot.{tenant_id}.{}-{}",
642 format_iso_basic(from),
643 format_iso_basic(to)
644 );
645 let snapshot_path = storage.write_atomic_parquet(tenant_id, &file_stem, &events)?;
646 let bytes_after = fs::metadata(&snapshot_path).map_or(0, |m| m.len());
647
648 for fi in &candidates {
653 if let Err(e) = fs::remove_file(&fi.path) {
654 tracing::error!(
655 file = %fi.path.display(),
656 "failed to remove pre-snapshot raw file: {e}"
657 );
658 }
659 }
660
661 let duration_ms = start_time.elapsed().as_millis() as u64;
662 tracing::info!(
663 tenant_id = tenant_id,
664 files_compacted = candidates.len(),
665 events = events.len(),
666 dropped_by_retention = dropped_by_retention,
667 mib_before = bytes_before as f64 / (1024.0 * 1024.0),
668 mib_after = bytes_after as f64 / (1024.0 * 1024.0),
669 duration_ms = duration_ms,
670 "tenant compaction complete"
671 );
672
673 Ok(CompactionResult {
674 files_compacted: candidates.len(),
675 bytes_before,
676 bytes_after,
677 events_compacted: events.len(),
678 duration_ms,
679 })
680 }
681
682 fn discover_tenants(&self) -> Result<Vec<String>> {
687 let Ok(entries) = fs::read_dir(&self.storage_dir) else {
688 return Ok(Vec::new());
689 };
690 let mut tenants: Vec<String> = entries
691 .filter_map(std::result::Result::ok)
692 .filter_map(|entry| {
693 let ft = entry.file_type().ok()?;
694 if !ft.is_dir() {
695 return None;
696 }
697 let name = entry.file_name().to_string_lossy().into_owned();
698 if name.starts_with('.') || name == "__system" {
701 return None;
702 }
703 Some(name)
704 })
705 .collect();
706 tenants.sort();
707 Ok(tenants)
708 }
709
710 pub fn stats(&self) -> CompactionStats {
712 (*self.stats.read()).clone()
713 }
714
715 pub fn config(&self) -> &CompactionConfig {
717 &self.config
718 }
719
720 #[cfg_attr(feature = "hotpath", hotpath::measure)]
722 pub fn compact_now(&self) -> Result<CompactionResult> {
723 tracing::info!("Manual compaction triggered");
724 self.compact()
725 }
726}
727
728#[derive(Debug, Clone, Default, Serialize)]
730pub struct CompactionResult {
731 pub files_compacted: usize,
732 pub bytes_before: u64,
733 pub bytes_after: u64,
734 pub events_compacted: usize,
735 pub duration_ms: u64,
736}
737
738pub(super) fn format_iso_basic(t: DateTime<Utc>) -> String {
744 t.format("%Y-%m-%dT%H%M%SZ").to_string()
745}
746
747pub struct CompactionTask {
749 manager: Arc<CompactionManager>,
750 interval: Duration,
751}
752
753impl CompactionTask {
754 pub fn new(manager: Arc<CompactionManager>, interval_seconds: u64) -> Self {
756 Self {
757 manager,
758 interval: Duration::from_secs(interval_seconds),
759 }
760 }
761
762 #[cfg_attr(feature = "hotpath", hotpath::measure)]
764 pub async fn run(self) {
765 let mut interval = tokio::time::interval(self.interval);
766
767 loop {
768 interval.tick().await;
769
770 if self.manager.should_compact() {
771 tracing::debug!("Auto-compaction check triggered");
772
773 match self.manager.compact() {
774 Ok(result) => {
775 if result.files_compacted > 0 {
776 tracing::info!(
777 "Auto-compaction succeeded: {} files, {:.2} MB saved",
778 result.files_compacted,
779 (result.bytes_before - result.bytes_after) as f64
780 / (1024.0 * 1024.0)
781 );
782 }
783 }
784 Err(e) => {
785 tracing::error!("Auto-compaction failed: {}", e);
786 }
787 }
788 }
789 }
790 }
791}
792
793#[cfg(test)]
794mod tests {
795 use super::*;
796 use tempfile::TempDir;
797
798 #[test]
799 fn test_compaction_manager_creation() {
800 let temp_dir = TempDir::new().unwrap();
801 let config = CompactionConfig::default();
802 let manager = CompactionManager::new(temp_dir.path(), config);
803
804 assert_eq!(manager.stats().total_compactions, 0);
805 }
806
807 #[test]
808 fn test_should_compact() {
809 let temp_dir = TempDir::new().unwrap();
810 let config = CompactionConfig {
811 auto_compact: true,
812 compaction_interval_seconds: 1,
813 ..Default::default()
814 };
815 let manager = CompactionManager::new(temp_dir.path(), config);
816
817 assert!(manager.should_compact());
819 }
820
821 #[test]
822 fn test_file_selection_size_based() {
823 let temp_dir = TempDir::new().unwrap();
824 let config = CompactionConfig {
825 small_file_threshold: 1024 * 1024, min_files_to_compact: 2,
827 strategy: CompactionStrategy::SizeBased,
828 ..Default::default()
829 };
830 let manager = CompactionManager::new(temp_dir.path(), config);
831
832 let files = vec![
833 FileInfo {
834 path: PathBuf::from("small1.parquet"),
835 size: 500_000, created: Utc::now(),
837 },
838 FileInfo {
839 path: PathBuf::from("small2.parquet"),
840 size: 600_000, created: Utc::now(),
842 },
843 FileInfo {
844 path: PathBuf::from("large.parquet"),
845 size: 10_000_000, created: Utc::now(),
847 },
848 ];
849
850 let selected = manager.select_files_for_compaction(&files);
851 assert_eq!(selected.len(), 2); }
853
854 #[test]
855 fn test_default_compaction_config() {
856 let config = CompactionConfig::default();
857 assert_eq!(config.min_files_to_compact, 3);
858 assert_eq!(config.target_file_size, 128 * 1024 * 1024);
859 assert_eq!(config.max_file_size, 256 * 1024 * 1024);
860 assert_eq!(config.small_file_threshold, 10 * 1024 * 1024);
861 assert_eq!(config.compaction_interval_seconds, 3600);
862 assert!(config.auto_compact);
863 assert_eq!(config.strategy, CompactionStrategy::SizeBased);
864 }
865
866 #[test]
867 fn test_should_compact_disabled() {
868 let temp_dir = TempDir::new().unwrap();
869 let config = CompactionConfig {
870 auto_compact: false,
871 ..Default::default()
872 };
873 let manager = CompactionManager::new(temp_dir.path(), config);
874
875 assert!(!manager.should_compact());
876 }
877
878 #[test]
879 fn test_compact_empty_directory() {
880 let temp_dir = TempDir::new().unwrap();
881 let config = CompactionConfig::default();
882 let manager = CompactionManager::new(temp_dir.path(), config);
883
884 let result = manager.compact().unwrap();
885 assert_eq!(result.files_compacted, 0);
886 assert_eq!(result.bytes_before, 0);
887 assert_eq!(result.bytes_after, 0);
888 assert_eq!(result.events_compacted, 0);
889 }
890
891 #[test]
892 fn test_compact_now() {
893 let temp_dir = TempDir::new().unwrap();
894 let config = CompactionConfig::default();
895 let manager = CompactionManager::new(temp_dir.path(), config);
896
897 let result = manager.compact_now().unwrap();
898 assert_eq!(result.files_compacted, 0);
899 }
900
901 #[test]
902 fn test_get_config() {
903 let temp_dir = TempDir::new().unwrap();
904 let config = CompactionConfig {
905 min_files_to_compact: 5,
906 ..Default::default()
907 };
908 let manager = CompactionManager::new(temp_dir.path(), config);
909
910 assert_eq!(manager.config().min_files_to_compact, 5);
911 }
912
913 #[test]
914 fn test_get_stats() {
915 let temp_dir = TempDir::new().unwrap();
916 let config = CompactionConfig::default();
917 let manager = CompactionManager::new(temp_dir.path(), config);
918
919 let stats = manager.stats();
920 assert_eq!(stats.total_compactions, 0);
921 assert_eq!(stats.total_files_compacted, 0);
922 assert_eq!(stats.total_bytes_before, 0);
923 assert_eq!(stats.total_bytes_after, 0);
924 assert_eq!(stats.total_events_compacted, 0);
925 assert_eq!(stats.last_compaction_duration_ms, 0);
926 assert_eq!(stats.space_saved_bytes, 0);
927 }
928
929 #[test]
930 fn test_file_selection_not_enough_small_files() {
931 let temp_dir = TempDir::new().unwrap();
932 let config = CompactionConfig {
933 small_file_threshold: 1024 * 1024,
934 min_files_to_compact: 3, strategy: CompactionStrategy::SizeBased,
936 ..Default::default()
937 };
938 let manager = CompactionManager::new(temp_dir.path(), config);
939
940 let files = vec![
941 FileInfo {
942 path: PathBuf::from("small1.parquet"),
943 size: 500_000,
944 created: Utc::now(),
945 },
946 FileInfo {
947 path: PathBuf::from("small2.parquet"),
948 size: 600_000,
949 created: Utc::now(),
950 },
951 ];
952
953 let selected = manager.select_files_for_compaction(&files);
954 assert_eq!(selected.len(), 0); }
956
957 #[test]
958 fn test_file_selection_time_based() {
959 let temp_dir = TempDir::new().unwrap();
960 let config = CompactionConfig {
961 min_files_to_compact: 2,
962 strategy: CompactionStrategy::TimeBased,
963 ..Default::default()
964 };
965 let manager = CompactionManager::new(temp_dir.path(), config);
966
967 let old_time = Utc::now() - chrono::Duration::hours(48);
968 let files = vec![
969 FileInfo {
970 path: PathBuf::from("old1.parquet"),
971 size: 1_000_000,
972 created: old_time,
973 },
974 FileInfo {
975 path: PathBuf::from("old2.parquet"),
976 size: 2_000_000,
977 created: old_time,
978 },
979 FileInfo {
980 path: PathBuf::from("new.parquet"),
981 size: 500_000,
982 created: Utc::now(),
983 },
984 ];
985
986 let selected = manager.select_files_for_compaction(&files);
987 assert_eq!(selected.len(), 2); }
989
990 #[test]
991 fn test_file_selection_time_based_not_enough() {
992 let temp_dir = TempDir::new().unwrap();
993 let config = CompactionConfig {
994 min_files_to_compact: 3,
995 strategy: CompactionStrategy::TimeBased,
996 ..Default::default()
997 };
998 let manager = CompactionManager::new(temp_dir.path(), config);
999
1000 let old_time = Utc::now() - chrono::Duration::hours(48);
1001 let files = vec![
1002 FileInfo {
1003 path: PathBuf::from("old1.parquet"),
1004 size: 1_000_000,
1005 created: old_time,
1006 },
1007 FileInfo {
1008 path: PathBuf::from("new.parquet"),
1009 size: 500_000,
1010 created: Utc::now(),
1011 },
1012 ];
1013
1014 let selected = manager.select_files_for_compaction(&files);
1015 assert_eq!(selected.len(), 0); }
1017
1018 #[test]
1019 fn test_file_selection_full_compaction() {
1020 let temp_dir = TempDir::new().unwrap();
1021 let config = CompactionConfig {
1022 strategy: CompactionStrategy::FullCompaction,
1023 ..Default::default()
1024 };
1025 let manager = CompactionManager::new(temp_dir.path(), config);
1026
1027 let files = vec![
1028 FileInfo {
1029 path: PathBuf::from("file1.parquet"),
1030 size: 1_000_000,
1031 created: Utc::now(),
1032 },
1033 FileInfo {
1034 path: PathBuf::from("file2.parquet"),
1035 size: 2_000_000,
1036 created: Utc::now(),
1037 },
1038 ];
1039
1040 let selected = manager.select_files_for_compaction(&files);
1041 assert_eq!(selected.len(), 2); }
1043
1044 #[test]
1045 fn test_compaction_strategy_serde() {
1046 let strategies = vec![
1047 CompactionStrategy::SizeBased,
1048 CompactionStrategy::TimeBased,
1049 CompactionStrategy::FullCompaction,
1050 ];
1051
1052 for strategy in strategies {
1053 let json = serde_json::to_string(&strategy).unwrap();
1054 let parsed: CompactionStrategy = serde_json::from_str(&json).unwrap();
1055 assert_eq!(parsed, strategy);
1056 }
1057 }
1058
1059 #[test]
1060 fn test_compaction_stats_default() {
1061 let stats = CompactionStats::default();
1062 assert_eq!(stats.total_compactions, 0);
1063 assert_eq!(stats.total_files_compacted, 0);
1064 }
1065
1066 #[test]
1067 fn test_compaction_stats_serde() {
1068 let stats = CompactionStats {
1069 total_compactions: 5,
1070 total_files_compacted: 20,
1071 total_bytes_before: 1000000,
1072 total_bytes_after: 500000,
1073 total_events_compacted: 10000,
1074 last_compaction_duration_ms: 500,
1075 space_saved_bytes: 500000,
1076 };
1077
1078 let json = serde_json::to_string(&stats).unwrap();
1079 assert!(json.contains("\"total_compactions\":5"));
1080 assert!(json.contains("\"space_saved_bytes\":500000"));
1081 }
1082
1083 #[test]
1084 fn test_compaction_result_serde() {
1085 let result = CompactionResult {
1086 files_compacted: 3,
1087 bytes_before: 1000000,
1088 bytes_after: 500000,
1089 events_compacted: 5000,
1090 duration_ms: 250,
1091 };
1092
1093 let json = serde_json::to_string(&result).unwrap();
1094 assert!(json.contains("\"files_compacted\":3"));
1095 assert!(json.contains("\"bytes_before\":1000000"));
1096 }
1097
1098 #[test]
1099 fn test_compaction_task_creation() {
1100 let temp_dir = TempDir::new().unwrap();
1101 let config = CompactionConfig::default();
1102 let manager = Arc::new(CompactionManager::new(temp_dir.path(), config));
1103
1104 let _task = CompactionTask::new(manager.clone(), 60);
1105 }
1107
1108 #[test]
1109 fn test_list_parquet_files_empty() {
1110 let temp_dir = TempDir::new().unwrap();
1111 let config = CompactionConfig::default();
1112 let manager = CompactionManager::new(temp_dir.path(), config);
1113
1114 let files = manager.list_parquet_files().unwrap();
1115 assert!(files.is_empty());
1116 }
1117
1118 #[test]
1119 fn test_list_parquet_files_with_non_parquet() {
1120 let temp_dir = TempDir::new().unwrap();
1121 let config = CompactionConfig::default();
1122 let manager = CompactionManager::new(temp_dir.path(), config);
1123
1124 std::fs::write(temp_dir.path().join("test.txt"), "test").unwrap();
1126 std::fs::write(temp_dir.path().join("data.json"), "{}").unwrap();
1127
1128 let files = manager.list_parquet_files().unwrap();
1129 assert!(files.is_empty()); }
1131
1132 fn ingest_and_flush_per_call(storage_dir: &std::path::Path, tenant: &str, count: usize) {
1137 for i in 0..count {
1141 let storage = ParquetStorage::with_config(
1142 storage_dir,
1143 crate::infrastructure::persistence::ParquetStorageConfig {
1144 batch_size: 1,
1145 ..Default::default()
1146 },
1147 )
1148 .unwrap();
1149 let event = crate::domain::entities::Event::from_strings(
1150 "test.event".to_string(),
1151 format!("{tenant}-{i}"),
1152 tenant.to_string(),
1153 serde_json::json!({"i": i}),
1154 None,
1155 )
1156 .unwrap();
1157 storage.append_event(event).unwrap();
1158 storage.flush().unwrap();
1159 }
1160 }
1161
1162 #[test]
1163 fn test_compact_tenant_emits_one_snapshot_and_removes_originals() {
1164 let temp_dir = TempDir::new().unwrap();
1165
1166 ingest_and_flush_per_call(temp_dir.path(), "alice", 4);
1168
1169 let config = CompactionConfig {
1170 min_files_to_compact: 2,
1171 small_file_threshold: 100 * 1024 * 1024,
1172 strategy: CompactionStrategy::SizeBased,
1173 ..Default::default()
1174 };
1175 let manager = CompactionManager::new(temp_dir.path(), config);
1176
1177 let result = manager.compact_tenant("alice").unwrap();
1178 assert_eq!(result.files_compacted, 4);
1179 assert_eq!(result.events_compacted, 4);
1180
1181 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1184 let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1185 assert_eq!(
1186 alice_files.len(),
1187 1,
1188 "expected exactly one snapshot file for alice"
1189 );
1190
1191 let name = alice_files[0]
1192 .file_name()
1193 .and_then(|n| n.to_str())
1194 .unwrap()
1195 .to_string();
1196 assert!(
1197 name.starts_with("snapshot.alice."),
1198 "expected snapshot prefix, got {name}"
1199 );
1200 assert!(name.ends_with(".parquet"));
1201
1202 let tmps: Vec<_> = std::fs::read_dir(alice_files[0].parent().unwrap())
1204 .unwrap()
1205 .filter_map(std::result::Result::ok)
1206 .filter(|e| e.path().to_string_lossy().ends_with(".tmp"))
1207 .collect();
1208 assert!(tmps.is_empty());
1209
1210 let loaded = storage.load_events_for_tenant("alice").unwrap();
1212 assert_eq!(loaded.len(), 4);
1213 for e in &loaded {
1214 assert_eq!(e.tenant_id_str(), "alice");
1215 }
1216 }
1217
1218 #[test]
1219 fn test_compact_tenant_skips_existing_snapshot_files() {
1220 let temp_dir = TempDir::new().unwrap();
1225 ingest_and_flush_per_call(temp_dir.path(), "alice", 4);
1226
1227 let config = CompactionConfig {
1228 min_files_to_compact: 2,
1229 small_file_threshold: 100 * 1024 * 1024,
1230 ..Default::default()
1231 };
1232 let manager = CompactionManager::new(temp_dir.path(), config);
1233
1234 let r1 = manager.compact_tenant("alice").unwrap();
1235 assert_eq!(r1.files_compacted, 4);
1236
1237 let r2 = manager.compact_tenant("alice").unwrap();
1238 assert_eq!(r2.files_compacted, 0, "snapshot must not be re-compacted");
1239
1240 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1242 let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1243 assert_eq!(alice_files.len(), 1);
1244 }
1245
1246 #[test]
1247 fn test_compact_tenant_below_threshold_is_a_noop() {
1248 let temp_dir = TempDir::new().unwrap();
1250 ingest_and_flush_per_call(temp_dir.path(), "alice", 1);
1251
1252 let manager = CompactionManager::new(temp_dir.path(), CompactionConfig::default());
1253 let result = manager.compact_tenant("alice").unwrap();
1254 assert_eq!(result.files_compacted, 0);
1255
1256 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1258 let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1259 assert_eq!(alice_files.len(), 1);
1260 let name = alice_files[0]
1261 .file_name()
1262 .unwrap()
1263 .to_string_lossy()
1264 .into_owned();
1265 assert!(
1266 !name.starts_with("snapshot."),
1267 "raw file must not be renamed"
1268 );
1269 }
1270
1271 #[test]
1272 fn test_compact_iterates_every_tenant() {
1273 let temp_dir = TempDir::new().unwrap();
1276 ingest_and_flush_per_call(temp_dir.path(), "alice", 3);
1277 ingest_and_flush_per_call(temp_dir.path(), "bob", 3);
1278
1279 let config = CompactionConfig {
1280 min_files_to_compact: 2,
1281 small_file_threshold: 100 * 1024 * 1024,
1282 strategy: CompactionStrategy::SizeBased,
1283 ..Default::default()
1284 };
1285 let manager = CompactionManager::new(temp_dir.path(), config);
1286
1287 let result = manager.compact().unwrap();
1288 assert_eq!(result.files_compacted, 6);
1289 assert_eq!(result.events_compacted, 6);
1290
1291 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1292 for tenant in ["alice", "bob"] {
1293 let files = storage.list_parquet_files_for_tenant(tenant).unwrap();
1294 assert_eq!(files.len(), 1, "{tenant} should have one snapshot");
1295 let name = files[0].file_name().unwrap().to_string_lossy().into_owned();
1296 assert!(name.starts_with(&format!("snapshot.{tenant}.")));
1297 }
1298 }
1299
1300 #[test]
1301 fn test_retention_drops_events_older_than_ttl() {
1302 let temp_dir = TempDir::new().unwrap();
1306
1307 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1312 let now = Utc::now();
1313 for i in 0..100 {
1314 let day_offset = 60 - (i * 60 / 99);
1316 let ts = now - chrono::Duration::days(i64::from(day_offset));
1317 let event = crate::domain::entities::Event::reconstruct_from_strings(
1318 uuid::Uuid::new_v4(),
1319 "test.event".to_string(),
1320 format!("e-{i}"),
1321 "alice".to_string(),
1322 serde_json::json!({"i": i}),
1323 ts,
1324 None,
1325 1,
1326 );
1327 storage.append_event(event).unwrap();
1328 if i % 10 == 9 {
1332 storage.flush().unwrap();
1333 }
1334 }
1335 storage.flush().unwrap();
1336
1337 let mut retention = RetentionConfig::default();
1339 retention.set("alice", Some(Duration::from_hours(30 * 24)));
1340 let config = CompactionConfig {
1341 min_files_to_compact: 2,
1342 small_file_threshold: 100 * 1024 * 1024,
1343 strategy: CompactionStrategy::SizeBased,
1344 retention,
1345 ..Default::default()
1346 };
1347 let manager = CompactionManager::new(temp_dir.path(), config);
1348
1349 let result = manager.compact_tenant("alice").unwrap();
1350 assert!(result.events_compacted > 0);
1351 assert!(
1352 result.events_compacted < 100,
1353 "retention should have dropped some events; kept {} of 100",
1354 result.events_compacted
1355 );
1356
1357 let storage2 = ParquetStorage::new(temp_dir.path()).unwrap();
1359 let loaded = storage2.load_events_for_tenant("alice").unwrap();
1360 assert_eq!(loaded.len(), result.events_compacted);
1361
1362 let cutoff = Utc::now() - chrono::Duration::days(30);
1365 for e in &loaded {
1366 assert!(
1367 e.timestamp >= cutoff - chrono::Duration::seconds(60),
1368 "event with ts {} survived retention but is older than cutoff {}",
1369 e.timestamp.to_rfc3339(),
1370 cutoff.to_rfc3339()
1371 );
1372 }
1373 }
1374
1375 #[test]
1376 fn test_retention_keeps_forever_by_default_for_non_system_tenants() {
1377 let temp_dir = TempDir::new().unwrap();
1380 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1381 let now = Utc::now();
1382 for i in 0..6 {
1383 let ts = now - chrono::Duration::days(i * 365);
1384 let event = crate::domain::entities::Event::reconstruct_from_strings(
1385 uuid::Uuid::new_v4(),
1386 "test.event".to_string(),
1387 format!("e-{i}"),
1388 "alice".to_string(),
1389 serde_json::json!({"i": i}),
1390 ts,
1391 None,
1392 1,
1393 );
1394 storage.append_event(event).unwrap();
1395 if i % 2 == 1 {
1396 storage.flush().unwrap();
1397 }
1398 }
1399 storage.flush().unwrap();
1400
1401 let config = CompactionConfig {
1403 min_files_to_compact: 2,
1404 small_file_threshold: 100 * 1024 * 1024,
1405 strategy: CompactionStrategy::SizeBased,
1406 ..Default::default()
1407 };
1408 let manager = CompactionManager::new(temp_dir.path(), config);
1409 let result = manager.compact_tenant("alice").unwrap();
1410 assert_eq!(result.events_compacted, 6, "no events should be dropped");
1411 }
1412
1413 #[test]
1414 fn test_retention_system_tenant_default_is_30_days() {
1415 let cfg = RetentionConfig::default();
1419 let ttl = cfg.ttl_for("system").unwrap();
1420 assert_eq!(ttl.as_secs(), 30 * 24 * 3600);
1421 assert!(cfg.ttl_for("acme").is_none());
1423 }
1424
1425 #[test]
1426 fn test_retention_drops_all_events_deletes_originals_without_snapshot() {
1427 let temp_dir = TempDir::new().unwrap();
1431 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1432 let very_old = Utc::now() - chrono::Duration::days(90);
1433 for i in 0..6 {
1434 let event = crate::domain::entities::Event::reconstruct_from_strings(
1435 uuid::Uuid::new_v4(),
1436 "test.event".to_string(),
1437 format!("e-{i}"),
1438 "alice".to_string(),
1439 serde_json::json!({"i": i}),
1440 very_old,
1441 None,
1442 1,
1443 );
1444 storage.append_event(event).unwrap();
1445 if i % 2 == 1 {
1446 storage.flush().unwrap();
1447 }
1448 }
1449 storage.flush().unwrap();
1450
1451 let mut retention = RetentionConfig::default();
1452 retention.set("alice", Some(Duration::from_hours(7 * 24)));
1453 let config = CompactionConfig {
1454 min_files_to_compact: 2,
1455 small_file_threshold: 100 * 1024 * 1024,
1456 strategy: CompactionStrategy::SizeBased,
1457 retention,
1458 ..Default::default()
1459 };
1460 let manager = CompactionManager::new(temp_dir.path(), config);
1461 let result = manager.compact_tenant("alice").unwrap();
1462 assert_eq!(result.events_compacted, 0);
1463 assert!(result.files_compacted >= 2); let storage2 = ParquetStorage::new(temp_dir.path()).unwrap();
1467 let alice_files = storage2.list_parquet_files_for_tenant("alice").unwrap();
1468 assert!(alice_files.is_empty(), "all originals should be deleted");
1469 }
1470
1471 #[test]
1472 fn test_compaction_with_simulated_crash_leaves_data_recoverable() {
1473 let temp_dir = TempDir::new().unwrap();
1479
1480 ingest_and_flush_per_call(temp_dir.path(), "alice", 3);
1482
1483 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1485 let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1486 let partition = alice_files[0].parent().unwrap().to_path_buf();
1487
1488 let crashed_tmp = partition.join("snapshot.alice.range.parquet.tmp");
1491 std::fs::write(&crashed_tmp, b"partial parquet bytes").unwrap();
1492 assert!(crashed_tmp.is_file());
1493
1494 let storage2 = ParquetStorage::new(temp_dir.path()).unwrap();
1496 assert!(
1497 !crashed_tmp.exists(),
1498 "stale tmp file should have been cleaned by ParquetStorage::new"
1499 );
1500
1501 let events = storage2.load_events_for_tenant("alice").unwrap();
1503 assert_eq!(events.len(), 3);
1504 }
1505
1506 #[test]
1507 fn test_cold_tier_archives_dropped_events_before_deletion() {
1508 use crate::infrastructure::persistence::cold_tier::LocalFsArchive;
1514
1515 let live_dir = TempDir::new().unwrap();
1516 let archive_dir = TempDir::new().unwrap();
1517
1518 let storage = ParquetStorage::new(live_dir.path()).unwrap();
1520 let now = Utc::now();
1521 for i in 0..50 {
1522 let day_offset = 60 - (i * 60 / 49);
1523 let ts = now - chrono::Duration::days(i64::from(day_offset));
1524 let event = crate::domain::entities::Event::reconstruct_from_strings(
1525 uuid::Uuid::new_v4(),
1526 "test.event".to_string(),
1527 format!("e-{i}"),
1528 "alice".to_string(),
1529 serde_json::json!({"i": i}),
1530 ts,
1531 None,
1532 1,
1533 );
1534 storage.append_event(event).unwrap();
1535 if i % 5 == 4 {
1536 storage.flush().unwrap();
1537 }
1538 }
1539 storage.flush().unwrap();
1540
1541 let mut retention = RetentionConfig::default();
1543 retention.set("alice", Some(Duration::from_hours(30 * 24)));
1544 let archive: Arc<dyn ArchiveTarget> =
1545 Arc::new(LocalFsArchive::new(archive_dir.path()).unwrap());
1546 let config = CompactionConfig {
1547 min_files_to_compact: 2,
1548 small_file_threshold: 100 * 1024 * 1024,
1549 strategy: CompactionStrategy::SizeBased,
1550 retention,
1551 archive: Some(archive),
1552 ..Default::default()
1553 };
1554 let manager = CompactionManager::new(live_dir.path(), config);
1555
1556 let result = manager.compact_tenant("alice").unwrap();
1557 assert!(result.events_compacted > 0, "some events kept");
1558 assert!(
1559 result.events_compacted < 50,
1560 "some events dropped to retention; kept {} of 50",
1561 result.events_compacted
1562 );
1563
1564 let live_after = ParquetStorage::new(live_dir.path())
1566 .unwrap()
1567 .load_events_for_tenant("alice")
1568 .unwrap();
1569 assert_eq!(live_after.len(), result.events_compacted);
1570
1571 let mut archive_files = vec![];
1574 let mut stack = vec![archive_dir.path().to_path_buf()];
1575 while let Some(d) = stack.pop() {
1576 for entry in std::fs::read_dir(&d).unwrap().flatten() {
1577 let p = entry.path();
1578 if p.is_dir() {
1579 stack.push(p);
1580 } else if p
1581 .file_name()
1582 .is_some_and(|n| n.to_string_lossy().starts_with("archive.alice."))
1583 {
1584 archive_files.push(p);
1585 }
1586 }
1587 }
1588 assert!(
1589 !archive_files.is_empty(),
1590 "archive directory must contain at least one archive.alice.* file"
1591 );
1592
1593 let archive_storage = ParquetStorage::new(archive_dir.path()).unwrap();
1596 let archived = archive_storage.load_events_for_tenant("alice").unwrap();
1597 assert_eq!(
1598 live_after.len() + archived.len(),
1599 50,
1600 "live + archived must equal original event count (live={}, archived={})",
1601 live_after.len(),
1602 archived.len()
1603 );
1604 }
1605
1606 #[test]
1607 fn test_cold_tier_failure_keeps_originals_on_disk() {
1608 let live_dir = TempDir::new().unwrap();
1612 let storage = ParquetStorage::new(live_dir.path()).unwrap();
1613 let now = Utc::now();
1614 for i in 0..20 {
1615 let ts = now - chrono::Duration::days(60 - i);
1616 let event = crate::domain::entities::Event::reconstruct_from_strings(
1617 uuid::Uuid::new_v4(),
1618 "test.event".to_string(),
1619 format!("e-{i}"),
1620 "alice".to_string(),
1621 serde_json::json!({"i": i}),
1622 ts,
1623 None,
1624 1,
1625 );
1626 storage.append_event(event).unwrap();
1627 if i % 5 == 4 {
1628 storage.flush().unwrap();
1629 }
1630 }
1631 storage.flush().unwrap();
1632
1633 let count_files = |dir: &std::path::Path| -> usize {
1636 let mut n = 0;
1637 let mut stack = vec![dir.to_path_buf()];
1638 while let Some(d) = stack.pop() {
1639 for entry in std::fs::read_dir(&d).unwrap().flatten() {
1640 let p = entry.path();
1641 if p.is_dir() {
1642 stack.push(p);
1643 } else if p.extension().is_some_and(|e| e == "parquet") {
1644 n += 1;
1645 }
1646 }
1647 }
1648 n
1649 };
1650 let before = count_files(live_dir.path());
1651 assert!(before > 0);
1652
1653 #[derive(Debug)]
1654 struct FailingArchive;
1655 impl ArchiveTarget for FailingArchive {
1656 fn archive(
1657 &self,
1658 _: &str,
1659 _: DateTime<Utc>,
1660 _: DateTime<Utc>,
1661 _: &[crate::domain::entities::Event],
1662 ) -> Result<()> {
1663 Err(AllSourceError::StorageError(
1664 "simulated archive outage".to_string(),
1665 ))
1666 }
1667 }
1668
1669 let mut retention = RetentionConfig::default();
1670 retention.set("alice", Some(Duration::from_hours(30 * 24)));
1671 let config = CompactionConfig {
1672 min_files_to_compact: 2,
1673 small_file_threshold: 100 * 1024 * 1024,
1674 strategy: CompactionStrategy::SizeBased,
1675 retention,
1676 archive: Some(Arc::new(FailingArchive) as Arc<dyn ArchiveTarget>),
1677 ..Default::default()
1678 };
1679 let manager = CompactionManager::new(live_dir.path(), config);
1680
1681 let result = manager.compact_tenant("alice");
1682 assert!(result.is_err(), "compaction must fail when archive fails");
1683
1684 let after = count_files(live_dir.path());
1686 assert_eq!(
1687 before, after,
1688 "no files should be removed after archive failure"
1689 );
1690
1691 let storage2 = ParquetStorage::new(live_dir.path()).unwrap();
1692 let loaded = storage2.load_events_for_tenant("alice").unwrap();
1693 assert_eq!(
1694 loaded.len(),
1695 20,
1696 "all 20 events still present after failed archive"
1697 );
1698 }
1699
1700 #[test]
1701 fn test_cold_tier_not_invoked_when_no_events_dropped() {
1702 let live_dir = TempDir::new().unwrap();
1706 let storage = ParquetStorage::new(live_dir.path()).unwrap();
1707 let now = Utc::now();
1708 for i in 0..10 {
1709 let ts = now - chrono::Duration::hours(i);
1710 let event = crate::domain::entities::Event::reconstruct_from_strings(
1711 uuid::Uuid::new_v4(),
1712 "test.event".to_string(),
1713 format!("e-{i}"),
1714 "alice".to_string(),
1715 serde_json::json!({"i": i}),
1716 ts,
1717 None,
1718 1,
1719 );
1720 storage.append_event(event).unwrap();
1721 if i % 3 == 2 {
1722 storage.flush().unwrap();
1723 }
1724 }
1725 storage.flush().unwrap();
1726
1727 #[derive(Debug)]
1728 struct PanickingArchive;
1729 impl ArchiveTarget for PanickingArchive {
1730 fn archive(
1731 &self,
1732 _: &str,
1733 _: DateTime<Utc>,
1734 _: DateTime<Utc>,
1735 _: &[crate::domain::entities::Event],
1736 ) -> Result<()> {
1737 panic!("archive must not be called when no events are dropped");
1738 }
1739 }
1740
1741 let config = CompactionConfig {
1743 min_files_to_compact: 2,
1744 small_file_threshold: 100 * 1024 * 1024,
1745 strategy: CompactionStrategy::SizeBased,
1746 archive: Some(Arc::new(PanickingArchive) as Arc<dyn ArchiveTarget>),
1747 ..Default::default()
1748 };
1749 let manager = CompactionManager::new(live_dir.path(), config);
1750 let result = manager.compact_tenant("alice").unwrap();
1751 assert_eq!(result.events_compacted, 10);
1752 }
1753
1754 #[test]
1755 fn test_discover_tenants_skips_system_and_hidden() {
1756 let temp_dir = TempDir::new().unwrap();
1759 std::fs::create_dir_all(temp_dir.path().join("alice")).unwrap();
1760 std::fs::create_dir_all(temp_dir.path().join("bob")).unwrap();
1761 std::fs::create_dir_all(temp_dir.path().join("__system")).unwrap();
1762 std::fs::create_dir_all(temp_dir.path().join(".hidden")).unwrap();
1763
1764 let manager = CompactionManager::new(temp_dir.path(), CompactionConfig::default());
1765 let tenants = manager.discover_tenants().unwrap();
1766 assert_eq!(tenants, vec!["alice".to_string(), "bob".to_string()]);
1767 }
1768}