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 match storage.load_events_from_file_path(&fi.path, &event_tenant) {
561 Ok(mut e) => events.append(&mut e),
562 Err(e) => {
563 tracing::error!(
564 file = %fi.path.display(),
565 "failed to read parquet file for compaction: {e}"
566 );
567 }
568 }
569 }
570
571 if events.is_empty() {
572 tracing::warn!(
573 tenant_id = tenant_id,
574 "candidate files had no readable events; skipping snapshot"
575 );
576 return Ok(CompactionResult::default());
577 }
578
579 let dropped_by_retention = if let Some(ttl) = self.config.retention.ttl_for(tenant_id) {
595 let cutoff = Utc::now()
596 - chrono::Duration::from_std(ttl).unwrap_or_else(|_| chrono::Duration::zero());
597 let before = events.len();
598
599 let (drained, kept): (Vec<_>, Vec<_>) = std::mem::take(&mut events)
604 .into_iter()
605 .partition(|e| e.timestamp < cutoff);
606 events = kept;
607 let dropped = before - events.len();
608
609 if dropped > 0 {
610 tracing::info!(
611 retention_tenant = tenant_id,
612 dropped = dropped,
613 kept = events.len(),
614 cutoff = %cutoff.to_rfc3339(),
615 ttl_secs = ttl.as_secs(),
616 "retention: dropped events older than TTL"
617 );
618
619 if let Some(archive) = self.config.archive.as_ref() {
620 let from = drained
621 .iter()
622 .map(|e| e.timestamp)
623 .min()
624 .expect("dropped > 0 guarantees non-empty drained");
625 let to = drained
626 .iter()
627 .map(|e| e.timestamp)
628 .max()
629 .expect("dropped > 0 guarantees non-empty drained");
630 archive.archive(tenant_id, from, to, &drained)?;
631 tracing::info!(
632 retention_tenant = tenant_id,
633 archived_to = %archive.description(),
634 archived = drained.len(),
635 "retention: dropped events archived to cold tier"
636 );
637 }
638 }
639 dropped
640 } else {
641 0
642 };
643
644 if events.is_empty() {
652 tracing::info!(
653 tenant_id = tenant_id,
654 files_dropped = candidates.len(),
655 events_dropped = dropped_by_retention,
656 "retention: every event aged out — deleting originals without snapshot"
657 );
658 for fi in &candidates {
659 if let Err(e) = fs::remove_file(&fi.path) {
660 tracing::error!(
661 file = %fi.path.display(),
662 "failed to remove fully-aged raw file: {e}"
663 );
664 }
665 }
666 return Ok(CompactionResult {
667 files_compacted: candidates.len(),
668 bytes_before,
669 bytes_after: 0,
670 events_compacted: 0,
671 duration_ms: start_time.elapsed().as_millis() as u64,
672 });
673 }
674
675 events.sort_by_key(|e| e.timestamp);
676 let from = events.first().expect("non-empty checked above").timestamp;
677 let to = events.last().expect("non-empty checked above").timestamp;
678
679 let file_stem = format!(
682 "snapshot.{tenant_id}.{}-{}",
683 format_iso_basic(from),
684 format_iso_basic(to)
685 );
686 let snapshot_path = storage.write_atomic_parquet(tenant_id, &file_stem, &events)?;
687 let bytes_after = fs::metadata(&snapshot_path).map_or(0, |m| m.len());
688
689 for fi in &candidates {
694 if let Err(e) = fs::remove_file(&fi.path) {
695 tracing::error!(
696 file = %fi.path.display(),
697 "failed to remove pre-snapshot raw file: {e}"
698 );
699 }
700 }
701
702 let duration_ms = start_time.elapsed().as_millis() as u64;
703 tracing::info!(
704 tenant_id = tenant_id,
705 files_compacted = candidates.len(),
706 events = events.len(),
707 dropped_by_retention = dropped_by_retention,
708 mib_before = bytes_before as f64 / (1024.0 * 1024.0),
709 mib_after = bytes_after as f64 / (1024.0 * 1024.0),
710 duration_ms = duration_ms,
711 "tenant compaction complete"
712 );
713
714 Ok(CompactionResult {
715 files_compacted: candidates.len(),
716 bytes_before,
717 bytes_after,
718 events_compacted: events.len(),
719 duration_ms,
720 })
721 }
722
723 fn discover_tenants(&self) -> Result<Vec<String>> {
728 let Ok(entries) = fs::read_dir(&self.storage_dir) else {
729 return Ok(Vec::new());
730 };
731 let mut tenants: Vec<String> = entries
732 .filter_map(std::result::Result::ok)
733 .filter_map(|entry| {
734 let ft = entry.file_type().ok()?;
735 if !ft.is_dir() {
736 return None;
737 }
738 let name = entry.file_name().to_string_lossy().into_owned();
739 if name.starts_with('.') || name == "__system" {
742 return None;
743 }
744 Some(name)
745 })
746 .collect();
747 tenants.sort();
748 Ok(tenants)
749 }
750
751 pub fn stats(&self) -> CompactionStats {
753 (*self.stats.read()).clone()
754 }
755
756 pub fn config(&self) -> &CompactionConfig {
758 &self.config
759 }
760
761 #[cfg_attr(feature = "hotpath", hotpath::measure)]
763 pub fn compact_now(&self) -> Result<CompactionResult> {
764 tracing::info!("Manual compaction triggered");
765 self.compact()
766 }
767}
768
769#[derive(Debug, Clone, Default, Serialize)]
771pub struct CompactionResult {
772 pub files_compacted: usize,
773 pub bytes_before: u64,
774 pub bytes_after: u64,
775 pub events_compacted: usize,
776 pub duration_ms: u64,
777}
778
779pub(super) fn format_iso_basic(t: DateTime<Utc>) -> String {
785 t.format("%Y-%m-%dT%H%M%SZ").to_string()
786}
787
788pub struct CompactionTask {
790 manager: Arc<CompactionManager>,
791 interval: Duration,
792}
793
794impl CompactionTask {
795 pub fn new(manager: Arc<CompactionManager>, interval_seconds: u64) -> Self {
797 Self {
798 manager,
799 interval: Duration::from_secs(interval_seconds),
800 }
801 }
802
803 #[cfg_attr(feature = "hotpath", hotpath::measure)]
805 pub async fn run(self) {
806 let mut interval = tokio::time::interval(self.interval);
807
808 loop {
809 interval.tick().await;
810
811 if self.manager.should_compact() {
812 tracing::debug!("Auto-compaction check triggered");
813
814 match self.manager.compact() {
815 Ok(result) => {
816 if result.files_compacted > 0 {
817 tracing::info!(
818 "Auto-compaction succeeded: {} files, {:.2} MB saved",
819 result.files_compacted,
820 (result.bytes_before - result.bytes_after) as f64
821 / (1024.0 * 1024.0)
822 );
823 }
824 }
825 Err(e) => {
826 tracing::error!("Auto-compaction failed: {}", e);
827 }
828 }
829 }
830 }
831 }
832}
833
834#[cfg(test)]
835mod tests {
836 use super::*;
837 use tempfile::TempDir;
838
839 #[test]
840 fn test_compaction_manager_creation() {
841 let temp_dir = TempDir::new().unwrap();
842 let config = CompactionConfig::default();
843 let manager = CompactionManager::new(temp_dir.path(), config);
844
845 assert_eq!(manager.stats().total_compactions, 0);
846 }
847
848 #[test]
849 fn test_should_compact() {
850 let temp_dir = TempDir::new().unwrap();
851 let config = CompactionConfig {
852 auto_compact: true,
853 compaction_interval_seconds: 1,
854 ..Default::default()
855 };
856 let manager = CompactionManager::new(temp_dir.path(), config);
857
858 assert!(manager.should_compact());
860 }
861
862 #[test]
863 fn test_file_selection_size_based() {
864 let temp_dir = TempDir::new().unwrap();
865 let config = CompactionConfig {
866 small_file_threshold: 1024 * 1024, min_files_to_compact: 2,
868 strategy: CompactionStrategy::SizeBased,
869 ..Default::default()
870 };
871 let manager = CompactionManager::new(temp_dir.path(), config);
872
873 let files = vec![
874 FileInfo {
875 path: PathBuf::from("small1.parquet"),
876 size: 500_000, created: Utc::now(),
878 },
879 FileInfo {
880 path: PathBuf::from("small2.parquet"),
881 size: 600_000, created: Utc::now(),
883 },
884 FileInfo {
885 path: PathBuf::from("large.parquet"),
886 size: 10_000_000, created: Utc::now(),
888 },
889 ];
890
891 let selected = manager.select_files_for_compaction(&files);
892 assert_eq!(selected.len(), 2); }
894
895 #[test]
896 fn test_default_compaction_config() {
897 let config = CompactionConfig::default();
898 assert_eq!(config.min_files_to_compact, 3);
899 assert_eq!(config.target_file_size, 128 * 1024 * 1024);
900 assert_eq!(config.max_file_size, 256 * 1024 * 1024);
901 assert_eq!(config.small_file_threshold, 10 * 1024 * 1024);
902 assert_eq!(config.compaction_interval_seconds, 3600);
903 assert!(config.auto_compact);
904 assert_eq!(config.strategy, CompactionStrategy::SizeBased);
905 }
906
907 #[test]
908 fn test_should_compact_disabled() {
909 let temp_dir = TempDir::new().unwrap();
910 let config = CompactionConfig {
911 auto_compact: false,
912 ..Default::default()
913 };
914 let manager = CompactionManager::new(temp_dir.path(), config);
915
916 assert!(!manager.should_compact());
917 }
918
919 #[test]
920 fn test_compact_empty_directory() {
921 let temp_dir = TempDir::new().unwrap();
922 let config = CompactionConfig::default();
923 let manager = CompactionManager::new(temp_dir.path(), config);
924
925 let result = manager.compact().unwrap();
926 assert_eq!(result.files_compacted, 0);
927 assert_eq!(result.bytes_before, 0);
928 assert_eq!(result.bytes_after, 0);
929 assert_eq!(result.events_compacted, 0);
930 }
931
932 #[test]
933 fn test_compact_now() {
934 let temp_dir = TempDir::new().unwrap();
935 let config = CompactionConfig::default();
936 let manager = CompactionManager::new(temp_dir.path(), config);
937
938 let result = manager.compact_now().unwrap();
939 assert_eq!(result.files_compacted, 0);
940 }
941
942 #[test]
943 fn test_get_config() {
944 let temp_dir = TempDir::new().unwrap();
945 let config = CompactionConfig {
946 min_files_to_compact: 5,
947 ..Default::default()
948 };
949 let manager = CompactionManager::new(temp_dir.path(), config);
950
951 assert_eq!(manager.config().min_files_to_compact, 5);
952 }
953
954 #[test]
955 fn test_get_stats() {
956 let temp_dir = TempDir::new().unwrap();
957 let config = CompactionConfig::default();
958 let manager = CompactionManager::new(temp_dir.path(), config);
959
960 let stats = manager.stats();
961 assert_eq!(stats.total_compactions, 0);
962 assert_eq!(stats.total_files_compacted, 0);
963 assert_eq!(stats.total_bytes_before, 0);
964 assert_eq!(stats.total_bytes_after, 0);
965 assert_eq!(stats.total_events_compacted, 0);
966 assert_eq!(stats.last_compaction_duration_ms, 0);
967 assert_eq!(stats.space_saved_bytes, 0);
968 }
969
970 #[test]
971 fn test_file_selection_not_enough_small_files() {
972 let temp_dir = TempDir::new().unwrap();
973 let config = CompactionConfig {
974 small_file_threshold: 1024 * 1024,
975 min_files_to_compact: 3, strategy: CompactionStrategy::SizeBased,
977 ..Default::default()
978 };
979 let manager = CompactionManager::new(temp_dir.path(), config);
980
981 let files = vec![
982 FileInfo {
983 path: PathBuf::from("small1.parquet"),
984 size: 500_000,
985 created: Utc::now(),
986 },
987 FileInfo {
988 path: PathBuf::from("small2.parquet"),
989 size: 600_000,
990 created: Utc::now(),
991 },
992 ];
993
994 let selected = manager.select_files_for_compaction(&files);
995 assert_eq!(selected.len(), 0); }
997
998 #[test]
999 fn test_file_selection_time_based() {
1000 let temp_dir = TempDir::new().unwrap();
1001 let config = CompactionConfig {
1002 min_files_to_compact: 2,
1003 strategy: CompactionStrategy::TimeBased,
1004 ..Default::default()
1005 };
1006 let manager = CompactionManager::new(temp_dir.path(), config);
1007
1008 let old_time = Utc::now() - chrono::Duration::hours(48);
1009 let files = vec![
1010 FileInfo {
1011 path: PathBuf::from("old1.parquet"),
1012 size: 1_000_000,
1013 created: old_time,
1014 },
1015 FileInfo {
1016 path: PathBuf::from("old2.parquet"),
1017 size: 2_000_000,
1018 created: old_time,
1019 },
1020 FileInfo {
1021 path: PathBuf::from("new.parquet"),
1022 size: 500_000,
1023 created: Utc::now(),
1024 },
1025 ];
1026
1027 let selected = manager.select_files_for_compaction(&files);
1028 assert_eq!(selected.len(), 2); }
1030
1031 #[test]
1032 fn test_file_selection_time_based_not_enough() {
1033 let temp_dir = TempDir::new().unwrap();
1034 let config = CompactionConfig {
1035 min_files_to_compact: 3,
1036 strategy: CompactionStrategy::TimeBased,
1037 ..Default::default()
1038 };
1039 let manager = CompactionManager::new(temp_dir.path(), config);
1040
1041 let old_time = Utc::now() - chrono::Duration::hours(48);
1042 let files = vec![
1043 FileInfo {
1044 path: PathBuf::from("old1.parquet"),
1045 size: 1_000_000,
1046 created: old_time,
1047 },
1048 FileInfo {
1049 path: PathBuf::from("new.parquet"),
1050 size: 500_000,
1051 created: Utc::now(),
1052 },
1053 ];
1054
1055 let selected = manager.select_files_for_compaction(&files);
1056 assert_eq!(selected.len(), 0); }
1058
1059 #[test]
1060 fn test_file_selection_full_compaction() {
1061 let temp_dir = TempDir::new().unwrap();
1062 let config = CompactionConfig {
1063 strategy: CompactionStrategy::FullCompaction,
1064 ..Default::default()
1065 };
1066 let manager = CompactionManager::new(temp_dir.path(), config);
1067
1068 let files = vec![
1069 FileInfo {
1070 path: PathBuf::from("file1.parquet"),
1071 size: 1_000_000,
1072 created: Utc::now(),
1073 },
1074 FileInfo {
1075 path: PathBuf::from("file2.parquet"),
1076 size: 2_000_000,
1077 created: Utc::now(),
1078 },
1079 ];
1080
1081 let selected = manager.select_files_for_compaction(&files);
1082 assert_eq!(selected.len(), 2); }
1084
1085 #[test]
1086 fn test_compaction_strategy_serde() {
1087 let strategies = vec![
1088 CompactionStrategy::SizeBased,
1089 CompactionStrategy::TimeBased,
1090 CompactionStrategy::FullCompaction,
1091 ];
1092
1093 for strategy in strategies {
1094 let json = serde_json::to_string(&strategy).unwrap();
1095 let parsed: CompactionStrategy = serde_json::from_str(&json).unwrap();
1096 assert_eq!(parsed, strategy);
1097 }
1098 }
1099
1100 #[test]
1101 fn test_compaction_stats_default() {
1102 let stats = CompactionStats::default();
1103 assert_eq!(stats.total_compactions, 0);
1104 assert_eq!(stats.total_files_compacted, 0);
1105 }
1106
1107 #[test]
1108 fn test_compaction_stats_serde() {
1109 let stats = CompactionStats {
1110 total_compactions: 5,
1111 total_files_compacted: 20,
1112 total_bytes_before: 1000000,
1113 total_bytes_after: 500000,
1114 total_events_compacted: 10000,
1115 last_compaction_duration_ms: 500,
1116 space_saved_bytes: 500000,
1117 };
1118
1119 let json = serde_json::to_string(&stats).unwrap();
1120 assert!(json.contains("\"total_compactions\":5"));
1121 assert!(json.contains("\"space_saved_bytes\":500000"));
1122 }
1123
1124 #[test]
1125 fn test_compaction_result_serde() {
1126 let result = CompactionResult {
1127 files_compacted: 3,
1128 bytes_before: 1000000,
1129 bytes_after: 500000,
1130 events_compacted: 5000,
1131 duration_ms: 250,
1132 };
1133
1134 let json = serde_json::to_string(&result).unwrap();
1135 assert!(json.contains("\"files_compacted\":3"));
1136 assert!(json.contains("\"bytes_before\":1000000"));
1137 }
1138
1139 #[test]
1140 fn test_compaction_task_creation() {
1141 let temp_dir = TempDir::new().unwrap();
1142 let config = CompactionConfig::default();
1143 let manager = Arc::new(CompactionManager::new(temp_dir.path(), config));
1144
1145 let _task = CompactionTask::new(manager.clone(), 60);
1146 }
1148
1149 #[test]
1150 fn test_list_parquet_files_empty() {
1151 let temp_dir = TempDir::new().unwrap();
1152 let config = CompactionConfig::default();
1153 let manager = CompactionManager::new(temp_dir.path(), config);
1154
1155 let files = manager.list_parquet_files().unwrap();
1156 assert!(files.is_empty());
1157 }
1158
1159 #[test]
1160 fn test_list_parquet_files_with_non_parquet() {
1161 let temp_dir = TempDir::new().unwrap();
1162 let config = CompactionConfig::default();
1163 let manager = CompactionManager::new(temp_dir.path(), config);
1164
1165 std::fs::write(temp_dir.path().join("test.txt"), "test").unwrap();
1167 std::fs::write(temp_dir.path().join("data.json"), "{}").unwrap();
1168
1169 let files = manager.list_parquet_files().unwrap();
1170 assert!(files.is_empty()); }
1172
1173 fn ingest_and_flush_per_call(storage_dir: &std::path::Path, tenant: &str, count: usize) {
1178 for i in 0..count {
1182 let storage = ParquetStorage::with_config(
1183 storage_dir,
1184 crate::infrastructure::persistence::ParquetStorageConfig {
1185 batch_size: 1,
1186 ..Default::default()
1187 },
1188 )
1189 .unwrap();
1190 let event = crate::domain::entities::Event::from_strings(
1191 "test.event".to_string(),
1192 format!("{tenant}-{i}"),
1193 tenant.to_string(),
1194 serde_json::json!({"i": i}),
1195 None,
1196 )
1197 .unwrap();
1198 storage.append_event(event).unwrap();
1199 storage.flush().unwrap();
1200 }
1201 }
1202
1203 #[test]
1204 fn test_compact_tenant_emits_one_snapshot_and_removes_originals() {
1205 let temp_dir = TempDir::new().unwrap();
1206
1207 ingest_and_flush_per_call(temp_dir.path(), "alice", 4);
1209
1210 let config = CompactionConfig {
1211 min_files_to_compact: 2,
1212 small_file_threshold: 100 * 1024 * 1024,
1213 strategy: CompactionStrategy::SizeBased,
1214 ..Default::default()
1215 };
1216 let manager = CompactionManager::new(temp_dir.path(), config);
1217
1218 let result = manager.compact_tenant("alice").unwrap();
1219 assert_eq!(result.files_compacted, 4);
1220 assert_eq!(result.events_compacted, 4);
1221
1222 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1225 let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1226 assert_eq!(
1227 alice_files.len(),
1228 1,
1229 "expected exactly one snapshot file for alice"
1230 );
1231
1232 let name = alice_files[0]
1233 .file_name()
1234 .and_then(|n| n.to_str())
1235 .unwrap()
1236 .to_string();
1237 assert!(
1238 name.starts_with("snapshot.alice."),
1239 "expected snapshot prefix, got {name}"
1240 );
1241 assert!(name.ends_with(".parquet"));
1242
1243 let tmps: Vec<_> = std::fs::read_dir(alice_files[0].parent().unwrap())
1245 .unwrap()
1246 .filter_map(std::result::Result::ok)
1247 .filter(|e| e.path().to_string_lossy().ends_with(".tmp"))
1248 .collect();
1249 assert!(tmps.is_empty());
1250
1251 let loaded = storage.load_events_for_tenant("alice").unwrap();
1253 assert_eq!(loaded.len(), 4);
1254 for e in &loaded {
1255 assert_eq!(e.tenant_id_str(), "alice");
1256 }
1257 }
1258
1259 #[test]
1260 fn test_compact_tenant_skips_existing_snapshot_files() {
1261 let temp_dir = TempDir::new().unwrap();
1266 ingest_and_flush_per_call(temp_dir.path(), "alice", 4);
1267
1268 let config = CompactionConfig {
1269 min_files_to_compact: 2,
1270 small_file_threshold: 100 * 1024 * 1024,
1271 ..Default::default()
1272 };
1273 let manager = CompactionManager::new(temp_dir.path(), config);
1274
1275 let r1 = manager.compact_tenant("alice").unwrap();
1276 assert_eq!(r1.files_compacted, 4);
1277
1278 let r2 = manager.compact_tenant("alice").unwrap();
1279 assert_eq!(r2.files_compacted, 0, "snapshot must not be re-compacted");
1280
1281 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1283 let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1284 assert_eq!(alice_files.len(), 1);
1285 }
1286
1287 #[test]
1288 fn test_compact_tenant_below_threshold_is_a_noop() {
1289 let temp_dir = TempDir::new().unwrap();
1291 ingest_and_flush_per_call(temp_dir.path(), "alice", 1);
1292
1293 let manager = CompactionManager::new(temp_dir.path(), CompactionConfig::default());
1294 let result = manager.compact_tenant("alice").unwrap();
1295 assert_eq!(result.files_compacted, 0);
1296
1297 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1299 let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1300 assert_eq!(alice_files.len(), 1);
1301 let name = alice_files[0]
1302 .file_name()
1303 .unwrap()
1304 .to_string_lossy()
1305 .into_owned();
1306 assert!(
1307 !name.starts_with("snapshot."),
1308 "raw file must not be renamed"
1309 );
1310 }
1311
1312 #[test]
1313 fn test_compact_iterates_every_tenant() {
1314 let temp_dir = TempDir::new().unwrap();
1317 ingest_and_flush_per_call(temp_dir.path(), "alice", 3);
1318 ingest_and_flush_per_call(temp_dir.path(), "bob", 3);
1319
1320 let config = CompactionConfig {
1321 min_files_to_compact: 2,
1322 small_file_threshold: 100 * 1024 * 1024,
1323 strategy: CompactionStrategy::SizeBased,
1324 ..Default::default()
1325 };
1326 let manager = CompactionManager::new(temp_dir.path(), config);
1327
1328 let result = manager.compact().unwrap();
1329 assert_eq!(result.files_compacted, 6);
1330 assert_eq!(result.events_compacted, 6);
1331
1332 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1333 for tenant in ["alice", "bob"] {
1334 let files = storage.list_parquet_files_for_tenant(tenant).unwrap();
1335 assert_eq!(files.len(), 1, "{tenant} should have one snapshot");
1336 let name = files[0].file_name().unwrap().to_string_lossy().into_owned();
1337 assert!(name.starts_with(&format!("snapshot.{tenant}.")));
1338 }
1339 }
1340
1341 #[test]
1342 fn test_retention_drops_events_older_than_ttl() {
1343 let temp_dir = TempDir::new().unwrap();
1347
1348 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1353 let now = Utc::now();
1354 for i in 0..100 {
1355 let day_offset = 60 - (i * 60 / 99);
1357 let ts = now - chrono::Duration::days(i64::from(day_offset));
1358 let event = crate::domain::entities::Event::reconstruct_from_strings(
1359 uuid::Uuid::new_v4(),
1360 "test.event".to_string(),
1361 format!("e-{i}"),
1362 "alice".to_string(),
1363 serde_json::json!({"i": i}),
1364 ts,
1365 None,
1366 1,
1367 );
1368 storage.append_event(event).unwrap();
1369 if i % 10 == 9 {
1373 storage.flush().unwrap();
1374 }
1375 }
1376 storage.flush().unwrap();
1377
1378 let mut retention = RetentionConfig::default();
1380 retention.set("alice", Some(Duration::from_hours(30 * 24)));
1381 let config = CompactionConfig {
1382 min_files_to_compact: 2,
1383 small_file_threshold: 100 * 1024 * 1024,
1384 strategy: CompactionStrategy::SizeBased,
1385 retention,
1386 ..Default::default()
1387 };
1388 let manager = CompactionManager::new(temp_dir.path(), config);
1389
1390 let result = manager.compact_tenant("alice").unwrap();
1391 assert!(result.events_compacted > 0);
1392 assert!(
1393 result.events_compacted < 100,
1394 "retention should have dropped some events; kept {} of 100",
1395 result.events_compacted
1396 );
1397
1398 let storage2 = ParquetStorage::new(temp_dir.path()).unwrap();
1400 let loaded = storage2.load_events_for_tenant("alice").unwrap();
1401 assert_eq!(loaded.len(), result.events_compacted);
1402
1403 let cutoff = Utc::now() - chrono::Duration::days(30);
1406 for e in &loaded {
1407 assert!(
1408 e.timestamp >= cutoff - chrono::Duration::seconds(60),
1409 "event with ts {} survived retention but is older than cutoff {}",
1410 e.timestamp.to_rfc3339(),
1411 cutoff.to_rfc3339()
1412 );
1413 }
1414 }
1415
1416 #[test]
1417 fn test_retention_keeps_forever_by_default_for_non_system_tenants() {
1418 let temp_dir = TempDir::new().unwrap();
1421 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1422 let now = Utc::now();
1423 for i in 0..6 {
1424 let ts = now - chrono::Duration::days(i * 365);
1425 let event = crate::domain::entities::Event::reconstruct_from_strings(
1426 uuid::Uuid::new_v4(),
1427 "test.event".to_string(),
1428 format!("e-{i}"),
1429 "alice".to_string(),
1430 serde_json::json!({"i": i}),
1431 ts,
1432 None,
1433 1,
1434 );
1435 storage.append_event(event).unwrap();
1436 if i % 2 == 1 {
1437 storage.flush().unwrap();
1438 }
1439 }
1440 storage.flush().unwrap();
1441
1442 let config = CompactionConfig {
1444 min_files_to_compact: 2,
1445 small_file_threshold: 100 * 1024 * 1024,
1446 strategy: CompactionStrategy::SizeBased,
1447 ..Default::default()
1448 };
1449 let manager = CompactionManager::new(temp_dir.path(), config);
1450 let result = manager.compact_tenant("alice").unwrap();
1451 assert_eq!(result.events_compacted, 6, "no events should be dropped");
1452 }
1453
1454 #[test]
1455 fn test_retention_system_tenant_default_is_30_days() {
1456 let cfg = RetentionConfig::default();
1460 let ttl = cfg.ttl_for("system").unwrap();
1461 assert_eq!(ttl.as_secs(), 30 * 24 * 3600);
1462 assert!(cfg.ttl_for("acme").is_none());
1464 }
1465
1466 #[test]
1467 fn test_retention_drops_all_events_deletes_originals_without_snapshot() {
1468 let temp_dir = TempDir::new().unwrap();
1472 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1473 let very_old = Utc::now() - chrono::Duration::days(90);
1474 for i in 0..6 {
1475 let event = crate::domain::entities::Event::reconstruct_from_strings(
1476 uuid::Uuid::new_v4(),
1477 "test.event".to_string(),
1478 format!("e-{i}"),
1479 "alice".to_string(),
1480 serde_json::json!({"i": i}),
1481 very_old,
1482 None,
1483 1,
1484 );
1485 storage.append_event(event).unwrap();
1486 if i % 2 == 1 {
1487 storage.flush().unwrap();
1488 }
1489 }
1490 storage.flush().unwrap();
1491
1492 let mut retention = RetentionConfig::default();
1493 retention.set("alice", Some(Duration::from_hours(7 * 24)));
1494 let config = CompactionConfig {
1495 min_files_to_compact: 2,
1496 small_file_threshold: 100 * 1024 * 1024,
1497 strategy: CompactionStrategy::SizeBased,
1498 retention,
1499 ..Default::default()
1500 };
1501 let manager = CompactionManager::new(temp_dir.path(), config);
1502 let result = manager.compact_tenant("alice").unwrap();
1503 assert_eq!(result.events_compacted, 0);
1504 assert!(result.files_compacted >= 2); let storage2 = ParquetStorage::new(temp_dir.path()).unwrap();
1508 let alice_files = storage2.list_parquet_files_for_tenant("alice").unwrap();
1509 assert!(alice_files.is_empty(), "all originals should be deleted");
1510 }
1511
1512 #[test]
1513 fn test_compaction_with_simulated_crash_leaves_data_recoverable() {
1514 let temp_dir = TempDir::new().unwrap();
1520
1521 ingest_and_flush_per_call(temp_dir.path(), "alice", 3);
1523
1524 let storage = ParquetStorage::new(temp_dir.path()).unwrap();
1526 let alice_files = storage.list_parquet_files_for_tenant("alice").unwrap();
1527 let partition = alice_files[0].parent().unwrap().to_path_buf();
1528
1529 let crashed_tmp = partition.join("snapshot.alice.range.parquet.tmp");
1532 std::fs::write(&crashed_tmp, b"partial parquet bytes").unwrap();
1533 assert!(crashed_tmp.is_file());
1534
1535 let storage2 = ParquetStorage::new(temp_dir.path()).unwrap();
1537 assert!(
1538 !crashed_tmp.exists(),
1539 "stale tmp file should have been cleaned by ParquetStorage::new"
1540 );
1541
1542 let events = storage2.load_events_for_tenant("alice").unwrap();
1544 assert_eq!(events.len(), 3);
1545 }
1546
1547 #[test]
1548 fn test_cold_tier_archives_dropped_events_before_deletion() {
1549 use crate::infrastructure::persistence::cold_tier::LocalFsArchive;
1555
1556 let live_dir = TempDir::new().unwrap();
1557 let archive_dir = TempDir::new().unwrap();
1558
1559 let storage = ParquetStorage::new(live_dir.path()).unwrap();
1561 let now = Utc::now();
1562 for i in 0..50 {
1563 let day_offset = 60 - (i * 60 / 49);
1564 let ts = now - chrono::Duration::days(i64::from(day_offset));
1565 let event = crate::domain::entities::Event::reconstruct_from_strings(
1566 uuid::Uuid::new_v4(),
1567 "test.event".to_string(),
1568 format!("e-{i}"),
1569 "alice".to_string(),
1570 serde_json::json!({"i": i}),
1571 ts,
1572 None,
1573 1,
1574 );
1575 storage.append_event(event).unwrap();
1576 if i % 5 == 4 {
1577 storage.flush().unwrap();
1578 }
1579 }
1580 storage.flush().unwrap();
1581
1582 let mut retention = RetentionConfig::default();
1584 retention.set("alice", Some(Duration::from_hours(30 * 24)));
1585 let archive: Arc<dyn ArchiveTarget> =
1586 Arc::new(LocalFsArchive::new(archive_dir.path()).unwrap());
1587 let config = CompactionConfig {
1588 min_files_to_compact: 2,
1589 small_file_threshold: 100 * 1024 * 1024,
1590 strategy: CompactionStrategy::SizeBased,
1591 retention,
1592 archive: Some(archive),
1593 ..Default::default()
1594 };
1595 let manager = CompactionManager::new(live_dir.path(), config);
1596
1597 let result = manager.compact_tenant("alice").unwrap();
1598 assert!(result.events_compacted > 0, "some events kept");
1599 assert!(
1600 result.events_compacted < 50,
1601 "some events dropped to retention; kept {} of 50",
1602 result.events_compacted
1603 );
1604
1605 let live_after = ParquetStorage::new(live_dir.path())
1607 .unwrap()
1608 .load_events_for_tenant("alice")
1609 .unwrap();
1610 assert_eq!(live_after.len(), result.events_compacted);
1611
1612 let mut archive_files = vec![];
1615 let mut stack = vec![archive_dir.path().to_path_buf()];
1616 while let Some(d) = stack.pop() {
1617 for entry in std::fs::read_dir(&d).unwrap().flatten() {
1618 let p = entry.path();
1619 if p.is_dir() {
1620 stack.push(p);
1621 } else if p
1622 .file_name()
1623 .is_some_and(|n| n.to_string_lossy().starts_with("archive.alice."))
1624 {
1625 archive_files.push(p);
1626 }
1627 }
1628 }
1629 assert!(
1630 !archive_files.is_empty(),
1631 "archive directory must contain at least one archive.alice.* file"
1632 );
1633
1634 let archive_storage = ParquetStorage::new(archive_dir.path()).unwrap();
1637 let archived = archive_storage.load_events_for_tenant("alice").unwrap();
1638 assert_eq!(
1639 live_after.len() + archived.len(),
1640 50,
1641 "live + archived must equal original event count (live={}, archived={})",
1642 live_after.len(),
1643 archived.len()
1644 );
1645 }
1646
1647 #[test]
1648 fn test_cold_tier_failure_keeps_originals_on_disk() {
1649 let live_dir = TempDir::new().unwrap();
1653 let storage = ParquetStorage::new(live_dir.path()).unwrap();
1654 let now = Utc::now();
1655 for i in 0..20 {
1656 let ts = now - chrono::Duration::days(60 - i);
1657 let event = crate::domain::entities::Event::reconstruct_from_strings(
1658 uuid::Uuid::new_v4(),
1659 "test.event".to_string(),
1660 format!("e-{i}"),
1661 "alice".to_string(),
1662 serde_json::json!({"i": i}),
1663 ts,
1664 None,
1665 1,
1666 );
1667 storage.append_event(event).unwrap();
1668 if i % 5 == 4 {
1669 storage.flush().unwrap();
1670 }
1671 }
1672 storage.flush().unwrap();
1673
1674 let count_files = |dir: &std::path::Path| -> usize {
1677 let mut n = 0;
1678 let mut stack = vec![dir.to_path_buf()];
1679 while let Some(d) = stack.pop() {
1680 for entry in std::fs::read_dir(&d).unwrap().flatten() {
1681 let p = entry.path();
1682 if p.is_dir() {
1683 stack.push(p);
1684 } else if p.extension().is_some_and(|e| e == "parquet") {
1685 n += 1;
1686 }
1687 }
1688 }
1689 n
1690 };
1691 let before = count_files(live_dir.path());
1692 assert!(before > 0);
1693
1694 #[derive(Debug)]
1695 struct FailingArchive;
1696 impl ArchiveTarget for FailingArchive {
1697 fn archive(
1698 &self,
1699 _: &str,
1700 _: DateTime<Utc>,
1701 _: DateTime<Utc>,
1702 _: &[crate::domain::entities::Event],
1703 ) -> Result<()> {
1704 Err(AllSourceError::StorageError(
1705 "simulated archive outage".to_string(),
1706 ))
1707 }
1708 }
1709
1710 let mut retention = RetentionConfig::default();
1711 retention.set("alice", Some(Duration::from_hours(30 * 24)));
1712 let config = CompactionConfig {
1713 min_files_to_compact: 2,
1714 small_file_threshold: 100 * 1024 * 1024,
1715 strategy: CompactionStrategy::SizeBased,
1716 retention,
1717 archive: Some(Arc::new(FailingArchive) as Arc<dyn ArchiveTarget>),
1718 ..Default::default()
1719 };
1720 let manager = CompactionManager::new(live_dir.path(), config);
1721
1722 let result = manager.compact_tenant("alice");
1723 assert!(result.is_err(), "compaction must fail when archive fails");
1724
1725 let after = count_files(live_dir.path());
1727 assert_eq!(
1728 before, after,
1729 "no files should be removed after archive failure"
1730 );
1731
1732 let storage2 = ParquetStorage::new(live_dir.path()).unwrap();
1733 let loaded = storage2.load_events_for_tenant("alice").unwrap();
1734 assert_eq!(
1735 loaded.len(),
1736 20,
1737 "all 20 events still present after failed archive"
1738 );
1739 }
1740
1741 #[test]
1742 fn test_cold_tier_not_invoked_when_no_events_dropped() {
1743 let live_dir = TempDir::new().unwrap();
1747 let storage = ParquetStorage::new(live_dir.path()).unwrap();
1748 let now = Utc::now();
1749 for i in 0..10 {
1750 let ts = now - chrono::Duration::hours(i);
1751 let event = crate::domain::entities::Event::reconstruct_from_strings(
1752 uuid::Uuid::new_v4(),
1753 "test.event".to_string(),
1754 format!("e-{i}"),
1755 "alice".to_string(),
1756 serde_json::json!({"i": i}),
1757 ts,
1758 None,
1759 1,
1760 );
1761 storage.append_event(event).unwrap();
1762 if i % 3 == 2 {
1763 storage.flush().unwrap();
1764 }
1765 }
1766 storage.flush().unwrap();
1767
1768 #[derive(Debug)]
1769 struct PanickingArchive;
1770 impl ArchiveTarget for PanickingArchive {
1771 fn archive(
1772 &self,
1773 _: &str,
1774 _: DateTime<Utc>,
1775 _: DateTime<Utc>,
1776 _: &[crate::domain::entities::Event],
1777 ) -> Result<()> {
1778 panic!("archive must not be called when no events are dropped");
1779 }
1780 }
1781
1782 let config = CompactionConfig {
1784 min_files_to_compact: 2,
1785 small_file_threshold: 100 * 1024 * 1024,
1786 strategy: CompactionStrategy::SizeBased,
1787 archive: Some(Arc::new(PanickingArchive) as Arc<dyn ArchiveTarget>),
1788 ..Default::default()
1789 };
1790 let manager = CompactionManager::new(live_dir.path(), config);
1791 let result = manager.compact_tenant("alice").unwrap();
1792 assert_eq!(result.events_compacted, 10);
1793 }
1794
1795 #[test]
1796 fn test_discover_tenants_skips_system_and_hidden() {
1797 let temp_dir = TempDir::new().unwrap();
1800 std::fs::create_dir_all(temp_dir.path().join("alice")).unwrap();
1801 std::fs::create_dir_all(temp_dir.path().join("bob")).unwrap();
1802 std::fs::create_dir_all(temp_dir.path().join("__system")).unwrap();
1803 std::fs::create_dir_all(temp_dir.path().join(".hidden")).unwrap();
1804
1805 let manager = CompactionManager::new(temp_dir.path(), CompactionConfig::default());
1806 let tenants = manager.discover_tenants().unwrap();
1807 assert_eq!(tenants, vec!["alice".to_string(), "bob".to_string()]);
1808 }
1809}