fix(trigger): atomic JSON writes to prevent mid-write corruption
All 11 JSON write locations across 7 files now use write-to-tmp + os.replace pattern via shared atomic_write_json() in config.py. Prevents file corruption if the process crashes mid-write. Updated 8 test files to provide the real atomic_write_json on mocked config modules. 370 tests passing. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.6
parent
083bdceec0
commit
538d3cb7e9
@@ -13,6 +13,9 @@ Provides package-relative paths for trigger data directories.
|
||||
Works in both pip-installed and development environments.
|
||||
"""
|
||||
|
||||
import json
|
||||
import os
|
||||
import tempfile
|
||||
from pathlib import Path
|
||||
from datetime import datetime, timezone
|
||||
|
||||
@@ -36,6 +39,33 @@ def _log_warning(message: str) -> None:
|
||||
AIPASS_PKG_ROOT = TRIGGER_ROOT.parent
|
||||
|
||||
|
||||
def atomic_write_json(path: Path, data, indent: int = 2, ensure_ascii: bool = True, encoding: str = 'utf-8') -> None:
|
||||
"""Write JSON data to a file atomically using write-to-tmp + os.replace.
|
||||
|
||||
Prevents file corruption from process crashes mid-write by writing to
|
||||
a temp file in the same directory first, then atomically renaming.
|
||||
|
||||
Args:
|
||||
path: Target file path
|
||||
data: JSON-serializable data
|
||||
indent: JSON indent level
|
||||
ensure_ascii: JSON ensure_ascii flag
|
||||
encoding: File encoding
|
||||
"""
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
fd, tmp_path = tempfile.mkstemp(dir=path.parent, suffix='.tmp')
|
||||
try:
|
||||
with os.fdopen(fd, 'w', encoding=encoding) as f:
|
||||
json.dump(data, f, indent=indent, ensure_ascii=ensure_ascii)
|
||||
os.replace(tmp_path, path)
|
||||
except BaseException:
|
||||
try:
|
||||
os.unlink(tmp_path)
|
||||
except OSError:
|
||||
pass
|
||||
raise
|
||||
|
||||
|
||||
def print_introspection():
|
||||
"""Display module introspection info."""
|
||||
try:
|
||||
|
||||
@@ -45,7 +45,7 @@ from typing import Any, Dict, List, Optional
|
||||
|
||||
|
||||
from aipass.prax.apps.modules.logger import get_direct_logger
|
||||
from aipass.trigger.apps.config import TRIGGER_ROOT
|
||||
from aipass.trigger.apps.config import TRIGGER_ROOT, atomic_write_json
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
logger = get_direct_logger()
|
||||
@@ -145,10 +145,7 @@ def _save_circuit_breaker_state() -> None:
|
||||
'opened_at': _circuit_breaker.opened_at,
|
||||
'cooldown_seconds': _circuit_breaker.cooldown_seconds,
|
||||
}
|
||||
TRIGGER_CONFIG_FILE.parent.mkdir(parents=True, exist_ok=True)
|
||||
TRIGGER_CONFIG_FILE.write_text(
|
||||
json.dumps(data, indent=2), encoding='utf-8'
|
||||
)
|
||||
atomic_write_json(TRIGGER_CONFIG_FILE, data)
|
||||
except Exception as exc:
|
||||
logger.warning("Failed to save circuit breaker state: %s", exc)
|
||||
|
||||
@@ -193,9 +190,7 @@ def _clear_circuit_breaker_state() -> None:
|
||||
data = json.loads(raw)
|
||||
if 'circuit_breaker' in data:
|
||||
del data['circuit_breaker']
|
||||
TRIGGER_CONFIG_FILE.write_text(
|
||||
json.dumps(data, indent=2), encoding='utf-8'
|
||||
)
|
||||
atomic_write_json(TRIGGER_CONFIG_FILE, data)
|
||||
except Exception as exc:
|
||||
logger.warning("Failed to clear circuit breaker state: %s", exc)
|
||||
|
||||
@@ -535,10 +530,7 @@ def _save_registry(data: dict) -> bool:
|
||||
"""
|
||||
try:
|
||||
data["metadata"]["last_updated"] = datetime.now().isoformat()
|
||||
REGISTRY_FILE.parent.mkdir(parents=True, exist_ok=True)
|
||||
REGISTRY_FILE.write_text(
|
||||
json.dumps(data, indent=2), encoding='utf-8'
|
||||
)
|
||||
atomic_write_json(REGISTRY_FILE, data)
|
||||
return True
|
||||
except Exception as exc:
|
||||
logger.warning("Failed to save error registry to %s: %s", REGISTRY_FILE, exc)
|
||||
|
||||
@@ -25,7 +25,7 @@ import json
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List
|
||||
from aipass.trigger.apps.config import TRIGGER_ROOT
|
||||
from aipass.trigger.apps.config import TRIGGER_ROOT, atomic_write_json
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
def _find_repo_root() -> Path:
|
||||
@@ -158,7 +158,7 @@ def _save_dashboard(branch_path: Path, dashboard: Dict) -> bool:
|
||||
try:
|
||||
dashboard_path = branch_path / "DASHBOARD.local.json"
|
||||
dashboard["last_refreshed"] = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||||
dashboard_path.write_text(json.dumps(dashboard, indent=2))
|
||||
atomic_write_json(dashboard_path, dashboard)
|
||||
return True
|
||||
except Exception as exc:
|
||||
_log_warning(f"save dashboard to {branch_path}: {exc}")
|
||||
|
||||
@@ -26,7 +26,7 @@ import re
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Optional
|
||||
from aipass.trigger.apps.config import TRIGGER_ROOT, AIPASS_PKG_ROOT
|
||||
from aipass.trigger.apps.config import TRIGGER_ROOT, AIPASS_PKG_ROOT, atomic_write_json
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
def _find_repo_root() -> Path:
|
||||
@@ -71,10 +71,7 @@ def _load_registry() -> dict:
|
||||
|
||||
def _save_registry(registry: dict) -> None:
|
||||
"""Save registry to JSON file"""
|
||||
import json
|
||||
FLOW_JSON_DIR.mkdir(parents=True, exist_ok=True)
|
||||
with open(REGISTRY_FILE, 'w', encoding='utf-8') as f:
|
||||
json.dump(registry, f, indent=2)
|
||||
atomic_write_json(REGISTRY_FILE, registry)
|
||||
|
||||
|
||||
def _get_plan_number(file_path: Path) -> Optional[str]:
|
||||
|
||||
@@ -23,7 +23,7 @@ import time
|
||||
from pathlib import Path
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import Any, Callable, Dict, List, Optional, Set
|
||||
from aipass.trigger.apps.config import TRIGGER_ROOT
|
||||
from aipass.trigger.apps.config import TRIGGER_ROOT, atomic_write_json
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
SYSTEM_LOGS_DIR = TRIGGER_ROOT.parent / "system_logs"
|
||||
@@ -89,9 +89,7 @@ def _load_trigger_data() -> Dict[str, Any]:
|
||||
def _save_trigger_data(data: Dict[str, Any]) -> None:
|
||||
"""Save trigger_data.json."""
|
||||
try:
|
||||
TRIGGER_DATA_FILE.parent.mkdir(parents=True, exist_ok=True)
|
||||
with open(TRIGGER_DATA_FILE, 'w') as f:
|
||||
json.dump(data, f, indent=2)
|
||||
atomic_write_json(TRIGGER_DATA_FILE, data)
|
||||
except Exception as exc:
|
||||
_log_warning(f"save trigger data failed: {exc}")
|
||||
return
|
||||
|
||||
@@ -13,6 +13,8 @@ from datetime import datetime, timezone
|
||||
from typing import Dict, List, Any, Optional
|
||||
import inspect
|
||||
|
||||
from aipass.trigger.apps.config import atomic_write_json
|
||||
|
||||
# Infrastructure — redirect to temp dir during tests
|
||||
_test_log_dir = os.environ.get("AIPASS_TEST_LOG_DIR")
|
||||
if _test_log_dir:
|
||||
@@ -132,8 +134,7 @@ def ensure_json_exists(module_name: str, json_type: str) -> bool:
|
||||
|
||||
template = _get_default_template(json_type, module_name)
|
||||
|
||||
with open(json_path, 'w', encoding='utf-8') as f:
|
||||
json.dump(template, f, indent=2, ensure_ascii=False)
|
||||
atomic_write_json(json_path, template, ensure_ascii=False)
|
||||
return True
|
||||
|
||||
|
||||
@@ -173,8 +174,7 @@ def save_json(module_name: str, json_type: str, data: Any) -> bool:
|
||||
if json_type == "data" and isinstance(data, dict):
|
||||
data["last_updated"] = datetime.now().date().isoformat()
|
||||
|
||||
with open(json_path, 'w', encoding='utf-8') as f:
|
||||
json.dump(data, f, indent=2, ensure_ascii=False)
|
||||
atomic_write_json(json_path, data, ensure_ascii=False)
|
||||
return True
|
||||
|
||||
|
||||
|
||||
@@ -29,7 +29,7 @@ import hashlib
|
||||
from datetime import datetime, timedelta
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, Set, Optional, Callable
|
||||
from aipass.trigger.apps.config import TRIGGER_ROOT, AIPASS_PKG_ROOT
|
||||
from aipass.trigger.apps.config import TRIGGER_ROOT, AIPASS_PKG_ROOT, atomic_write_json
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
from aipass.prax.apps.modules.logger import get_direct_logger
|
||||
@@ -150,10 +150,7 @@ def _save_seen_hashes() -> None:
|
||||
if TRIGGER_DATA_FILE.exists():
|
||||
data = json.loads(TRIGGER_DATA_FILE.read_text(encoding='utf-8'))
|
||||
data['seen_error_hashes'] = list(_seen_error_hashes)
|
||||
TRIGGER_DATA_FILE.parent.mkdir(parents=True, exist_ok=True)
|
||||
TRIGGER_DATA_FILE.write_text(
|
||||
json.dumps(data, indent=2), encoding='utf-8'
|
||||
)
|
||||
atomic_write_json(TRIGGER_DATA_FILE, data)
|
||||
except Exception as exc:
|
||||
logger.warning("Failed to save seen hashes: %s", exc)
|
||||
return # Write failure - hashes remain in memory only
|
||||
@@ -195,10 +192,7 @@ def _save_log_positions(positions: Dict[str, int]) -> None:
|
||||
if TRIGGER_DATA_FILE.exists():
|
||||
data = json.loads(TRIGGER_DATA_FILE.read_text(encoding='utf-8'))
|
||||
data['log_positions'] = positions
|
||||
TRIGGER_DATA_FILE.parent.mkdir(parents=True, exist_ok=True)
|
||||
TRIGGER_DATA_FILE.write_text(
|
||||
json.dumps(data, indent=2), encoding='utf-8'
|
||||
)
|
||||
atomic_write_json(TRIGGER_DATA_FILE, data)
|
||||
except Exception as exc:
|
||||
logger.warning("Failed to save log positions: %s", exc)
|
||||
return # Write failure - positions remain in memory only
|
||||
|
||||
@@ -22,7 +22,7 @@ from pathlib import Path
|
||||
from typing import Any, Dict, List
|
||||
|
||||
from aipass.prax.apps.modules.logger import get_direct_logger
|
||||
from aipass.trigger.apps.config import TRIGGER_ROOT
|
||||
from aipass.trigger.apps.config import TRIGGER_ROOT, atomic_write_json
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
logger = get_direct_logger()
|
||||
@@ -59,10 +59,7 @@ def write_config(data: dict) -> bool:
|
||||
True on success, False on failure
|
||||
"""
|
||||
try:
|
||||
TRIGGER_CONFIG_FILE.parent.mkdir(parents=True, exist_ok=True)
|
||||
TRIGGER_CONFIG_FILE.write_text(
|
||||
json.dumps(data, indent=2), encoding='utf-8'
|
||||
)
|
||||
atomic_write_json(TRIGGER_CONFIG_FILE, data)
|
||||
return True
|
||||
except Exception as exc:
|
||||
logger.warning("write_config failed: %s", exc)
|
||||
|
||||
@@ -64,9 +64,11 @@ def _mock_infrastructure(monkeypatch):
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.handlers.log_watcher", mock_log_watcher)
|
||||
|
||||
# -- trigger config -----------------------------------------------------
|
||||
from aipass.trigger.apps.config import atomic_write_json
|
||||
mock_config = MagicMock()
|
||||
mock_config.TRIGGER_ROOT = "/fake/trigger"
|
||||
mock_config.AIPASS_PKG_ROOT = "/fake/aipass"
|
||||
mock_config.atomic_write_json = atomic_write_json
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.config", mock_config)
|
||||
|
||||
# -- CLI console --------------------------------------------------------
|
||||
|
||||
@@ -46,8 +46,10 @@ def _mock_infrastructure(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> Non
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.handlers.json.json_handler", json_mod)
|
||||
|
||||
# -- trigger config (TRIGGER_ROOT) --------------------------------------
|
||||
from aipass.trigger.apps.config import atomic_write_json
|
||||
mock_config = MagicMock()
|
||||
mock_config.TRIGGER_ROOT = tmp_path
|
||||
mock_config.atomic_write_json = atomic_write_json
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.config", mock_config)
|
||||
|
||||
# Force re-import so the module picks up mocked sys.modules
|
||||
|
||||
@@ -63,7 +63,10 @@ def _mock_infrastructure(monkeypatch):
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.handlers.json.json_handler", json_mod)
|
||||
|
||||
# -- trigger config (needed by error_registry import chain) -------------
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.config", MagicMock())
|
||||
from aipass.trigger.apps.config import atomic_write_json
|
||||
config_mock = MagicMock()
|
||||
config_mock.atomic_write_json = atomic_write_json
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.config", config_mock)
|
||||
|
||||
# -- Force re-import so mocks take effect -------------------------------
|
||||
monkeypatch.delitem(sys.modules, "aipass.trigger.apps.handlers.error_reporter", raising=False)
|
||||
|
||||
@@ -51,9 +51,11 @@ def _mock_infrastructure(monkeypatch):
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.handlers.watchers.log_watcher", mock_log_watcher)
|
||||
|
||||
# -- trigger config -----------------------------------------------------
|
||||
from aipass.trigger.apps.config import atomic_write_json
|
||||
mock_config = MagicMock()
|
||||
mock_config.TRIGGER_ROOT = "/fake/trigger"
|
||||
mock_config.AIPASS_PKG_ROOT = "/fake/aipass"
|
||||
mock_config.atomic_write_json = atomic_write_json
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.config", mock_config)
|
||||
|
||||
# -- CLI console --------------------------------------------------------
|
||||
|
||||
@@ -50,9 +50,11 @@ def _mock_infrastructure(monkeypatch):
|
||||
)
|
||||
|
||||
# -- trigger config (TRIGGER_ROOT, AIPASS_PKG_ROOT) ---------------------
|
||||
from aipass.trigger.apps.config import atomic_write_json
|
||||
mock_config = MagicMock()
|
||||
mock_config.TRIGGER_ROOT = Path("/tmp/fake_trigger_root")
|
||||
mock_config.AIPASS_PKG_ROOT = Path("/tmp/fake_aipass_pkg")
|
||||
mock_config.atomic_write_json = atomic_write_json
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.config", mock_config)
|
||||
|
||||
# -- error_registry (report) -------------------------------------------
|
||||
|
||||
@@ -45,8 +45,10 @@ def _mock_infrastructure(monkeypatch):
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.handlers.json.json_handler", json_mod)
|
||||
|
||||
# -- trigger config (TRIGGER_ROOT) --------------------------------------
|
||||
from aipass.trigger.apps.config import atomic_write_json
|
||||
config_mod = MagicMock()
|
||||
config_mod.TRIGGER_ROOT = Path("/tmp/fake_trigger_root")
|
||||
config_mod.atomic_write_json = atomic_write_json
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.config", config_mod)
|
||||
|
||||
# -- Force re-import so mocks take effect -------------------------------
|
||||
|
||||
@@ -18,8 +18,10 @@ def _mock_infrastructure(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> Non
|
||||
"""Mock heavy infrastructure imports."""
|
||||
import sys
|
||||
|
||||
from aipass.trigger.apps.config import atomic_write_json
|
||||
mock_config = MagicMock()
|
||||
mock_config.TRIGGER_ROOT = tmp_path
|
||||
mock_config.atomic_write_json = atomic_write_json
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.config", mock_config)
|
||||
|
||||
mock_json_handler = MagicMock()
|
||||
|
||||
@@ -51,8 +51,10 @@ def _mock_infrastructure(monkeypatch):
|
||||
)
|
||||
|
||||
# -- trigger config (TRIGGER_ROOT) --------------------------------------
|
||||
from aipass.trigger.apps.config import atomic_write_json
|
||||
mock_config = MagicMock()
|
||||
mock_config.TRIGGER_ROOT = Path("/tmp/fake_trigger_root")
|
||||
mock_config.atomic_write_json = atomic_write_json
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.config", mock_config)
|
||||
|
||||
# -- trigger core (trigger.fire) ----------------------------------------
|
||||
|
||||
Reference in New Issue
Block a user