"""Mass Storage Vault — auto-resizing storage for skills, memory, and projects. The vault automatically manages disk space: - Monitors available disk space - Auto-resizes: when storage grows, it checks disk space and cleans up - Compresses old/cold data to save space - Tracks storage usage across all components - Provides a unified storage API for all system components Storage tiers: - Hot: frequently accessed data (SQLite DBs, active skills) - Warm: occasionally accessed (episodic memory, completed goals) - Cold: rarely accessed (old conversations, completed projects) — compressed Auto-resize strategy: 1. Check disk space before writing 2. If disk is > 80% full, trigger cleanup: a. Compress cold data (gzip old entries) b. Archive completed projects c. Prune low-value cache entries d. Vacuum SQLite databases 3. If still > 90% full, escalate: a. Delete expired entries b. Compress warm data to cold c. Reduce cache sizes """ from __future__ import annotations import gzip import json import logging import os import shutil import sqlite3 import time from typing import Any logger = logging.getLogger(__name__) class StorageVault: """Auto-resizing mass storage vault. Manages all on-disk storage for the LLM system: - SQLite databases (memory, goals, cache, links) - Skill storage - Model files - Agent state - Project artifacts Automatically monitors disk space and cleans up when needed. Compresses cold data to save space. Grows dynamically. """ DISK_WARNING_THRESHOLD = 0.80 # 80% disk usage DISK_CRITICAL_THRESHOLD = 0.90 # 90% disk usage CLEANUP_INTERVAL_S = 300.0 # check every 5 minutes COLD_DATA_AGE_DAYS = 7 # data older than 7 days → cold (compressed) ARCHIVE_AGE_DAYS = 30 # data older than 30 days → archived def __init__(self, data_dir: str) -> None: self.data_dir = data_dir os.makedirs(data_dir, exist_ok=True) # Storage subdirectories self.db_dir = os.path.join(data_dir, "db") self.cache_dir = os.path.join(data_dir, "cache") self.archive_dir = os.path.join(data_dir, "archive") self.compressed_dir = os.path.join(data_dir, "compressed") self.artifacts_dir = os.path.join(data_dir, "artifacts") for d in [self.db_dir, self.cache_dir, self.archive_dir, self.compressed_dir, self.artifacts_dir]: os.makedirs(d, exist_ok=True) self._last_cleanup = 0.0 self._stats = { "total_storage_bytes": 0, "db_storage_bytes": 0, "cache_storage_bytes": 0, "archive_storage_bytes": 0, "compressed_storage_bytes": 0, "artifacts_storage_bytes": 0, "disk_free_bytes": 0, "disk_total_bytes": 0, "disk_usage_percent": 0.0, "cleanups_performed": 0, "items_compressed": 0, "items_archived": 0, "items_deleted": 0, "dbs_vacuumed": 0, "auto_resize_enabled": True, } self._update_storage_stats() def get_db_path(self, name: str) -> str: """Get path for a named SQLite database.""" return os.path.join(self.db_dir, f"{name}.db") def get_cache_path(self, name: str) -> str: """Get path for a cache file.""" return os.path.join(self.cache_dir, name) def get_artifact_path(self, name: str) -> str: """Get path for a project artifact.""" return os.path.join(self.artifacts_dir, name) def store_artifact(self, name: str, data: bytes) -> str: """Store a project artifact (code, images, etc.).""" path = self.get_artifact_path(name) self._check_and_cleanup() with open(path, "wb") as f: f.write(data) self._update_storage_stats() return path def store_artifact_text(self, name: str, text: str) -> str: """Store a text artifact.""" return self.store_artifact(name, text.encode()) def load_artifact(self, name: str) -> bytes | None: """Load an artifact.""" path = self.get_artifact_path(name) if not os.path.exists(path): return None with open(path, "rb") as f: return f.read() def list_artifacts(self) -> list[dict[str, Any]]: """List all artifacts with metadata.""" artifacts = [] if not os.path.exists(self.artifacts_dir): return artifacts for name in sorted(os.listdir(self.artifacts_dir)): path = os.path.join(self.artifacts_dir, name) if os.path.isfile(path): stat = os.stat(path) artifacts.append({ "name": name, "size_bytes": stat.st_size, "created_at": stat.st_ctime, "modified_at": stat.st_mtime, }) return artifacts def delete_artifact(self, name: str) -> bool: """Delete an artifact.""" path = self.get_artifact_path(name) if os.path.exists(path): os.remove(path) self._update_storage_stats() return True return False def _check_and_cleanup(self) -> bool: """Check disk space and cleanup if needed. Returns True if cleanup ran.""" if not self._stats["auto_resize_enabled"]: return False now = time.time() if now - self._last_cleanup < self.CLEANUP_INTERVAL_S: return False self._last_cleanup = now self._update_storage_stats() usage = self._stats["disk_usage_percent"] if usage < self.DISK_WARNING_THRESHOLD: return False logger.info("Disk usage %.1f%% — triggering cleanup", usage * 100) self._stats["cleanups_performed"] += 1 # Phase 1: Compress cold data self._compress_cold_data() # Phase 2: Archive old completed projects self._archive_old_data() # Phase 3: Vacuum SQLite databases self._vacuum_dbs() # Phase 4: If still critical, delete expired entries self._update_storage_stats() if self._stats["disk_usage_percent"] >= self.DISK_CRITICAL_THRESHOLD: self._delete_expired_data() self._update_storage_stats() logger.info("Cleanup complete — disk usage now %.1f%%", self._stats["disk_usage_percent"] * 100) return True def _compress_cold_data(self) -> None: """Compress old data files to save space.""" cutoff = time.time() - (self.COLD_DATA_AGE_DAYS * 86400) # Compress old cache files if os.path.exists(self.cache_dir): for name in os.listdir(self.cache_dir): path = os.path.join(self.cache_dir, name) if os.path.isfile(path): stat = os.stat(path) if stat.st_mtime < cutoff and not name.endswith(".gz"): self._gzip_file(path) self._stats["items_compressed"] += 1 # Compress old artifacts if os.path.exists(self.artifacts_dir): for name in os.listdir(self.artifacts_dir): path = os.path.join(self.artifacts_dir, name) if os.path.isfile(path): stat = os.stat(path) if stat.st_mtime < cutoff and not name.endswith(".gz"): self._gzip_file(path) self._stats["items_compressed"] += 1 def _archive_old_data(self) -> None: """Archive old data to the archive directory.""" cutoff = time.time() - (self.ARCHIVE_AGE_DAYS * 86400) # Move old compressed files to archive for directory in [self.cache_dir, self.compressed_dir]: if not os.path.exists(directory): continue for name in os.listdir(directory): path = os.path.join(directory, name) if os.path.isfile(path): stat = os.stat(path) if stat.st_mtime < cutoff: dest = os.path.join(self.archive_dir, name) shutil.move(path, dest) self._stats["items_archived"] += 1 def _vacuum_dbs(self) -> None: """Vacuum SQLite databases to reclaim space.""" if not os.path.exists(self.db_dir): return for name in os.listdir(self.db_dir): if not name.endswith(".db"): continue path = os.path.join(self.db_dir, name) try: with sqlite3.connect(path) as conn: conn.execute("VACUUM") self._stats["dbs_vacuumed"] += 1 logger.debug("Vacuumed: %s", name) except Exception as e: logger.debug("Vacuum failed for %s: %s", name, e) def _delete_expired_data(self) -> None: """Delete expired archive data when disk is critical.""" if not os.path.exists(self.archive_dir): return cutoff = time.time() - (self.ARCHIVE_AGE_DAYS * 2 * 86400) # 60 days for name in os.listdir(self.archive_dir): path = os.path.join(self.archive_dir, name) if os.path.isfile(path): stat = os.stat(path) if stat.st_mtime < cutoff: os.remove(path) self._stats["items_deleted"] += 1 logger.info("Deleted expired archive: %s", name) def _gzip_file(self, path: str) -> None: """Compress a file with gzip.""" gz_path = path + ".gz" try: with open(path, "rb") as f_in: with gzip.open(gz_path, "wb") as f_out: shutil.copyfileobj(f_in, f_out) os.remove(path) # remove original after compression except Exception as e: logger.debug("Compression failed for %s: %s", path, e) def _update_storage_stats(self) -> None: """Update storage statistics.""" usage = shutil.disk_usage(self.data_dir) self._stats["disk_free_bytes"] = usage.free self._stats["disk_total_bytes"] = usage.total self._stats["disk_usage_percent"] = round(1 - (usage.free / usage.total), 4) self._stats["db_storage_bytes"] = self._dir_size(self.db_dir) self._stats["cache_storage_bytes"] = self._dir_size(self.cache_dir) self._stats["archive_storage_bytes"] = self._dir_size(self.archive_dir) self._stats["compressed_storage_bytes"] = self._dir_size(self.compressed_dir) self._stats["artifacts_storage_bytes"] = self._dir_size(self.artifacts_dir) self._stats["total_storage_bytes"] = ( self._stats["db_storage_bytes"] + self._stats["cache_storage_bytes"] + self._stats["archive_storage_bytes"] + self._stats["compressed_storage_bytes"] + self._stats["artifacts_storage_bytes"] ) @staticmethod def _dir_size(path: str) -> int: """Get total size of a directory.""" if not os.path.exists(path): return 0 total = 0 for dirpath, _, filenames in os.walk(path): for f in filenames: fp = os.path.join(dirpath, f) if not os.path.islink(fp): total += os.path.getsize(fp) return total def force_cleanup(self) -> dict[str, Any]: """Force a cleanup cycle.""" self._last_cleanup = 0.0 # reset to allow immediate cleanup self._check_and_cleanup() return self.get_stats() def get_stats(self) -> dict[str, Any]: self._update_storage_stats() return { **self._stats, "total_storage_mb": round(self._stats["total_storage_bytes"] / 1e6, 2), "disk_free_gb": round(self._stats["disk_free_bytes"] / 1e9, 2), "disk_total_gb": round(self._stats["disk_total_bytes"] / 1e9, 2), "db_storage_mb": round(self._stats["db_storage_bytes"] / 1e6, 2), "cache_storage_mb": round(self._stats["cache_storage_bytes"] / 1e6, 2), "archive_storage_mb": round(self._stats["archive_storage_bytes"] / 1e6, 2), "artifacts_storage_mb": round(self._stats["artifacts_storage_bytes"] / 1e6, 2), }