From 5bba1b930e1163dfc4929fe6496e91c42348200a Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Sun, 7 Jun 2026 06:01:12 -0600 Subject: [PATCH] feat(ci): add journal_io raw-mechanic gate; migrate 4 journal-data stragglers MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Add scripts/check_journal_io_mechanic.py: a static AST check that flags raw durable-write / locked-RMW mechanics (os.replace, filesystem Path.replace, mkstemp+manual-replace, flock LOCK_EX RMW) against owner journal data anywhere outside solstone/think/journal_io/. Unlike the import-binding access lint, this gate SCANS owners — they may import journal_io primitives but must carry no hand-rolled mechanics. Seeds with an empty allowlist; wired into make ci via install-checks alongside the access and layer-hygiene checks. To seed the gate empty and green, migrate every remaining journal-data straggler to journal_io (byte-identical output, equivalent locking): - link/paths.py: state/token/totp writes -> write_json (token/totp keep 0o600) - convey/chat_stream.py: chat.jsonl -> atomic_replace (ensure_ascii=False house pattern); declared the L2 owner of chronicle chat//chat.jsonl - activities.py: _write_jsonl_records -> atomic_replace; locked_modify -> hold_lock (drops the bare-LOCK_EX blocking RMW + max_retries retry loop) - importers/shared.py: chronicle import./imported.jsonl -> atomic_replace (a 4th straggler the original inventory missed; surfaced in prep) locked_modify now raises journal_io.LockTimeout on contention; surface it in sol's voice at the activities CLI write boundary via a new ACTIVITIES_BUSY reason, mirroring the entities ENTITY_BUSY pattern. Tests reassert locking and byte-fidelity invariants instead of the old flock/mkstemp internals. Co-Authored-By: Claude Opus 4.8 (1M context) --- AGENTS.md | 1 + Makefile | 9 +- scripts/check_journal_io_access.py | 1 + scripts/check_journal_io_mechanic.py | 427 ++++++++++++++++++++++++ solstone/apps/activities/call.py | 66 ++-- solstone/convey/chat_stream.py | 18 +- solstone/convey/reasons.py | 5 + solstone/think/activities.py | 65 +--- solstone/think/importers/shared.py | 13 +- solstone/think/link/paths.py | 39 +-- tests/link/test_paths.py | 8 +- tests/test_activities_locking.py | 11 +- tests/test_check_journal_io_mechanic.py | 245 ++++++++++++++ 13 files changed, 773 insertions(+), 135 deletions(-) create mode 100644 scripts/check_journal_io_mechanic.py create mode 100644 tests/test_check_journal_io_mechanic.py diff --git a/AGENTS.md b/AGENTS.md index 35244f92d..13f016d26 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -195,6 +195,7 @@ Each domain has exactly **one** write-owning module (or one tightly-scoped famil | Awareness activity state (`awareness/activity_state.json`) | `solstone/think/thinking.py` | | Identity (`identity/*.md`, `identity/history.jsonl` audit log) | `solstone/think/identity.py` | | Todos (`facets/*/todos/*.jsonl`) | `solstone/apps/todos/todo.py` + `solstone/apps/todos/call.py` | +| Chronicle chat stream (`chronicle/**/chat//chat.jsonl`) | `solstone/convey/chat_stream.py` | | Config (`config/journal.json`) | `solstone/think/journal_config.py` | | Push devices (`config/push_devices.json`) | `solstone/think/push/devices.py` | | Convey config (`config/convey.json`) | `solstone/convey/config.py` + `solstone/think/facets.py` | diff --git a/Makefile b/Makefile index 88d04d90a..82d0cc738 100644 --- a/Makefile +++ b/Makefile @@ -14,7 +14,7 @@ export TMPDIR := /var/tmp PYTEST_BASETEMP_INIT := BASETEMP=$$(mktemp -d /var/tmp/solstone-pytest-XXXXXX); trap 'rm -rf "$$BASETEMP"' EXIT INT TERM; PYTEST_BASETEMP_FLAG := --basetemp "$$BASETEMP" -.PHONY: install uninstall test test-cov test-app test-only format format-check install-checks ci clean clean-install coverage watch versions update update-prices preflight pre-commit skills dev all sandbox sandbox-stop install-pinchtab install-models parakeet-helper parakeet-helper-clean wheel-macos wheel-macos-clean verify-browser update-browser-baselines review verify verify-api update-api-baselines service-logs check-layer-hygiene check-api-conventions check-journal-io-access smoke-cogitate release release-test FORCE +.PHONY: install uninstall test test-cov test-app test-only format format-check install-checks ci clean clean-install coverage watch versions update update-prices preflight pre-commit skills dev all sandbox sandbox-stop install-pinchtab install-models parakeet-helper parakeet-helper-clean wheel-macos wheel-macos-clean verify-browser update-browser-baselines review verify verify-api update-api-baselines service-logs check-layer-hygiene check-api-conventions check-journal-io-access check-journal-io-mechanic smoke-cogitate release release-test FORCE # Default target - install package in editable mode all: install @@ -461,6 +461,9 @@ install-checks: .installed @echo "=== Running journal-io access check ===" @$(MAKE) check-journal-io-access @echo "" + @echo "=== Running journal-io mechanic check ===" + @$(MAKE) check-journal-io-mechanic + @echo "" @echo "=== Checking extras consistency ===" @$(VENV_BIN)/python scripts/check_extras_consistency.py @echo "" @@ -531,6 +534,10 @@ check-api-conventions: .installed check-journal-io-access: .installed $(VENV_BIN)/python scripts/check_journal_io_access.py +# Journal raw-mechanic check (see AGENTS.md §7 L2) +check-journal-io-mechanic: .installed + $(VENV_BIN)/python scripts/check_journal_io_mechanic.py + # Re-run the live four-backend integrated-façade cogitate smoke. Spawns the # archived runner (extro `vpe/workspace/archived/`) against this venv so the # real openhands-sdk Agent path is exercised end-to-end. Requires real API diff --git a/scripts/check_journal_io_access.py b/scripts/check_journal_io_access.py index 0604671e0..78bdb7488 100644 --- a/scripts/check_journal_io_access.py +++ b/scripts/check_journal_io_access.py @@ -133,6 +133,7 @@ OWNER_FILES: frozenset[str] = frozenset( "solstone/think/importers/cli.py", "solstone/think/importers/documents.py", "solstone/think/importers/shared.py", + "solstone/convey/chat_stream.py", # imports/** bundle + sync-cursor writers (local/CLI import flows). "solstone/think/importers/plaud.py", # streamed imported-audio install. "solstone/think/importers/sync.py", diff --git a/scripts/check_journal_io_mechanic.py b/scripts/check_journal_io_mechanic.py new file mode 100644 index 000000000..e3ec8f49c --- /dev/null +++ b/scripts/check_journal_io_mechanic.py @@ -0,0 +1,427 @@ +#!/usr/bin/env python3 +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Journal raw-mechanic lint. + +This check flags raw durable-write and locked read-modify-write mechanics used +against owner journal data anywhere outside ``solstone.think.journal_io``. +Unlike ``check_journal_io_access.py``, owner modules are intentionally scanned: +owners may import journal_io primitives, but after migration they must not carry +their own raw replacement or bare exclusive-flock mechanics. + +Flagged mechanics: + + D1 - ``os.replace``. Calls are resolved through ``import os`` / aliases and + ``from os import replace`` bindings. Bare attribute-name matching is forbidden. + + D2 - ``Path.replace`` heuristic. The detector flags one-argument + ``.replace(...)`` calls only when the receiver looks like a temp path: + direct ``.with_suffix/.with_name/.with_stem(...)`` receiver, a receiver name + containing ``tmp``/``temp``, or a receiver name bound from known temp/path + creation patterns. This deliberately avoids string/date ``replace`` calls. + Limit: a non-temp-named Path variable with one non-string argument is a false + negative; in this repository the ``os.replace`` rule and temp-name tracking + cover every real durable-write site. The temp-substring receiver-name match + can also over-match names like ``template`` or ``attempt``. + + D3 - ``flock(LOCK_EX)``. ``fcntl.flock`` calls with a lock-mode subtree + containing ``LOCK_EX`` and not ``LOCK_NB`` are flagged. ``LOCK_UN`` and + ``LOCK_EX | LOCK_NB`` are not violations. + +The check ships green with an empty allowlist. The allowlist is keyed by +``(file, kind)`` with an allowed count so it can ratchet down, matching the +existing journal_io access check. +""" + +from __future__ import annotations + +import argparse +import ast +import sys +from pathlib import Path + +ROOT = Path(__file__).resolve().parent.parent + +PATH_TEMP_METHODS: frozenset[str] = frozenset({"with_suffix", "with_name", "with_stem"}) +VIOLATION_KINDS: frozenset[str] = frozenset( + {"os.replace", "Path.replace", "flock(LOCK_EX)"} +) + +EXCLUDED_FILES: frozenset[str] = frozenset( + { + # Ops/runtime state and single-process guards. + "solstone/think/scheduler.py", + "solstone/think/supervisor.py", + "solstone/think/readiness.py", + "solstone/think/providers/state.py", + "solstone/think/providers/local_install.py", + "solstone/think/providers/mlx_install.py", + "solstone/think/services/scout.py", + "solstone/think/services/spl.py", + "solstone/think/steward.py", + "solstone/talent/steward.py", + "solstone/think/install_guard.py", + "solstone/think/install_models.py", + "solstone/think/setup.py", + "solstone/think/user_config.py", + "solstone/think/voice/brain.py", + "solstone/think/sync_check.py", + "solstone/think/runner.py", + "solstone/think/providers_cli.py", + "solstone/think/start.py", + "solstone/talent/daily_schedule.py", + "solstone/think/routines.py", + "solstone/apps/sol/maint/005_migrate_dream_to_think_schedules.py", + "solstone/apps/timeline/maint/001_register_schedules.py", + # App-storage and temporary upload/transcription files. + "solstone/apps/import/routes.py", + "solstone/apps/support/routes.py", + "solstone/observe/transcribe/whisper.py", + "solstone/observe/transcribe/_parakeet_coreml.py", + "solstone/observe/transcribe/revai.py", + "solstone/think/journal_export.py", + # UI/pipeline runtime state, not owner journal content. + "solstone/apps/home/routes.py", + "solstone/apps/transcripts/routes.py", + "solstone/think/data_state.py", + "solstone/think/skills_cli.py", + } +) + +EXCLUDED_PREFIXES: tuple[str, ...] = () + +# Empty by design. If this check flags owner journal-data, migrate the site. If +# it flags non-owner runtime storage, classify it in EXCLUDED_FILES. +ALLOWLIST: dict[tuple[str, str], int] = {} + + +def _is_test_file(rel: Path) -> bool: + return ( + "tests" in rel.parts + or rel.name == "conftest.py" + or (rel.name.startswith("test_") and rel.suffix == ".py") + ) + + +def _is_excluded(rel: Path) -> bool: + rel_str = rel.as_posix() + return rel_str in EXCLUDED_FILES or any( + rel_str.startswith(prefix) for prefix in EXCLUDED_PREFIXES + ) + + +def discover_modules(root: Path) -> list[Path]: + """Return posix-relative scanned modules under ``solstone/``.""" + scope = root / "solstone" + if not scope.is_dir(): + return [] + + found: list[Path] = [] + for path in sorted(scope.rglob("*.py")): + rel = path.relative_to(root) + rel_str = rel.as_posix() + if "__pycache__" in rel.parts: + continue + if rel_str.startswith("solstone/think/journal_io/"): + continue + if _is_test_file(rel): + continue + if _is_excluded(rel): + continue + found.append(rel) + return found + + +def _call_detail(func: ast.expr) -> str: + try: + return ast.unparse(func) + except Exception: + return "" + + +def _is_attr_call(expr: ast.AST, attrs: frozenset[str]) -> bool: + return ( + isinstance(expr, ast.Call) + and isinstance(expr.func, ast.Attribute) + and expr.func.attr in attrs + ) + + +def _is_tempfile_call( + expr: ast.AST, + name: str, + tempfile_aliases: set[str], + direct_names: set[str], +) -> bool: + if not isinstance(expr, ast.Call): + return False + func = expr.func + if isinstance(func, ast.Name): + return func.id in direct_names + return ( + isinstance(func, ast.Attribute) + and func.attr == name + and isinstance(func.value, ast.Name) + and func.value.id in tempfile_aliases + ) + + +def _iter_target_names(target: ast.AST) -> list[str]: + if isinstance(target, ast.Name): + return [target.id] + if isinstance(target, (ast.Tuple, ast.List)): + names: list[str] = [] + for elt in target.elts: + names.extend(_iter_target_names(elt)) + return names + return [] + + +def _collect_bindings( + tree: ast.AST, +) -> tuple[set[str], set[str], set[str], set[str], set[str]]: + os_aliases: set[str] = set() + os_replace_names: set[str] = set() + fcntl_aliases: set[str] = set() + flock_names: set[str] = set() + tempfile_aliases: set[str] = set() + mkstemp_names: set[str] = set() + named_temp_names: set[str] = set() + temp_names: set[str] = set() + + for node in ast.walk(tree): + if isinstance(node, ast.Import): + for alias in node.names: + bound = alias.asname or alias.name + if alias.name == "os": + os_aliases.add(bound) + elif alias.name == "fcntl": + fcntl_aliases.add(bound) + elif alias.name == "tempfile": + tempfile_aliases.add(bound) + elif isinstance(node, ast.ImportFrom): + module = node.module or "" + for alias in node.names: + bound = alias.asname or alias.name + if module == "os" and alias.name == "replace": + os_replace_names.add(bound) + elif module == "fcntl" and alias.name == "flock": + flock_names.add(bound) + elif module == "tempfile" and alias.name == "mkstemp": + mkstemp_names.add(bound) + elif module == "tempfile" and alias.name == "NamedTemporaryFile": + named_temp_names.add(bound) + + for node in ast.walk(tree): + if isinstance(node, (ast.Assign, ast.AnnAssign)): + value = node.value + targets = node.targets if isinstance(node, ast.Assign) else [node.target] + if value is None: + continue + for target in targets: + if _is_attr_call(value, PATH_TEMP_METHODS) or _is_tempfile_call( + value, + "NamedTemporaryFile", + tempfile_aliases, + named_temp_names, + ): + temp_names.update(_iter_target_names(target)) + if _is_tempfile_call(value, "mkstemp", tempfile_aliases, mkstemp_names): + if ( + isinstance(target, (ast.Tuple, ast.List)) + and len(target.elts) >= 2 + and isinstance(target.elts[1], ast.Name) + ): + temp_names.add(target.elts[1].id) + elif isinstance(node, ast.With): + for item in node.items: + if item.optional_vars is None: + continue + if _is_tempfile_call( + item.context_expr, + "NamedTemporaryFile", + tempfile_aliases, + named_temp_names, + ): + temp_names.update(_iter_target_names(item.optional_vars)) + + return os_aliases, os_replace_names, fcntl_aliases, flock_names, temp_names + + +def _is_os_replace_call( + func: ast.expr, + os_aliases: set[str], + os_replace_names: set[str], +) -> bool: + if isinstance(func, ast.Name): + return func.id in os_replace_names + return ( + isinstance(func, ast.Attribute) + and func.attr == "replace" + and isinstance(func.value, ast.Name) + and func.value.id in os_aliases + ) + + +def _is_temp_name(name: str, temp_names: set[str]) -> bool: + lowered = name.lower() + return "tmp" in lowered or "temp" in lowered or name in temp_names + + +def _is_path_like_receiver(receiver: ast.expr, temp_names: set[str]) -> bool: + if _is_attr_call(receiver, PATH_TEMP_METHODS): + return True + return isinstance(receiver, ast.Name) and _is_temp_name(receiver.id, temp_names) + + +def _is_path_replace_call(node: ast.Call, temp_names: set[str]) -> bool: + if not isinstance(node.func, ast.Attribute) or node.func.attr != "replace": + return False + if len(node.args) != 1 or node.keywords: + return False + if any(isinstance(arg, ast.Starred) for arg in node.args): + return False + arg = node.args[0] + if isinstance(arg, ast.Constant) and isinstance(arg.value, str): + return False + return _is_path_like_receiver(node.func.value, temp_names) + + +def _is_flock_call( + func: ast.expr, + fcntl_aliases: set[str], + flock_names: set[str], +) -> bool: + if isinstance(func, ast.Name): + return func.id in flock_names + return ( + isinstance(func, ast.Attribute) + and func.attr == "flock" + and isinstance(func.value, ast.Name) + and func.value.id in fcntl_aliases + ) + + +def _lock_mode_names(node: ast.AST) -> set[str]: + names: set[str] = set() + for child in ast.walk(node): + if isinstance(child, ast.Name): + names.add(child.id) + elif isinstance(child, ast.Attribute): + names.add(child.attr) + return names + + +def _is_bare_lock_ex(node: ast.Call) -> bool: + if len(node.args) < 2: + return False + lock_names = _lock_mode_names(node.args[1]) + return "LOCK_EX" in lock_names and "LOCK_NB" not in lock_names + + +def scan_source(source: str, filename: str = "") -> list[tuple[int, str, str]]: + """Return ``(lineno, kind, detail)`` mechanic violations for source.""" + tree = ast.parse(source, filename=filename) + os_aliases, os_replace_names, fcntl_aliases, flock_names, temp_names = ( + _collect_bindings(tree) + ) + + findings: list[tuple[int, str, str]] = [] + for node in ast.walk(tree): + if not isinstance(node, ast.Call): + continue + + if _is_os_replace_call(node.func, os_aliases, os_replace_names): + findings.append((node.lineno, "os.replace", _call_detail(node.func))) + continue + + if _is_path_replace_call(node, temp_names): + findings.append((node.lineno, "Path.replace", _call_detail(node.func))) + continue + + if _is_flock_call(node.func, fcntl_aliases, flock_names) and _is_bare_lock_ex( + node + ): + findings.append((node.lineno, "flock(LOCK_EX)", _call_detail(node.func))) + + findings.sort() + return findings + + +def scan_file(path: Path) -> list[tuple[int, str, str]]: + return scan_source(path.read_text(encoding="utf-8"), filename=str(path)) + + +def count_violations(root: Path) -> dict[tuple[str, str], int]: + """Map ``(posix-relpath, kind)`` -> occurrence count across the tree.""" + counts: dict[tuple[str, str], int] = {} + for rel in discover_modules(root): + for _lineno, kind, _detail in scan_file(root / rel): + key = (rel.as_posix(), kind) + counts[key] = counts.get(key, 0) + 1 + return counts + + +def evaluate( + root: Path, + allowlist: dict[tuple[str, str], int], +) -> tuple[list[str], list[str]]: + """Return ``(new_violations, tracked)`` human-readable lines.""" + new: list[str] = [] + tracked: list[str] = [] + for rel in discover_modules(root): + rel_str = rel.as_posix() + findings = scan_file(root / rel) + by_kind: dict[str, list[int]] = {} + for lineno, kind, _detail in findings: + by_kind.setdefault(kind, []).append(lineno) + for kind, linenos in sorted(by_kind.items()): + count = len(linenos) + allowed = allowlist.get((rel_str, kind), 0) + if count > allowed: + lines = ", ".join(str(n) for n in sorted(linenos)) + new.append( + f"{rel_str}: raw journal I/O mechanic {kind} " + f"({count} occurrence(s), allowed {allowed}) at line(s) {lines}" + ) + elif allowed: + tracked.append(f"{rel_str}: {count}/{allowed} {kind} (allowlisted)") + return new, tracked + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description="Journal raw-mechanic lint") + parser.add_argument( + "--root", + type=Path, + default=ROOT, + help="Repository root to scan (defaults to the checkout root).", + ) + args = parser.parse_args(argv) + + new, tracked = evaluate(args.root, ALLOWLIST) + + if tracked: + print("journal-io-mechanic: known violations (allowlisted, ratcheting down):") + for line in tracked: + print(f" {line}") + print() + + if new: + print("journal-io-mechanic: NEW violations:", file=sys.stderr) + for line in new: + print(f" {line}", file=sys.stderr) + print(file=sys.stderr) + print( + "Route durable journal writes through solstone.think.journal_io.", + file=sys.stderr, + ) + return 1 + + print("journal-io-mechanic: pass") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/solstone/apps/activities/call.py b/solstone/apps/activities/call.py index 841ba8ea9..b8f838c4b 100644 --- a/solstone/apps/activities/call.py +++ b/solstone/apps/activities/call.py @@ -15,6 +15,7 @@ from typing import Any import typer +from solstone.convey.reasons import ACTIVITIES_BUSY from solstone.think.activities import ( append_activity_record, append_edit, @@ -30,6 +31,7 @@ from solstone.think.activities import ( from solstone.think.entities.loading import load_entities from solstone.think.entities.matching import find_matching_entity from solstone.think.facets import get_facets, log_call_action +from solstone.think.journal_io import LockTimeout from solstone.think.utils import ( get_sol_facet, now_ms, @@ -466,7 +468,13 @@ def create_record( note="created", ) - if not append_activity_record(resolved_facet, resolved_day, record): + try: + created = append_activity_record(resolved_facet, resolved_day, record) + except LockTimeout: + typer.echo(ACTIVITIES_BUSY.message, err=True) + raise typer.Exit(1) + + if not created: typer.echo(f"Error: activity already exists: {span_id}", err=True) raise typer.Exit(1) @@ -523,14 +531,18 @@ def update_record_command( raise typer.Exit(1) note_text = note or f"updated fields: {', '.join(sorted(patch))}" - updated = update_activity_record( - resolved_facet, - resolved_day, - span_id, - patch, - actor="cli:update", - note=note_text, - ) + try: + updated = update_activity_record( + resolved_facet, + resolved_day, + span_id, + patch, + actor="cli:update", + note=note_text, + ) + except LockTimeout: + typer.echo(ACTIVITIES_BUSY.message, err=True) + raise typer.Exit(1) if updated is None: typer.echo(f"activity not found: {span_id}", err=True) raise typer.Exit(1) @@ -569,13 +581,17 @@ def mute_record( """Hide an activity record without deleting it.""" resolved_facet = resolve_sol_facet(facet) resolved_day = resolve_sol_day(day) - record = mute_activity_record( - resolved_facet, - resolved_day, - span_id, - actor="cli:mute", - reason=reason, - ) + try: + record = mute_activity_record( + resolved_facet, + resolved_day, + span_id, + actor="cli:mute", + reason=reason, + ) + except LockTimeout: + typer.echo(ACTIVITIES_BUSY.message, err=True) + raise typer.Exit(1) if record is None: typer.echo(f"activity not found: {span_id}", err=True) raise typer.Exit(1) @@ -614,13 +630,17 @@ def unmute_record( """Restore a previously hidden activity record.""" resolved_facet = resolve_sol_facet(facet) resolved_day = resolve_sol_day(day) - record = unmute_activity_record( - resolved_facet, - resolved_day, - span_id, - actor="cli:unmute", - reason=reason, - ) + try: + record = unmute_activity_record( + resolved_facet, + resolved_day, + span_id, + actor="cli:unmute", + reason=reason, + ) + except LockTimeout: + typer.echo(ACTIVITIES_BUSY.message, err=True) + raise typer.Exit(1) if record is None: typer.echo(f"activity not found: {span_id}", err=True) raise typer.Exit(1) diff --git a/solstone/convey/chat_stream.py b/solstone/convey/chat_stream.py index 22913934d..4001ad23a 100644 --- a/solstone/convey/chat_stream.py +++ b/solstone/convey/chat_stream.py @@ -5,7 +5,6 @@ from __future__ import annotations import json import logging -import os import sys import threading import time @@ -15,6 +14,7 @@ from typing import Any from solstone.think.callosum import callosum_send from solstone.think.indexer.journal import index_file +from solstone.think.journal_io import atomic_replace from solstone.think.streams import update_stream, write_segment_stream from solstone.think.utils import ( day_path, @@ -413,17 +413,5 @@ def _read_events_file(path: Path) -> list[dict[str, Any]]: def _write_events_file(path: Path, events: list[dict[str, Any]]) -> None: - tmp_path = path.with_suffix(f".{os.getpid()}-{threading.get_ident()}.tmp") - try: - with open(tmp_path, "w", encoding="utf-8") as handle: - for event in events: - handle.write(json.dumps(event, ensure_ascii=False)) - handle.write("\n") - os.replace(tmp_path, path) - except Exception: - try: - if tmp_path.exists(): - tmp_path.unlink() - except OSError: - pass - raise + body = "".join(json.dumps(event, ensure_ascii=False) + "\n" for event in events) + atomic_replace(path, body) diff --git a/solstone/convey/reasons.py b/solstone/convey/reasons.py index 9e4aa29f5..f70020487 100644 --- a/solstone/convey/reasons.py +++ b/solstone/convey/reasons.py @@ -316,6 +316,11 @@ ENTITY_BUSY = Reason( "I couldn't update that entity right now because it was busy. Try again in a moment.", 503, ) +ACTIVITIES_BUSY = Reason( + "activities_busy", + "I couldn't update activities right now because they were busy. Try again in a moment.", + 503, +) # identity IDENTITY_BUSY = Reason( diff --git a/solstone/think/activities.py b/solstone/think/activities.py index 7de27fb47..3b8fc4e49 100644 --- a/solstone/think/activities.py +++ b/solstone/think/activities.py @@ -11,18 +11,15 @@ stored as facets/{facet}/activities/{day}.jsonl. """ import difflib -import fcntl import json import logging import os -import random import re -import tempfile -import time from datetime import UTC, datetime from pathlib import Path from typing import Any +from solstone.think.journal_io import atomic_replace, hold_lock from solstone.think.utils import get_journal, segment_parse logger = logging.getLogger(__name__) @@ -784,17 +781,10 @@ def _read_jsonl_records(path: Path) -> list[dict[str, Any]]: def _write_jsonl_records(path: Path, records: list[dict[str, Any]]) -> None: """Atomically write JSONL entries to *path*.""" - path.parent.mkdir(parents=True, exist_ok=True) - fd, tmp_name = tempfile.mkstemp(dir=path.parent, suffix=".tmp") - try: - with os.fdopen(fd, "w", encoding="utf-8") as handle: - for record in records: - handle.write(json.dumps(record, ensure_ascii=False) + "\n") - os.replace(tmp_name, path) - except BaseException: - if os.path.exists(tmp_name): - os.unlink(tmp_name) - raise + atomic_replace( + path, + "".join(json.dumps(record, ensure_ascii=False) + "\n" for record in records), + ) def _fallback_activity_title(record: dict[str, Any]) -> str: @@ -835,42 +825,21 @@ def locked_modify( modify_fn: Any, *, create_if_missing: bool = False, - max_retries: int = 3, ) -> None: """Perform a locked load-modify-save cycle on a JSONL file.""" - lock_path = path.parent / f"{path.name}.lock" - - last_error: OSError | None = None - for attempt in range(max_retries): - try: - path.parent.mkdir(parents=True, exist_ok=True) - with open(lock_path, "w", encoding="utf-8") as lock_file: - fcntl.flock(lock_file, fcntl.LOCK_EX) - try: - existed = path.exists() - if not existed and not create_if_missing: - raise FileNotFoundError(path) - current = _read_jsonl_records(path) if existed else [] - updated = modify_fn([dict(item) for item in current]) - if not isinstance(updated, list): - raise TypeError("modify_fn must return list[dict]") - if not existed and not updated: - return - if existed and updated == current: - return - _write_jsonl_records(path, updated) - finally: - fcntl.flock(lock_file, fcntl.LOCK_UN) + with hold_lock(path): + existed = path.exists() + if not existed and not create_if_missing: + raise FileNotFoundError(path) + current = _read_jsonl_records(path) if existed else [] + updated = modify_fn([dict(item) for item in current]) + if not isinstance(updated, list): + raise TypeError("modify_fn must return list[dict]") + if not existed and not updated: + return + if existed and updated == current: return - except (FileNotFoundError, TypeError, ValueError): - raise - except OSError as exc: - last_error = exc - if attempt < max_retries - 1: - time.sleep(random.uniform(0.05, 0.3) * (attempt + 1)) - - if last_error is not None: - raise last_error + _write_jsonl_records(path, updated) def append_edit( diff --git a/solstone/think/importers/shared.py b/solstone/think/importers/shared.py index a88f1a740..b80462cfb 100644 --- a/solstone/think/importers/shared.py +++ b/solstone/think/importers/shared.py @@ -436,7 +436,6 @@ def write_structured_import( Returns: List of created file paths (absolute) """ - import tempfile from collections import defaultdict # Group entries by day (YYYYMMDD extracted from ts) @@ -505,17 +504,7 @@ def write_structured_import( lines.append(json.dumps(entry)) content = "\n".join(lines) + "\n" - # Atomic write: write to temp file, then rename - fd, tmp_path = tempfile.mkstemp( - dir=str(import_dir), suffix=".tmp", prefix="imported_" - ) - try: - with os.fdopen(fd, "w", encoding="utf-8") as f: - f.write(content) - os.replace(tmp_path, str(out_path)) - except BaseException: - os.unlink(tmp_path) - raise + atomic_replace(out_path, content) created.append(str(out_path)) logger.info("Wrote %d entries to %s", len(day_entries), out_path) diff --git a/solstone/think/link/paths.py b/solstone/think/link/paths.py index 655d0373f..7ce680158 100644 --- a/solstone/think/link/paths.py +++ b/solstone/think/link/paths.py @@ -31,6 +31,7 @@ import uuid from dataclasses import dataclass from pathlib import Path +from solstone.think.journal_io import write_json from solstone.think.utils import get_journal # Production spl-relay endpoint. Single source of truth — self-hosters @@ -146,19 +147,11 @@ class LinkState: return None def save(self) -> None: - path = state_path() - path.parent.mkdir(parents=True, exist_ok=True) - tmp = path.with_suffix(".json.tmp") - with open(tmp, "w", encoding="utf-8") as f: - json.dump( - {"instance_id": self.instance_id, "home_label": self.home_label}, - f, - indent=2, - ) - f.write("\n") - f.flush() - os.fsync(f.fileno()) - os.replace(tmp, path) + write_json( + state_path(), + {"instance_id": self.instance_id, "home_label": self.home_label}, + indent=2, + ) def load_service_token() -> str | None: @@ -178,15 +171,7 @@ def load_service_token() -> str | None: def save_service_token(token: str) -> None: """Persist the service token atomically with mode 0600.""" path = service_token_path() - path.parent.mkdir(parents=True, exist_ok=True) - tmp = path.with_suffix(".json.tmp") - with open(tmp, "w", encoding="utf-8") as f: - json.dump({"service_token": token}, f, indent=2) - f.write("\n") - f.flush() - os.fsync(f.fileno()) - os.chmod(tmp, 0o600) - os.replace(tmp, path) + write_json(path, {"service_token": token}, indent=2, mode=0o600) def generate_totp_secret() -> str: @@ -209,12 +194,4 @@ def load_totp_secret() -> str | None: def save_totp_secret(secret: str) -> None: """Persist the relay pairing TOTP secret atomically with mode 0600.""" path = totp_secret_path() - path.parent.mkdir(parents=True, exist_ok=True) - tmp = path.with_suffix(".json.tmp") - with open(tmp, "w", encoding="utf-8") as f: - json.dump({"totp_secret": secret}, f, indent=2) - f.write("\n") - f.flush() - os.fsync(f.fileno()) - os.chmod(tmp, 0o600) - os.replace(tmp, path) + write_json(path, {"totp_secret": secret}, indent=2, mode=0o600) diff --git a/tests/link/test_paths.py b/tests/link/test_paths.py index 15653864c..d3b75cae2 100644 --- a/tests/link/test_paths.py +++ b/tests/link/test_paths.py @@ -192,7 +192,9 @@ def test_save_service_token_is_atomic( token_path = service_token_path() assert token_path.exists() - assert not any(path.name.endswith(".tmp") for path in token_path.parent.iterdir()) + assert json.loads(token_path.read_text("utf-8")) == {"service_token": "tok.123"} + assert token_path.stat().st_mode & 0o777 == 0o600 + assert {path.name for path in token_path.parent.iterdir()} == {token_path.name} def test_load_service_token_reads_legacy_account_key( @@ -245,7 +247,9 @@ def test_save_totp_secret_is_atomic( secret_path = totp_secret_path() assert secret_path.exists() - assert not any(path.name.endswith(".tmp") for path in secret_path.parent.iterdir()) + assert json.loads(secret_path.read_text("utf-8")) == {"totp_secret": "SECRET"} + assert secret_path.stat().st_mode & 0o777 == 0o600 + assert {path.name for path in secret_path.parent.iterdir()} == {secret_path.name} def test_generate_totp_secret_shape() -> None: diff --git a/tests/test_activities_locking.py b/tests/test_activities_locking.py index 963a8fa11..13d52af9d 100644 --- a/tests/test_activities_locking.py +++ b/tests/test_activities_locking.py @@ -58,7 +58,9 @@ def test_locked_modify_serializes_concurrent_edits(tmp_path, monkeypatch): target=worker, args=("cli:update", "first writer"), kwargs={"hold_lock": True} ) second = threading.Thread( - target=worker, args=("cli:mute", "second writer"), kwargs={"hold_lock": False} + target=worker, + args=("cli:mute", "second writer café"), + kwargs={"hold_lock": False}, ) first.start() @@ -75,9 +77,12 @@ def test_locked_modify_serializes_concurrent_edits(tmp_path, monkeypatch): assert len(records) == 1 assert [edit["note"] for edit in records[0]["edits"]] == [ "first writer", - "second writer", + "second writer café", ] - raw_lines = record_path.read_text(encoding="utf-8").splitlines() + raw_text = record_path.read_text(encoding="utf-8") + assert "café" in raw_text + assert "\\u00e9" not in raw_text + raw_lines = raw_text.splitlines() assert len(raw_lines) == 1 assert json.loads(raw_lines[0])["edits"][1]["actor"] == "cli:mute" diff --git a/tests/test_check_journal_io_mechanic.py b/tests/test_check_journal_io_mechanic.py new file mode 100644 index 000000000..3e7b5aaa4 --- /dev/null +++ b/tests/test_check_journal_io_mechanic.py @@ -0,0 +1,245 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Self-test for scripts/check_journal_io_mechanic.py.""" + +from __future__ import annotations + +import importlib.util +import subprocess +import sys +from pathlib import Path + +import pytest + +REPO_ROOT = Path(__file__).resolve().parents[1] +SCRIPT = REPO_ROOT / "scripts" / "check_journal_io_mechanic.py" + + +def _load_checker(): + spec = importlib.util.spec_from_file_location("check_journal_io_mechanic", SCRIPT) + assert spec and spec.loader + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + +cjm = _load_checker() + + +BAD_OS_REPLACE = """\ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc +import os + + +def persist(tmp, path): + os.replace(tmp, path) +""" + +OWNER_OS_REPLACE = """\ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc +import os + + +def persist(tmp, path): + os.replace(tmp, path) +""" + +HOME_OS_REPLACE = """\ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc +import os + + +def atomic(tmp, path): + os.replace(tmp, path) +""" + +EXCLUDED_OS_REPLACE = """\ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc +import os + + +def persist(tmp, path): + os.replace(tmp, path) +""" + +TEST_OS_REPLACE = """\ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc +import os + + +def test_raw_replace(tmp_path): + os.replace(tmp_path / "a", tmp_path / "b") +""" + + +def _write_file(root: Path, rel: str, content: str) -> None: + path = root / rel + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(content, encoding="utf-8") + + +@pytest.fixture +def bad_root(tmp_path) -> Path: + root = tmp_path / "bad" + _write_file(root, "solstone/apps/badapp/routes.py", BAD_OS_REPLACE) + return root + + +@pytest.fixture +def owner_root(tmp_path) -> Path: + root = tmp_path / "owner" + _write_file(root, "solstone/think/entities/saving.py", OWNER_OS_REPLACE) + return root + + +@pytest.fixture +def home_root(tmp_path) -> Path: + root = tmp_path / "home" + _write_file(root, "solstone/think/journal_io/custom.py", HOME_OS_REPLACE) + return root + + +@pytest.fixture +def excluded_root(tmp_path) -> Path: + root = tmp_path / "excluded" + _write_file(root, "solstone/think/scheduler.py", EXCLUDED_OS_REPLACE) + return root + + +@pytest.fixture +def test_root(tmp_path) -> Path: + root = tmp_path / "test" + _write_file(root, "solstone/apps/badapp/test_raw.py", TEST_OS_REPLACE) + return root + + +def _run(root: Path) -> subprocess.CompletedProcess: + return subprocess.run( + [sys.executable, str(SCRIPT), "--root", str(root)], + capture_output=True, + text=True, + ) + + +@pytest.mark.parametrize( + ("source", "kind"), + [ + ( + "import os\n\ndef persist(tmp, path):\n os.replace(tmp, path)\n", + "os.replace", + ), + ( + "import os as ops\n\ndef persist(tmp, path):\n ops.replace(tmp, path)\n", + "os.replace", + ), + ( + "from os import replace as swap\n\ndef persist(tmp, path):\n swap(tmp, path)\n", + "os.replace", + ), + ("def persist(tmp, path):\n tmp.replace(path)\n", "Path.replace"), + ( + "def persist(dest):\n dest.with_suffix('.tmp').replace(dest)\n", + "Path.replace", + ), + ( + "import tempfile\n\n" + "def persist(path):\n" + " fd, tmp_path = tempfile.mkstemp(dir=path.parent)\n" + " tmp_path.replace(path)\n", + "Path.replace", + ), + ( + "import fcntl\n\ndef lock(f):\n fcntl.flock(f, fcntl.LOCK_EX)\n", + "flock(LOCK_EX)", + ), + ( + "from fcntl import flock, LOCK_EX\n\ndef lock(f):\n flock(f, LOCK_EX)\n", + "flock(LOCK_EX)", + ), + ], +) +def test_scan_source_flags_raw_mechanics(source: str, kind: str) -> None: + findings = cjm.scan_source(source) + assert [finding[1] for finding in findings] == [kind] + + +@pytest.mark.parametrize( + "source", + [ + "def clean(s):\n return s.replace('a', 'b')\n", + "def clean(dt, tz):\n return dt.replace(tzinfo=tz)\n", + "def clean(path):\n return path.name.replace('_', ' ')\n", + "import fcntl\n\ndef lock(f):\n fcntl.flock(f, fcntl.LOCK_EX | fcntl.LOCK_NB)\n", + "import fcntl\n\ndef unlock(f):\n fcntl.flock(f, fcntl.LOCK_UN)\n", + ( + "import tempfile\n" + "from solstone.think.journal_io import install_file\n\n" + "def persist(dest):\n" + " with tempfile.NamedTemporaryFile(delete=False) as tmp:\n" + " tmp.write(b'ok')\n" + " install_file(tmp.name, dest)\n" + ), + ], +) +def test_scan_source_ignores_false_positive_surfaces(source: str) -> None: + assert cjm.scan_source(source) == [] + + +def test_bad_module_exits_one_and_names_file_and_kind(bad_root): + result = _run(bad_root) + assert result.returncode == 1, result.stdout + result.stderr + assert "solstone/apps/badapp/routes.py" in result.stderr + assert "os.replace" in result.stderr + + +def test_owner_path_is_scanned(owner_root): + result = _run(owner_root) + assert result.returncode == 1, result.stdout + result.stderr + assert "solstone/think/entities/saving.py" in result.stderr + assert "os.replace" in result.stderr + + +def test_journal_io_home_is_skipped(home_root): + result = _run(home_root) + assert result.returncode == 0, result.stdout + result.stderr + assert "journal-io-mechanic: pass" in result.stdout + + +def test_excluded_ops_path_is_skipped(excluded_root): + result = _run(excluded_root) + assert result.returncode == 0, result.stdout + result.stderr + assert "journal-io-mechanic: pass" in result.stdout + + +def test_test_modules_are_skipped(test_root): + result = _run(test_root) + assert result.returncode == 0, result.stdout + result.stderr + assert "journal-io-mechanic: pass" in result.stdout + + +def test_ratchet_by_file_kind_count(bad_root): + new, tracked = cjm.evaluate(bad_root, {}) + assert new + assert tracked == [] + + counts = cjm.count_violations(bad_root) + new_exact, tracked_exact = cjm.evaluate(bad_root, counts) + assert new_exact == [] + assert tracked_exact + + key = next(iter(counts)) + ratcheted = dict(counts) + ratcheted[key] = counts[key] - 1 + new_over, _ = cjm.evaluate(bad_root, ratcheted) + assert new_over + + +def test_repo_tree_is_green(): + result = _run(REPO_ROOT) + assert result.returncode == 0, result.stdout + result.stderr -- 2.51.2