From 1ff5f609bc8bd5992dddb7bdac398e536b22b9d9 Mon Sep 17 00:00:00 2001 From: Caleb Date: Wed, 23 Sep 2026 13:36:33 +0800 Subject: [PATCH 1/8] fix: harden temporary file cleanup --- core/lifecycle.py | 2 +- core/temp_monitor.py | 130 ++++++++++++++++++++++++++----------- tests/test_temp_monitor.py | 126 +++++++++++++++++++++++++++++++++++ 3 files changed, 220 insertions(+), 38 deletions(-) create mode 100644 tests/test_temp_monitor.py diff --git a/core/lifecycle.py b/core/lifecycle.py index 8294a70a..12b01461 100644 --- a/core/lifecycle.py +++ b/core/lifecycle.py @@ -246,7 +246,7 @@ async def init_and_run_system(self): self.temp_monitor = AsyncTempMonitor( folder_path=str(temp_folder), kira_config=self.kira_config, - check_interval=10, + check_interval=5 * 60, batch_size=20, ) diff --git a/core/temp_monitor.py b/core/temp_monitor.py index 69123aec..759054bd 100644 --- a/core/temp_monitor.py +++ b/core/temp_monitor.py @@ -1,9 +1,10 @@ import asyncio import heapq +import stat import time from pathlib import Path from watchfiles import awatch -from typing import Dict, List, Tuple, Optional +from typing import Dict, List, Literal, Optional, Tuple from typing import TYPE_CHECKING @@ -186,25 +187,56 @@ async def _get_files_exceeding_limit(self) -> List[Tuple[str, int, float, float] return [] return heapq.nsmallest(min(excess_count, len(eligible_files)), eligible_files, key=lambda x: x[3]) - async def _delete_file(self, path_str: str) -> Optional[int]: - """Delete a single file asynchronously""" - loop = asyncio.get_event_loop() + async def _delete_file(self, path_str: str) -> Tuple[Literal["deleted", "missing", "failed"], int, Optional[str]]: + """Delete a single file, retrying once after clearing read-only mode.""" + loop = asyncio.get_running_loop() def delete(): + file_path = Path(path_str) try: - file_path = Path(path_str) - if file_path.exists() and file_path.is_file(): - size = file_path.stat().st_size + if not file_path.exists() or not file_path.is_file(): + return "missing", 0, None + + size = file_path.stat().st_size + try: file_path.unlink() - return size - except Exception as e: - logger.error(f"Failed to delete {path_str}: {e}") - return None + except PermissionError: + # Git objects are commonly read-only on Windows. Make the + # file owner-writable, then retry exactly once. + mode = file_path.stat().st_mode + file_path.chmod(mode | stat.S_IWRITE) + file_path.unlink() + return "deleted", size, None + except FileNotFoundError: + return "missing", 0, None + except OSError as e: + return "failed", 0, f"{type(e).__name__}: {e}" return await loop.run_in_executor(None, delete) + async def _cleanup_empty_dirs(self): + """Remove empty subdirectories from deepest to shallowest.""" + loop = asyncio.get_running_loop() + + def cleanup_empty_dirs(): + if not self.folder_path.exists(): + return + + directories = ( + path for path in self.folder_path.rglob('*') if path.is_dir() + ) + for directory in sorted( + directories, key=lambda path: len(path.parts), reverse=True + ): + try: + directory.rmdir() + except OSError: + continue + + await loop.run_in_executor(None, cleanup_empty_dirs) + async def cleanup(self): - """Execute cleanup asynchronously""" + """Execute cleanup asynchronously.""" # Use lock to prevent concurrent cleanup async with self._cleanup_lock: current_time = time.time() @@ -235,6 +267,31 @@ async def cleanup(self): deleted_count = 0 freed_space = 0 + attempted_paths = set() + failed_deletions: Dict[str, str] = {} + + async def delete_candidate(path_str: str): + nonlocal deleted_count, freed_space + + # A file can qualify as expired, excessive, and oversized in + # the same cleanup. Never retry a failure within one cycle. + if path_str in attempted_paths: + return None + attempted_paths.add(path_str) + + status, deleted_size, error = await self._delete_file(path_str) + if status == "failed": + failed_deletions[path_str] = error or "unknown error" + return status + + cached = self.file_cache.pop(path_str, None) + if cached is not None: + self.total_size = max(0, self.total_size - cached[0]) + + if status == "deleted": + deleted_count += 1 + freed_space += deleted_size + return status # Phase 1: Delete expired files (by max_age_hours) if expired_files: @@ -250,15 +307,10 @@ async def cleanup(self): f"{Path(path_str).name} (age: {file_age:.1f}s < " f"protection: {min_expired_protection}s)") continue - deleted_size = await self._delete_file(path_str) - if deleted_size is not None: - self.total_size -= deleted_size - if path_str in self.file_cache: - del self.file_cache[path_str] - deleted_count += 1 - freed_space += deleted_size + status = await delete_candidate(path_str) + if status == "deleted": file_age_hours = (current_time - creation_time) / 3600 - logger.debug(f"DELETED expired: {Path(path_str).name} (age: {file_age_hours:.1f}h, size: {deleted_size / 1024:.2f}KB)") + logger.debug(f"DELETED expired: {Path(path_str).name} (age: {file_age_hours:.1f}h, size: {size / 1024:.2f}KB)") # Phase 2: Delete files exceeding max_files limit # Protection period already filtered inside _get_files_exceeding_limit @@ -266,14 +318,9 @@ async def cleanup(self): if excess_files: logger.debug(f"Found {len(excess_files)} excess files (limit: {self.max_files})") for path_str, size, mtime, creation_time in excess_files: - deleted_size = await self._delete_file(path_str) - if deleted_size is not None: - self.total_size -= deleted_size - if path_str in self.file_cache: - del self.file_cache[path_str] - deleted_count += 1 - freed_space += deleted_size - logger.debug(f"DELETED excess: {Path(path_str).name} (size: {deleted_size / 1024:.2f}KB)") + status = await delete_candidate(path_str) + if status == "deleted": + logger.debug(f"DELETED excess: {Path(path_str).name} (size: {size / 1024:.2f}KB)") # Phase 3: Delete oldest files if still over size limit if self.total_size > self.max_size_bytes: @@ -288,14 +335,11 @@ async def cleanup(self): logger.error(f"ATTEMPTED TO DELETE PROTECTED FILE (age: {file_age:.2f}s): {Path(path_str).name} - SKIPPING") continue - deleted_size = await self._delete_file(path_str) - if deleted_size is not None: - self.total_size -= deleted_size - if path_str in self.file_cache: - del self.file_cache[path_str] - deleted_count += 1 - freed_space += deleted_size - logger.debug(f"DELETED: {Path(path_str).name} (age: {file_age:.2f}s, size: {deleted_size / 1024:.2f}KB)") + status = await delete_candidate(path_str) + if status == "deleted": + logger.debug(f"DELETED: {Path(path_str).name} (age: {file_age:.2f}s, size: {size / 1024:.2f}KB)") + + await self._cleanup_empty_dirs() if deleted_count > 0: logger.info(f"Cleanup completed: deleted {deleted_count} files, " @@ -303,6 +347,19 @@ async def cleanup(self): f"remaining: {len(self.file_cache)} files, " f"{self.total_size / 1024 / 1024:.2f}MB") + if failed_deletions: + failure_items = list(failed_deletions.items()) + details = "; ".join( + f"{path}: {error}" for path, error in failure_items[:3] + ) + omitted = len(failure_items) - 3 + if omitted > 0: + details += f"; and {omitted} more" + logger.warning( + f"Cleanup completed with {len(failed_deletions)} deletion failure(s); " + f"will retry on the next cleanup cycle. {details}" + ) + async def _process_changes(self, changes): """Process file change events""" @@ -323,7 +380,6 @@ async def _periodic_cleanup_loop(self): await asyncio.sleep(self.check_interval) if self._stop_event.is_set(): break - self.last_check_time = 0 # Reset to allow cleanup await self.cleanup() except asyncio.CancelledError: break diff --git a/tests/test_temp_monitor.py b/tests/test_temp_monitor.py new file mode 100644 index 00000000..8646f898 --- /dev/null +++ b/tests/test_temp_monitor.py @@ -0,0 +1,126 @@ +import os +import stat +import time +from pathlib import Path + +import pytest + +from core import temp_monitor +from core.temp_monitor import AsyncTempMonitor + + +class FakeConfig: + def __init__(self, cache_config): + self.cache_config = cache_config + + def get_config(self, key, default=None): + if key == "bot_config.cache": + return self.cache_config + return default + + +def make_monitor(tmp_path, **cache_overrides): + cache_config = { + "max_size_mb": 50, + "max_files": 50, + "max_age_hours": 24, + } + cache_config.update(cache_overrides) + return AsyncTempMonitor( + str(tmp_path), + FakeConfig(cache_config), + check_interval=0, + file_protection_seconds=0, + ) + + +@pytest.mark.asyncio +async def test_delete_file_retries_once_after_permission_error(tmp_path, monkeypatch): + target = tmp_path / "readonly.bin" + target.write_bytes(b"content") + original_unlink = Path.unlink + original_chmod = Path.chmod + unlink_calls = 0 + chmod_modes = [] + + def flaky_unlink(path, *args, **kwargs): + nonlocal unlink_calls + if path == target: + unlink_calls += 1 + if unlink_calls == 1: + raise PermissionError("read-only") + return original_unlink(path, *args, **kwargs) + + def record_chmod(path, mode, *args, **kwargs): + if path == target: + chmod_modes.append(mode) + return original_chmod(path, mode, *args, **kwargs) + + monkeypatch.setattr(Path, "unlink", flaky_unlink) + monkeypatch.setattr(Path, "chmod", record_chmod) + + monitor = make_monitor(tmp_path) + status, deleted_size, error = await monitor._delete_file(str(target)) + + assert status == "deleted" + assert deleted_size == len(b"content") + assert error is None + assert unlink_calls == 2 + assert len(chmod_modes) == 1 + assert chmod_modes[0] & stat.S_IWRITE + assert not target.exists() + + +@pytest.mark.asyncio +async def test_cleanup_summarizes_failures_once_and_retries_next_cycle( + tmp_path, monkeypatch +): + target = tmp_path / "locked.bin" + target.write_bytes(b"content") + old_time = time.time() - 3600 + os.utime(target, (old_time, old_time)) + + monitor = make_monitor(tmp_path, max_size_mb=0, max_age_hours=0) + await monitor._build_cache() + + delete_calls = [] + + async def fail_delete(path_str): + delete_calls.append(path_str) + return "failed", 0, "PermissionError: locked" + + warnings = [] + monkeypatch.setattr(monitor, "_delete_file", fail_delete) + monkeypatch.setattr(temp_monitor.logger, "warning", warnings.append) + + await monitor.cleanup() + + assert delete_calls == [str(target)] + assert len(warnings) == 1 + assert "1 deletion failure(s)" in warnings[0] + assert "next cleanup cycle" in warnings[0] + assert str(target) in monitor.file_cache + + monitor.last_check_time = 0 + await monitor.cleanup() + + assert delete_calls == [str(target), str(target)] + assert len(warnings) == 2 + + +@pytest.mark.asyncio +async def test_cleanup_removes_empty_directories(tmp_path): + nested_dir = tmp_path / "job" / "repo" / ".git" / "objects" + nested_dir.mkdir(parents=True) + target = nested_dir / "old.bin" + target.write_bytes(b"content") + old_time = time.time() - 3600 + os.utime(target, (old_time, old_time)) + + monitor = make_monitor(tmp_path, max_age_hours=0) + await monitor._build_cache() + await monitor.cleanup() + + assert not target.exists() + assert not (tmp_path / "job").exists() + assert tmp_path.exists() From a7f90d370b0d1d8482ee9af2e7fd78df5341fddd Mon Sep 17 00:00:00 2001 From: Caleb Date: Wed, 23 Sep 2026 14:03:32 +0800 Subject: [PATCH 2/8] feat: configure temp cleanup interval --- core/config/default.py | 3 ++- core/temp_monitor.py | 36 ++++++++++++++++++++----- tests/test_temp_monitor.py | 21 +++++++++++++++ webui/frontend/src/i18n/en.ts | 2 ++ webui/frontend/src/i18n/zh.ts | 2 ++ webui/frontend/src/views/ConfigView.vue | 1 + webui/routes/config.py | 4 +++ 7 files changed, 62 insertions(+), 7 deletions(-) diff --git a/core/config/default.py b/core/config/default.py index 6cae2e36..d210350d 100644 --- a/core/config/default.py +++ b/core/config/default.py @@ -23,7 +23,8 @@ "cache": { "max_size_mb": 50, "max_files": 50, - "max_age_hours": 24 + "max_age_hours": 24, + "check_interval_minutes": 5 }, "image_compression": { "enabled": False, diff --git a/core/temp_monitor.py b/core/temp_monitor.py index 759054bd..7a751b1c 100644 --- a/core/temp_monitor.py +++ b/core/temp_monitor.py @@ -32,6 +32,7 @@ def __init__(self, folder_path: str, kira_config: 'KiraConfig', """ self.folder_path = Path(folder_path) self.kira_config = kira_config + self._default_check_interval = check_interval self.check_interval = check_interval self.batch_size = batch_size self.file_protection_seconds = file_protection_seconds @@ -47,9 +48,10 @@ def __init__(self, folder_path: str, kira_config: 'KiraConfig', self.last_check_time = 0 self._cleanup_lock = asyncio.Lock() self._stop_event = asyncio.Event() + self._config_changed_event = asyncio.Event() def _refresh_config(self): - """Read latest config values from KiraConfig to support runtime changes""" + """Read latest config values from KiraConfig to support runtime changes.""" cache_config = self.kira_config.get_config("bot_config.cache", {}) or {} max_size_mb = cache_config.get("max_size_mb", 50) self.max_size_bytes = max_size_mb * 1024 * 1024 @@ -57,6 +59,21 @@ def _refresh_config(self): max_age_hours = cache_config.get("max_age_hours", 24) self.max_age_seconds = max_age_hours * 3600 + interval_minutes = cache_config.get("check_interval_minutes") + if ( + isinstance(interval_minutes, (int, float)) + and not isinstance(interval_minutes, bool) + and interval_minutes > 0 + ): + self.check_interval = int(interval_minutes * 60) + else: + self.check_interval = self._default_check_interval + + def notify_config_changed(self): + """Apply runtime config changes and wake the periodic scheduler.""" + self._refresh_config() + self._config_changed_event.set() + async def _build_cache(self): """Build file cache asynchronously""" loop = asyncio.get_running_loop() @@ -239,6 +256,7 @@ async def cleanup(self): """Execute cleanup asynchronously.""" # Use lock to prevent concurrent cleanup async with self._cleanup_lock: + self._refresh_config() current_time = time.time() # Control check frequency (inside lock to avoid TOCTOU race) @@ -247,9 +265,6 @@ async def cleanup(self): self.last_check_time = current_time - # Refresh config to pick up runtime changes - self._refresh_config() - # Check for expired files first expired_files = await self._get_expired_files() @@ -374,10 +389,19 @@ async def _process_changes(self, changes): await self.cleanup() async def _periodic_cleanup_loop(self): - """Periodically check for expired files even without file changes""" + """Periodically check for expired files even without file changes.""" while not self._stop_event.is_set(): try: - await asyncio.sleep(self.check_interval) + try: + await asyncio.wait_for( + self._config_changed_event.wait(), + timeout=self.check_interval, + ) + self._config_changed_event.clear() + continue + except asyncio.TimeoutError: + pass + if self._stop_event.is_set(): break await self.cleanup() diff --git a/tests/test_temp_monitor.py b/tests/test_temp_monitor.py index 8646f898..824b1fe4 100644 --- a/tests/test_temp_monitor.py +++ b/tests/test_temp_monitor.py @@ -124,3 +124,24 @@ async def test_cleanup_removes_empty_directories(tmp_path): assert not target.exists() assert not (tmp_path / "job").exists() assert tmp_path.exists() + + +@pytest.mark.asyncio +async def test_configured_check_interval_updates_at_runtime(tmp_path): + config = FakeConfig( + { + "max_size_mb": 50, + "max_files": 50, + "max_age_hours": 24, + "check_interval_minutes": 5, + } + ) + monitor = AsyncTempMonitor(str(tmp_path), config, check_interval=60) + + assert monitor.check_interval == 5 * 60 + + config.cache_config["check_interval_minutes"] = 2 + monitor.notify_config_changed() + + assert monitor.check_interval == 2 * 60 + assert monitor._config_changed_event.is_set() diff --git a/webui/frontend/src/i18n/en.ts b/webui/frontend/src/i18n/en.ts index e5b57235..f54f6bd9 100644 --- a/webui/frontend/src/i18n/en.ts +++ b/webui/frontend/src/i18n/en.ts @@ -679,6 +679,7 @@ export default { max_size_mb: 'Maximum disk space allowed for the cache folder', max_files: 'Maximum number of files allowed in the cache folder', max_age_hours: 'Cache files older than this will be automatically cleaned up', + check_interval_minutes: 'How often the temporary cache is checked for cleanup. Changes take effect immediately', log_level: 'Minimum log level displayed in the terminal', log_file_path: 'Path to the log file, leave empty for default', log_file_max_size: 'Maximum size of a single log file in megabytes', @@ -742,6 +743,7 @@ export default { max_size_mb: 'Max Storage (MB)', max_files: 'Max Files', max_age_hours: 'Max Cache Age (hours)', + check_interval_minutes: 'Cleanup Check Interval (minutes)', log_level: 'Log Level', log_file_path: 'Log File Path', log_file_max_size: 'Max Log File Size (MB)', diff --git a/webui/frontend/src/i18n/zh.ts b/webui/frontend/src/i18n/zh.ts index 41da74aa..b9960177 100644 --- a/webui/frontend/src/i18n/zh.ts +++ b/webui/frontend/src/i18n/zh.ts @@ -679,6 +679,7 @@ export default { max_size_mb: '缓存文件夹允许占用的最大磁盘空间', max_files: '缓存文件夹中允许的最大文件数量', max_age_hours: '超过此时间的缓存文件将被自动清理', + check_interval_minutes: '检查临时缓存并执行清理的时间间隔,修改后立即生效', log_level: '终端显示的最低日志等级', log_file_path: '日志文件路径,留空使用默认路径', log_file_max_size: '单个日志文件的最大体积(MB)', @@ -742,6 +743,7 @@ export default { max_size_mb: '最大存储(MB)', max_files: '最大文件数', max_age_hours: '最大缓存时间(小时)', + check_interval_minutes: '清理检查间隔(分钟)', log_level: '日志等级', log_file_path: '日志文件路径', log_file_max_size: '日志文件最大体积(MB)', diff --git a/webui/frontend/src/views/ConfigView.vue b/webui/frontend/src/views/ConfigView.vue index d17ec4c6..ed57de81 100644 --- a/webui/frontend/src/views/ConfigView.vue +++ b/webui/frontend/src/views/ConfigView.vue @@ -543,6 +543,7 @@ const allGroups: ConfigGroup[] = [ { key: 'bot_config.cache.max_size_mb', labelKey: 'configuration.message.max_size_mb', labelFallback: 'Max Storage (MB)', hintKey: 'configuration.hints.max_size_mb', hintFallback: 'Maximum disk space allowed for the cache folder', type: 'integer', default: 50, validation: { min: 1, max: 10240, required: true } }, { key: 'bot_config.cache.max_files', labelKey: 'configuration.message.max_files', labelFallback: 'Max Files', hintKey: 'configuration.hints.max_files', hintFallback: 'Maximum number of files allowed in the cache folder', type: 'integer', default: 50, validation: { min: 1, max: 100000, required: true } }, { key: 'bot_config.cache.max_age_hours', labelKey: 'configuration.message.max_age_hours', labelFallback: 'Max Cache Age (hours)', hintKey: 'configuration.hints.max_age_hours', hintFallback: 'Cache files older than this will be automatically cleaned up', type: 'integer', default: 24, validation: { min: 1, max: 8760, required: true } }, + { key: 'bot_config.cache.check_interval_minutes', labelKey: 'configuration.message.check_interval_minutes', labelFallback: 'Cleanup Check Interval (minutes)', hintKey: 'configuration.hints.check_interval_minutes', hintFallback: 'How often the temporary cache is checked for files to clean up', type: 'integer', default: 5, validation: { min: 1, max: 1440, required: true } }, ], }, { diff --git a/webui/routes/config.py b/webui/routes/config.py index 7559faf0..12ee90d7 100644 --- a/webui/routes/config.py +++ b/webui/routes/config.py @@ -93,6 +93,10 @@ async def update_configuration(self, payload: Dict): updated = False if "bot_config" in payload: config["bot_config"] = bot_config + temp_monitor = getattr(self.lifecycle, "temp_monitor", None) + notify_config_changed = getattr(temp_monitor, "notify_config_changed", None) + if callable(notify_config_changed): + notify_config_changed() updated = True if "models" in payload: config["models"] = models From d270c4938ead67453cd0407b3bdf47406613c44d Mon Sep 17 00:00:00 2001 From: Caleb Date: Wed, 23 Sep 2026 14:24:14 +0800 Subject: [PATCH 3/8] fix: address temp cleanup review findings --- core/temp_monitor.py | 9 ++++++++- tests/test_temp_monitor.py | 25 +++++++++++++++++++++++++ 2 files changed, 33 insertions(+), 1 deletion(-) diff --git a/core/temp_monitor.py b/core/temp_monitor.py index 7a751b1c..113de3ee 100644 --- a/core/temp_monitor.py +++ b/core/temp_monitor.py @@ -65,7 +65,7 @@ def _refresh_config(self): and not isinstance(interval_minutes, bool) and interval_minutes > 0 ): - self.check_interval = int(interval_minutes * 60) + self.check_interval = max(1, int(interval_minutes * 60)) else: self.check_interval = self._default_check_interval @@ -239,6 +239,7 @@ def cleanup_empty_dirs(): if not self.folder_path.exists(): return + now = time.time() directories = ( path for path in self.folder_path.rglob('*') if path.is_dir() ) @@ -246,6 +247,11 @@ def cleanup_empty_dirs(): directories, key=lambda path: len(path.parts), reverse=True ): try: + if ( + self.file_protection_seconds > 0 + and now - directory.stat().st_mtime < self.file_protection_seconds + ): + continue directory.rmdir() except OSError: continue @@ -275,6 +281,7 @@ async def cleanup(self): ) if not needs_cleanup: + await self._cleanup_empty_dirs() return logger.info(f"CLEANUP TRIGGERED - Files: {len(self.file_cache)}/{self.max_files}, " diff --git a/tests/test_temp_monitor.py b/tests/test_temp_monitor.py index 824b1fe4..d148b959 100644 --- a/tests/test_temp_monitor.py +++ b/tests/test_temp_monitor.py @@ -145,3 +145,28 @@ async def test_configured_check_interval_updates_at_runtime(tmp_path): assert monitor.check_interval == 2 * 60 assert monitor._config_changed_event.is_set() + + config.cache_config["check_interval_minutes"] = 0.01 + monitor.notify_config_changed() + + assert monitor.check_interval == 1 + + +@pytest.mark.asyncio +async def test_cleanup_empty_dirs_preserves_recent_directories(tmp_path): + recent_dir = tmp_path / "recent" + old_dir = tmp_path / "old" + recent_dir.mkdir() + old_dir.mkdir() + old_time = time.time() - 120 + os.utime(old_dir, (old_time, old_time)) + + monitor = AsyncTempMonitor( + str(tmp_path), + FakeConfig({}), + file_protection_seconds=60, + ) + await monitor.cleanup() + + assert recent_dir.exists() + assert not old_dir.exists() From 08c4f80c65fecc88cbd7c71a362b9db72114a2df Mon Sep 17 00:00:00 2001 From: Caleb Date: Wed, 23 Sep 2026 14:56:41 +0800 Subject: [PATCH 4/8] refactor: rescan temp directory during cleanup --- core/temp_monitor.py | 527 +++++++++++++++++++++---------------- requirements.txt | 1 - tests/test_temp_monitor.py | 101 +++++-- 3 files changed, 377 insertions(+), 252 deletions(-) diff --git a/core/temp_monitor.py b/core/temp_monitor.py index 113de3ee..3f633963 100644 --- a/core/temp_monitor.py +++ b/core/temp_monitor.py @@ -1,12 +1,10 @@ import asyncio import heapq +import os import stat import time from pathlib import Path -from watchfiles import awatch -from typing import Dict, List, Literal, Optional, Tuple - -from typing import TYPE_CHECKING +from typing import TYPE_CHECKING, Dict, List, Literal, Optional, Tuple from core.logging_manager import get_logger @@ -15,21 +13,24 @@ logger = get_logger("atm", "yellow") +FileKey = Tuple[int, int] +FileVersion = Tuple[int, int, int, int, int] +FileCacheEntry = Tuple[int, float, float] +FileCandidate = Tuple[str, int, float, float] +DirectoryCandidate = Tuple[Path, FileKey] +DeleteStatus = Literal["deleted", "missing", "changed", "failed"] + class AsyncTempMonitor: - def __init__(self, folder_path: str, kira_config: 'KiraConfig', - check_interval: int = 60, batch_size: int = 20, - file_protection_seconds: int = 60): - """ - Asynchronous temporary folder monitor - - Args: - folder_path: Path of the folder to monitor - kira_config: KiraConfig instance for dynamic config reading - check_interval: Minimum check interval in seconds - batch_size: Maximum number of files to clean up at once - file_protection_seconds: File protection period in seconds, new files won't be deleted within this period - """ + def __init__( + self, + folder_path: str, + kira_config: 'KiraConfig', + check_interval: int = 60, + batch_size: int = 20, + file_protection_seconds: int = 60, + ): + """Initialize the periodic temporary-directory cleaner.""" self.folder_path = Path(folder_path) self.kira_config = kira_config self._default_check_interval = check_interval @@ -38,18 +39,35 @@ def __init__(self, folder_path: str, kira_config: 'KiraConfig', self.file_protection_seconds = file_protection_seconds self.folder_path.mkdir(parents=True, exist_ok=True) - # Initialize from config self._refresh_config() - # Cache file information: path -> (size, mtime, creation_time) - # creation_time is when the file was first seen by the monitor - self.file_cache: Dict[str, Tuple[int, float, float]] = {} + # The file cache is a fresh snapshot rebuilt before every cleanup. + self.file_cache: Dict[str, FileCacheEntry] = {} self.total_size = 0 - self.last_check_time = 0 + self._file_versions: Dict[str, FileVersion] = {} + self._first_seen: Dict[str, Tuple[FileKey, float]] = {} + self._pending_retries: Dict[str, FileKey] = {} + self._has_scanned = False self._cleanup_lock = asyncio.Lock() self._stop_event = asyncio.Event() self._config_changed_event = asyncio.Event() + @staticmethod + def _file_key(file_stat: os.stat_result) -> FileKey: + """Return the stable identity used to distinguish path replacements.""" + return file_stat.st_dev, file_stat.st_ino + + @staticmethod + def _file_version(file_stat: os.stat_result) -> FileVersion: + """Return the scanned version used to detect changes before deletion.""" + return ( + file_stat.st_dev, + file_stat.st_ino, + file_stat.st_size, + file_stat.st_mtime_ns, + file_stat.st_ctime_ns, + ) + def _refresh_config(self): """Read latest config values from KiraConfig to support runtime changes.""" cache_config = self.kira_config.get_config("bot_config.cache", {}) or {} @@ -74,154 +92,192 @@ def notify_config_changed(self): self._refresh_config() self._config_changed_event.set() - async def _build_cache(self): - """Build file cache asynchronously""" + async def _scan_folder(self) -> List[DirectoryCandidate]: + """Rebuild the file snapshot and capture protected directory state.""" loop = asyncio.get_running_loop() + previous_first_seen = self._first_seen.copy() + has_scanned = self._has_scanned + protection_seconds = self.file_protection_seconds def scan_folder(): - cache = {} + cache: Dict[str, FileCacheEntry] = {} + versions: Dict[str, FileVersion] = {} + first_seen_entries: Dict[str, Tuple[FileKey, float]] = {} + eligible_directories: List[DirectoryCandidate] = [] total = 0 current_time = time.time() - for file_path in self.folder_path.rglob('*'): - if file_path.is_file(): - stat = file_path.stat() - # For existing files, use mtime as creation_time - cache[str(file_path)] = (stat.st_size, stat.st_mtime, stat.st_mtime) - total += stat.st_size - return cache, total - - # Execute blocking IO operations in thread pool - self.file_cache, self.total_size = await loop.run_in_executor(None, scan_folder) - - logger.info(f"Cache initialization completed: {len(self.file_cache)} files, " - f"total size: {self.total_size / 1024 / 1024:.2f}MB") - - async def _update_cache(self, change_type: int, file_path: str): - """Update cache asynchronously""" - path_str = str(file_path) - loop = asyncio.get_event_loop() - current_time = time.time() - if change_type == 1: # Added - # Only add if not already in cache - if path_str not in self.file_cache: - def add_file(): - p = Path(file_path) - if p.exists() and p.is_file(): - stat = p.stat() - return stat.st_size, stat.st_mtime - return None, None - - result = await loop.run_in_executor(None, add_file) - if result[0] is not None: - size, mtime = result - # Use current_time as creation_time for new files - self.file_cache[path_str] = (size, mtime, current_time) - self.total_size += size - logger.debug(f"File added to cache: {Path(path_str).name}, size: {size}, mtime: {mtime}, creation_time: {current_time}") - - elif change_type == 2: # Modified - if path_str in self.file_cache: - old_size, old_mtime, creation_time = self.file_cache[path_str] - - def modify_file(): - p = Path(file_path) - if p.exists() and p.is_file(): - stat = p.stat() - return stat.st_size, stat.st_mtime - return None, None - - result = await loop.run_in_executor(None, modify_file) - if result[0] is not None: - new_size, new_mtime = result - # Keep the original creation_time - self.file_cache[path_str] = (new_size, new_mtime, creation_time) - self.total_size = self.total_size - old_size + new_size - - elif change_type == 3: # Deleted - if path_str in self.file_cache: - old_size, _, _ = self.file_cache[path_str] - del self.file_cache[path_str] - self.total_size -= old_size - - async def _get_oldest_files(self, limit: int = 10) -> List[Tuple[str, int, float, float]]: - """Get oldest files based on creation_time, skip files within protection period""" + def on_walk_error(error: OSError): + logger.debug(f"Skip temporary directory during scan: {error}") + + for root, directory_names, file_names in os.walk( + self.folder_path, + onerror=on_walk_error, + followlinks=False, + ): + root_path = Path(root) + + for directory_name in directory_names: + directory = root_path / directory_name + try: + directory_stat = directory.lstat() + if not stat.S_ISDIR(directory_stat.st_mode): + continue + if ( + protection_seconds > 0 + and current_time - directory_stat.st_mtime + < protection_seconds + ): + continue + eligible_directories.append( + (directory, self._file_key(directory_stat)) + ) + except OSError as e: + logger.debug( + f"Skip temporary directory {directory} due to stat error: {e}" + ) + + for file_name in file_names: + file_path = root_path / file_name + try: + file_stat = file_path.stat() + if not stat.S_ISREG(file_stat.st_mode): + continue + except OSError as e: + logger.debug( + f"Skip temporary file {file_path} due to stat error: {e}" + ) + continue + + path_str = str(file_path) + file_key = self._file_key(file_stat) + previous = previous_first_seen.get(path_str) + if previous is not None and previous[0] == file_key: + first_seen = previous[1] + elif has_scanned: + first_seen = current_time + else: + # Preserve the previous startup behavior for files that + # existed before the cleaner started. + first_seen = file_stat.st_mtime + + cache[path_str] = ( + file_stat.st_size, + file_stat.st_mtime, + first_seen, + ) + versions[path_str] = self._file_version(file_stat) + first_seen_entries[path_str] = (file_key, first_seen) + total += file_stat.st_size + + return ( + cache, + versions, + first_seen_entries, + eligible_directories, + total, + ) + + ( + self.file_cache, + self._file_versions, + self._first_seen, + eligible_directories, + self.total_size, + ) = await loop.run_in_executor(None, scan_folder) + self._has_scanned = True + + # A replaced or removed path is not the same failed deletion and must + # not inherit its pending retry state. + self._pending_retries = { + path_str: file_key + for path_str, file_key in self._pending_retries.items() + if path_str in self._file_versions + and self._file_versions[path_str][:2] == file_key + } + + logger.debug( + f"Temporary folder scan completed: {len(self.file_cache)} files, " + f"total size: {self.total_size / 1024 / 1024:.2f}MB" + ) + return eligible_directories + + async def _get_oldest_files(self, limit: int = 10) -> List[FileCandidate]: + """Get oldest files while skipping files inside the protection period.""" current_time = time.time() eligible_files = [] - protected_count = 0 - for path_str, (size, mtime, creation_time) in self.file_cache.items(): - # Use creation_time to check protection period - file_age = current_time - creation_time - if file_age < self.file_protection_seconds: - protected_count += 1 + for path_str, (size, mtime, first_seen) in self.file_cache.items(): + if current_time - first_seen < self.file_protection_seconds: continue - - eligible_files.append((path_str, size, mtime, creation_time)) + eligible_files.append((path_str, size, mtime, first_seen)) if not eligible_files: logger.warning("No eligible files for deletion (all files are protected)") return [] - # Use nsmallest to get files with smallest creation_time (oldest files) - oldest_files = heapq.nsmallest(limit, eligible_files, key=lambda x: x[3]) - - return oldest_files + return heapq.nsmallest(limit, eligible_files, key=lambda item: item[3]) - async def _get_expired_files(self) -> List[Tuple[str, int, float, float]]: - """Get files that have exceeded the max age""" + async def _get_expired_files(self) -> List[FileCandidate]: + """Get files that have exceeded the configured maximum age.""" current_time = time.time() expired_files = [] - for path_str, (size, mtime, creation_time) in self.file_cache.items(): - # Check if file age exceeds max age - file_age = current_time - creation_time - if file_age > self.max_age_seconds: - expired_files.append((path_str, size, mtime, creation_time)) + for path_str, (size, mtime, first_seen) in self.file_cache.items(): + if current_time - first_seen > self.max_age_seconds: + expired_files.append((path_str, size, mtime, first_seen)) - # Sort by age (oldest first) - expired_files.sort(key=lambda x: x[3]) + expired_files.sort(key=lambda item: item[3]) return expired_files - async def _get_files_exceeding_limit(self) -> List[Tuple[str, int, float, float]]: - """Get files when total file count exceeds max_files limit, skipping protected files""" + async def _get_files_exceeding_limit(self) -> List[FileCandidate]: + """Get oldest eligible files exceeding the configured count limit.""" if len(self.file_cache) <= self.max_files: return [] current_time = time.time() - # Filter out files within protection period, consistent with _get_oldest_files eligible_files = [] - for path_str, (size, mtime, creation_time) in self.file_cache.items(): - file_age = current_time - creation_time - if file_age < self.file_protection_seconds: + for path_str, (size, mtime, first_seen) in self.file_cache.items(): + if current_time - first_seen < self.file_protection_seconds: continue - eligible_files.append((path_str, size, mtime, creation_time)) + eligible_files.append((path_str, size, mtime, first_seen)) - # Recalculate excess count based on eligible files plus protected ones - # that will remain regardless excess_count = len(self.file_cache) - self.max_files if excess_count <= 0 or not eligible_files: return [] - return heapq.nsmallest(min(excess_count, len(eligible_files)), eligible_files, key=lambda x: x[3]) + return heapq.nsmallest( + min(excess_count, len(eligible_files)), + eligible_files, + key=lambda item: item[3], + ) - async def _delete_file(self, path_str: str) -> Tuple[Literal["deleted", "missing", "failed"], int, Optional[str]]: - """Delete a single file, retrying once after clearing read-only mode.""" + async def _delete_file( + self, + path_str: str, + expected_version: Optional[FileVersion] = None, + ) -> Tuple[DeleteStatus, int, Optional[str]]: + """Delete one unchanged file, retrying read-only failures once.""" loop = asyncio.get_running_loop() def delete(): file_path = Path(path_str) try: - if not file_path.exists() or not file_path.is_file(): - return "missing", 0, None - - size = file_path.stat().st_size + file_stat = file_path.stat() + if not stat.S_ISREG(file_stat.st_mode): + return "changed", 0, None + if ( + expected_version is not None + and self._file_version(file_stat) != expected_version + ): + return "changed", 0, None + + size = file_stat.st_size try: file_path.unlink() except PermissionError: # Git objects are commonly read-only on Windows. Make the # file owner-writable, then retry exactly once. - mode = file_path.stat().st_mode - file_path.chmod(mode | stat.S_IWRITE) + file_path.chmod(file_stat.st_mode | stat.S_IWRITE) file_path.unlink() return "deleted", size, None except FileNotFoundError: @@ -231,26 +287,22 @@ def delete(): return await loop.run_in_executor(None, delete) - async def _cleanup_empty_dirs(self): - """Remove empty subdirectories from deepest to shallowest.""" + async def _cleanup_empty_dirs( + self, + candidates: List[DirectoryCandidate], + ) -> None: + """Remove empty directories that were unprotected at scan time.""" loop = asyncio.get_running_loop() def cleanup_empty_dirs(): - if not self.folder_path.exists(): - return - - now = time.time() - directories = ( - path for path in self.folder_path.rglob('*') if path.is_dir() - ) - for directory in sorted( - directories, key=lambda path: len(path.parts), reverse=True + for directory, expected_key in sorted( + candidates, + key=lambda item: len(item[0].parts), + reverse=True, ): try: - if ( - self.file_protection_seconds > 0 - and now - directory.stat().st_mtime < self.file_protection_seconds - ): + directory_stat = directory.lstat() + if self._file_key(directory_stat) != expected_key: continue directory.rmdir() except OSError: @@ -259,33 +311,35 @@ def cleanup_empty_dirs(): await loop.run_in_executor(None, cleanup_empty_dirs) async def cleanup(self): - """Execute cleanup asynchronously.""" - # Use lock to prevent concurrent cleanup + """Scan the temporary directory and execute one cleanup cycle.""" async with self._cleanup_lock: self._refresh_config() + directory_candidates = await self._scan_folder() current_time = time.time() - - # Control check frequency (inside lock to avoid TOCTOU race) - if current_time - self.last_check_time < self.check_interval: - return - - self.last_check_time = current_time - - # Check for expired files first expired_files = await self._get_expired_files() + pending_files = [ + (path_str, *self.file_cache[path_str]) + for path_str in self._pending_retries + if path_str in self.file_cache + ] needs_cleanup = ( self.total_size > self.max_size_bytes or len(self.file_cache) > self.max_files or expired_files + or pending_files ) if not needs_cleanup: - await self._cleanup_empty_dirs() + await self._cleanup_empty_dirs(directory_candidates) return - logger.info(f"CLEANUP TRIGGERED - Files: {len(self.file_cache)}/{self.max_files}, " - f"Size: {self.total_size / 1024 / 1024:.2f}MB/{self.max_size_bytes / 1024 / 1024:.2f}MB") + logger.info( + f"CLEANUP TRIGGERED - Files: {len(self.file_cache)}/{self.max_files}, " + f"Size: {self.total_size / 1024 / 1024:.2f}MB/" + f"{self.max_size_bytes / 1024 / 1024:.2f}MB, " + f"Pending retries: {len(pending_files)}" + ) deleted_count = 0 freed_space = 0 @@ -295,18 +349,32 @@ async def cleanup(self): async def delete_candidate(path_str: str): nonlocal deleted_count, freed_space - # A file can qualify as expired, excessive, and oversized in - # the same cleanup. Never retry a failure within one cycle. if path_str in attempted_paths: return None attempted_paths.add(path_str) - status, deleted_size, error = await self._delete_file(path_str) + expected_version = self._file_versions.get(path_str) + status, deleted_size, error = await self._delete_file( + path_str, + expected_version, + ) if status == "failed": failed_deletions[path_str] = error or "unknown error" + if expected_version is not None: + self._pending_retries[path_str] = expected_version[:2] + return status + + if status == "changed": + self._pending_retries.pop(path_str, None) + logger.debug( + f"Skipping changed temporary file until the next scan: {path_str}" + ) return status cached = self.file_cache.pop(path_str, None) + self._file_versions.pop(path_str, None) + self._first_seen.pop(path_str, None) + self._pending_retries.pop(path_str, None) if cached is not None: self.total_size = max(0, self.total_size - cached[0]) @@ -315,59 +383,83 @@ async def delete_candidate(path_str: str): freed_space += deleted_size return status - # Phase 1: Delete expired files (by max_age_hours) + # Retry files that failed during an earlier cleanup cycle first. + for path_str, size, mtime, first_seen in pending_files: + status = await delete_candidate(path_str) + if status == "deleted": + logger.debug( + f"DELETED pending retry: {Path(path_str).name} " + f"(size: {size / 1024:.2f}KB)" + ) + if expired_files: - logger.debug(f"Found {len(expired_files)} expired files (older than {self.max_age_seconds / 3600:.1f}h)") - # Use a reduced protection period for expired files to avoid - # deleting files that are still actively in use when max_age - # is changed at runtime + logger.debug( + f"Found {len(expired_files)} expired files " + f"(older than {self.max_age_seconds / 3600:.1f}h)" + ) min_expired_protection = self.file_protection_seconds // 4 - for path_str, size, mtime, creation_time in expired_files: - file_age = current_time - creation_time + for path_str, size, mtime, first_seen in expired_files: + file_age = current_time - first_seen if file_age < min_expired_protection: - logger.warning(f"Skipping recently created expired file: " - f"{Path(path_str).name} (age: {file_age:.1f}s < " - f"protection: {min_expired_protection}s)") + logger.warning( + f"Skipping recently created expired file: " + f"{Path(path_str).name} (age: {file_age:.1f}s < " + f"protection: {min_expired_protection}s)" + ) continue status = await delete_candidate(path_str) if status == "deleted": - file_age_hours = (current_time - creation_time) / 3600 - logger.debug(f"DELETED expired: {Path(path_str).name} (age: {file_age_hours:.1f}h, size: {size / 1024:.2f}KB)") + logger.debug( + f"DELETED expired: {Path(path_str).name} " + f"(age: {file_age / 3600:.1f}h, " + f"size: {size / 1024:.2f}KB)" + ) - # Phase 2: Delete files exceeding max_files limit - # Protection period already filtered inside _get_files_exceeding_limit excess_files = await self._get_files_exceeding_limit() if excess_files: - logger.debug(f"Found {len(excess_files)} excess files (limit: {self.max_files})") - for path_str, size, mtime, creation_time in excess_files: + logger.debug( + f"Found {len(excess_files)} excess files (limit: {self.max_files})" + ) + for path_str, size, mtime, first_seen in excess_files: status = await delete_candidate(path_str) if status == "deleted": - logger.debug(f"DELETED excess: {Path(path_str).name} (size: {size / 1024:.2f}KB)") + logger.debug( + f"DELETED excess: {Path(path_str).name} " + f"(size: {size / 1024:.2f}KB)" + ) - # Phase 3: Delete oldest files if still over size limit if self.total_size > self.max_size_bytes: oldest_files = await self._get_oldest_files(limit=self.batch_size) - for path_str, size, mtime, creation_time in oldest_files: + for path_str, size, mtime, first_seen in oldest_files: if self.total_size <= self.max_size_bytes: break - # Double-check protection period before deletion - file_age = current_time - creation_time + file_age = current_time - first_seen if file_age < self.file_protection_seconds: - logger.error(f"ATTEMPTED TO DELETE PROTECTED FILE (age: {file_age:.2f}s): {Path(path_str).name} - SKIPPING") + logger.error( + f"ATTEMPTED TO DELETE PROTECTED FILE " + f"(age: {file_age:.2f}s): " + f"{Path(path_str).name} - SKIPPING" + ) continue status = await delete_candidate(path_str) if status == "deleted": - logger.debug(f"DELETED: {Path(path_str).name} (age: {file_age:.2f}s, size: {size / 1024:.2f}KB)") + logger.debug( + f"DELETED: {Path(path_str).name} " + f"(age: {file_age:.2f}s, " + f"size: {size / 1024:.2f}KB)" + ) - await self._cleanup_empty_dirs() + await self._cleanup_empty_dirs(directory_candidates) if deleted_count > 0: - logger.info(f"Cleanup completed: deleted {deleted_count} files, " - f"freed {freed_space / 1024 / 1024:.2f}MB, " - f"remaining: {len(self.file_cache)} files, " - f"{self.total_size / 1024 / 1024:.2f}MB") + logger.info( + f"Cleanup completed: deleted {deleted_count} files, " + f"freed {freed_space / 1024 / 1024:.2f}MB, " + f"remaining: {len(self.file_cache)} files, " + f"{self.total_size / 1024 / 1024:.2f}MB" + ) if failed_deletions: failure_items = list(failed_deletions.items()) @@ -378,25 +470,13 @@ async def delete_candidate(path_str: str): if omitted > 0: details += f"; and {omitted} more" logger.warning( - f"Cleanup completed with {len(failed_deletions)} deletion failure(s); " - f"will retry on the next cleanup cycle. {details}" + f"Cleanup completed with {len(failed_deletions)} " + f"deletion failure(s); will retry on the next cleanup cycle. " + f"{details}" ) - async def _process_changes(self, changes): - """Process file change events""" - - for change_type, file_path in changes: - await self._update_cache(change_type, file_path) - - # Check if cleanup is needed; cleanup() handles its own rate limiting - # and _refresh_config() inside the lock to avoid race conditions - if self.total_size > self.max_size_bytes or len(self.file_cache) > self.max_files: - logger.debug(f"Cleanup needed - Files: {len(self.file_cache)}/{self.max_files}, " - f"Size: {self.total_size / 1024 / 1024:.2f}MB/{self.max_size_bytes / 1024 / 1024:.2f}MB") - await self.cleanup() - async def _periodic_cleanup_loop(self): - """Periodically check for expired files even without file changes.""" + """Run cleanup periodically and react to runtime interval changes.""" while not self._stop_event.is_set(): try: try: @@ -418,44 +498,27 @@ async def _periodic_cleanup_loop(self): logger.error(f"Error in periodic cleanup: {e}") async def start_monitoring(self): - """Start asynchronous monitoring""" - logger.info("="*60) - logger.info(f"Starting monitoring folder: {self.folder_path}") + """Start periodic temporary-directory cleanup.""" + logger.info("=" * 60) + logger.info(f"Starting periodic cleanup for folder: {self.folder_path}") logger.info(f"Max size: {self.max_size_bytes / 1024 / 1024:.2f}MB") logger.info(f"Max files: {self.max_files}") logger.info(f"Max age: {self.max_age_seconds / 3600:.1f} hours") - logger.info(f"File protection period: {self.file_protection_seconds} seconds") - logger.info(f"Check interval: {self.check_interval} seconds") - logger.info("="*60) - - # Build cache - await self._build_cache() - - # Cleanup once on startup (reset check interval so it runs) - self.last_check_time = 0 - await self.cleanup() - - periodic_task = asyncio.create_task( - self._periodic_cleanup_loop(), - name="temp_periodic_cleanup" + logger.info( + f"File protection period: {self.file_protection_seconds} seconds" ) + logger.info(f"Check interval: {self.check_interval} seconds") + logger.info("=" * 60) try: - # Use awatch for asynchronous monitoring - async for changes in awatch(str(self.folder_path)): - if changes and not self._stop_event.is_set(): - await self._process_changes(changes) - + await self.cleanup() + await self._periodic_cleanup_loop() except asyncio.CancelledError: - logger.info("Monitoring cancelled") + logger.info("Periodic cleanup cancelled") finally: - periodic_task.cancel() - try: - await periodic_task - except asyncio.CancelledError: - pass - logger.info("Monitoring stopped") + logger.info("Periodic cleanup stopped") async def stop_monitoring(self): - """Stop monitoring""" + """Stop periodic cleanup and wake the scheduler immediately.""" self._stop_event.set() + self._config_changed_event.set() diff --git a/requirements.txt b/requirements.txt index d588214f..1eeb963f 100644 --- a/requirements.txt +++ b/requirements.txt @@ -15,7 +15,6 @@ psutil>=5.9.0 fastmcp~=2.14.4 httpx>=0.28.1 jinja2>=3.1.0 -watchfiles>=1.1.1 pycryptodome>=3.20.0 qrcode>=7.4.2 PyYAML>=6.0.2 diff --git a/tests/test_temp_monitor.py b/tests/test_temp_monitor.py index d148b959..36bafbb5 100644 --- a/tests/test_temp_monitor.py +++ b/tests/test_temp_monitor.py @@ -1,3 +1,4 @@ +import asyncio import os import stat import time @@ -72,53 +73,65 @@ def record_chmod(path, mode, *args, **kwargs): @pytest.mark.asyncio -async def test_cleanup_summarizes_failures_once_and_retries_next_cycle( +async def test_cleanup_retries_failed_file_after_size_returns_below_limit( tmp_path, monkeypatch ): - target = tmp_path / "locked.bin" - target.write_bytes(b"content") + locked = tmp_path / "locked.bin" + removable = tmp_path / "removable.bin" + locked.write_bytes(b"a" * 600) + removable.write_bytes(b"b" * 600) old_time = time.time() - 3600 - os.utime(target, (old_time, old_time)) - - monitor = make_monitor(tmp_path, max_size_mb=0, max_age_hours=0) - await monitor._build_cache() + os.utime(locked, (old_time, old_time)) + os.utime(removable, (old_time + 10, old_time + 10)) + monitor = make_monitor(tmp_path, max_size_mb=0.001) + original_delete = monitor._delete_file delete_calls = [] - async def fail_delete(path_str): + async def fail_locked(path_str, expected_version=None): delete_calls.append(path_str) - return "failed", 0, "PermissionError: locked" + if path_str == str(locked): + if delete_calls.count(path_str) == 1: + os.utime(locked, None) + return "failed", 0, "PermissionError: locked" + return await original_delete(path_str, expected_version) warnings = [] - monkeypatch.setattr(monitor, "_delete_file", fail_delete) + monkeypatch.setattr(monitor, "_delete_file", fail_locked) monkeypatch.setattr(temp_monitor.logger, "warning", warnings.append) await monitor.cleanup() - assert delete_calls == [str(target)] + assert locked.exists() + assert not removable.exists() + assert monitor.total_size < monitor.max_size_bytes + assert str(locked) in monitor._pending_retries assert len(warnings) == 1 - assert "1 deletion failure(s)" in warnings[0] - assert "next cleanup cycle" in warnings[0] - assert str(target) in monitor.file_cache - monitor.last_check_time = 0 await monitor.cleanup() - assert delete_calls == [str(target), str(target)] + assert delete_calls.count(str(locked)) == 2 assert len(warnings) == 2 + assert "next cleanup cycle" in warnings[-1] @pytest.mark.asyncio -async def test_cleanup_removes_empty_directories(tmp_path): +async def test_cleanup_removes_nested_empty_directories_in_one_cycle(tmp_path): nested_dir = tmp_path / "job" / "repo" / ".git" / "objects" nested_dir.mkdir(parents=True) target = nested_dir / "old.bin" target.write_bytes(b"content") old_time = time.time() - 3600 os.utime(target, (old_time, old_time)) + for directory in [ + nested_dir, + nested_dir.parent, + nested_dir.parent.parent, + tmp_path / "job", + ]: + os.utime(directory, (old_time, old_time)) monitor = make_monitor(tmp_path, max_age_hours=0) - await monitor._build_cache() await monitor.cleanup() assert not target.exists() @@ -153,7 +166,7 @@ async def test_configured_check_interval_updates_at_runtime(tmp_path): @pytest.mark.asyncio -async def test_cleanup_empty_dirs_preserves_recent_directories(tmp_path): +async def test_cleanup_preserves_recent_empty_directories(tmp_path): recent_dir = tmp_path / "recent" old_dir = tmp_path / "old" recent_dir.mkdir() @@ -170,3 +183,53 @@ async def test_cleanup_empty_dirs_preserves_recent_directories(tmp_path): assert recent_dir.exists() assert not old_dir.exists() + + +@pytest.mark.asyncio +async def test_cleanup_rescans_files_created_after_startup(tmp_path): + monitor = make_monitor(tmp_path, max_size_mb=0) + await monitor.cleanup() + + target = tmp_path / "new.bin" + target.write_bytes(b"content") + + await monitor.cleanup() + + assert not target.exists() + assert monitor.file_cache == {} + assert monitor.total_size == 0 + + +@pytest.mark.asyncio +async def test_cleanup_skips_file_changed_after_scan(tmp_path, monkeypatch): + target = tmp_path / "changing.bin" + target.write_bytes(b"old") + monitor = make_monitor(tmp_path, max_size_mb=0) + original_scan = monitor._scan_folder + + async def scan_then_change(): + directory_candidates = await original_scan() + target.write_bytes(b"replacement") + return directory_candidates + + monkeypatch.setattr(monitor, "_scan_folder", scan_then_change) + + await monitor.cleanup() + + assert target.read_bytes() == b"replacement" + assert str(target) not in monitor._pending_retries + + +@pytest.mark.asyncio +async def test_stop_monitoring_wakes_periodic_scheduler(tmp_path): + monitor = AsyncTempMonitor( + str(tmp_path), + FakeConfig({"check_interval_minutes": 5}), + ) + task = asyncio.create_task(monitor._periodic_cleanup_loop()) + await asyncio.sleep(0) + + await monitor.stop_monitoring() + await asyncio.wait_for(task, timeout=1) + + assert task.done() From 9e510768b4262f621cb707767db341008d718d77 Mon Sep 17 00:00:00 2001 From: Caleb Date: Wed, 23 Sep 2026 16:11:38 +0800 Subject: [PATCH 5/8] fix: harden temp cleanup retries --- core/temp_monitor.py | 72 ++++++++++++++---------- tests/test_temp_monitor.py | 109 ++++++++++++++++++++++++++++++++++--- 2 files changed, 144 insertions(+), 37 deletions(-) diff --git a/core/temp_monitor.py b/core/temp_monitor.py index 3f633963..e25ba91e 100644 --- a/core/temp_monitor.py +++ b/core/temp_monitor.py @@ -19,6 +19,7 @@ FileCandidate = Tuple[str, int, float, float] DirectoryCandidate = Tuple[Path, FileKey] DeleteStatus = Literal["deleted", "missing", "changed", "failed"] +DeleteResult = Tuple[DeleteStatus, int, Optional[str], Optional[FileVersion]] class AsyncTempMonitor: @@ -46,7 +47,7 @@ def __init__( self.total_size = 0 self._file_versions: Dict[str, FileVersion] = {} self._first_seen: Dict[str, Tuple[FileKey, float]] = {} - self._pending_retries: Dict[str, FileKey] = {} + self._pending_retries: Dict[str, FileVersion] = {} self._has_scanned = False self._cleanup_lock = asyncio.Lock() self._stop_event = asyncio.Event() @@ -187,13 +188,12 @@ def on_walk_error(error: OSError): ) = await loop.run_in_executor(None, scan_folder) self._has_scanned = True - # A replaced or removed path is not the same failed deletion and must - # not inherit its pending retry state. + # A removed, replaced, or modified path is not the same failed + # deletion and must not inherit its pending retry state. self._pending_retries = { - path_str: file_key - for path_str, file_key in self._pending_retries.items() - if path_str in self._file_versions - and self._file_versions[path_str][:2] == file_key + path_str: pending_version + for path_str, pending_version in self._pending_retries.items() + if self._file_versions.get(path_str) == pending_version } logger.debug( @@ -231,7 +231,7 @@ async def _get_expired_files(self) -> List[FileCandidate]: return expired_files async def _get_files_exceeding_limit(self) -> List[FileCandidate]: - """Get oldest eligible files exceeding the configured count limit.""" + """Get count-limit candidates ordered from oldest to newest.""" if len(self.file_cache) <= self.max_files: return [] @@ -242,21 +242,15 @@ async def _get_files_exceeding_limit(self) -> List[FileCandidate]: continue eligible_files.append((path_str, size, mtime, first_seen)) - excess_count = len(self.file_cache) - self.max_files - if excess_count <= 0 or not eligible_files: - return [] - return heapq.nsmallest( - min(excess_count, len(eligible_files)), - eligible_files, - key=lambda item: item[3], - ) + eligible_files.sort(key=lambda item: item[3]) + return eligible_files async def _delete_file( self, path_str: str, expected_version: Optional[FileVersion] = None, - ) -> Tuple[DeleteStatus, int, Optional[str]]: - """Delete one unchanged file, retrying read-only failures once.""" + ) -> DeleteResult: + """Delete one unchanged file and return its remaining version on failure.""" loop = asyncio.get_running_loop() def delete(): @@ -264,12 +258,12 @@ def delete(): try: file_stat = file_path.stat() if not stat.S_ISREG(file_stat.st_mode): - return "changed", 0, None + return "changed", 0, None, None if ( expected_version is not None and self._file_version(file_stat) != expected_version ): - return "changed", 0, None + return "changed", 0, None, None size = file_stat.st_size try: @@ -279,11 +273,27 @@ def delete(): # file owner-writable, then retry exactly once. file_path.chmod(file_stat.st_mode | stat.S_IWRITE) file_path.unlink() - return "deleted", size, None + return "deleted", size, None, None except FileNotFoundError: - return "missing", 0, None + return "missing", 0, None, None except OSError as e: - return "failed", 0, f"{type(e).__name__}: {e}" + error = f"{type(e).__name__}: {e}" + try: + remaining_stat = file_path.stat() + remaining_version = ( + self._file_version(remaining_stat) + if stat.S_ISREG(remaining_stat.st_mode) + else None + ) + if ( + expected_version is not None + and remaining_version is not None + and remaining_version[:4] != expected_version[:4] + ): + return "changed", 0, None, None + except OSError: + remaining_version = None + return "failed", 0, error, remaining_version return await loop.run_in_executor(None, delete) @@ -354,14 +364,15 @@ async def delete_candidate(path_str: str): attempted_paths.add(path_str) expected_version = self._file_versions.get(path_str) - status, deleted_size, error = await self._delete_file( - path_str, - expected_version, + status, deleted_size, error, remaining_version = ( + await self._delete_file(path_str, expected_version) ) if status == "failed": failed_deletions[path_str] = error or "unknown error" - if expected_version is not None: - self._pending_retries[path_str] = expected_version[:2] + if remaining_version is not None: + self._pending_retries[path_str] = remaining_version + else: + self._pending_retries.pop(path_str, None) return status if status == "changed": @@ -418,9 +429,12 @@ async def delete_candidate(path_str: str): excess_files = await self._get_files_exceeding_limit() if excess_files: logger.debug( - f"Found {len(excess_files)} excess files (limit: {self.max_files})" + f"Found {len(excess_files)} count-limit candidates " + f"(limit: {self.max_files})" ) for path_str, size, mtime, first_seen in excess_files: + if len(self.file_cache) <= self.max_files: + break status = await delete_candidate(path_str) if status == "deleted": logger.debug( diff --git a/tests/test_temp_monitor.py b/tests/test_temp_monitor.py index 36bafbb5..375570ef 100644 --- a/tests/test_temp_monitor.py +++ b/tests/test_temp_monitor.py @@ -61,11 +61,14 @@ def record_chmod(path, mode, *args, **kwargs): monkeypatch.setattr(Path, "chmod", record_chmod) monitor = make_monitor(tmp_path) - status, deleted_size, error = await monitor._delete_file(str(target)) + status, deleted_size, error, remaining_version = await monitor._delete_file( + str(target) + ) assert status == "deleted" assert deleted_size == len(b"content") assert error is None + assert remaining_version is None assert unlink_calls == 2 assert len(chmod_modes) == 1 assert chmod_modes[0] & stat.S_IWRITE @@ -73,7 +76,35 @@ def record_chmod(path, mode, *args, **kwargs): @pytest.mark.asyncio -async def test_cleanup_retries_failed_file_after_size_returns_below_limit( +async def test_delete_file_does_not_retry_replacement_seen_after_failure( + tmp_path, monkeypatch +): + target = tmp_path / "changing.bin" + target.write_bytes(b"old") + monitor = make_monitor(tmp_path) + expected_version = monitor._file_version(target.stat()) + + def replace_then_fail(path, *args, **kwargs): + if path == target: + target.write_bytes(b"replacement") + raise OSError("locked") + return Path.unlink(path, *args, **kwargs) + + monkeypatch.setattr(Path, "unlink", replace_then_fail) + + status, deleted_size, error, remaining_version = await monitor._delete_file( + str(target), expected_version + ) + + assert status == "changed" + assert deleted_size == 0 + assert error is None + assert remaining_version is None + assert target.read_bytes() == b"replacement" + + +@pytest.mark.asyncio +async def test_cleanup_retries_unchanged_file_after_size_returns_below_limit( tmp_path, monkeypatch ): locked = tmp_path / "locked.bin" @@ -90,10 +121,8 @@ async def test_cleanup_retries_failed_file_after_size_returns_below_limit( async def fail_locked(path_str, expected_version=None): delete_calls.append(path_str) - if path_str == str(locked): - if delete_calls.count(path_str) == 1: - os.utime(locked, None) - return "failed", 0, "PermissionError: locked" + if path_str == str(locked) and delete_calls.count(path_str) == 1: + return "failed", 0, "PermissionError: locked", expected_version return await original_delete(path_str, expected_version) warnings = [] @@ -111,8 +140,72 @@ async def fail_locked(path_str, expected_version=None): await monitor.cleanup() assert delete_calls.count(str(locked)) == 2 - assert len(warnings) == 2 - assert "next cleanup cycle" in warnings[-1] + assert not locked.exists() + assert str(locked) not in monitor._pending_retries + assert len(warnings) == 1 + + +@pytest.mark.asyncio +async def test_cleanup_drops_pending_retry_when_failed_file_changes( + tmp_path, monkeypatch +): + locked = tmp_path / "locked.bin" + removable = tmp_path / "removable.bin" + locked.write_bytes(b"a" * 600) + removable.write_bytes(b"b" * 600) + old_time = time.time() - 3600 + os.utime(locked, (old_time, old_time)) + os.utime(removable, (old_time + 10, old_time + 10)) + + monitor = make_monitor(tmp_path, max_size_mb=0.001) + original_delete = monitor._delete_file + delete_calls = [] + + async def fail_locked(path_str, expected_version=None): + delete_calls.append(path_str) + if path_str == str(locked) and delete_calls.count(path_str) == 1: + return "failed", 0, "PermissionError: locked", expected_version + return await original_delete(path_str, expected_version) + + monkeypatch.setattr(monitor, "_delete_file", fail_locked) + + await monitor.cleanup() + locked.write_bytes(b"replacement") + await monitor.cleanup() + + assert locked.read_bytes() == b"replacement" + assert delete_calls.count(str(locked)) == 1 + assert str(locked) not in monitor._pending_retries + + +@pytest.mark.asyncio +async def test_count_cleanup_continues_after_oldest_file_fails( + tmp_path, monkeypatch +): + locked = tmp_path / "locked.bin" + removable = tmp_path / "removable.bin" + locked.write_bytes(b"locked") + removable.write_bytes(b"removable") + old_time = time.time() - 3600 + os.utime(locked, (old_time, old_time)) + os.utime(removable, (old_time + 10, old_time + 10)) + + monitor = make_monitor(tmp_path, max_files=1) + original_delete = monitor._delete_file + + async def fail_locked(path_str, expected_version=None): + if path_str == str(locked): + return "failed", 0, "PermissionError: locked", expected_version + return await original_delete(path_str, expected_version) + + monkeypatch.setattr(monitor, "_delete_file", fail_locked) + + await monitor.cleanup() + + assert locked.exists() + assert not removable.exists() + assert len(monitor.file_cache) == monitor.max_files + assert str(locked) in monitor._pending_retries @pytest.mark.asyncio From 31c897eb8464990eaa78bffde41af948660e6550 Mon Sep 17 00:00:00 2001 From: Caleb Date: Wed, 23 Sep 2026 16:30:19 +0800 Subject: [PATCH 6/8] fix: keep size cleanup progressing --- core/temp_monitor.py | 23 ++++++++++++++++------- tests/test_temp_monitor.py | 32 ++++++++++++++++++++++++++++++++ 2 files changed, 48 insertions(+), 7 deletions(-) diff --git a/core/temp_monitor.py b/core/temp_monitor.py index e25ba91e..72b51ddc 100644 --- a/core/temp_monitor.py +++ b/core/temp_monitor.py @@ -202,12 +202,18 @@ def on_walk_error(error: OSError): ) return eligible_directories - async def _get_oldest_files(self, limit: int = 10) -> List[FileCandidate]: - """Get oldest files while skipping files inside the protection period.""" + async def _get_oldest_files( + self, + limit: int = 10, + exclude: Optional[set] = None, + ) -> List[FileCandidate]: + """Get oldest unattempted files outside the protection period.""" current_time = time.time() eligible_files = [] for path_str, (size, mtime, first_seen) in self.file_cache.items(): + if exclude and path_str in exclude: + continue if current_time - first_seen < self.file_protection_seconds: continue eligible_files.append((path_str, size, mtime, first_seen)) @@ -395,7 +401,7 @@ async def delete_candidate(path_str: str): return status # Retry files that failed during an earlier cleanup cycle first. - for path_str, size, mtime, first_seen in pending_files: + for path_str, size, _mtime, first_seen in pending_files: status = await delete_candidate(path_str) if status == "deleted": logger.debug( @@ -409,7 +415,7 @@ async def delete_candidate(path_str: str): f"(older than {self.max_age_seconds / 3600:.1f}h)" ) min_expired_protection = self.file_protection_seconds // 4 - for path_str, size, mtime, first_seen in expired_files: + for path_str, size, _mtime, first_seen in expired_files: file_age = current_time - first_seen if file_age < min_expired_protection: logger.warning( @@ -432,7 +438,7 @@ async def delete_candidate(path_str: str): f"Found {len(excess_files)} count-limit candidates " f"(limit: {self.max_files})" ) - for path_str, size, mtime, first_seen in excess_files: + for path_str, size, _mtime, first_seen in excess_files: if len(self.file_cache) <= self.max_files: break status = await delete_candidate(path_str) @@ -443,8 +449,11 @@ async def delete_candidate(path_str: str): ) if self.total_size > self.max_size_bytes: - oldest_files = await self._get_oldest_files(limit=self.batch_size) - for path_str, size, mtime, first_seen in oldest_files: + oldest_files = await self._get_oldest_files( + limit=self.batch_size, + exclude=attempted_paths, + ) + for path_str, size, _mtime, first_seen in oldest_files: if self.total_size <= self.max_size_bytes: break diff --git a/tests/test_temp_monitor.py b/tests/test_temp_monitor.py index 375570ef..76417102 100644 --- a/tests/test_temp_monitor.py +++ b/tests/test_temp_monitor.py @@ -178,6 +178,38 @@ async def fail_locked(path_str, expected_version=None): assert str(locked) not in monitor._pending_retries +@pytest.mark.asyncio +async def test_size_cleanup_excludes_already_attempted_pending_files( + tmp_path, monkeypatch +): + locked = tmp_path / "locked.bin" + removable = tmp_path / "removable.bin" + locked.write_bytes(b"a" * 600) + removable.write_bytes(b"b" * 600) + old_time = time.time() - 3600 + os.utime(locked, (old_time, old_time)) + os.utime(removable, (old_time + 10, old_time + 10)) + + monitor = make_monitor(tmp_path, max_size_mb=0.001) + monitor.batch_size = 1 + monitor._pending_retries[str(locked)] = monitor._file_version(locked.stat()) + original_delete = monitor._delete_file + + async def fail_locked(path_str, expected_version=None): + if path_str == str(locked): + return "failed", 0, "PermissionError: locked", expected_version + return await original_delete(path_str, expected_version) + + monkeypatch.setattr(monitor, "_delete_file", fail_locked) + + await monitor.cleanup() + + assert locked.exists() + assert not removable.exists() + assert monitor.total_size < monitor.max_size_bytes + assert str(locked) in monitor._pending_retries + + @pytest.mark.asyncio async def test_count_cleanup_continues_after_oldest_file_fails( tmp_path, monkeypatch From 3b8c71dccffc41045b0940d748f54ff4522b8144 Mon Sep 17 00:00:00 2001 From: Caleb Date: Wed, 23 Sep 2026 18:10:23 +0800 Subject: [PATCH 7/8] refactor: reduce temp cleanup log noise --- core/temp_monitor.py | 38 ++++++++------------------------------ 1 file changed, 8 insertions(+), 30 deletions(-) diff --git a/core/temp_monitor.py b/core/temp_monitor.py index 72b51ddc..333656d7 100644 --- a/core/temp_monitor.py +++ b/core/temp_monitor.py @@ -401,13 +401,8 @@ async def delete_candidate(path_str: str): return status # Retry files that failed during an earlier cleanup cycle first. - for path_str, size, _mtime, first_seen in pending_files: - status = await delete_candidate(path_str) - if status == "deleted": - logger.debug( - f"DELETED pending retry: {Path(path_str).name} " - f"(size: {size / 1024:.2f}KB)" - ) + for path_str, _size, _mtime, _first_seen in pending_files: + await delete_candidate(path_str) if expired_files: logger.debug( @@ -415,7 +410,7 @@ async def delete_candidate(path_str: str): f"(older than {self.max_age_seconds / 3600:.1f}h)" ) min_expired_protection = self.file_protection_seconds // 4 - for path_str, size, _mtime, first_seen in expired_files: + for path_str, _size, _mtime, first_seen in expired_files: file_age = current_time - first_seen if file_age < min_expired_protection: logger.warning( @@ -424,13 +419,7 @@ async def delete_candidate(path_str: str): f"protection: {min_expired_protection}s)" ) continue - status = await delete_candidate(path_str) - if status == "deleted": - logger.debug( - f"DELETED expired: {Path(path_str).name} " - f"(age: {file_age / 3600:.1f}h, " - f"size: {size / 1024:.2f}KB)" - ) + await delete_candidate(path_str) excess_files = await self._get_files_exceeding_limit() if excess_files: @@ -438,22 +427,17 @@ async def delete_candidate(path_str: str): f"Found {len(excess_files)} count-limit candidates " f"(limit: {self.max_files})" ) - for path_str, size, _mtime, first_seen in excess_files: + for path_str, _size, _mtime, _first_seen in excess_files: if len(self.file_cache) <= self.max_files: break - status = await delete_candidate(path_str) - if status == "deleted": - logger.debug( - f"DELETED excess: {Path(path_str).name} " - f"(size: {size / 1024:.2f}KB)" - ) + await delete_candidate(path_str) if self.total_size > self.max_size_bytes: oldest_files = await self._get_oldest_files( limit=self.batch_size, exclude=attempted_paths, ) - for path_str, size, _mtime, first_seen in oldest_files: + for path_str, _size, _mtime, first_seen in oldest_files: if self.total_size <= self.max_size_bytes: break @@ -466,13 +450,7 @@ async def delete_candidate(path_str: str): ) continue - status = await delete_candidate(path_str) - if status == "deleted": - logger.debug( - f"DELETED: {Path(path_str).name} " - f"(age: {file_age:.2f}s, " - f"size: {size / 1024:.2f}KB)" - ) + await delete_candidate(path_str) await self._cleanup_empty_dirs(directory_candidates) From 107a760998ca90c9a9852692e2b15ddb97cb6d4b Mon Sep 17 00:00:00 2001 From: Caleb Date: Wed, 23 Sep 2026 18:23:24 +0800 Subject: [PATCH 8/8] fix: avoid misleading cleanup warning --- core/temp_monitor.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/temp_monitor.py b/core/temp_monitor.py index 333656d7..3d408102 100644 --- a/core/temp_monitor.py +++ b/core/temp_monitor.py @@ -218,7 +218,7 @@ async def _get_oldest_files( continue eligible_files.append((path_str, size, mtime, first_seen)) - if not eligible_files: + if not eligible_files and not exclude: logger.warning("No eligible files for deletion (all files are protected)") return []