refactor(persistence): 24 hand-rolled atomic JSON/text writers go through utils.atomic_json_write / atomic_write_text

Each copy re-implemented temp+replace by hand and lacked one or more of
fsync, symlink preservation, atomic_replace's Windows-contention retry and
EXDEV/bind-mount fallback, mode preservation, or interrupt-safe temp
cleanup. Three (gateway/session_persistence, cron/suggestions,
agent/shell_hooks) were verbatim inlines of utils._atomic_write; two
modules defined their own directory-fsync helper, now utils.fsync_directory.
plugins/google_meet/_jsonfile.write_json_atomic is deleted (callers use the
canonical helper directly).

Behavior change: every one of these writers now fsyncs the payload, keeps a
pre-existing target's mode, cleans its temp file on BaseException, and
survives Windows AV/indexer contention and cross-device renames the way
config writes already did. cron/suggestions.json is 0600 from creation
(previously chmod'ed after the replace). Skipped on purpose: cron/jobs.py
two-phase staging, gateway/status._write_json_excl (create-only lock),
kanban_transfer staging (not atomic writers); tools/skill_usage.
_write_suppressed_names lives inside a PLUGIN-COMPAT block.
This commit is contained in:
teknium1
2026-09-12 20:12:12 -07:00
committed by Teknium
parent 2be8e6147a
commit 3ef8b384a9
33 changed files with 151 additions and 300 deletions
+2 -12
View File
@@ -13,7 +13,6 @@ import os
import re
import subprocess
import sys
import tempfile
import threading
import time
from contextlib import ExitStack, contextmanager, suppress
@@ -31,7 +30,7 @@ except ImportError: # pragma: no cover
fcntl = None # type: ignore[assignment]
from hermes_constants import get_hermes_home
from utils import atomic_replace
from utils import atomic_json_write
logger = logging.getLogger(__name__)
@@ -469,16 +468,7 @@ def save_allowlist(data: Dict[str, Any]) -> None:
"""Atomic write; on OSError log and keep the in-process approval."""
p = allowlist_path()
try:
p.parent.mkdir(parents=True, exist_ok=True)
fd, tmp_path = tempfile.mkstemp(prefix=f"{p.name}.", suffix=".tmp", dir=str(p.parent))
try:
with os.fdopen(fd, "w", encoding="utf-8") as fh:
fh.write(json.dumps(data, indent=2, sort_keys=True))
atomic_replace(tmp_path, p)
except Exception:
with suppress(OSError):
os.unlink(tmp_path)
raise
atomic_json_write(p, data, sort_keys=True)
except OSError as exc:
logger.warning("Failed to persist shell hook allowlist to %s: %s. The approval is in-memory for this run, "
"but the next startup will re-prompt (or skip registration on non-TTY runs without "
+3 -26
View File
@@ -12,8 +12,6 @@ from __future__ import annotations
import json
import logging
import os
import tempfile
import threading
import uuid
from pathlib import Path
@@ -21,7 +19,7 @@ from typing import Any, Dict, List, Optional
from hermes_constants import get_hermes_home
from hermes_time import now as _hermes_now
from utils import atomic_replace
from utils import atomic_json_write
logger = logging.getLogger(__name__)
@@ -48,13 +46,6 @@ def _current_suggestions_file() -> Path:
return SUGGESTIONS_FILE or (get_hermes_home().resolve() / "cron" / "suggestions.json")
def _secure_file(path: Path) -> None:
try:
os.chmod(path, 0o600)
except OSError:
pass
def _ensure_dir() -> None:
from cron.jobs import _ensure_cron_dir
@@ -81,22 +72,8 @@ def _load_raw() -> Dict[str, Any]:
def _save_raw(suggestions: List[Dict[str, Any]]) -> None:
_ensure_dir()
suggestions_file = _current_suggestions_file()
fd, tmp_path = tempfile.mkstemp(dir=str(suggestions_file.parent), suffix=".tmp", prefix=".sugg_")
try:
with os.fdopen(fd, "w", encoding="utf-8") as f:
payload = {"suggestions": suggestions, "updated_at": _hermes_now().isoformat()}
json.dump(payload, f, indent=2)
f.flush()
os.fsync(f.fileno())
atomic_replace(tmp_path, suggestions_file)
_secure_file(suggestions_file)
except BaseException:
try:
os.unlink(tmp_path)
except OSError:
pass
raise
payload = {"suggestions": suggestions, "updated_at": _hermes_now().isoformat()}
atomic_json_write(_current_suggestions_file(), payload, mode=0o600)
def load_suggestions() -> List[Dict[str, Any]]:
+2 -2
View File
@@ -48,7 +48,7 @@ _ROOM_GRANT_SECRET_FILE = ".room-link-grant-secret"
@lru_cache(maxsize=32)
def _gateway_room_grant_secret_for_home(home_value: str) -> bytes:
"""Load one restart-scoped grant secret for an exact installation root."""
from hermes_cli.install_identity import _fsync_directory
from utils import fsync_directory
(home := Path(home_value)).mkdir(parents=True, exist_ok=True)
path = home / _ROOM_GRANT_SECRET_FILE
def _read() -> bytes:
@@ -75,7 +75,7 @@ def _gateway_room_grant_secret_for_home(home_value: str) -> bytes:
except FileExistsError:
material = _read()
else:
_fsync_directory(home)
fsync_directory(home)
finally:
temporary.unlink(missing_ok=True)
return hmac.new(material, b"hermes-hosted-room-installation-grant-v1", hashlib.sha256).digest()
+2 -4
View File
@@ -15,6 +15,7 @@ import json
import os
import time
from typing import Optional
from utils import atomic_json_write
_MAX_ENTRIES = 1000
_MAX_TEXT_CHARS = 2000
@@ -47,10 +48,7 @@ def _update(chat_id, message_id, fields: dict) -> None:
if len(data) > _MAX_ENTRIES: # trim oldest by timestamp
for k, _ in sorted(data.items(), key=lambda kv: kv[1].get("ts", 0))[: len(data) - _MAX_ENTRIES]:
data.pop(k, None)
tmp = f"{path}.tmp.{os.getpid()}"
with open(tmp, "w", encoding="utf-8") as fh:
json.dump(data, fh, ensure_ascii=False)
os.replace(tmp, path) # atomic; tolerates concurrent writers racing
atomic_json_write(path, data, indent=None) # atomic; tolerates concurrent writers racing
except Exception:
return
+2 -19
View File
@@ -7,12 +7,10 @@ from __future__ import annotations
import contextlib
import logging
import json
import os
import tempfile
import threading
from pathlib import Path
from typing import TYPE_CHECKING, Any, Dict, Optional
from utils import atomic_replace
from utils import atomic_json_write
if TYPE_CHECKING:
from gateway.session import SessionEntry
@@ -475,22 +473,7 @@ class SessionPersistenceMixin:
def _save_sessions_json(self, data: Dict[str, Any]) -> None:
"""Write the legacy sessions.json mirror of the routing index (atomic + fsync)."""
self.sessions_dir.mkdir(parents=True, exist_ok=True)
sessions_file = self.sessions_dir / "sessions.json"
data = {"_README": _SESSIONS_JSON_README, **data}
fd, tmp_path = tempfile.mkstemp(dir=str(self.sessions_dir), suffix=".tmp", prefix=".sessions_")
try:
with os.fdopen(fd, "w", encoding="utf-8") as f:
json.dump(data, f, indent=2)
f.flush()
os.fsync(f.fileno())
atomic_replace(tmp_path, sessions_file)
except BaseException:
try:
os.unlink(tmp_path)
except OSError as e:
logger.debug("Could not remove temp file %s: %s", tmp_path, e)
raise
atomic_json_write(self.sessions_dir / "sessions.json", {"_README": _SESSIONS_JSON_README, **data})
def _save_entries(self) -> None:
"""Snapshot latest state under ``_lock`` and persist after releasing it."""
+2 -11
View File
@@ -20,6 +20,7 @@ from pathlib import Path
from typing import Any, Iterator, Optional
from hermes_constants import get_default_hermes_root, get_hermes_home
from utils import atomic_json_write
logger = logging.getLogger(__name__)
@@ -287,17 +288,7 @@ def _valid_process_start(v: Any) -> bool:
def _write_entries(path: Path, entries: list[dict[str, Any]]) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
tmp = path.with_name(f"{path.name}.{os.getpid()}.{uuid.uuid4().hex}.tmp")
try:
with open(tmp, "w", encoding="utf-8") as fh:
json.dump({"entries": entries}, fh, sort_keys=True)
os.replace(tmp, path)
finally:
try:
tmp.unlink(missing_ok=True)
except OSError:
pass
atomic_json_write(path, {"entries": entries}, indent=None, sort_keys=True)
def _process_start_time(pid: int) -> Optional[float]:
+2 -7
View File
@@ -681,13 +681,8 @@ def save_banner_snapshot(tools: List[dict], enabled_toolsets: List[str], availab
}
def _write():
import tempfile
path = _banner_snapshot_path()
path.parent.mkdir(parents=True, exist_ok=True)
fd, tmp = tempfile.mkstemp(dir=str(path.parent), prefix=".banner_snap.")
with os.fdopen(fd, "w", encoding="utf-8") as fh:
json.dump(payload, fh)
os.replace(tmp, path)
from utils import atomic_json_write
atomic_json_write(_banner_snapshot_path(), payload, indent=None)
_quiet(_write)
+3 -15
View File
@@ -374,21 +374,9 @@ def _build_hermes_tools_mcp_entry() -> dict:
def _write_atomic(target: Path, text: str) -> None:
"""Write via a same-directory temp file + rename (atomic on POSIX, ReplaceFile on Windows) so a
crash mid-write never leaves a half-written config.toml that codex would refuse to load."""
import tempfile
tmp_fd, tmp_path_str = tempfile.mkstemp(prefix=".config.toml.", dir=str(target.parent))
tmp_path = Path(tmp_path_str)
try:
with os.fdopen(tmp_fd, "w", encoding="utf-8") as fh:
fh.write(text)
tmp_path.replace(target)
except Exception:
try:
tmp_path.unlink(missing_ok=True)
except Exception:
pass
raise
"""Atomic rewrite so a crash mid-write never leaves a half-written config.toml codex refuses to load."""
from utils import atomic_write_text
atomic_write_text(target, text, tmp_prefix=".config.toml.", preserve_mode=True)
def migrate(
+2 -6
View File
@@ -16,7 +16,7 @@ from types import SimpleNamespace
from typing import Optional
from hermes_constants import get_hermes_home
from utils import atomic_replace
from utils import atomic_json_write
logger = logging.getLogger(__name__)
@@ -54,12 +54,8 @@ def _load_pending() -> list[dict]:
def _save_pending(entries: list[dict]) -> None:
path = _pending_file()
try:
path.parent.mkdir(parents=True, exist_ok=True)
tmp = path.with_suffix(".json.tmp")
tmp.write_text(json.dumps(entries, indent=2), encoding="utf-8")
atomic_replace(tmp, path)
atomic_json_write(_pending_file(), entries)
except OSError:
pass # non-fatal — worst case the user runs ``hermes debug delete`` manually
+2 -29
View File
@@ -6,12 +6,12 @@ import contextlib
import os
from pathlib import Path
import re
import tempfile
import threading
from typing import Optional
import uuid
from hermes_constants import get_default_hermes_root
from utils import atomic_write_text
_INSTALL_ID_FILENAME = "install_id"
_INSTALL_ID_RE = re.compile(r"^[0-9a-f]{32}$")
@@ -47,22 +47,6 @@ def _install_id_file_lock(root: Path):
os.close(fd)
def _fsync_directory(path: Path) -> None:
"""Best-effort durability for the directory entry after replace."""
if os.name == "nt":
return
try:
fd = os.open(path, os.O_RDONLY | getattr(os, "O_DIRECTORY", 0))
except OSError:
return
try:
os.fsync(fd)
except OSError:
pass
finally:
os.close(fd)
def _read_existing(path: Path) -> tuple[Optional[str], bool]:
"""``(valid id or None, mint?)`` — mint on a missing or malformed file, never on a read failure."""
try:
@@ -92,18 +76,7 @@ def read_or_create_install_id(root: Path | None = None) -> Optional[str]:
existing, mint = _read_existing(path)
if not mint:
return existing
fd, tmp_name = tempfile.mkstemp(dir=str(root), prefix=".install_id-")
try:
with os.fdopen(fd, "w", encoding="utf-8") as handle:
handle.write(uuid.uuid4().hex + "\n")
handle.flush()
os.fsync(handle.fileno())
os.replace(tmp_name, path)
_fsync_directory(root)
except BaseException:
with contextlib.suppress(OSError):
os.unlink(tmp_name)
raise
atomic_write_text(path, uuid.uuid4().hex + "\n", tmp_prefix=".install_id-", fsync_dir=True)
committed = path.read_text(encoding="utf-8").strip().lower()
return committed if _INSTALL_ID_RE.fullmatch(committed) else None
except OSError:
+2 -11
View File
@@ -155,17 +155,8 @@ def generate_presets(models_dir: Path, budget: HardwareBudget, preset_path: Path
body = "\n".join(f"{k} = {v}" for k, v in entry.keys.items())
sections.append(f"[{entry.model_id}]\n{body}\n")
preset_path.parent.mkdir(parents=True, exist_ok=True)
import os
import tempfile
fd, tmp = tempfile.mkstemp(prefix=preset_path.name, suffix=".tmp", dir=preset_path.parent)
try:
with os.fdopen(fd, "w", encoding="utf-8") as stream:
stream.write("\n".join(sections))
os.replace(tmp, preset_path)
finally:
Path(tmp).unlink(missing_ok=True)
from utils import atomic_write_text
atomic_write_text(preset_path, "\n".join(sections), tmp_prefix=f".{preset_path.name}_")
logger.info("wrote %d preset sections to %s", sum(e.keys is not None for e in entries), preset_path)
return entries
+2 -8
View File
@@ -18,7 +18,7 @@ from pathlib import Path
from typing import Any
from hermes_cli import __version__ as _HERMES_VERSION
from utils import atomic_replace
from utils import atomic_json_write
logger = logging.getLogger(__name__)
@@ -154,14 +154,8 @@ def _read_disk_cache() -> tuple[dict[str, Any] | None, float]:
def _write_disk_cache(data: dict[str, Any]) -> None:
path = _cache_path()
try:
path.parent.mkdir(parents=True, exist_ok=True)
tmp = path.with_suffix(path.suffix + ".tmp")
with open(tmp, "w", encoding="utf-8") as fh:
json.dump(data, fh, indent=2)
fh.write("\n")
atomic_replace(tmp, path)
atomic_json_write(_cache_path(), data)
except OSError as exc:
logger.info("model catalog cache write failed: %s", exc)
+2 -4
View File
@@ -27,6 +27,7 @@ import warnings
from dataclasses import dataclass
from pathlib import Path
from typing import Dict, Iterable, List, Optional, Tuple
from utils import atomic_json_write
COMPAT_REMOVAL_DATE = _dt.date(2026, 9, 14)
COMPAT_REMOVAL = COMPAT_REMOVAL_DATE.isoformat()
@@ -247,10 +248,7 @@ def _write_report_file(report: Dict[str, List[Hit]]) -> None:
"written_at": _dt.datetime.now(_dt.timezone.utc).isoformat(timespec="seconds"),
"plugins": {k: [h.__dict__ for h in v] for k, v in report.items()},
"lines": summary_lines(report)}
p.parent.mkdir(parents=True, exist_ok=True)
tmp = p.with_suffix(".tmp")
tmp.write_text(json.dumps(payload, indent=1), encoding="utf-8")
os.replace(tmp, p)
atomic_json_write(p, payload, indent=1)
except Exception:
pass
+2 -7
View File
@@ -1621,18 +1621,13 @@ def _plugin_toolset_keys_cache_path() -> Path:
def _persist_plugin_toolset_keys() -> None:
"""Persist discovered plugin toolset keys + portable MCP names (best-effort)."""
try:
import tempfile
from utils import atomic_json_write
keys = sorted({ts_key for ts_key, _, _ in get_plugin_toolsets()})
try:
portable = sorted(get_plugin_manager().get_portable_mcp_servers())
except Exception:
portable = []
path = _plugin_toolset_keys_cache_path()
path.parent.mkdir(parents=True, exist_ok=True)
fd, tmp = tempfile.mkstemp(dir=str(path.parent), prefix=".pt_keys.")
with os.fdopen(fd, "w", encoding="utf-8") as fh:
json.dump({"toolset_keys": keys, "portable_mcp": portable}, fh)
os.replace(tmp, path)
atomic_json_write(_plugin_toolset_keys_cache_path(), {"toolset_keys": keys, "portable_mcp": portable}, indent=None)
except Exception:
logger.debug("plugin toolset key persist failed", exc_info=True)
+3 -6
View File
@@ -1591,15 +1591,12 @@ def import_profile(archive_path: str, name: Optional[str] = None) -> Path:
# Rename
def _atomic_write_json(path: Path, data: dict) -> bool:
"""Write *data* to *path* via a sibling ``.tmp`` + rename. Returns False (tmp cleaned) on OSError."""
tmp = path.with_suffix(path.suffix + ".tmp")
"""Atomic rewrite of a third-party JSON config; False on OSError (nothing partially written)."""
from utils import atomic_json_write
try:
tmp.write_text(json.dumps(data, indent=2, ensure_ascii=False) + "\n", encoding="utf-8")
tmp.replace(path)
atomic_json_write(path, data)
return True
except OSError:
with contextlib.suppress(OSError):
tmp.unlink(missing_ok=True)
return False
+2 -3
View File
@@ -11,6 +11,7 @@ import sys
import time
from pathlib import Path
from typing import Optional
from utils import atomic_json_write
# Multiplexer / terminal-emulator identity env vars, checked in order when no real tty path is
# available (e.g. stdin piped but stdout still a pty owned by a known terminal).
@@ -86,9 +87,7 @@ def write_breadcrumb(session_id: str, cwd: Optional[str] = None) -> None:
directory.mkdir(parents=True, exist_ok=True)
now = time.time()
payload = {"session_id": session_id, "cwd": cwd or os.getcwd(), "ts": now}
tmp = directory / f".{terminal_id}.tmp"
tmp.write_text(json.dumps(payload), encoding="utf-8")
os.replace(tmp, directory / terminal_id)
atomic_json_write(directory / terminal_id, payload, indent=None)
_prune_stale(directory, now)
except Exception:
pass
+2 -19
View File
@@ -4,15 +4,14 @@
import asyncio
import logging
import ipaddress
import json
import os
import subprocess
import sys
import tempfile
import threading
import time
from pathlib import Path
from typing import TYPE_CHECKING, Any, Optional
from utils import atomic_json_write
if TYPE_CHECKING: # pragma: no cover - annotation only
import uvicorn
@@ -212,25 +211,9 @@ def _write_dashboard_ready_file(actual_port: int) -> None:
if not target:
return
tmp_name = ""
try:
path = Path(target)
path.parent.mkdir(parents=True, exist_ok=True)
payload = json.dumps({"port": int(actual_port)}, separators=(",", ":"))
with tempfile.NamedTemporaryFile(
"w", encoding="utf-8", dir=str(path.parent), prefix=f"{path.name}.", suffix=".tmp", delete=False
) as fh:
fh.write(payload)
fh.flush()
os.fsync(fh.fileno())
tmp_name = fh.name
os.replace(tmp_name, path)
atomic_json_write(Path(target), {"port": int(actual_port)}, indent=None, separators=(",", ":"))
except Exception as exc:
if tmp_name:
try:
Path(tmp_name).unlink(missing_ok=True)
except Exception:
pass
_log.warning("Failed to write dashboard ready file %r: %s", target, exc)
+2 -11
View File
@@ -19,6 +19,7 @@ from pathlib import Path
from typing import Dict, Optional
from hermes_constants import get_hermes_home
from utils import atomic_json_write
logger = logging.getLogger("cli")
@@ -458,21 +459,11 @@ def _load_worktree_merge_cache() -> Dict[str, bool]:
def _save_worktree_merge_cache(verdicts: Dict[str, bool]) -> None:
"""Atomically persist the newest ``_WORKTREE_MERGE_CACHE_MAX`` verdicts. Never raises."""
path = _worktree_merge_cache_path()
tmp = None
try:
items = list(verdicts.items())[-_WORKTREE_MERGE_CACHE_MAX:]
path.parent.mkdir(parents=True, exist_ok=True)
tmp = path.with_suffix(f".{os.getpid()}.tmp")
tmp.write_text(json.dumps({"version": 1, "verdicts": dict(items)}), encoding="utf-8")
os.replace(str(tmp), str(path))
atomic_json_write(_worktree_merge_cache_path(), {"version": 1, "verdicts": dict(items)}, indent=None)
except Exception as e:
logger.debug("Could not persist worktree merge cache: %s", e)
if tmp is not None:
try:
tmp.unlink()
except Exception:
pass
def _worktree_commits_all_merged_upstream(
+1 -14
View File
@@ -1,8 +1,7 @@
"""Tiny JSON file helpers shared by the bot, process manager, node registry and node server."""
"""Tiny JSON read helper shared by the bot, process manager, node registry and node server."""
from __future__ import annotations
import contextlib
import json
from pathlib import Path
from typing import Any, Optional
@@ -16,15 +15,3 @@ def read_json(path: Path) -> Optional[Any]:
return json.loads(path.read_text(encoding="utf-8"))
except (OSError, ValueError):
return None
def write_json_atomic(path: Path, data: Any, mode: Optional[int] = None) -> None:
"""Write ``json.dumps(data, indent=2)`` via a ``.json.tmp`` sibling + rename; *mode* (e.g.
``0o600``) is applied to the temp file so the final file never exists with looser perms."""
path.parent.mkdir(parents=True, exist_ok=True)
tmp = path.with_suffix(".json.tmp")
tmp.write_text(json.dumps(data, indent=2), encoding="utf-8")
if mode is not None:
with contextlib.suppress(OSError, NotImplementedError): # best-effort on non-POSIX filesystems
tmp.chmod(mode)
tmp.replace(path)
+2 -2
View File
@@ -22,7 +22,7 @@ from pathlib import Path
from types import SimpleNamespace
from typing import Optional
from plugins.google_meet._jsonfile import write_json_atomic
from utils import atomic_json_write
# Short three-segment code, a lookup URL, or /new. Anything else is rejected.
MEET_URL_RE = re.compile(
@@ -93,7 +93,7 @@ class _BotState:
def _flush(self) -> None:
data = {key: getattr(self, attr) if attr else None for key, attr, _ in _STATUS_FIELDS}
data.update(transcriptPath=str(self.transcript_path), pid=os.getpid()) # keeps table key order
write_json_atomic(self.status_path, data)
atomic_json_write(self.status_path, data)
def set(self, **kwargs) -> None:
self.__dict__.update(kwargs)
+3 -2
View File
@@ -13,7 +13,8 @@ from typing import Any, Dict, List, Optional
from hermes_constants import get_hermes_home
from plugins.google_meet._jsonfile import read_json, write_json_atomic
from plugins.google_meet._jsonfile import read_json
from utils import atomic_json_write
def _default_path() -> Path:
@@ -33,7 +34,7 @@ class NodeRegistry:
return nodes if isinstance(nodes, dict) else {}
def _save(self, nodes: Dict[str, Dict[str, Any]]) -> None:
write_json_atomic(self.path, {"nodes": nodes})
atomic_json_write(self.path, {"nodes": nodes})
def get(self, name: str) -> Optional[Dict[str, Any]]:
entry = self._load().get(name)
+3 -2
View File
@@ -17,7 +17,8 @@ from pathlib import Path
from typing import Any, Dict, Optional
from hermes_constants import get_hermes_home
from plugins.google_meet._jsonfile import read_json, write_json_atomic
from plugins.google_meet._jsonfile import read_json
from utils import atomic_json_write
from plugins.google_meet.node import protocol as _proto
_START_BOT_KEYS = ("url", "guest_name", "duration", "headed", "auth_state", "session_id", "out_dir")
@@ -80,7 +81,7 @@ class NodeServer:
if not (isinstance(tok, str) and tok):
tok = secrets.token_hex(16) # 32 hex chars
# Owner-only: the token grants full RPC access to the meet bot.
write_json_atomic(self.token_path, {"token": tok, "generated_at": time.time()}, mode=0o600)
atomic_json_write(self.token_path, {"token": tok, "generated_at": time.time()}, mode=0o600)
self._token = tok
return tok
+3 -2
View File
@@ -21,7 +21,8 @@ from typing import Any, Dict, Optional
from hermes_constants import get_hermes_home
from plugins.google_meet._jsonfile import read_json, write_json_atomic
from plugins.google_meet._jsonfile import read_json
from utils import atomic_json_write
def _root() -> Path:
@@ -33,7 +34,7 @@ def _read_active() -> Optional[Dict[str, Any]]:
def _write_active(data: Dict[str, Any]) -> None:
write_json_atomic(_root() / ".active.json", data)
atomic_json_write(_root() / ".active.json", data)
def _pid_alive(pid: int) -> bool:
+3 -2
View File
@@ -452,14 +452,15 @@ class TestAllowlistConcurrency:
p.parent.mkdir(parents=True, exist_ok=True)
tmp_paths_seen: list = []
real_mkstemp = shell_hooks.tempfile.mkstemp
import utils
real_mkstemp = utils.tempfile.mkstemp
def spying_mkstemp(*args, **kwargs):
fd, path = real_mkstemp(*args, **kwargs)
tmp_paths_seen.append(path)
return fd, path
monkeypatch.setattr(shell_hooks.tempfile, "mkstemp", spying_mkstemp)
monkeypatch.setattr(utils.tempfile, "mkstemp", spying_mkstemp)
shell_hooks.save_allowlist({"approvals": [{"event": "a", "command": "x"}]})
shell_hooks.save_allowlist({"approvals": [{"event": "b", "command": "y"}]})
@@ -79,13 +79,12 @@ class TestTomlValueFormatter:
"""If rename fails partway through (out of disk, permissions,
crash), the temp file must be cleaned up. Otherwise repeated
failed migrations would pile up .config.toml.* files."""
from pathlib import Path as _Path
original_replace = _Path.replace
import utils
def failing_replace(self, target):
def failing_replace(tmp, target):
raise OSError("simulated disk full")
monkeypatch.setattr(_Path, "replace", failing_replace)
monkeypatch.setattr(utils, "atomic_replace", failing_replace)
report = migrate(
{"mcp_servers": {"x": {"command": "y"}}},
codex_home=tmp_path,
+3 -2
View File
@@ -21,14 +21,15 @@ def _race_first_install_id(
if start_barrier is not None:
start_barrier.wait(timeout=10)
if writer_entered is not None:
original_mkstemp = install_identity.tempfile.mkstemp
import utils
original_mkstemp = utils.tempfile.mkstemp
def held_mkstemp(*args, **kwargs):
writer_entered.set()
assert release_writer.wait(timeout=10)
return original_mkstemp(*args, **kwargs)
install_identity.tempfile.mkstemp = held_mkstemp
utils.tempfile.mkstemp = held_mkstemp
results.put(read_or_create_install_id(root))
+69
View File
@@ -0,0 +1,69 @@
"""Invariant for the non-secret atomic JSON writers that used to inline ``utils._atomic_write``.
Contract under test for ``gateway/session_persistence``, ``cron/suggestions`` and
``agent/shell_hooks``: a failed replace leaves the previous file byte-identical AND leaves no temp
file behind (the interrupt-safe cleanup only the canonical helper guarantees). A hand-rolled copy
that skips the cleanup, or writes through the target instead of a sibling temp, fails this.
"""
from __future__ import annotations
import json
from pathlib import Path
import pytest
def _failing_replace(tmp, target):
raise OSError("simulated disk full")
def _leftovers(directory: Path, keep: str) -> list[str]:
return sorted(p.name for p in directory.iterdir() if p.name != keep)
@pytest.fixture
def broken_replace(monkeypatch):
import utils
monkeypatch.setattr(utils, "atomic_replace", _failing_replace)
def test_sessions_json_failed_replace_keeps_old_bytes_and_no_temp(tmp_path, monkeypatch, broken_replace):
from gateway.session_persistence import SessionPersistenceMixin
sessions_dir = tmp_path / "sessions"
sessions_dir.mkdir()
target = sessions_dir / "sessions.json"
target.write_text('{"old": true}', encoding="utf-8")
store = SessionPersistenceMixin()
store.sessions_dir = sessions_dir
with pytest.raises(OSError):
store._save_sessions_json({"k": "v"})
assert target.read_text(encoding="utf-8") == '{"old": true}'
assert _leftovers(sessions_dir, "sessions.json") == []
def test_suggestions_failed_replace_keeps_old_bytes_and_no_temp(tmp_path, monkeypatch, broken_replace):
from cron import suggestions
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
target = suggestions._current_suggestions_file()
target.parent.mkdir(parents=True, exist_ok=True)
target.write_text('{"old": true}', encoding="utf-8")
with pytest.raises(OSError):
suggestions._save_raw([{"id": 1}])
assert target.read_text(encoding="utf-8") == '{"old": true}'
assert _leftovers(target.parent, target.name) == []
def test_shell_hooks_allowlist_survives_failed_replace_without_temp(tmp_path, monkeypatch, broken_replace):
from agent import shell_hooks
monkeypatch.setenv("HERMES_HOME", str(tmp_path / "home"))
target = shell_hooks.allowlist_path()
target.parent.mkdir(parents=True, exist_ok=True)
target.write_text(json.dumps({"approvals": []}), encoding="utf-8")
shell_hooks.save_allowlist({"approvals": [{"event": "a", "command": "x"}]}) # logs, never raises
assert json.loads(target.read_text(encoding="utf-8")) == {"approvals": []}
assert _leftovers(target.parent, target.name) == []
+5 -24
View File
@@ -10,10 +10,11 @@ from __future__ import annotations
import json
import os
import re
import tempfile
import time
import uuid
from contextlib import contextmanager
from utils import atomic_json_write, fsync_directory
from pathlib import Path
from typing import Any
@@ -73,25 +74,14 @@ def _root(home: Path | str) -> Path:
return Path(home).resolve() / "runtime" / DELIVERY_DIR_NAME
def _fsync_dir(path: Path) -> None:
# Windows cannot open directories with os.open; file fsync still applies.
if os.name == "nt":
return
fd = os.open(path, os.O_RDONLY)
try:
os.fsync(fd)
finally:
os.close(fd)
@contextmanager
def _locked(home: Path | str):
root = _root(home)
root.parent.mkdir(parents=True, exist_ok=True)
root.mkdir(mode=0o700, exist_ok=True)
root.chmod(0o700)
_fsync_dir(root.parent)
_fsync_dir(root.parent.parent)
fsync_directory(root.parent)
fsync_directory(root.parent.parent)
lock = root / ".lock"
fd = os.open(lock, os.O_CREAT | os.O_WRONLY, 0o600)
os.close(fd)
@@ -107,16 +97,7 @@ def _read(path: Path) -> dict[str, Any] | None:
def _write(path: Path, record: dict[str, Any]) -> None:
fd, temporary = tempfile.mkstemp(dir=path.parent, prefix=".delivery-")
try:
with os.fdopen(fd, "w", encoding="utf-8") as stream:
json.dump(record, stream, ensure_ascii=False, sort_keys=True)
stream.flush()
os.fsync(stream.fileno())
os.replace(temporary, path)
_fsync_dir(path.parent)
finally:
Path(temporary).unlink(missing_ok=True)
atomic_json_write(path, record, indent=None, sort_keys=True, fsync_dir=True)
def deliver_to_live_owner(
+3 -4
View File
@@ -407,9 +407,8 @@ def _run_local_turn(argv: list[str], dm_file: str, *, env: Optional[dict[str, st
def _admit_live_dm(profile_home: Path | None, dm_file: str, author: Optional[dict] = None) -> dict | None:
"""Pin intent before admission; retries may inspect, never change transport."""
from tools.bot_live_delivery import (
_fsync_dir, deliver_to_live_owner, find_canonical_live_owner, read_delivery_result,
)
from tools.bot_live_delivery import deliver_to_live_owner, find_canonical_live_owner, read_delivery_result
from utils import fsync_directory
intent: dict[str, Any]
intent_path = Path(dm_file + ".live.json")
@@ -432,7 +431,7 @@ def _admit_live_dm(profile_home: Path | None, dm_file: str, author: Optional[dic
json.dump(intent, stream)
stream.flush()
os.fsync(stream.fileno())
_fsync_dir(intent_path.parent)
fsync_directory(intent_path.parent)
home = intent["owner"]["profile_home"]
record = read_delivery_result(home, intent["delivery_id"])
if record is None:
+6 -17
View File
@@ -21,13 +21,13 @@ import re
import shlex
import shutil
import sys
import tempfile
import time
import uuid
from pathlib import Path
from typing import Any, Iterator, Optional
from tools.bot_mode_probe import _default_home, _hermes_root
from utils import atomic_json_write
logger = logging.getLogger(__name__)
@@ -90,17 +90,8 @@ def _ensure_dirs(root: Path | str) -> Path:
return base
def _atomic_write_json(target: Path, payload: Any, *, prefix: str, sort_keys: bool = False) -> None:
"""tempfile + os.replace so readers never see a partial file; tempfile removed on failure."""
fd, tmp = tempfile.mkstemp(dir=str(target.parent), prefix=prefix, suffix=".tmp")
try:
with os.fdopen(fd, "w", encoding="utf-8") as f:
json.dump(payload, f, ensure_ascii=False, sort_keys=sort_keys)
os.replace(tmp, target)
except Exception:
with contextlib.suppress(OSError):
os.unlink(tmp)
raise
def _atomic_write_json(target: Path, payload: Any, *, sort_keys: bool = False) -> None:
atomic_json_write(target, payload, indent=None, sort_keys=sort_keys)
def _bot_mode_cfg(key: str, *, loader: str) -> Any:
@@ -145,8 +136,7 @@ def write_remote_roster(root: Path | str, rows: Any) -> int:
for norm in filter(None, map(_normalize_roster_row, rows if isinstance(rows, list) else [])):
by_key.setdefault((norm["connection_id"], norm["profile"]), norm)
cleaned = [by_key[k] for k in sorted(by_key)]
_atomic_write_json(base / ROSTER_FILE, {"updated_at": int(time.time()), "agents": cleaned},
prefix=".roster-", sort_keys=True)
_atomic_write_json(base / ROSTER_FILE, {"updated_at": int(time.time()), "agents": cleaned}, sort_keys=True)
return len(cleaned)
@@ -229,7 +219,7 @@ def enqueue_envelope(root: Path | str, *, target: dict, message: str, sender_pro
"target_connection": target["connection_id"], "target_profile": target["profile"],
"target_handle": target["handle"], "message": message,
}
_atomic_write_json(base / OUTBOX_DIR / f"{envelope['id']}.json", envelope, prefix=".env-")
_atomic_write_json(base / OUTBOX_DIR / f"{envelope['id']}.json", envelope)
return envelope
@@ -289,8 +279,7 @@ def write_reply(root: Path | str, envelope_id: str, *, reply: str = "", error: s
code = classify_agent_error(err)
path = base / REPLIES_DIR / f"{safe}.json"
_atomic_write_json(path, {"id": safe, "at": int(time.time()), "reply": str(reply or ""), "error": err, "reason": code},
prefix=".rep-")
_atomic_write_json(path, {"id": safe, "at": int(time.time()), "reply": str(reply or ""), "error": err, "reason": code})
return path
+4 -6
View File
@@ -10,7 +10,6 @@ Lives here, not in tool dispatch, so hits sit *after* every safety check and ski
import hashlib
import json
import logging
import os
import re
import threading
import time
@@ -18,6 +17,7 @@ from contextlib import suppress
from pathlib import Path
from typing import Dict, Optional, Tuple
from urllib.parse import urlparse
from utils import atomic_json_write
logger = logging.getLogger(__name__)
@@ -183,11 +183,9 @@ def _save_index(index: dict) -> None:
if len(index) > _INDEX_MAX_ENTRIES:
newest = sorted(index.items(), key=lambda kv: kv[1].get("fetched_at", 0), reverse=True)
index = dict(newest[:_INDEX_MAX_ENTRIES])
# Per-process tmp name: CLI, gateway, cron, and subagents all write this index; a shared tmp name
# would let concurrent writers truncate each other. os.replace is atomic: worst case is a lost insert.
tmp = path.with_suffix(f".tmp.{os.getpid()}")
tmp.write_text(json.dumps(index), encoding="utf-8")
tmp.replace(path)
# CLI, gateway, cron, and subagents all write this index; the replace is atomic, so the worst case
# under concurrent writers is a lost insert, never a truncated index.
atomic_json_write(path, index, indent=None)
except Exception as exc: # noqa: BLE001
logger.debug("Failed to save web extract cache index: %s", exc)
+2 -6
View File
@@ -14,7 +14,6 @@ from __future__ import annotations
import difflib
import json
import logging
import os
import re
import time
import uuid
@@ -24,6 +23,7 @@ from pathlib import Path
from typing import Any, Dict, List, Optional
from hermes_constants import get_hermes_home
from utils import atomic_json_write
logger = logging.getLogger(__name__)
@@ -82,11 +82,7 @@ def stage_write(subsystem: str, payload: Dict[str, Any], *, summary: str, origin
"created_at": time.time(), "payload": payload,
}
try:
path = _pending_path(subsystem, pid)
path.parent.mkdir(parents=True, exist_ok=True)
tmp = path.with_suffix(".json.tmp")
tmp.write_text(json.dumps(record, ensure_ascii=False, indent=2), encoding="utf-8")
os.replace(tmp, path)
atomic_json_write(_pending_path(subsystem, pid), record)
except Exception as e: # pragma: no cover - disk failure path
logger.error("Failed to stage pending %s write: %s", subsystem, e, exc_info=True)
return record
+2 -13
View File
@@ -8,15 +8,13 @@ marker" instead of raising."""
from __future__ import annotations
import contextlib
import json
import logging
import os
import tempfile
import threading
import time
from pathlib import Path
from typing import Any
from utils import atomic_json_write
logger = logging.getLogger(__name__)
@@ -59,16 +57,7 @@ def _store(path: Path, entries: dict[str, dict]) -> None:
if not entries:
path.unlink(missing_ok=True)
return
path.parent.mkdir(parents=True, exist_ok=True)
fd, tmp = tempfile.mkstemp(dir=path.parent, prefix=".turn-marker-")
try:
with os.fdopen(fd, "w", encoding="utf-8") as f:
json.dump(entries, f)
os.replace(tmp, path)
except Exception:
with contextlib.suppress(OSError):
os.unlink(tmp)
raise
atomic_json_write(path, entries, indent=None)
def _update(home: Path | str, session_key: str, mutate, what: str) -> None: