Volume Compaction Implementation Plan

For Claude: REQUIRED SUB-SKILL: Use superpowers:executing-plans to implement this plan task-by-task.

Goal: Implement production-grade volume garbage collection that reclaims dead space from append-only volume files, integrated as a background worker in s4-server.

Architecture: The compactor scans volume files, identifies dead blobs (those with no matching DedupEntry), copies live blobs to new volumes, atomically updates IndexRecords and DedupEntries via fjall batch writes, then safely deletes old volumes. It runs as a periodic background tokio task inside s4-server, using the s4-compactor crate as a library.

Tech Stack: Rust, fjall (LSM-tree), tokio, bincode, crc32fast, tracing


Overview

What is Volume Compaction?

S4 uses append-only volume files. When objects are deleted or overwritten, the old blob data remains as dead space. Volume compaction:

  1. Scans each volume file, reading every blob header
  2. Identifies dead blobs — those with no DedupEntry pointing to (volume_id, offset)
  3. Copies live blobs to a new volume
  4. Atomically updates all IndexRecords and DedupEntries to point to new locations
  5. Deletes old volume files after verification

Key Invariants

  • Never lose confirmed data — only delete old volumes after ALL live blobs are verified relocated
  • Crash-safe — incomplete compaction must be recoverable (idempotent restart)
  • Non-blocking — compaction runs in background, normal reads/writes continue
  • Dedup-aware — update DedupEntry locations, not just IndexRecords

File Map

File Action Description
s4-compactor/src/compactor.rs Rewrite Core compaction logic (~400 lines)
s4-compactor/src/scrubber.rs Rewrite CRC verification scanner (~150 lines)
s4-compactor/src/lib.rs Create Library exports
s4-compactor/src/main.rs Update Keep as standalone CLI wrapper
s4-compactor/Cargo.toml Update Add crc32fast dependency
s4-server/src/compaction_worker.rs Create Background worker (~120 lines)
s4-server/src/config.rs Modify Add CompactionConfig
s4-server/src/app.rs Modify Spawn compaction worker
s4-server/src/lib.rs Modify Export compaction_worker module
s4-server/Cargo.toml Modify Add s4-compactor dependency
s4-core/src/storage/dedup.rs Modify Add iter_entries() and update_location() methods
s4-core/src/storage/index.rs Modify Add scan_objects_by_location() method
s4-compactor/tests/compaction_integration.rs Create Integration tests
s4-core/tests/compaction_e2e.rs Create E2E tests with full engine
ARCHITECTURE.md Update Document compaction
README.md Update Add compaction config vars
docs/04-features/compaction.md Create Feature documentation

Task 1: Add Dedup Iteration Support to s4-core

Files: - Modify: s4-core/src/storage/dedup.rs

Step 1: Write the failing test

Add to the existing tests module in dedup.rs:

#[test]
fn test_iter_entries() {
    let (_temp, dedup) = create_test_dedup();

    let h1 = Deduplicator::compute_hash(b"data1");
    let h2 = Deduplicator::compute_hash(b"data2");
    let h3 = Deduplicator::compute_hash(b"data3");

    dedup.register_content(h1, 0, 0).unwrap();
    dedup.register_content(h2, 0, 100).unwrap();
    dedup.register_content(h3, 1, 0).unwrap();

    let entries: Vec<_> = dedup.iter_entries().unwrap().collect();
    assert_eq!(entries.len(), 3);

    // Filter by volume 0
    let vol0: Vec<_> = dedup.iter_entries().unwrap()
        .filter(|(_, e)| e.volume_id == 0)
        .collect();
    assert_eq!(vol0.len(), 2);
}

#[test]
fn test_update_location() {
    let (_temp, dedup) = create_test_dedup();
    let hash = Deduplicator::compute_hash(b"movable");

    dedup.register_content(hash, 0, 100).unwrap();

    // Update location
    dedup.update_location(&hash, 5, 200).unwrap();

    let entry = dedup.get_entry(&hash).unwrap().unwrap();
    assert_eq!(entry.volume_id, 5);
    assert_eq!(entry.offset, 200);
    assert_eq!(entry.ref_count, 1); // ref_count unchanged
}

Step 2: Run test to verify it fails

Run: cargo test -p s4-core -- test_iter_entries test_update_location Expected: FAIL — methods don't exist yet

Step 3: Write minimal implementation

Add to Deduplicator impl block in dedup.rs:

/// Iterates all dedup entries in this placement group.
///
/// Returns an iterator of `(content_hash, DedupEntry)` pairs.
/// Used by the compactor to find which blobs are live in a given volume.
///
/// # Performance
///
/// Full keyspace scan. Call from a blocking thread for large datasets.
pub fn iter_entries(&self) -> Result<Vec<([u8; 32], DedupEntry)>, StorageError> {
    let pg_prefix = self.pg_id.to_be_bytes();
    let mut results = Vec::new();

    for guard in self.partition.prefix(&pg_prefix) {
        let (key_bytes, value_bytes) =
            guard.into_inner().map_err(|e| StorageError::Database(e.to_string()))?;
        if key_bytes.len() != 36 {
            continue; // skip malformed entries
        }
        let mut hash = [0u8; 32];
        hash.copy_from_slice(&key_bytes[4..36]);
        let entry: DedupEntry = bincode::deserialize(&value_bytes)
            .map_err(|e| StorageError::Serialization(e.to_string()))?;
        results.push((hash, entry));
    }

    Ok(results)
}

/// Updates the physical location of a dedup entry without changing ref_count.
///
/// Used by the compactor after copying a blob to a new volume.
/// The content hash stays the same — only volume_id and offset change.
pub fn update_location(
    &self,
    content_hash: &[u8; 32],
    new_volume_id: u32,
    new_offset: u64,
) -> Result<(), StorageError> {
    let pg_key = self.make_key(content_hash);
    let existing = self.get_entry(content_hash)?
        .ok_or_else(|| StorageError::InvalidData(
            "Cannot update location: dedup entry not found".to_string()
        ))?;
    let updated = DedupEntry {
        volume_id: new_volume_id,
        offset: new_offset,
        ref_count: existing.ref_count,
    };
    let value = bincode::serialize(&updated)
        .map_err(|e| StorageError::Serialization(e.to_string()))?;
    self.partition
        .insert(&pg_key[..], &value[..])
        .map_err(|e| StorageError::Database(e.to_string()))?;
    Ok(())
}

/// Creates a BatchOp to update the physical location of a dedup entry.
///
/// Used by the compactor for atomic batch updates when relocating blobs.
pub fn make_update_location_op(
    &self,
    content_hash: &[u8; 32],
    new_volume_id: u32,
    new_offset: u64,
) -> Result<BatchOp, StorageError> {
    let pg_key = self.make_key(content_hash);
    let existing = self.get_entry(content_hash)?
        .ok_or_else(|| StorageError::InvalidData(
            "Cannot update location: dedup entry not found".to_string()
        ))?;
    let updated = DedupEntry {
        volume_id: new_volume_id,
        offset: new_offset,
        ref_count: existing.ref_count,
    };
    let value = bincode::serialize(&updated)
        .map_err(|e| StorageError::Serialization(e.to_string()))?;
    Ok(BatchOp {
        keyspace: KeyspaceId::Dedup,
        action: BatchAction::Put(pg_key, value),
    })
}

Step 4: Run tests to verify they pass

Run: cargo test -p s4-core -- test_iter_entries test_update_location Expected: PASS

Step 5: Commit

git add s4-core/src/storage/dedup.rs
git commit -m "feat(dedup): add iter_entries and update_location for compaction"

Task 2: Add Index Location Scan to s4-core

Files: - Modify: s4-core/src/storage/index.rs

Step 1: Write the failing test

Add to the existing tests module in index.rs:

#[tokio::test]
async fn test_scan_objects_by_volume() {
    let temp = tempfile::TempDir::new().unwrap();
    let db = IndexDb::new(temp.path()).unwrap();

    // Insert records in different volumes
    let mut r1 = IndexRecord::new(0, 0, 100, [1u8; 32], "e1".into(), "text/plain".into());
    let mut r2 = IndexRecord::new(0, 200, 100, [2u8; 32], "e2".into(), "text/plain".into());
    let mut r3 = IndexRecord::new(1, 0, 100, [3u8; 32], "e3".into(), "text/plain".into());

    db.put("bucket/key1", &r1).await.unwrap();
    db.put("bucket/key2", &r2).await.unwrap();
    db.put("bucket/key3", &r3).await.unwrap();

    let vol0_records = db.scan_objects_by_volume(0).await.unwrap();
    assert_eq!(vol0_records.len(), 2);

    let vol1_records = db.scan_objects_by_volume(1).await.unwrap();
    assert_eq!(vol1_records.len(), 1);
}

Step 2: Run test to verify it fails

Run: cargo test -p s4-core -- test_scan_objects_by_volume Expected: FAIL

Step 3: Write minimal implementation

Add to IndexDb impl block:

/// Scans the objects keyspace and returns all records stored in a given volume.
///
/// Used by the compactor to find IndexRecords that need updating after
/// blobs are relocated to a new volume.
///
/// # Returns
///
/// Vec of (key, IndexRecord) for records where `record.file_id == volume_id`.
pub async fn scan_objects_by_volume(
    &self,
    volume_id: u32,
) -> Result<Vec<(String, IndexRecord)>, StorageError> {
    let objects = self.objects.clone();

    task::spawn_blocking(move || {
        let mut results = Vec::new();
        for guard in objects.iter() {
            let (key_bytes, value_bytes) =
                guard.into_inner().map_err(|e| StorageError::Database(e.to_string()))?;
            let key = String::from_utf8(key_bytes.to_vec())
                .map_err(|e| StorageError::InvalidData(e.to_string()))?;
            let record: IndexRecord = bincode::deserialize(&value_bytes)
                .map_err(|e| StorageError::Serialization(e.to_string()))?;
            if record.file_id == volume_id {
                results.push((key, record));
            }
        }
        Ok(results)
    })
    .await
    .map_err(|e| StorageError::Database(e.to_string()))?
}

Step 4: Run test to verify it passes

Run: cargo test -p s4-core -- test_scan_objects_by_volume Expected: PASS

Step 5: Commit

git add s4-core/src/storage/index.rs
git commit -m "feat(index): add scan_objects_by_volume for compaction"

Task 3: Implement Core Compaction Logic in s4-compactor

Files: - Rewrite: s4-compactor/src/compactor.rs - Create: s4-compactor/src/lib.rs - Modify: s4-compactor/Cargo.toml

Step 1: Update Cargo.toml

Add crc32fast to dependencies:

[dependencies]
s4-core = { path = "../s4-core" }
tokio = { workspace = true }
serde = { workspace = true }
config = { workspace = true }
thiserror = { workspace = true }
tracing = { workspace = true }
tracing-subscriber = { workspace = true }
anyhow = { workspace = true }
crc32fast = { workspace = true }
bincode = { workspace = true }

Check if crc32fast and bincode are already workspace deps — if not, add to root Cargo.toml.

Step 2: Create s4-compactor/src/lib.rs

// Copyright 2026 S4Core Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// ...Apache 2.0 header...

//! S4 Volume Compactor — background garbage collection for append-only volumes.
//!
//! This crate provides the compaction logic used by s4-server's background
//! worker. It can also be used as a standalone CLI tool for maintenance.
//!
//! # How It Works
//!
//! 1. **Scan**: Read each volume file sequentially, deserializing blob headers.
//! 2. **Classify**: For each blob, check the dedup keyspace — if no entry
//!    points to (volume_id, offset), the blob is dead.
//! 3. **Plan**: Calculate fragmentation ratio per volume. Volumes above the
//!    threshold are candidates for compaction.
//! 4. **Compact**: Copy live blobs to a new volume via VolumeWriter.
//! 5. **Update**: Atomically update IndexRecords and DedupEntries to new locations.
//! 6. **Delete**: Remove old volume file after all references are updated.

pub mod compactor;
pub mod scrubber;

pub use compactor::{CompactionConfig, CompactionResult, CompactionStats, VolumeCompactor};
pub use scrubber::{ScrubResult, VolumeScrubber};

Step 3: Rewrite s4-compactor/src/compactor.rs

// Copyright 2026 S4Core Team
// ...Apache 2.0 header...

//! Volume compaction — reclaims dead space from append-only volume files.

use s4_core::error::StorageError;
use s4_core::storage::{
    BatchAction, BatchOp, Deduplicator, IndexDb, KeyspaceId, VolumeReader, VolumeWriter,
};
use s4_core::types::{BlobHeader, IndexRecord};
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use tokio::sync::RwLock;
use std::sync::Arc;
use tracing::{debug, info, warn};

/// Configuration for the compaction process.
#[derive(Debug, Clone)]
pub struct CompactionConfig {
    /// Minimum fragmentation ratio (0.0–1.0) to trigger compaction.
    /// Default: 0.3 (30% dead space).
    pub fragmentation_threshold: f64,
    /// Minimum dead bytes in a volume to consider it for compaction.
    /// Default: 10MB. Prevents compacting nearly-empty volumes.
    pub min_dead_bytes: u64,
    /// Maximum number of volumes to compact in a single run.
    /// Default: 10. Prevents long-running compaction.
    pub max_volumes_per_run: usize,
    /// Maximum volume size (bytes) for new compacted volumes.
    pub max_volume_size: u64,
    /// Whether to perform a dry run (report only, don't compact).
    pub dry_run: bool,
}

impl Default for CompactionConfig {
    fn default() -> Self {
        Self {
            fragmentation_threshold: 0.3,
            min_dead_bytes: 10 * 1024 * 1024, // 10MB
            max_volumes_per_run: 10,
            max_volume_size: 1024 * 1024 * 1024, // 1GB
            dry_run: false,
        }
    }
}

/// Result of compacting a single volume.
#[derive(Debug, Clone)]
pub struct CompactionResult {
    /// Volume ID that was compacted.
    pub source_volume_id: u32,
    /// Number of live blobs copied.
    pub live_blobs_copied: u64,
    /// Number of dead blobs skipped.
    pub dead_blobs_skipped: u64,
    /// Bytes reclaimed (dead space freed).
    pub bytes_reclaimed: u64,
    /// New volume ID where live blobs were written.
    pub target_volume_id: Option<u32>,
}

/// Aggregate statistics for a compaction run.
#[derive(Debug, Clone, Default)]
pub struct CompactionStats {
    /// Total volumes scanned.
    pub volumes_scanned: u64,
    /// Volumes that were compacted.
    pub volumes_compacted: u64,
    /// Volumes skipped (below threshold).
    pub volumes_skipped: u64,
    /// Total live blobs copied.
    pub total_live_blobs: u64,
    /// Total dead blobs removed.
    pub total_dead_blobs: u64,
    /// Total bytes reclaimed.
    pub total_bytes_reclaimed: u64,
    /// Errors encountered (non-fatal).
    pub errors: u64,
}

/// Per-volume fragmentation analysis.
#[derive(Debug, Clone)]
pub struct VolumeAnalysis {
    /// Volume ID.
    pub volume_id: u32,
    /// Total bytes used by blobs (live + dead).
    pub total_bytes: u64,
    /// Bytes used by live blobs.
    pub live_bytes: u64,
    /// Bytes used by dead blobs.
    pub dead_bytes: u64,
    /// Number of live blobs.
    pub live_count: u64,
    /// Number of dead blobs.
    pub dead_count: u64,
    /// Fragmentation ratio (dead_bytes / total_bytes).
    pub fragmentation: f64,
}

/// The volume compactor.
///
/// Holds references to the storage components needed for compaction.
/// Designed to be called from a background worker or CLI tool.
pub struct VolumeCompactor {
    volumes_dir: PathBuf,
    index_db: Arc<IndexDb>,
    deduplicator: Deduplicator,
    volume_writer: Arc<RwLock<VolumeWriter>>,
    config: CompactionConfig,
}

impl VolumeCompactor {
    /// Creates a new compactor.
    pub fn new(
        volumes_dir: PathBuf,
        index_db: Arc<IndexDb>,
        deduplicator: Deduplicator,
        volume_writer: Arc<RwLock<VolumeWriter>>,
        config: CompactionConfig,
    ) -> Self {
        Self {
            volumes_dir,
            index_db,
            deduplicator,
            volume_writer,
            config,
        }
    }

    /// Runs a full compaction cycle.
    ///
    /// 1. Discovers all volume files
    /// 2. Analyzes fragmentation for each
    /// 3. Compacts volumes above the threshold
    ///
    /// Returns aggregate statistics.
    pub async fn run(&self) -> Result<CompactionStats, StorageError> {
        let mut stats = CompactionStats::default();

        // Step 1: Discover volume files
        let volume_ids = self.discover_volumes().await?;
        info!("Compactor: discovered {} volume files", volume_ids.len());

        // Step 2: Build dedup location index (volume_id, offset) -> content_hash
        // This avoids O(n) dedup scan per blob
        let dedup_index = self.build_dedup_index()?;
        info!("Compactor: built dedup index with {} entries", dedup_index.len());

        // Step 3: Get current writer volume to skip it (it's being written to)
        let current_volume_id = {
            let writer = self.volume_writer.read().await;
            writer.current_volume_id()
        };

        // Step 4: Analyze and compact
        let mut compacted = 0u64;
        for &vol_id in &volume_ids {
            if vol_id == current_volume_id {
                debug!("Compactor: skipping active volume {}", vol_id);
                continue;
            }

            stats.volumes_scanned += 1;

            match self.analyze_volume(vol_id, &dedup_index).await {
                Ok(analysis) => {
                    if analysis.fragmentation >= self.config.fragmentation_threshold
                        && analysis.dead_bytes >= self.config.min_dead_bytes
                    {
                        if self.config.dry_run {
                            info!(
                                "Compactor [DRY-RUN]: would compact volume {} \
                                 (frag: {:.1}%, dead: {} bytes, live: {} blobs)",
                                vol_id,
                                analysis.fragmentation * 100.0,
                                analysis.dead_bytes,
                                analysis.live_count,
                            );
                            stats.volumes_skipped += 1;
                        } else {
                            match self.compact_volume(vol_id, &dedup_index).await {
                                Ok(result) => {
                                    info!(
                                        "Compactor: compacted volume {} -> {}: \
                                         {} live blobs copied, {} dead skipped, \
                                         {} bytes reclaimed",
                                        vol_id,
                                        result.target_volume_id.unwrap_or(0),
                                        result.live_blobs_copied,
                                        result.dead_blobs_skipped,
                                        result.bytes_reclaimed,
                                    );
                                    stats.total_live_blobs += result.live_blobs_copied;
                                    stats.total_dead_blobs += result.dead_blobs_skipped;
                                    stats.total_bytes_reclaimed += result.bytes_reclaimed;
                                    stats.volumes_compacted += 1;
                                    compacted += 1;
                                }
                                Err(e) => {
                                    warn!("Compactor: failed to compact volume {}: {}", vol_id, e);
                                    stats.errors += 1;
                                }
                            }
                        }
                    } else {
                        debug!(
                            "Compactor: skipping volume {} (frag: {:.1}%, dead: {} bytes)",
                            vol_id,
                            analysis.fragmentation * 100.0,
                            analysis.dead_bytes,
                        );
                        stats.volumes_skipped += 1;
                    }
                }
                Err(e) => {
                    warn!("Compactor: failed to analyze volume {}: {}", vol_id, e);
                    stats.errors += 1;
                }
            }

            if compacted >= self.config.max_volumes_per_run as u64 {
                info!("Compactor: reached max volumes per run ({})", self.config.max_volumes_per_run);
                break;
            }
        }

        Ok(stats)
    }

    /// Discovers all volume files in the volumes directory.
    async fn discover_volumes(&self) -> Result<Vec<u32>, StorageError> {
        let mut volume_ids = Vec::new();
        let mut entries = tokio::fs::read_dir(&self.volumes_dir).await?;

        while let Some(entry) = entries.next_entry().await? {
            let name = entry.file_name();
            let name_str = name.to_string_lossy();
            if let Some(id_str) = name_str.strip_prefix("volume_").and_then(|s| s.strip_suffix(".dat")) {
                if let Ok(id) = id_str.parse::<u32>() {
                    volume_ids.push(id);
                }
            }
        }

        volume_ids.sort();
        Ok(volume_ids)
    }

    /// Builds an in-memory index of (volume_id, offset) -> content_hash
    /// from the dedup keyspace for fast blob liveness checks.
    fn build_dedup_index(&self) -> Result<HashMap<(u32, u64), [u8; 32]>, StorageError> {
        let entries = self.deduplicator.iter_entries()?;
        let mut index = HashMap::with_capacity(entries.len());
        for (hash, entry) in entries {
            index.insert((entry.volume_id, entry.offset), hash);
        }
        Ok(index)
    }

    /// Analyzes a single volume for fragmentation.
    async fn analyze_volume(
        &self,
        volume_id: u32,
        dedup_index: &HashMap<(u32, u64), [u8; 32]>,
    ) -> Result<VolumeAnalysis, StorageError> {
        let reader = VolumeReader::new(&self.volumes_dir);
        let volume_path = self.volumes_dir.join(format!("volume_{:06}.dat", volume_id));
        let file_size = tokio::fs::metadata(&volume_path).await?.len();

        let mut offset = 0u64;
        let mut live_bytes = 0u64;
        let mut dead_bytes = 0u64;
        let mut live_count = 0u64;
        let mut dead_count = 0u64;

        while offset < file_size {
            match reader.read_blob(volume_id, offset).await {
                Ok((header, key, data)) => {
                    let header_size = header.serialized_size()
                        .map_err(|e| StorageError::Serialization(e.to_string()))? as u64;
                    let blob_total = header_size + header.key_len as u64 + header.blob_len;

                    if dedup_index.contains_key(&(volume_id, offset)) {
                        live_bytes += blob_total;
                        live_count += 1;
                    } else {
                        dead_bytes += blob_total;
                        dead_count += 1;
                    }

                    offset += blob_total;
                }
                Err(_) => break, // End of readable data
            }
        }

        let total_bytes = live_bytes + dead_bytes;
        let fragmentation = if total_bytes > 0 {
            dead_bytes as f64 / total_bytes as f64
        } else {
            0.0
        };

        Ok(VolumeAnalysis {
            volume_id,
            total_bytes,
            live_bytes,
            dead_bytes,
            live_count,
            dead_count,
            fragmentation,
        })
    }

    /// Compacts a single volume: copies live blobs to a new volume,
    /// updates index and dedup, then deletes the old volume.
    async fn compact_volume(
        &self,
        volume_id: u32,
        dedup_index: &HashMap<(u32, u64), [u8; 32]>,
    ) -> Result<CompactionResult, StorageError> {
        let reader = VolumeReader::new(&self.volumes_dir);
        let volume_path = self.volumes_dir.join(format!("volume_{:06}.dat", volume_id));
        let file_size = tokio::fs::metadata(&volume_path).await?.len();

        let mut live_blobs_copied = 0u64;
        let mut dead_blobs_skipped = 0u64;
        let mut bytes_reclaimed = 0u64;
        let mut target_volume_id = None;

        // Collect relocations: (old_offset, content_hash, key, new_volume_id, new_offset)
        let mut relocations: Vec<(u64, [u8; 32], String, u32, u64)> = Vec::new();

        let mut offset = 0u64;
        while offset < file_size {
            match reader.read_blob(volume_id, offset).await {
                Ok((header, key, data)) => {
                    let header_size = header.serialized_size()
                        .map_err(|e| StorageError::Serialization(e.to_string()))? as u64;
                    let blob_total = header_size + header.key_len as u64 + header.blob_len;

                    if let Some(&content_hash) = dedup_index.get(&(volume_id, offset)) {
                        // Live blob — copy to new volume via the shared writer
                        let (new_vol_id, new_offset) = {
                            let mut writer = self.volume_writer.write().await;
                            writer.write_blob(&header, &key, &data).await?
                        };

                        target_volume_id = Some(new_vol_id);
                        relocations.push((offset, content_hash, key, new_vol_id, new_offset));
                        live_blobs_copied += 1;
                    } else {
                        dead_blobs_skipped += 1;
                        bytes_reclaimed += blob_total;
                    }

                    offset += blob_total;
                }
                Err(_) => break,
            }
        }

        // Atomically update all IndexRecords and DedupEntries in a single batch
        if !relocations.is_empty() {
            let mut ops = Vec::with_capacity(relocations.len() * 2);

            for (old_offset, content_hash, key, new_vol_id, new_offset) in &relocations {
                // Update DedupEntry location
                let dedup_op = self.deduplicator.make_update_location_op(
                    &content_hash,
                    *new_vol_id,
                    *new_offset,
                )?;
                ops.push(dedup_op);

                // Find and update all IndexRecords pointing to old location
                let records = self.index_db.scan_objects_by_volume(volume_id).await?;
                for (rec_key, mut record) in records {
                    if record.offset == *old_offset && record.file_id == volume_id {
                        record.file_id = *new_vol_id;
                        record.offset = *new_offset;
                        let value = bincode::serialize(&record)
                            .map_err(|e| StorageError::Serialization(e.to_string()))?;
                        ops.push(BatchOp {
                            keyspace: KeyspaceId::Objects,
                            action: BatchAction::Put(rec_key.into_bytes(), value),
                        });
                    }
                }
            }

            // Single atomic commit for all relocations
            self.index_db.batch_write(ops).await?;

            // Sync the volume writer to ensure new blobs are durable
            {
                let mut writer = self.volume_writer.write().await;
                writer.sync().await?;
            }
        }

        // Delete old volume file ONLY after all updates are committed and synced
        if live_blobs_copied > 0 || dead_blobs_skipped > 0 {
            // Rename to .compacted first (recoverable if crash between rename and delete)
            let compacted_path = volume_path.with_extension("dat.compacted");
            tokio::fs::rename(&volume_path, &compacted_path).await?;
            tokio::fs::remove_file(&compacted_path).await?;
            info!("Compactor: deleted old volume {}", volume_id);
        }

        Ok(CompactionResult {
            source_volume_id: volume_id,
            live_blobs_copied,
            dead_blobs_skipped,
            bytes_reclaimed,
            target_volume_id,
        })
    }
}

Step 4: Run cargo check

Run: cargo check -p s4-compactor Expected: PASS

Step 5: Commit

git add s4-compactor/
git commit -m "feat(compactor): implement volume compaction engine"

Task 4: Implement Volume Scrubber

Files: - Rewrite: s4-compactor/src/scrubber.rs

Step 1: Write the scrubber

// Copyright 2026 S4Core Team
// ...Apache 2.0 header...

//! Volume scrubber — CRC32 integrity verification for stored blobs.
//!
//! Scans volume files and verifies each blob's CRC32 checksum against its
//! stored data. Reports any corruption detected.

use s4_core::error::StorageError;
use s4_core::storage::VolumeReader;
use std::path::{Path, PathBuf};
use tracing::{debug, info, warn};

/// Result of scrubbing a single blob.
#[derive(Debug, Clone)]
pub struct ScrubError {
    /// Volume ID where corruption was found.
    pub volume_id: u32,
    /// Offset of the corrupted blob.
    pub offset: u64,
    /// Object key.
    pub key: String,
    /// Expected CRC32 (from header).
    pub expected_crc: u32,
    /// Actual CRC32 (computed from data).
    pub actual_crc: u32,
}

/// Aggregate scrub results.
#[derive(Debug, Clone, Default)]
pub struct ScrubResult {
    /// Total blobs checked.
    pub blobs_checked: u64,
    /// Blobs with valid CRC.
    pub blobs_ok: u64,
    /// Blobs with CRC mismatch.
    pub blobs_corrupted: u64,
    /// Volume files scanned.
    pub volumes_scanned: u64,
    /// Details of corrupted blobs.
    pub errors: Vec<ScrubError>,
}

/// Volume integrity scrubber.
pub struct VolumeScrubber {
    volumes_dir: PathBuf,
}

impl VolumeScrubber {
    /// Creates a new scrubber.
    pub fn new(volumes_dir: impl Into<PathBuf>) -> Self {
        Self { volumes_dir: volumes_dir.into() }
    }

    /// Scrubs all volume files for CRC32 integrity.
    pub async fn scrub_all(&self) -> Result<ScrubResult, StorageError> {
        let mut result = ScrubResult::default();
        let mut entries = tokio::fs::read_dir(&self.volumes_dir).await?;

        let mut volume_ids = Vec::new();
        while let Some(entry) = entries.next_entry().await? {
            let name = entry.file_name();
            let name_str = name.to_string_lossy();
            if let Some(id_str) = name_str.strip_prefix("volume_").and_then(|s| s.strip_suffix(".dat")) {
                if let Ok(id) = id_str.parse::<u32>() {
                    volume_ids.push(id);
                }
            }
        }
        volume_ids.sort();

        for vol_id in volume_ids {
            match self.scrub_volume(vol_id).await {
                Ok(vol_result) => {
                    result.blobs_checked += vol_result.blobs_checked;
                    result.blobs_ok += vol_result.blobs_ok;
                    result.blobs_corrupted += vol_result.blobs_corrupted;
                    result.errors.extend(vol_result.errors);
                    result.volumes_scanned += 1;
                }
                Err(e) => {
                    warn!("Scrubber: failed to scrub volume {}: {}", vol_id, e);
                }
            }
        }

        Ok(result)
    }

    /// Scrubs a single volume file.
    pub async fn scrub_volume(&self, volume_id: u32) -> Result<ScrubResult, StorageError> {
        let reader = VolumeReader::new(&self.volumes_dir);
        let volume_path = self.volumes_dir.join(format!("volume_{:06}.dat", volume_id));
        let file_size = tokio::fs::metadata(&volume_path).await?.len();

        let mut result = ScrubResult::default();
        let mut offset = 0u64;

        while offset < file_size {
            match reader.read_blob(volume_id, offset).await {
                Ok((header, key, data)) => {
                    result.blobs_checked += 1;

                    let actual_crc = crc32fast::hash(&data);
                    if actual_crc == header.crc {
                        result.blobs_ok += 1;
                    } else {
                        result.blobs_corrupted += 1;
                        result.errors.push(ScrubError {
                            volume_id,
                            offset,
                            key,
                            expected_crc: header.crc,
                            actual_crc,
                        });
                    }

                    let header_size = header.serialized_size()
                        .map_err(|e| StorageError::Serialization(e.to_string()))? as u64;
                    offset += header_size + header.key_len as u64 + header.blob_len;
                }
                Err(_) => break,
            }
        }

        result.volumes_scanned = 1;
        Ok(result)
    }
}

Step 2: Run cargo check

Run: cargo check -p s4-compactor Expected: PASS

Step 3: Commit

git add s4-compactor/src/scrubber.rs
git commit -m "feat(scrubber): implement CRC32 volume integrity verification"

Task 5: Integrate Compactor as Background Worker in s4-server

Files: - Modify: s4-server/Cargo.toml — add s4-compactor dependency - Create: s4-server/src/compaction_worker.rs - Modify: s4-server/src/config.rs — add CompactionConfig - Modify: s4-server/src/app.rs — spawn compaction worker - Modify: s4-server/src/lib.rs — export module

Step 1: Update s4-server/Cargo.toml

Add s4-compactor = { path = "../s4-compactor" } to dependencies.

Step 2: Add CompactionConfig to config.rs

/// Volume compaction configuration.
///
/// The compactor runs as a background task and periodically reclaims
/// dead space from append-only volume files.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CompactionConfig {
    /// Enable compaction worker.
    /// Can be set via S4_COMPACTION_ENABLED environment variable.
    pub enabled: bool,
    /// Compaction check interval in hours (default: 6).
    /// Can be set via S4_COMPACTION_INTERVAL_HOURS environment variable.
    pub interval_hours: u64,
    /// Minimum fragmentation ratio to trigger compaction (0.0-1.0, default: 0.3).
    /// Can be set via S4_COMPACTION_THRESHOLD environment variable.
    pub fragmentation_threshold: f64,
    /// Dry-run mode - analyze without compacting.
    /// Can be set via S4_COMPACTION_DRY_RUN environment variable.
    pub dry_run: bool,
}

Add Default impl reading from env vars (same pattern as LifecycleConfig).

Add pub compaction: CompactionConfig field to Config struct.

Step 3: Create compaction_worker.rs

Follow the exact same pattern as lifecycle_worker.rs:

// Copyright 2026 S4Core Team
// ...Apache 2.0 header...

//! Background volume compaction worker.
//!
//! Runs periodically to scan volumes for dead space and reclaim it
//! by copying live blobs to new volumes and deleting old ones.

use crate::config::CompactionConfig;
use s4_compactor::{CompactionConfig as CoreCompactionConfig, VolumeCompactor};
use s4_core::storage::BitcaskStorageEngine;
use std::sync::Arc;
use tokio::sync::RwLock;
use tokio::time::{interval, Duration};
use tracing::{error, info};

/// Background compaction worker.
pub struct CompactionWorker {
    storage: Arc<RwLock<BitcaskStorageEngine>>,
    config: CompactionConfig,
}

impl CompactionWorker {
    pub fn new(storage: Arc<RwLock<BitcaskStorageEngine>>, config: CompactionConfig) -> Self {
        Self { storage, config }
    }

    pub fn spawn(self) -> tokio::task::JoinHandle<()> {
        tokio::spawn(async move { self.run_loop().await })
    }

    async fn run_loop(&self) {
        let interval_duration = Duration::from_secs(self.config.interval_hours * 3600);
        let mut timer = interval(interval_duration);

        info!(
            "Compaction worker started (interval: {} hours, threshold: {:.0}%, dry_run: {})",
            self.config.interval_hours,
            self.config.fragmentation_threshold * 100.0,
            self.config.dry_run,
        );

        // Skip first tick
        timer.tick().await;

        loop {
            timer.tick().await;
            info!("Starting volume compaction cycle");

            let start = std::time::Instant::now();
            // NOTE: The actual compaction needs access to internal engine fields
            // (volumes_dir, index_db, deduplicator, volume_writer).
            // We need to expose these from BitcaskStorageEngine or provide
            // a run_compaction() method on the engine itself.
            // See Task 6 for adding this method.

            let storage = self.storage.read().await;
            match storage.run_compaction(self.config.fragmentation_threshold, self.config.dry_run).await {
                Ok(stats) => {
                    info!(
                        "Compaction completed in {:?}: {} volumes compacted, \
                         {} bytes reclaimed, {} errors",
                        start.elapsed(),
                        stats.volumes_compacted,
                        stats.total_bytes_reclaimed,
                        stats.errors,
                    );
                }
                Err(e) => {
                    error!("Compaction failed: {:?}", e);
                }
            }
        }
    }
}

Step 4: Update app.rs

Add compaction worker spawn next to lifecycle worker (same pattern).

Step 5: Commit

git add s4-server/
git commit -m "feat(server): add background compaction worker"

Task 6: Expose run_compaction() on BitcaskStorageEngine

Files: - Modify: s4-core/src/storage/bitcask.rs

Add a public method that creates a VolumeCompactor from the engine's internal state and runs it:

/// Runs a volume compaction cycle.
///
/// Scans all volumes for dead space and compacts those above the
/// fragmentation threshold. This method is called by the background
/// compaction worker in s4-server.
///
/// # Arguments
///
/// * `fragmentation_threshold` - Minimum dead ratio (0.0–1.0) to trigger compaction
/// * `dry_run` - If true, analyze only (don't compact)
///
/// # Returns
///
/// Aggregate statistics from the compaction run.
pub async fn run_compaction(
    &self,
    fragmentation_threshold: f64,
    dry_run: bool,
) -> Result<s4_compactor::CompactionStats, StorageError> {
    let config = s4_compactor::CompactionConfig {
        fragmentation_threshold,
        dry_run,
        ..Default::default()
    };
    let compactor = s4_compactor::VolumeCompactor::new(
        self.volumes_dir.clone(),
        self.index_db.clone(),
        self.deduplicator.clone(),
        self.volume_writer.clone(),
        config,
    );
    compactor.run().await
}

Also add s4-compactor to s4-core's Cargo.toml... Actually NO — that creates a circular dependency. Instead, keep compaction logic in s4-compactor and expose internal fields from BitcaskStorageEngine:

/// Returns the volumes directory path.
pub fn volumes_dir(&self) -> &Path { &self.volumes_dir }

/// Returns the index database (shared Arc).
pub fn index_db(&self) -> &Arc<IndexDb> { &self.index_db }

/// Returns the deduplicator.
pub fn deduplicator(&self) -> &Deduplicator { &self.deduplicator }

/// Returns the volume writer (shared Arc<RwLock>).
pub fn volume_writer(&self) -> &Arc<RwLock<VolumeWriter>> { &self.volume_writer }

Then the compaction worker constructs VolumeCompactor directly from these accessors, avoiding circular deps.

Commit:

git add s4-core/src/storage/bitcask.rs
git commit -m "feat(bitcask): expose internal fields for compaction worker"

Task 7: Unit Tests for Compactor

Files: - Create: s4-compactor/tests/compaction_integration.rs

Write integration tests that: 1. Create a temp BitcaskStorageEngine 2. Write objects, delete some (creating dead space) 3. Run compaction 4. Verify: all live objects still readable with correct data 5. Verify: dead space reclaimed (volume file deleted) 6. Verify: IndexRecords point to new locations 7. Verify: DedupEntries updated

Commit:

git add s4-compactor/tests/
git commit -m "test(compactor): add integration tests for volume compaction"

Task 8: E2E Test via Storage Engine

Files: - Create: s4-core/tests/compaction_e2e.rs

Full end-to-end test: 1. Create engine, write 100 objects across multiple volumes 2. Delete 60% of objects 3. Verify fragmentation is detected 4. Run compaction 5. Verify all remaining objects return correct SHA-256 hashes 6. Verify old volumes are gone 7. Restart engine (simulate crash recovery), verify all data intact

Commit:

git add s4-core/tests/compaction_e2e.rs
git commit -m "test(e2e): add compaction end-to-end test with crash recovery"

Task 9: Update Documentation

Files: - Update: ARCHITECTURE.md — add Compaction section - Update: README.md — add compaction env vars to table - Update: AGENTS.md — mention compaction in architecture - Create: docs/04-features/compaction.md — full feature doc

Key additions: - S4_COMPACTION_ENABLED (default: true) - S4_COMPACTION_INTERVAL_HOURS (default: 6) - S4_COMPACTION_THRESHOLD (default: 0.3) - S4_COMPACTION_DRY_RUN (default: false)

Commit:

git add ARCHITECTURE.md README.md AGENTS.md docs/
git commit -m "docs: add volume compaction documentation"

Task 10: Final Verification

Step 1: Run full test suite:

cargo fmt --check
cargo clippy --workspace --all-targets -- -D warnings
cargo test --workspace

Step 2: Manual smoke test:

cargo build --release
# Start server, upload objects, delete some, wait for compaction cycle
# Or trigger manually via lower interval

Step 3: Commit any fixes.


Decision Log

Decision Alternatives Rationale
Integrated background worker (not separate binary) Standalone s4-compactor binary Shares engine state; no fjall lock conflict; simpler deployment
s4-compactor stays as library crate Merge into s4-core Separation of concerns; s4-core stays focused on storage engine
Build in-memory dedup index for scan Query dedup per-blob during scan O(1) lookup per blob vs O(log n); compaction scans millions of blobs
Rename .dat to .dat.compacted before delete Direct delete Crash-safe: if crash between rename and delete, old volume still discoverable
Atomic batch for all relocations per volume Per-blob atomic update Fewer fjall persists; all-or-nothing consistency for entire volume
Skip active (current) volume Compact all volumes Current volume is being written to; compacting it requires extra coordination
Fragmentation threshold + min dead bytes Just threshold Prevents wasting I/O compacting small volumes with few dead bytes