diff --git a/AGENTS.md b/AGENTS.md index 0d801dcfb..b2745ee83 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -205,6 +205,9 @@ Each domain has exactly **one** write-owning module (or one tightly-scoped famil | Schedules (`config/schedules.json`) | `solstone/think/schedule_config.py` | | Push devices (`config/push_devices.json`) | `solstone/think/push/devices.py` | | Local inference operational telemetry (`health/local-inference/YYYYMMDD.jsonl`) | `solstone/think/providers/local_admission.py` | +| Provider install status records and proof cache (`health/providers/{local,parakeet}.json`, `health/providers/{local,parakeet}.proof-cache.json`) | `solstone/think/providers/install_state.py` + `solstone/think/providers/artifact_proof.py` | +| Provider install leases (`health/providers/{local,parakeet}.lease`) | `solstone/think/providers/install_lease.py` | +| Provider artifact manifests (`cache/providers/**/.solstone-provider-manifest.json`, `cache/providers/local/mlx/**/*.manifest.json`) | `solstone/think/providers/artifact_proof.py` | | Media offload ledger (`health/offload/.jsonl`) | `solstone/think/offload_ledger.py` | | Parakeet server placement record (`health/parakeet-cpp.placement`) | `solstone/think/providers/parakeet_server.py` | | Hosted backup binding (`backup/hosted/binding.json`) | `solstone/think/backup/hosted.py` | diff --git a/Makefile b/Makefile index 7090911ac..ade9dfb6d 100644 --- a/Makefile +++ b/Makefile @@ -483,6 +483,9 @@ install-checks: .installed @echo "=== Running journal-config owner check ===" @$(MAKE) check-journal-config-owner @echo "" + @echo "=== Running provider-install owner check ===" + @$(MAKE) check-provider-install-owner + @echo "" @echo "=== Running call-http-only check ===" @$(MAKE) check-call-http-only @echo "" @@ -612,6 +615,10 @@ check-journal-io-mechanic: .installed check-journal-config-owner: .installed $(VENV_BIN)/python scripts/check_journal_config_owner.py +# Provider install ownership boundary gate +check-provider-install-owner: .installed + $(VENV_BIN)/python scripts/check_provider_install_owner.py + # sol call HTTP-only gate (call.py reaches the journal only over HTTP) check-call-http-only: .installed $(VENV_BIN)/python scripts/check_call_http_only.py diff --git a/docs/PROVIDERS.md b/docs/PROVIDERS.md index d8b51d423..66c80573b 100644 --- a/docs/PROVIDERS.md +++ b/docs/PROVIDERS.md @@ -99,7 +99,7 @@ config, UI, credential flow, or validation path for enterprise integrations. `local.py` remains a thin product-policy wrapper rather than a second general cloud adapter. It owns guarantees that OpenHands alone does not provide: -- bundled runtime installation and readiness; +- bundled runtime installation and manifest-backed readiness; - context-budget fitting and local schema preparation; - Qwen sampling and chat-template controls; - cross-process local admission and bounded retry; @@ -107,8 +107,9 @@ cloud adapter. It owns guarantees that OpenHands alone does not provide: - confidential egress/attestation gates; - stable local error classification. -Bundled local posts to the supervisor-owned loopback server. A configured -endpoint uses: +Bundled local posts to the supervisor-owned loopback server. Its install status +lives under `health/providers/`, while artifact truth lives in provider manifests +and the affirmative proof cache. A configured endpoint uses: - `providers.local.endpoint_url` - `providers.local.served_model_id` @@ -165,4 +166,12 @@ also: - moves `providers.contexts` enable/extract controls to `talent_overrides`; - moves Rev.ai/Plaud validation state to `service_key_validation`. +The next Thinking maintenance task moves legacy provider install truth out of +`providers.bundled`. It promotes only artifacts that can be proven against the +current pins, writes provider-owned status and manifests, and then removes the +retired operational fields. Missing or mismatched proof exits successfully +without promotion and is repaired by the provider installer under the provider +lease. Unreadable proof exits successfully without promotion and is preserved +until the owner fixes the underlying access or I/O problem. + There are no runtime compatibility shims for the retired shapes. diff --git a/scripts/check_provider_install_owner.py b/scripts/check_provider_install_owner.py new file mode 100755 index 000000000..5ef8a4924 --- /dev/null +++ b/scripts/check_provider_install_owner.py @@ -0,0 +1,693 @@ +#!/usr/bin/env python3 +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Provider install ownership lint. + +Provider install state, leases, proof caches, and artifact manifests are +provider-owned operational domains. Production code must route status writes +through ``solstone.think.providers.install_state``, lease acquisition through +``install_lease``, and manifest/proof-cache mechanics through +``artifact_proof``. This gate catches raw writes, raw locks, private-owner +wrappers, and operational access to retired ``providers.bundled`` install state. + +The production allowlist is intentionally empty. The one numbered migration calls +the owner migration API and is clean by construction, so a new production +violation should fail immediately. +""" + +from __future__ import annotations + +import argparse +import ast +import sys +from pathlib import Path + +ROOT = Path(__file__).resolve().parent.parent + +OWNER_FILES: frozenset[str] = frozenset( + { + "solstone/think/providers/artifact_proof.py", + "solstone/think/providers/install_lease.py", + "solstone/think/providers/install_state.py", + } +) +ALLOWLIST: dict[tuple[str, str], int] = {} +PRIVATE_OWNER_SYMBOLS = { + "_cleanup_legacy_provider_install_config", + "_read_current_unlocked", + "_write_affirmative_cache", +} + + +def _is_owner(rel: Path) -> bool: + return rel.as_posix() in OWNER_FILES + + +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 discover_modules(root: Path) -> list[Path]: + scope = root / "solstone" + if not scope.is_dir(): + return [] + found: list[Path] = [] + for path in sorted(scope.rglob("*.py")): + rel = path.relative_to(root) + if "__pycache__" in rel.parts: + continue + if _is_test_file(rel) or _is_owner(rel): + continue + found.append(rel) + return found + + +class _Bindings: + def __init__(self) -> None: + self.atomic_replace_names: set[str] = set() + self.hold_lock_names: set[str] = set() + self.flock_names: set[str] = set() + self.open_names: set[str] = {"open"} + self.os_modules: set[str] = set() + self.os_open_names: set[str] = set() + self.os_replace_names: set[str] = set() + self.fcntl_modules: set[str] = set() + self.journal_io_modules: set[str] = set() + self.owner_modules: set[str] = set() + self.private_owner_names: set[str] = set() + self.status_path_names: set[str] = set() + self.proof_cache_path_names: set[str] = set() + self.lease_path_names: set[str] = set() + self.manifest_path_names: set[str] = set() + + +def _collect_bindings(tree: ast.AST) -> _Bindings: + bindings = _Bindings() + 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 in { + "solstone.think.journal_io", + "solstone.think.journal_io.atomic", + "solstone.think.journal_io.locking", + }: + bindings.journal_io_modules.add(bound) + elif alias.name in { + "solstone.think.providers.artifact_proof", + "solstone.think.providers.install_lease", + "solstone.think.providers.install_state", + }: + bindings.owner_modules.add(bound) + elif alias.name == "os": + bindings.os_modules.add(bound) + elif alias.name == "fcntl": + bindings.fcntl_modules.add(bound) + elif isinstance(node, ast.ImportFrom): + module = node.module or "" + for alias in node.names: + bound = alias.asname or alias.name + if ( + module + in { + "solstone.think.journal_io", + "solstone.think.journal_io.atomic", + } + and alias.name == "atomic_replace" + ): + bindings.atomic_replace_names.add(bound) + elif ( + module + in { + "solstone.think.journal_io", + "solstone.think.journal_io.locking", + } + and alias.name == "hold_lock" + ): + bindings.hold_lock_names.add(bound) + elif module == "fcntl" and alias.name == "flock": + bindings.flock_names.add(bound) + elif module == "os": + if alias.name == "open": + bindings.os_open_names.add(bound) + elif alias.name == "replace": + bindings.os_replace_names.add(bound) + elif module == "builtins" and alias.name == "open": + bindings.open_names.add(bound) + elif module == "solstone.think.providers.install_state": + if alias.name == "provider_status_path": + bindings.status_path_names.add(bound) + elif alias.name in PRIVATE_OWNER_SYMBOLS: + bindings.private_owner_names.add(bound) + elif module == "solstone.think.providers.artifact_proof": + if alias.name == "proof_cache_path": + bindings.proof_cache_path_names.add(bound) + elif alias.name in { + "artifact_manifest_path", + "mlx_snapshot_manifest_path", + "mlx_variant_manifest_path", + }: + bindings.manifest_path_names.add(bound) + elif alias.name in PRIVATE_OWNER_SYMBOLS: + bindings.private_owner_names.add(bound) + elif module == "solstone.think.providers.install_lease": + if alias.name == "lease_path": + bindings.lease_path_names.add(bound) + elif alias.name in PRIVATE_OWNER_SYMBOLS: + bindings.private_owner_names.add(bound) + return bindings + + +def _dotted_name(node: ast.AST) -> str | None: + if isinstance(node, ast.Name): + return node.id + if isinstance(node, ast.Attribute): + base = _dotted_name(node.value) + if base: + return f"{base}.{node.attr}" + return None + + +def _called_attr(func: ast.expr, modules: set[str], attr: str, full_name: str) -> bool: + if not isinstance(func, ast.Attribute) or func.attr != attr: + return False + if isinstance(func.value, ast.Name) and func.value.id in modules: + return True + return _dotted_name(func) == full_name + + +def _constant_path_parts(node: ast.AST) -> list[str]: + if isinstance(node, ast.Constant) and isinstance(node.value, str): + return [node.value] + if isinstance(node, ast.BinOp) and isinstance(node.op, ast.Div): + return _constant_path_parts(node.left) + _constant_path_parts(node.right) + if isinstance(node, ast.BinOp) and isinstance(node.op, ast.Add): + left = "".join(_constant_path_parts(node.left)) + right = "".join(_constant_path_parts(node.right)) + if left or right: + return [left + right] + if isinstance(node, ast.Call): + parts: list[str] = [] + for arg in node.args: + parts.extend(_constant_path_parts(arg)) + return parts + return [] + + +def _normalized_parts(parts: list[str]) -> list[str]: + normalized: list[str] = [] + for part in parts: + normalized.extend( + piece for piece in part.replace("\\", "/").split("/") if piece + ) + return normalized + + +def _contains_health_provider(parts: list[str]) -> bool: + normalized = _normalized_parts(parts) + return any( + normalized[index : index + 2] == ["health", "providers"] + for index in range(len(normalized) - 1) + ) + + +def _contains_manifest(parts: list[str]) -> bool: + return any( + part.endswith(".solstone-provider-manifest.json") + or part.endswith("snapshot.manifest.json") + or part.endswith("variant-solstone-budget1120.manifest.json") + for part in _normalized_parts(parts) + ) + + +def _contains_provider_status(parts: list[str]) -> bool: + if not _contains_health_provider(parts): + return False + return any( + part in {"local.json", "parakeet.json"} for part in _normalized_parts(parts) + ) + + +def _contains_proof_cache(parts: list[str]) -> bool: + if not _contains_health_provider(parts): + return False + return any(part.endswith(".proof-cache.json") for part in _normalized_parts(parts)) + + +def _contains_lease(parts: list[str]) -> bool: + if not _contains_health_provider(parts): + return False + return any(part.endswith(".lease") for part in _normalized_parts(parts)) + + +def _is_status_path_call(node: ast.AST, bindings: _Bindings) -> bool: + return _is_path_helper_call( + node, + bindings.status_path_names, + "provider_status_path", + "solstone.think.providers.install_state.provider_status_path", + bindings, + ) + + +def _is_proof_cache_path_call(node: ast.AST, bindings: _Bindings) -> bool: + return _is_path_helper_call( + node, + bindings.proof_cache_path_names, + "proof_cache_path", + "solstone.think.providers.artifact_proof.proof_cache_path", + bindings, + ) + + +def _is_lease_path_call(node: ast.AST, bindings: _Bindings) -> bool: + return _is_path_helper_call( + node, + bindings.lease_path_names, + "lease_path", + "solstone.think.providers.install_lease.lease_path", + bindings, + ) + + +def _is_manifest_path_call(node: ast.AST, bindings: _Bindings) -> bool: + if not isinstance(node, ast.Call): + return False + func = node.func + if isinstance(func, ast.Name) and func.id in bindings.manifest_path_names: + return True + dotted = _dotted_name(func) + return dotted in { + "solstone.think.providers.artifact_proof.artifact_manifest_path", + "solstone.think.providers.artifact_proof.mlx_snapshot_manifest_path", + "solstone.think.providers.artifact_proof.mlx_variant_manifest_path", + } + + +def _is_path_helper_call( + node: ast.AST, + direct_names: set[str], + attr: str, + full_name: str, + bindings: _Bindings, +) -> bool: + if not isinstance(node, ast.Call): + return False + func = node.func + if isinstance(func, ast.Name) and func.id in direct_names: + return True + return _called_attr(func, bindings.owner_modules, attr, full_name) + + +def _assigned_path_names( + tree: ast.AST, bindings: _Bindings +) -> tuple[set[str], set[str], set[str], set[str]]: + status_names: set[str] = set() + cache_names: set[str] = set() + lease_names: set[str] = set() + manifest_names: set[str] = set() + changed = True + while changed: + changed = False + for node in ast.walk(tree): + value: ast.AST | None = None + targets: list[ast.expr] = [] + if isinstance(node, ast.Assign): + value = node.value + targets = list(node.targets) + elif isinstance(node, ast.AnnAssign): + value = node.value + targets = [node.target] + if value is None: + continue + is_status = _is_owned_path_expr(value, bindings, status_names, "status") + is_cache = _is_owned_path_expr(value, bindings, cache_names, "cache") + is_lease = _is_owned_path_expr(value, bindings, lease_names, "lease") + is_manifest = _is_owned_path_expr( + value, bindings, manifest_names, "manifest" + ) + for target in targets: + if not isinstance(target, ast.Name): + continue + if is_status and target.id not in status_names: + status_names.add(target.id) + changed = True + if is_cache and target.id not in cache_names: + cache_names.add(target.id) + changed = True + if is_lease and target.id not in lease_names: + lease_names.add(target.id) + changed = True + if is_manifest and target.id not in manifest_names: + manifest_names.add(target.id) + changed = True + return status_names, cache_names, lease_names, manifest_names + + +def _is_owned_path_expr( + node: ast.AST, + bindings: _Bindings, + known_names: set[str], + kind: str, +) -> bool: + if isinstance(node, ast.Name): + return node.id in known_names + if kind == "status" and _is_status_path_call(node, bindings): + return True + if kind == "cache" and _is_proof_cache_path_call(node, bindings): + return True + if kind == "lease" and _is_lease_path_call(node, bindings): + return True + if kind == "manifest" and _is_manifest_path_call(node, bindings): + return True + parts = _constant_path_parts(node) + if kind == "status": + return _contains_provider_status(parts) + if kind == "cache": + return _contains_proof_cache(parts) + if kind == "lease": + return _contains_lease(parts) + return _contains_manifest(parts) + + +def _called_atomic_replace(func: ast.expr, bindings: _Bindings) -> bool: + if isinstance(func, ast.Name): + return func.id in bindings.atomic_replace_names + return _called_attr( + func, + bindings.journal_io_modules, + "atomic_replace", + "solstone.think.journal_io.atomic_replace", + ) + + +def _called_hold_lock(func: ast.expr, bindings: _Bindings) -> bool: + if isinstance(func, ast.Name): + return func.id in bindings.hold_lock_names + return _called_attr( + func, + bindings.journal_io_modules, + "hold_lock", + "solstone.think.journal_io.hold_lock", + ) + + +def _called_flock(func: ast.expr, bindings: _Bindings) -> bool: + if isinstance(func, ast.Name): + return func.id in bindings.flock_names + return _called_attr(func, bindings.fcntl_modules, "flock", "fcntl.flock") + + +def _called_os_open(func: ast.expr, bindings: _Bindings) -> bool: + if isinstance(func, ast.Name): + return func.id in bindings.os_open_names + return _called_attr(func, bindings.os_modules, "open", "os.open") + + +def _called_os_replace(func: ast.expr, bindings: _Bindings) -> bool: + if isinstance(func, ast.Name): + return func.id in bindings.os_replace_names + return _called_attr(func, bindings.os_modules, "replace", "os.replace") + + +def _called_open(func: ast.expr, bindings: _Bindings) -> bool: + return isinstance(func, ast.Name) and func.id in bindings.open_names + + +def _is_path_replace_call(node: ast.Call) -> bool: + return isinstance(node.func, ast.Attribute) and node.func.attr == "replace" + + +def _is_path_write_call(node: ast.Call) -> bool: + return isinstance(node.func, ast.Attribute) and node.func.attr in { + "write_text", + "write_bytes", + "open", + "unlink", + } + + +def _owned_kind( + node: ast.AST, + bindings: _Bindings, + status_names: set[str], + cache_names: set[str], + lease_names: set[str], + manifest_names: set[str], +) -> str | None: + if _is_owned_path_expr(node, bindings, status_names, "status"): + return "provider_status" + if _is_owned_path_expr(node, bindings, cache_names, "cache"): + return "proof_cache" + if _is_owned_path_expr(node, bindings, lease_names, "lease"): + return "provider_lease" + if _is_owned_path_expr(node, bindings, manifest_names, "manifest"): + return "provider_manifest" + return None + + +def _bundled_access(node: ast.AST) -> bool: + if isinstance(node, ast.Subscript): + slc = node.slice + return isinstance(slc, ast.Constant) and slc.value == "bundled" + if isinstance(node, ast.Call) and isinstance(node.func, ast.Attribute): + return ( + node.func.attr == "get" + and bool(node.args) + and isinstance(node.args[0], ast.Constant) + and node.args[0].value == "bundled" + ) + return False + + +def _is_private_owner_expr(node: ast.AST, bindings: _Bindings) -> bool: + if isinstance(node, ast.Name): + return node.id in bindings.private_owner_names + if isinstance(node, ast.Attribute) and node.attr in PRIVATE_OWNER_SYMBOLS: + if isinstance(node.value, ast.Name) and node.value.id in bindings.owner_modules: + return True + dotted = _dotted_name(node) + return dotted in { + f"solstone.think.providers.install_state.{node.attr}", + f"solstone.think.providers.artifact_proof.{node.attr}", + f"solstone.think.providers.install_lease.{node.attr}", + } + return False + + +def scan_source(source: str, filename: str = "") -> list[tuple[int, str, str]]: + tree = ast.parse(source, filename=filename) + bindings = _collect_bindings(tree) + status_names, cache_names, lease_names, manifest_names = _assigned_path_names( + tree, bindings + ) + findings: list[tuple[int, str, str]] = [] + + for node in ast.walk(tree): + if isinstance(node, ast.ImportFrom) and node.module in { + "solstone.think.providers.install_state", + "solstone.think.providers.artifact_proof", + "solstone.think.providers.install_lease", + }: + for alias in node.names: + if alias.name in PRIVATE_OWNER_SYMBOLS: + findings.append((node.lineno, "private_owner_symbol", alias.name)) + elif isinstance(node, (ast.Assign, ast.AnnAssign)): + value = node.value + if value is not None and _is_private_owner_expr(value, bindings): + findings.append( + ( + node.lineno, + "private_owner_wrapper", + "private provider-install owner symbol alias or re-export", + ) + ) + elif isinstance(node, ast.FunctionDef): + for child in ast.walk(node): + if isinstance(child, ast.Call) and _is_private_owner_expr( + child.func, bindings + ): + findings.append( + ( + child.lineno, + "private_owner_wrapper", + f"{node.name} wraps private provider-install owner symbol", + ) + ) + elif _bundled_access(node): + findings.append( + ( + node.lineno, + "providers_bundled_operational", + "providers.bundled access outside migration owner", + ) + ) + if not isinstance(node, ast.Call): + continue + if _called_atomic_replace(node.func, bindings) and node.args: + kind = _owned_kind( + node.args[0], + bindings, + status_names, + cache_names, + lease_names, + manifest_names, + ) + if kind: + findings.append((node.lineno, f"{kind}_replace", "atomic_replace")) + elif _called_os_replace(node.func, bindings) and len(node.args) >= 2: + kind = _owned_kind( + node.args[1], + bindings, + status_names, + cache_names, + lease_names, + manifest_names, + ) + if kind: + findings.append((node.lineno, f"{kind}_replace", "os.replace")) + elif _is_path_replace_call(node) and node.args: + kind = _owned_kind( + node.args[0], + bindings, + status_names, + cache_names, + lease_names, + manifest_names, + ) + if kind: + findings.append((node.lineno, f"{kind}_replace", "Path.replace")) + elif _called_open(node.func, bindings) and node.args: + kind = _owned_kind( + node.args[0], + bindings, + status_names, + cache_names, + lease_names, + manifest_names, + ) + if kind: + findings.append((node.lineno, f"{kind}_raw_open", "open")) + elif _called_os_open(node.func, bindings) and node.args: + kind = _owned_kind( + node.args[0], + bindings, + status_names, + cache_names, + lease_names, + manifest_names, + ) + if kind: + findings.append((node.lineno, f"{kind}_raw_open", "os.open")) + elif _is_path_write_call(node): + path_expr = node.func.value + kind = _owned_kind( + path_expr, + bindings, + status_names, + cache_names, + lease_names, + manifest_names, + ) + if kind: + findings.append((node.lineno, f"{kind}_write", node.func.attr)) + elif _called_hold_lock(node.func, bindings) and node.args: + kind = _owned_kind( + node.args[0], + bindings, + status_names, + cache_names, + lease_names, + manifest_names, + ) + if kind: + findings.append((node.lineno, "second_provider_install_lock", kind)) + elif _called_flock(node.func, bindings): + if lease_names or any( + _is_owned_path_expr(arg, bindings, lease_names, "lease") + for arg in node.args + ): + findings.append( + ( + node.lineno, + "second_provider_install_lock", + "fcntl.flock targets provider lease", + ) + ) + + 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]: + 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], list[str]]: + over: list[str] = [] + stale: list[str] = [] + tracked: list[str] = [] + counts = count_violations(root) + for key in sorted(set(counts) | set(allowlist)): + count = counts.get(key, 0) + allowed = allowlist.get(key, 0) + rel, kind = key + if count > allowed: + over.append(f"{rel}: {kind} count {count} exceeds allowed {allowed}") + elif count < allowed: + stale.append(f"{rel}: {kind} count {count} below allowed {allowed}") + elif allowed: + tracked.append(f"{rel}: {count}/{allowed} {kind} (allowlisted)") + return over, stale, tracked + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description="Provider install owner lint") + parser.add_argument( + "--root", + type=Path, + default=ROOT, + help="Repository root to scan (defaults to the checkout root).", + ) + args = parser.parse_args(argv) + over, stale, tracked = evaluate(args.root, ALLOWLIST) + if tracked: + print("provider-install-owner: known violations (allowlisted):") + for line in tracked: + print(f" {line}") + print() + if over or stale: + print("provider-install-owner: violations:", file=sys.stderr) + for line in over: + print(f" {line}", file=sys.stderr) + for line in stale: + print(f" stale allowlist: {line}", file=sys.stderr) + print( + "Route provider install state, leases, manifests, and proof caches " + "through their owner APIs.", + file=sys.stderr, + ) + return 1 + print("provider-install-owner: pass") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/solstone/apps/thinking/maint/001_migrate_provider_install_state.py b/solstone/apps/thinking/maint/001_migrate_provider_install_state.py new file mode 100644 index 000000000..3737921fa --- /dev/null +++ b/solstone/apps/thinking/maint/001_migrate_provider_install_state.py @@ -0,0 +1,47 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Move provider install truth to provider-owned status and manifest records.""" + +from __future__ import annotations + +from pathlib import Path +from typing import Any + +from solstone.think.providers.install_state import ( + migrate_legacy_provider_artifact_truth, +) +from solstone.think.utils import get_journal + +MAINT_RETRY_ON_NEXT_START = True +MAINT_BLOCKS_SUPERVISOR_START = True + + +def migrate(config: dict[str, Any], journal: Path) -> bool: + result = migrate_legacy_provider_artifact_truth(journal_path=journal) + return bool( + result["actions"] or result["cleanup"]["removed"] or result["cleanup"]["moved"] + ) + + +def main() -> None: + journal = Path(get_journal()) + result = migrate_legacy_provider_artifact_truth(journal_path=journal) + changed = bool( + result["actions"] or result["cleanup"]["removed"] or result["cleanup"]["moved"] + ) + if not changed: + print("Provider install state already uses provider-owned records.") + return + + for action in result["actions"]: + print(action["message"]) + cleanup = result["cleanup"] + if cleanup["moved"]: + print("Moved local Vulkan device override to providers.local.") + if cleanup["removed"]: + print("Removed legacy provider install state from providers.bundled.") + + +if __name__ == "__main__": + main() diff --git a/solstone/apps/thinking/tests/test_provider_install_state_migration.py b/solstone/apps/thinking/tests/test_provider_install_state_migration.py new file mode 100644 index 000000000..5b5e82ff2 --- /dev/null +++ b/solstone/apps/thinking/tests/test_provider_install_state_migration.py @@ -0,0 +1,175 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import importlib +import json +from types import SimpleNamespace + +from solstone.think.providers.artifact_proof import ReadinessOutcome +from solstone.think.providers.install_lease import acquire_install_lease +from solstone.think.providers.install_state import ( + migrate_legacy_provider_artifact_truth, + read_install_status, +) + + +def _write_config(journal, bundled_local): + path = journal / "config" / "journal.json" + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text( + json.dumps({"providers": {"bundled": {"local": bundled_local}}}, indent=2) + + "\n", + encoding="utf-8", + ) + + +def _read_config(journal): + return json.loads((journal / "config" / "journal.json").read_text("utf-8")) + + +def _readiness(status: str, reason_code: str = "ready") -> ReadinessOutcome: + return ReadinessOutcome( + provider="local", + status=status, # type: ignore[arg-type] + reason_code=reason_code, + target={}, + install={ + "install_state": "idle", + "install_error": None, + "error_code": None, + "attempt_id": None, + "progress_bytes_received": None, + "progress_bytes_total": None, + "last_transition_at": None, + "last_progress_at": None, + }, + host={"gpu_available": True, "gpu_probe_ok": True, "ram_sufficient": True}, + artifacts={ + "binary_installed": status == "ready", + "model_installed": status == "ready", + }, + proof={ + "binary": { + "status": status, + "reason_code": reason_code, + "cache_hit": False, + }, + "model": {"status": status, "reason_code": reason_code, "cache_hit": False}, + }, + ) + + +def test_migration_task_declares_retry_and_startup_blocking(): + module = importlib.import_module( + "solstone.apps.thinking.maint.001_migrate_provider_install_state" + ) + + assert module.MAINT_RETRY_ON_NEXT_START is True + assert module.MAINT_BLOCKS_SUPERVISOR_START is True + + +def test_ready_legacy_state_publishes_status_and_cleans_config(tmp_path, monkeypatch): + legacy = { + "install_state": "installed", + "vulkan_device_index": 1, + "binary_artifact": "old", + "binary_sha256": "old", + "binary_path": "old", + } + _write_config(tmp_path, legacy) + from solstone.think.providers import local_install + + monkeypatch.setattr( + local_install, "inspect_readiness", lambda _model: _readiness("ready") + ) + monkeypatch.setattr( + local_install, + "target_fingerprint", + lambda _model: {"provider": "local", "unit": "test"}, + ) + + result = migrate_legacy_provider_artifact_truth(journal_path=tmp_path) + + assert result["actions"][0]["status"] == "ready" + config = _read_config(tmp_path) + assert config["providers"]["local"]["vulkan_device_index"] == 1 + assert config["providers"]["bundled"]["local"] == {} + status = read_install_status(name="local", journal_path=tmp_path) + assert status["install_state"] == "installed" + + +def test_fedora_shape_no_manifest_exits_zero_without_cleanup(tmp_path, monkeypatch): + _write_config(tmp_path, {"install_state": "installed"}) + from solstone.think.providers import local_cuda, local_install + + monkeypatch.setattr( + local_install, + "inspect_readiness", + lambda _model: _readiness("missing-or-mismatched", "manifest_missing"), + ) + monkeypatch.setattr( + local_install, + "target_fingerprint", + lambda _model: {"provider": "local", "unit": "test"}, + ) + monkeypatch.setattr( + local_cuda, + "resolve_local_backend", + lambda _pin: SimpleNamespace(backend="vulkan", reason="test"), + ) + + result = migrate_legacy_provider_artifact_truth(journal_path=tmp_path) + + action = result["actions"][0] + assert action["status"] == "missing-or-mismatched" + assert action["cleanup"] is False + assert "reinstall will rebuild the proof" in action["message"] + assert ( + _read_config(tmp_path)["providers"]["bundled"]["local"]["install_state"] + == "installed" + ) + assert not (tmp_path / "health" / "providers" / "local.json").exists() + + +def test_proof_unavailable_is_non_destructive(tmp_path, monkeypatch): + _write_config(tmp_path, {"install_state": "installed", "binary_path": "old"}) + from solstone.think.providers import local_install + + monkeypatch.setattr( + local_install, + "inspect_readiness", + lambda _model: _readiness("proof-unavailable", "manifest_io_error"), + ) + monkeypatch.setattr( + local_install, + "target_fingerprint", + lambda _model: {"provider": "local", "unit": "test"}, + ) + + result = migrate_legacy_provider_artifact_truth(journal_path=tmp_path) + + assert result["actions"][0]["status"] == "proof-unavailable" + assert ( + _read_config(tmp_path)["providers"]["bundled"]["local"]["binary_path"] == "old" + ) + assert not (tmp_path / "health" / "providers" / "local.json").exists() + + +def test_busy_lease_defers_to_next_start(tmp_path): + _write_config(tmp_path, {"install_state": "installed"}) + lease = acquire_install_lease("local", journal_path=tmp_path) + assert lease is not None + + try: + result = migrate_legacy_provider_artifact_truth(journal_path=tmp_path) + finally: + lease.release() + + assert result["actions"][0]["status"] == "busy" + assert result["actions"][0]["reason_code"] == "install_busy" + assert ( + _read_config(tmp_path)["providers"]["bundled"]["local"]["install_state"] + == "installed" + ) diff --git a/solstone/convey/maint_cli.py b/solstone/convey/maint_cli.py index 53d219fde..889bf77e4 100644 --- a/solstone/convey/maint_cli.py +++ b/solstone/convey/maint_cli.py @@ -77,30 +77,11 @@ def show_task_details(journal: Path, task_name: str) -> None: status, exit_code, ran_ts = get_task_status(journal, task.app, task.name) state_file = get_state_file(journal, task.app, task.name) - duration_ms = None - log_lines: list[str] = [] - errors: list[str] = [] + attempts: list[dict] = [] if status != "pending" and state_file.exists(): - with open(state_file, "r") as f: - for raw_line in f: - line = raw_line.strip() - if not line: - continue - try: - event = json.loads(line) - except json.JSONDecodeError: - continue - - event_type = event.get("event") - if event_type == "line": - text = event.get("line") - if isinstance(text, str): - log_lines.append(text) - elif event_type == "exit": - if isinstance(event.get("duration_ms"), int): - duration_ms = event["duration_ms"] - if event.get("error"): - errors.append(str(event["error"])) + attempts = _read_attempt_logs(state_file) + latest = attempts[-1] if attempts else {} + duration_ms = latest.get("duration_ms") print(task.qualified_name) if task.description: @@ -133,11 +114,49 @@ def show_task_details(journal: Path, task_name: str) -> None: print("Task has not been run yet.") return - for line in log_lines: - print(line) - - for error in errors: - print(f"Error: {error}") + for index, attempt in enumerate(reversed(attempts), start=1): + if index > 1: + print() + print(f"Prior attempt {index}:") + for line in attempt["lines"]: + print(line) + + for error in attempt["errors"]: + print(f"Error: {error}") + + +def _read_attempt_logs(state_file: Path) -> list[dict]: + attempts: list[dict] = [] + current: dict | None = None + with open(state_file, "r") as f: + for raw_line in f: + line = raw_line.strip() + if not line: + continue + try: + event = json.loads(line) + except json.JSONDecodeError: + continue + event_type = event.get("event") + if event_type == "exec": + if current is not None: + attempts.append(current) + current = {"lines": [], "errors": [], "duration_ms": None} + continue + if current is None: + current = {"lines": [], "errors": [], "duration_ms": None} + if event_type == "line": + text = event.get("line") + if isinstance(text, str): + current["lines"].append(text) + elif event_type == "exit": + if isinstance(event.get("duration_ms"), int): + current["duration_ms"] = event["duration_ms"] + if event.get("error"): + current["errors"].append(str(event["error"])) + if current is not None: + attempts.append(current) + return attempts def main() -> None: @@ -238,7 +257,9 @@ Examples: print() # Run pending tasks - ran, succeeded = run_pending_tasks(journal) + results = run_pending_tasks(journal) + ran = len(results) + succeeded = sum(1 for result in results if result.success) if ran == 0: print("No pending maintenance tasks.") else: diff --git a/solstone/think/maint.py b/solstone/think/maint.py index f95dcc3c8..b735154fe 100644 --- a/solstone/think/maint.py +++ b/solstone/think/maint.py @@ -24,6 +24,7 @@ Execution: from __future__ import annotations +import ast import json import logging import queue @@ -32,6 +33,7 @@ import subprocess import sys import threading import time +import uuid from dataclasses import dataclass from pathlib import Path from typing import Optional @@ -49,6 +51,8 @@ class MaintTask: name: str script_path: Path description: str = "" + retry_on_next_start: bool = False + blocks_supervisor_start: bool = False @property def qualified_name(self) -> str: @@ -56,6 +60,16 @@ class MaintTask: return f"{self.app}:{self.name}" +@dataclass(frozen=True) +class MaintTaskResult: + """Result of a maintenance task run.""" + + task: MaintTask + success: bool + exit_code: int + state_file: Path + + def discover_tasks() -> list[MaintTask]: """Discover all maint tasks from apps/*/maint/*.py. @@ -80,8 +94,10 @@ def discover_tasks() -> list[MaintTask]: if script.name.startswith("_"): continue - # Extract description from module docstring + # Extract description from module docstring and opt-ins without import. description = "" + retry_on_next_start = False + blocks_supervisor_start = False try: content = script.read_text() # Find first docstring (handles files with license headers) @@ -94,6 +110,9 @@ def discover_tasks() -> list[MaintTask]: content[start + 3 : end].strip().split("\n")[0] ) break + opt_ins = _parse_task_opt_ins(content, filename=str(script)) + retry_on_next_start = opt_ins["retry_on_next_start"] + blocks_supervisor_start = opt_ins["blocks_supervisor_start"] except Exception: pass @@ -103,6 +122,8 @@ def discover_tasks() -> list[MaintTask]: name=script.stem, script_path=script, description=description, + retry_on_next_start=retry_on_next_start, + blocks_supervisor_start=blocks_supervisor_start, ) ) @@ -110,6 +131,33 @@ def discover_tasks() -> list[MaintTask]: return tasks +def _parse_task_opt_ins(source: str, *, filename: str) -> dict[str, bool]: + tree = ast.parse(source, filename=filename) + opt_ins = { + "retry_on_next_start": False, + "blocks_supervisor_start": False, + } + names = { + "MAINT_RETRY_ON_NEXT_START": "retry_on_next_start", + "MAINT_BLOCKS_SUPERVISOR_START": "blocks_supervisor_start", + } + for node in tree.body: + target: ast.expr | None = None + value: ast.expr | None = None + if isinstance(node, ast.Assign) and len(node.targets) == 1: + target = node.targets[0] + value = node.value + elif isinstance(node, ast.AnnAssign): + target = node.target + value = node.value + if not isinstance(target, ast.Name) or target.id not in names: + continue + if not isinstance(value, ast.Constant) or not isinstance(value.value, bool): + continue + opt_ins[names[target.id]] = value.value + return opt_ins + + def get_state_file(journal: Path, app: str, task: str) -> Path: """Get path to task state file.""" return journal / "maint" / app / f"{task}.jsonl" @@ -126,30 +174,24 @@ def _parse_state_file(state_file: Path) -> dict: return default try: + latest = _latest_attempt_events(state_file) + if not latest: + return default + duration_ms = None line_count = 0 exec_ts = None exit_ts = None - - with open(state_file, "r") as f: - for raw_line in f: - line = raw_line.strip() - if not line: - continue - - try: - event = json.loads(line) - except json.JSONDecodeError: - continue - event_type = event.get("event") - if event_type == "exec" and exec_ts is None: - exec_ts = event.get("ts") - elif event_type == "line": - line_count += 1 - elif event_type == "exit": - exit_ts = event.get("ts") - if isinstance(event.get("duration_ms"), int): - duration_ms = event["duration_ms"] + for event in latest: + event_type = event.get("event") + if event_type == "exec" and exec_ts is None: + exec_ts = event.get("ts") + elif event_type == "line": + line_count += 1 + elif event_type == "exit": + exit_ts = event.get("ts") + if isinstance(event.get("duration_ms"), int): + duration_ms = event["duration_ms"] return { "duration_ms": duration_ms, @@ -179,24 +221,14 @@ def get_task_status( if not state_file.exists(): return "pending", None, None - # Track the first exec event and the last event. try: + events = _latest_attempt_events(state_file) exec_ts = None last_event = None - - with open(state_file, "r") as f: - for raw_line in f: - line = raw_line.strip() - if not line: - continue - - try: - event = json.loads(line) - except json.JSONDecodeError: - continue - if event.get("event") == "exec" and exec_ts is None: - exec_ts = event.get("ts") - last_event = event + for event in events: + if event.get("event") == "exec" and exec_ts is None: + exec_ts = event.get("ts") + last_event = event if last_event and last_event.get("event") == "exit": ts = last_event.get("ts") @@ -214,6 +246,59 @@ def get_task_status( return "in_progress", None, None +def _read_events(state_file: Path) -> list[dict]: + events: list[dict] = [] + with open(state_file, "r") as f: + for raw_line in f: + line = raw_line.strip() + if not line: + continue + try: + event = json.loads(line) + except json.JSONDecodeError: + continue + if isinstance(event, dict): + events.append(event) + return events + + +def _attempt_blocks(events: list[dict]) -> list[list[dict]]: + blocks: list[list[dict]] = [] + current: list[dict] = [] + current_attempt_id: str | None = None + for event in events: + event_type = event.get("event") + if event_type == "exec": + if current: + blocks.append(current) + current = [event] + attempt_id = event.get("attempt_id") + current_attempt_id = attempt_id if isinstance(attempt_id, str) else None + continue + if not current: + current = [event] + continue + event_attempt_id = event.get("attempt_id") + if ( + current_attempt_id is not None + and isinstance(event_attempt_id, str) + and event_attempt_id != current_attempt_id + ): + blocks.append(current) + current = [event] + current_attempt_id = event_attempt_id + continue + current.append(event) + if current: + blocks.append(current) + return blocks + + +def _latest_attempt_events(state_file: Path) -> list[dict]: + blocks = _attempt_blocks(_read_events(state_file)) + return blocks[-1] if blocks else [] + + def _write_event(f, event: dict) -> None: """Write a JSONL event to file.""" f.write(json.dumps(event) + "\n") @@ -260,6 +345,7 @@ def run_task( state_file = get_state_file(journal, task.app, task.name) start_time = time.time() start_ts = int(start_time * 1000) + attempt_id = uuid.uuid4().hex # Build command to run the task cmd = [sys.executable, "-m", f"solstone.apps.{task.app}.maint.{task.name}"] @@ -277,12 +363,13 @@ def run_task( ) try: - with open(state_file, "w") as f: + with open(state_file, "a") as f: # Write exec event _write_event( f, { "event": "exec", + "attempt_id": attempt_id, "ts": start_ts, "app": task.app, "task": task.name, @@ -383,6 +470,7 @@ def run_task( f, { "event": "line", + "attempt_id": attempt_id, "ts": now_ms(), "line": line, }, @@ -399,6 +487,7 @@ def run_task( f, { "event": "exit", + "attempt_id": attempt_id, "ts": now_ms(), "exit_code": exit_code, "duration_ms": duration_ms, @@ -412,6 +501,7 @@ def run_task( f, { "event": "exit", + "attempt_id": attempt_id, "ts": now_ms(), "exit_code": exit_code, "duration_ms": duration_ms, @@ -469,6 +559,7 @@ def run_task( f, { "event": "exit", + "attempt_id": attempt_id, "ts": now_ms(), "exit_code": -1, "error": str(e), @@ -491,7 +582,7 @@ def run_task( return False, -1 -def run_pending_tasks(journal: Path, emit_fn=None) -> tuple[int, int]: +def run_pending_tasks(journal: Path, emit_fn=None) -> list[MaintTaskResult]: """Run all pending maintenance tasks. Args: @@ -499,31 +590,36 @@ def run_pending_tasks(journal: Path, emit_fn=None) -> tuple[int, int]: emit_fn: Optional function to emit Callosum events Returns: - Tuple of (tasks_run, tasks_succeeded) + Per-task run results, in execution order. """ tasks = discover_tasks() pending = [] for task in tasks: status, _, _ = get_task_status(journal, task.app, task.name) - if status == "pending": + if status == "pending" or (status == "failed" and task.retry_on_next_start): pending.append(task) if not pending: - return 0, 0 + return [] logger.info(f"Found {len(pending)} pending maintenance task(s)") - ran = 0 - succeeded = 0 + results: list[MaintTaskResult] = [] for task in pending: - ran += 1 success, _ = run_task(journal, task, emit_fn) - if success: - succeeded += 1 + _status, exit_code, _ran_ts = get_task_status(journal, task.app, task.name) + results.append( + MaintTaskResult( + task=task, + success=success, + exit_code=exit_code if exit_code is not None else -1, + state_file=get_state_file(journal, task.app, task.name), + ) + ) - return ran, succeeded + return results def list_tasks(journal: Path) -> list[dict]: diff --git a/solstone/think/providers/install_state.py b/solstone/think/providers/install_state.py index b89cc22b7..b5fb0d1e6 100644 --- a/solstone/think/providers/install_state.py +++ b/solstone/think/providers/install_state.py @@ -13,7 +13,11 @@ from datetime import datetime, timezone from pathlib import Path from typing import Any, Callable, Literal, TypedDict, cast, get_args -from solstone.think.journal_config import JournalConfigMutation, mutate_journal_config +from solstone.think.journal_config import ( + JournalConfigMutation, + mutate_journal_config, + read_journal_config, +) from solstone.think.journal_io.atomic import atomic_replace from solstone.think.journal_io.locking import hold_lock from solstone.think.utils import get_journal @@ -408,11 +412,54 @@ def record_interrupted_install( ) +def migrate_legacy_provider_artifact_truth( + *, + journal_path: str | Path | None = None, +) -> dict[str, Any]: + """Promote trustworthy legacy provider state into manifests/status records.""" + journal = Path(journal_path) if journal_path is not None else Path(get_journal()) + config = read_journal_config(journal) + providers = config.get("providers") + bundled = providers.get("bundled") if isinstance(providers, dict) else None + if not isinstance(bundled, dict): + cleanup = _cleanup_legacy_provider_install_config( + clean_providers=frozenset(), journal_path=journal + ) + return {"actions": [], "cleanup": cleanup} + + actions: list[dict[str, Any]] = [] + clean_providers: set[str] = set() + for provider in ("local", "parakeet"): + legacy = bundled.get(provider) + if not _legacy_provider_has_operational_state(legacy): + continue + action = _migrate_legacy_provider(provider, legacy, journal) + actions.append(action) + if action.get("cleanup"): + clean_providers.add(provider) + + cleanup = _cleanup_legacy_provider_install_config( + clean_providers=frozenset(clean_providers), journal_path=journal + ) + return {"actions": actions, "cleanup": cleanup} + + def migrate_legacy_provider_install_state( *, journal_path: str | Path | None = None, ) -> dict[str, int]: """Remove legacy provider install operational fields from journal config.""" + return _cleanup_legacy_provider_install_config( + clean_providers=PROVIDERS, journal_path=journal_path + ) + + +def _cleanup_legacy_provider_install_config( + *, + clean_providers: frozenset[str], + journal_path: str | Path | None = None, +) -> dict[str, int]: + """Remove legacy provider install fields after owner proof is established.""" def apply(config: dict[str, Any]) -> JournalConfigMutation[dict[str, int]]: removed = 0 @@ -429,7 +476,7 @@ def migrate_legacy_provider_install_state( owner_config["vulkan_device_index"] = value moved += 1 removed += 1 - for provider in PROVIDERS: + for provider in clean_providers: record = bundled.get(provider) if not isinstance(record, dict): continue @@ -445,6 +492,269 @@ def migrate_legacy_provider_install_state( return mutate_journal_config(apply, journal_path=journal_path).value +def _legacy_provider_has_operational_state(value: object) -> bool: + return isinstance(value, dict) and any( + key in value for key in _LEGACY_OPERATIONAL_KEYS | {"vulkan_device_index"} + ) + + +def _migrate_legacy_provider( + provider: str, legacy: object, journal: Path +) -> dict[str, Any]: + from solstone.think.providers.install_lease import acquire_install_lease + + lease = acquire_install_lease(provider, journal_path=journal) + if lease is None: + return { + "provider": provider, + "status": "busy", + "reason_code": "install_busy", + "cleanup": False, + "message": ( + f"{provider} install is in progress; legacy provider state will be " + "retried on the next start." + ), + } + try: + if provider == "local": + return _migrate_legacy_local(legacy, journal) + return _migrate_legacy_parakeet(legacy, journal) + finally: + lease.release() + + +def _migrate_legacy_local(legacy: object, journal: Path) -> dict[str, Any]: + from solstone.think.models import LOCAL_MODEL + from solstone.think.providers import local_cuda, local_install, mlx_install + + readiness = local_install.inspect_readiness(LOCAL_MODEL) + fingerprint = local_install.target_fingerprint(LOCAL_MODEL) + if readiness.ready: + _publish_installed_status("local", fingerprint, journal) + return _ready_action("local", "already-ready") + if readiness.status in {"proof-unavailable", "host-ineligible"}: + return _not_promoted_action("local", readiness.status, readiness.reason_code) + if not isinstance(legacy, dict) or legacy.get("install_state") != "installed": + return _not_promoted_action( + "local", + "missing-or-mismatched", + "legacy_status_not_installed", + ) + if legacy.get("mlx_model_id") or legacy.get("mlx_snapshot_dir"): + try: + spec = mlx_install.resolve_model_spec(str(legacy.get("mlx_model_id") or "")) + except Exception: + spec = None + if ( + spec is not None + and legacy.get("mlx_revision") == spec.revision + and mlx_install.inspect_readiness(spec.name).ready + ): + _publish_installed_status( + "local", mlx_install.target_fingerprint(spec.name), journal + ) + return _ready_action("local", "already-ready") + return _not_promoted_action( + "local", + "missing-or-mismatched", + "manifest_missing", + message=( + "Existing local MLX artifacts cannot be trusted because there is " + "no Solstone manifest for the current pin. A reinstall will rebuild " + "the proof rather than trusting the old tree." + ), + ) + + choice = local_cuda.resolve_local_backend(local_install.CUDA_SERVER_PIN) + if choice.backend != "vulkan": + return _not_promoted_action( + "local", "missing-or-mismatched", "manifest_missing" + ) + try: + _verify_legacy_local_llama(legacy) + local_install._write_vulkan_manifest( + artifact_key=local_install.llama_server_artifact_key(), + pin=local_install.pin_for_current_platform(), + attempt_status=None, + fingerprint=fingerprint, + ) + local_install._write_model_manifest( + model_id=LOCAL_MODEL, + attempt_status=None, + fingerprint=fingerprint, + ) + except OSError: + return _not_promoted_action("local", "proof-unavailable", "legacy_io_error") + except Exception as exc: + return _not_promoted_action( + "local", + "missing-or-mismatched", + getattr(exc, "reason_code", "manifest_missing"), + message=( + "Existing local artifacts cannot be trusted because there is no " + "Solstone manifest for the current pin. A reinstall will rebuild " + "the proof rather than trusting the old tree." + ), + ) + final = local_install.inspect_readiness(LOCAL_MODEL) + if not final.ready: + return _not_promoted_action("local", final.status, final.reason_code) + _publish_installed_status("local", fingerprint, journal) + return _ready_action("local", "promoted") + + +def _migrate_legacy_parakeet(legacy: object, journal: Path) -> dict[str, Any]: + from solstone.think.providers import parakeet_install + + readiness = parakeet_install.inspect_readiness(journal) + fingerprint = parakeet_install.target_fingerprint(journal_path=journal) + if readiness.ready: + _publish_installed_status("parakeet", fingerprint, journal) + return _ready_action("parakeet", "already-ready") + if readiness.status in {"proof-unavailable", "host-ineligible"}: + return _not_promoted_action("parakeet", readiness.status, readiness.reason_code) + if not isinstance(legacy, dict) or legacy.get("install_state") != "installed": + return _not_promoted_action( + "parakeet", + "missing-or-mismatched", + "legacy_status_not_installed", + ) + try: + _verify_legacy_parakeet(legacy, journal) + for backend in ("cpu", "vulkan"): + parakeet_install._write_binary_manifest( + backend=backend, + attempt_status=None, + fingerprint=fingerprint, + journal_path=journal, + ) + parakeet_install._write_model_manifest( + attempt_status=None, + fingerprint=fingerprint, + journal_path=journal, + ) + except OSError: + return _not_promoted_action("parakeet", "proof-unavailable", "legacy_io_error") + except Exception as exc: + return _not_promoted_action( + "parakeet", + "missing-or-mismatched", + getattr(exc, "reason_code", "manifest_missing"), + ) + final = parakeet_install.inspect_readiness(journal) + if not final.ready: + return _not_promoted_action("parakeet", final.status, final.reason_code) + _publish_installed_status("parakeet", fingerprint, journal) + return _ready_action("parakeet", "promoted") + + +def _verify_legacy_local_llama(legacy: dict[str, Any]) -> None: + from solstone.think.models import LOCAL_MODEL + from solstone.think.providers import local_install + from solstone.think.providers.local import LOCAL_MODEL_SPECS + + artifact_key = local_install.llama_server_artifact_key() + pin = local_install.pin_for_current_platform() + spec = LOCAL_MODEL_SPECS[LOCAL_MODEL] + expected_binary = local_install.binary_path_for_pin(artifact_key, pin) + if legacy.get("binary_artifact") != artifact_key: + raise ValueError("legacy_binary_artifact_mismatch") + if legacy.get("binary_sha256") != pin["sha256"]: + raise ValueError("legacy_binary_pin_mismatch") + if legacy.get("binary_path") != str(expected_binary): + raise ValueError("legacy_binary_path_mismatch") + if not expected_binary.is_file() or not (expected_binary.stat().st_mode & 0o111): + raise ValueError("legacy_binary_missing") + if legacy.get("model_id") != spec.model_id: + raise ValueError("legacy_model_id_mismatch") + if legacy.get("model_path") != str(local_install.model_path(spec.model_id)): + raise ValueError("legacy_model_path_mismatch") + local_install._verify_sha256(local_install.model_path(spec.model_id), spec.sha256) + if spec.mmproj_sha256: + projector = local_install.mmproj_path(spec.model_id) + if projector is None or legacy.get("mmproj_path") != str(projector): + raise ValueError("legacy_projector_path_mismatch") + local_install._verify_sha256(projector, spec.mmproj_sha256) + + +def _verify_legacy_parakeet(legacy: dict[str, Any], journal: Path) -> None: + from solstone.think import parakeet_readiness + from solstone.think.providers import parakeet_install + + artifact_key = parakeet_install.parakeet_server_artifact_key() + for backend in ("cpu", "vulkan"): + pin = parakeet_install._pin_for_backend(artifact_key, backend) + if legacy.get(f"binary_artifact_{backend}") != artifact_key: + raise ValueError(f"legacy_binary_artifact_{backend}_mismatch") + if legacy.get(f"binary_sha256_{backend}") != pin["sha256"]: + raise ValueError(f"legacy_binary_pin_{backend}_mismatch") + if legacy.get(f"binary_path_{backend}") != str( + parakeet_install.binary_path(backend, journal) + ): + raise ValueError(f"legacy_binary_path_{backend}_mismatch") + parakeet_readiness.check_parakeet_cpp_files(journal) + spec = parakeet_install.PARAKEET_MODEL_SPEC + if legacy.get("model_repo") != spec.repo: + raise ValueError("legacy_model_repo_mismatch") + if legacy.get("model_filename") != spec.filename: + raise ValueError("legacy_model_filename_mismatch") + if legacy.get("model_revision") != spec.revision: + raise ValueError("legacy_model_revision_mismatch") + if legacy.get("model_path") != str(parakeet_install.model_path(journal)): + raise ValueError("legacy_model_path_mismatch") + parakeet_install._verify_sha256(parakeet_install.model_path(journal), spec.sha256) + + +def _publish_installed_status( + provider: str, + fingerprint: dict[str, Any], + journal: Path, +) -> InstallStatus: + current = read_install_status(name=provider, journal_path=journal) + if current["install_state"] in IN_FLIGHT_STATES: + raise InstallStatusConflictError("cannot migrate while install is in flight") + fingerprint_json = canonical_fingerprint(fingerprint) + current["target_fingerprint_json"] = fingerprint_json + current["target_fingerprint_sha256"] = fingerprint_sha256(fingerprint_json) + current["owner"] = {"entry": "legacy_provider_install_state_migration"} + return write_install_status( + transition_state(current, new_state="installed"), + journal_path=journal, + ) + + +def _ready_action(provider: str, action: str) -> dict[str, Any]: + return { + "provider": provider, + "status": "ready", + "reason_code": "ready", + "action": action, + "cleanup": True, + "message": f"{provider} provider install state migrated.", + } + + +def _not_promoted_action( + provider: str, + status: str, + reason_code: str, + *, + message: str | None = None, +) -> dict[str, Any]: + if message is None: + message = ( + f"{provider} provider legacy install state was not promoted: {reason_code}." + ) + return { + "provider": provider, + "status": status, + "reason_code": reason_code, + "action": "not-promoted", + "cleanup": False, + "message": message, + } + + def _read_current_unlocked(path: Path, provider: ProviderName) -> InstallStatus: if not path.exists(): return make_idle_status(provider) @@ -663,6 +973,7 @@ __all__ = [ "canonical_fingerprint", "fingerprint_sha256", "make_idle_status", + "migrate_legacy_provider_artifact_truth", "migrate_legacy_provider_install_state", "now_iso", "observe_install_attempt", diff --git a/solstone/think/supervisor.py b/solstone/think/supervisor.py index 92e17c941..3e3c78d75 100644 --- a/solstone/think/supervisor.py +++ b/solstone/think/supervisor.py @@ -3705,7 +3705,9 @@ def main() -> None: # Run pending journal-maintenance tasks before spawning any writer children. # Callosum isn't up yet (emit_fn=None); migrations log through supervisor's logger only. try: - ran, succeeded = run_pending_tasks(journal_path, emit_fn=None) + maint_results = run_pending_tasks(journal_path, emit_fn=None) + ran = len(maint_results) + succeeded = sum(1 for result in maint_results if result.success) if ran > 0: print(f" Ran {ran} maintenance task(s)", flush=True) if ran == succeeded: @@ -3716,6 +3718,23 @@ def main() -> None: succeeded, ran, ) + blocking_failures = [ + result + for result in maint_results + if not result.success and result.task.blocks_supervisor_start + ] + if blocking_failures: + failure = blocking_failures[0] + message = ( + "Startup blocked by maintenance task " + f"{failure.task.qualified_name} " + f"(exit {failure.exit_code}). Log: {failure.state_file}. " + "This task is retry-on-next-start; fix the error and start " + "the supervisor again." + ) + logging.error(message) + print(f" {message}", file=sys.stderr, flush=True) + sys.exit(1) except Exception: logging.exception("Maintenance runner raised; continuing startup") diff --git a/tests/test_check_provider_install_owner.py b/tests/test_check_provider_install_owner.py new file mode 100644 index 000000000..5a8d0d616 --- /dev/null +++ b/tests/test_check_provider_install_owner.py @@ -0,0 +1,191 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import importlib.util +import subprocess +import sys +from pathlib import Path +from types import ModuleType + +SCRIPT = ( + Path(__file__).resolve().parents[1] / "scripts" / "check_provider_install_owner.py" +) + + +def _load_checker() -> ModuleType: + spec = importlib.util.spec_from_file_location( + "check_provider_install_owner", SCRIPT + ) + assert spec is not None + assert spec.loader is not None + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + +def _write_file(root: Path, rel: str, source: str) -> Path: + path = root / rel + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(source, encoding="utf-8") + return path + + +def _run(root: Path) -> subprocess.CompletedProcess[str]: + return subprocess.run( + [sys.executable, str(SCRIPT), "--root", str(root)], + check=False, + capture_output=True, + text=True, + ) + + +def _kinds(findings: list[tuple[int, str, str]]) -> set[str]: + return {kind for _lineno, kind, _detail in findings} + + +checker = _load_checker() + + +def test_scan_flags_direct_status_replace() -> None: + findings = checker.scan_source( + "from pathlib import Path\n" + "from solstone.think.journal_io.atomic import atomic_replace\n" + "status = Path('journal') / 'health' / 'providers' / 'local.json'\n" + "atomic_replace(status, '{}')\n" + ) + + assert "provider_status_replace" in _kinds(findings) + + +def test_scan_flags_proof_cache_write() -> None: + findings = checker.scan_source( + "from solstone.think.providers.artifact_proof import proof_cache_path\n" + "path = proof_cache_path('local')\n" + "path.write_text('{}')\n" + ) + + assert "proof_cache_write" in _kinds(findings) + + +def test_scan_flags_raw_lease_open_and_flock() -> None: + findings = checker.scan_source( + "import fcntl\n" + "import os\n" + "from pathlib import Path\n" + "lease = Path('journal') / 'health' / 'providers' / 'local.lease'\n" + "fd = os.open(lease, os.O_RDWR)\n" + "fcntl.flock(fd, fcntl.LOCK_EX)\n" + ) + + kinds = _kinds(findings) + assert "provider_lease_raw_open" in kinds + assert "second_provider_install_lock" in kinds + + +def test_scan_flags_manifest_write() -> None: + findings = checker.scan_source( + "from pathlib import Path\n" + "manifest = Path('cache') / '.solstone-provider-manifest.json'\n" + "manifest.write_text('{}')\n" + ) + + assert "provider_manifest_write" in _kinds(findings) + + +def test_scan_flags_private_owner_alias() -> None: + findings = checker.scan_source( + "from solstone.think.providers.install_state import _read_current_unlocked as raw\n" + "writer = raw\n" + ) + + kinds = _kinds(findings) + assert "private_owner_symbol" in kinds + assert "private_owner_wrapper" in kinds + + +def test_scan_flags_providers_bundled_access() -> None: + findings = checker.scan_source( + "def f(config):\n" + " providers = config.get('providers', {})\n" + " return providers.get('bundled')\n" + ) + + assert "providers_bundled_operational" in _kinds(findings) + + +def test_scan_allows_owner_api_calls() -> None: + findings = checker.scan_source( + "from solstone.think.providers.install_lease import acquire_install_lease\n" + "from solstone.think.providers.install_state import write_install_status\n" + "lease = acquire_install_lease('local')\n" + "write_install_status(status)\n" + ) + + assert findings == [] + + +def test_e2e_flags_violation(tmp_path: Path) -> None: + _write_file( + tmp_path, + "solstone/bad.py", + "from pathlib import Path\n" + "p = Path('journal') / 'health' / 'providers' / 'local.json'\n" + "p.write_text('{}')\n", + ) + + result = _run(tmp_path) + + assert result.returncode == 1 + assert "provider-install-owner: violations:" in result.stderr + assert "provider_status_write" in result.stderr + + +def test_e2e_clean_source_passes(tmp_path: Path) -> None: + _write_file( + tmp_path, + "solstone/good.py", + "from solstone.think.providers.install_state import write_install_status\n" + "def f(status):\n" + " return write_install_status(status)\n", + ) + + result = _run(tmp_path) + + assert result.returncode == 0 + assert "provider-install-owner: pass" in result.stdout + assert result.stderr == "" + + +def test_allowlist_ratchet_and_stale_entry(tmp_path: Path) -> None: + _write_file( + tmp_path, + "solstone/bad.py", + "from pathlib import Path\n" + "p = Path('journal') / 'health' / 'providers' / 'local.json'\n" + "p.write_text('{}')\n", + ) + counts = checker.count_violations(tmp_path) + + over, stale, tracked = checker.evaluate(tmp_path, counts) + assert over == [] + assert stale == [] + assert tracked + + ratcheted = {next(iter(counts)): 0} + over, stale, _tracked = checker.evaluate(tmp_path, ratcheted) + assert over + assert stale == [] + + stale_allowlist = {("solstone/missing.py", "provider_status_write"): 1} + over, stale, _tracked = checker.evaluate(tmp_path, stale_allowlist) + assert over + assert stale + + +def test_landed_tree_is_clean() -> None: + over, stale, _tracked = checker.evaluate(checker.ROOT, checker.ALLOWLIST) + + assert over == [] + assert stale == [] diff --git a/tests/test_maint.py b/tests/test_maint.py index 8298c5b18..b527ef9f3 100644 --- a/tests/test_maint.py +++ b/tests/test_maint.py @@ -15,6 +15,7 @@ import pytest from solstone.think.maint import ( MaintTask, _terminate_with_grace, + discover_tasks, get_state_file, get_task_status, list_tasks, @@ -105,6 +106,24 @@ class TestStatusTracking: assert exit_code is None assert ran_ts == 1000 + def test_latest_attempt_status_wins(self, temp_journal): + state_dir = temp_journal / "maint" / "chat" + state_dir.mkdir(parents=True) + state_file = state_dir / "retry_task.jsonl" + state_file.write_text( + '{"event": "exec", "attempt_id": "a1", "ts": 1000}\n' + '{"event": "exit", "attempt_id": "a1", "ts": 2000, "exit_code": 1}\n' + '{"event": "exec", "attempt_id": "a2", "ts": 3000}\n' + '{"event": "line", "attempt_id": "a2", "ts": 3500, "line": "retry"}\n' + '{"event": "exit", "attempt_id": "a2", "ts": 4000, "exit_code": 0}\n' + ) + + status, exit_code, ran_ts = get_task_status(temp_journal, "chat", "retry_task") + + assert status == "success" + assert exit_code == 0 + assert ran_ts == 4000 + class TestListTasks: """Tests for listing tasks with status metadata.""" @@ -275,11 +294,92 @@ class TestRunPendingTasks: encoding="utf-8", ) - ran, succeeded = run_pending_tasks(journal) + results = run_pending_tasks(journal) - assert (ran, succeeded) == (0, 0) + assert results == [] assert json.loads(config_path.read_text("utf-8")) == payload + def test_failed_task_retries_only_when_opted_in(self, monkeypatch, tmp_path): + journal = tmp_path / "journal" + retry_task = MaintTask( + app="test_app", + name="retry", + script_path=Path("/dummy/retry.py"), + retry_on_next_start=True, + ) + no_retry_task = MaintTask( + app="test_app", + name="no_retry", + script_path=Path("/dummy/no_retry.py"), + ) + in_progress_task = MaintTask( + app="test_app", + name="in_progress", + script_path=Path("/dummy/in_progress.py"), + retry_on_next_start=True, + ) + monkeypatch.setattr( + "solstone.think.maint.discover_tasks", + lambda: [retry_task, no_retry_task, in_progress_task], + ) + for task, body in ( + ( + retry_task, + '{"event": "exec", "ts": 1000}\n' + '{"event": "exit", "ts": 2000, "exit_code": 1}\n', + ), + ( + no_retry_task, + '{"event": "exec", "ts": 1000}\n' + '{"event": "exit", "ts": 2000, "exit_code": 1}\n', + ), + (in_progress_task, '{"event": "exec", "ts": 1000}\n'), + ): + state_file = get_state_file(journal, task.app, task.name) + state_file.parent.mkdir(parents=True, exist_ok=True) + state_file.write_text(body, encoding="utf-8") + calls = [] + + def fake_run_task(_journal, task, emit_fn=None): + calls.append(task.name) + state_file = get_state_file(_journal, task.app, task.name) + with state_file.open("a", encoding="utf-8") as handle: + handle.write( + '{"event": "exec", "attempt_id": "retry", "ts": 3000}\n' + '{"event": "exit", "attempt_id": "retry", "ts": 4000, ' + '"exit_code": 0}\n' + ) + return True, 0 + + monkeypatch.setattr("solstone.think.maint.run_task", fake_run_task) + + results = run_pending_tasks(journal) + + assert calls == ["retry"] + assert [result.task.name for result in results] == ["retry"] + assert results[0].success is True + + def test_existing_tasks_keep_default_opt_ins(self): + tasks = discover_tasks() + existing = [ + task + for task in tasks + if task.qualified_name != "thinking:001_migrate_provider_install_state" + ] + + assert len(existing) == 22 + assert all(not task.retry_on_next_start for task in existing) + assert all(not task.blocks_supervisor_start for task in existing) + + migration = [ + task + for task in tasks + if task.qualified_name == "thinking:001_migrate_provider_install_state" + ] + assert len(migration) == 1 + assert migration[0].retry_on_next_start is True + assert migration[0].blocks_supervisor_start is True + class TestRunTask: """Tests for running individual tasks.""" @@ -478,6 +578,7 @@ class TestRunTask: last_event = events[-1] assert list(last_event.keys()) == [ "event", + "attempt_id", "ts", "exit_code", "duration_ms", diff --git a/tests/test_maintenance.py b/tests/test_maintenance.py index ab683b45a..8e7cd1d2b 100644 --- a/tests/test_maintenance.py +++ b/tests/test_maintenance.py @@ -574,7 +574,7 @@ def test_supervisor_registers_maintenance_before_scheduler_init(tmp_path, monkey "argv", ["supervisor", "0", "--no-daily", "--no-convey", "--no-cortex", "--no-spl"], ) - monkeypatch.setattr(mod, "run_pending_tasks", lambda *a, **k: (0, 0)) + monkeypatch.setattr(mod, "run_pending_tasks", lambda *a, **k: []) monkeypatch.setattr(mod, "_sweep_orphaned_sol_processes", lambda *_a, **_k: 0) monkeypatch.setattr( mod, diff --git a/tests/test_supervisor.py b/tests/test_supervisor.py index c578a30fc..424391ec2 100644 --- a/tests/test_supervisor.py +++ b/tests/test_supervisor.py @@ -20,6 +20,7 @@ from unittest.mock import MagicMock import psutil import pytest +from solstone.think.maint import MaintTask, MaintTaskResult from solstone.think.processing import ( DISPLAY_POWERSAVE_UNAVAILABLE, DisplayPowersaveReading, @@ -291,7 +292,7 @@ def test_graceful_shutdown_calls_stop_process_for_each_managed_proc( "argv", ["supervisor", "0", "--no-daily", "--no-schedule"], ) - monkeypatch.setattr(mod, "run_pending_tasks", lambda *a, **k: (0, 0)) + monkeypatch.setattr(mod, "run_pending_tasks", lambda *a, **k: []) monkeypatch.setattr(mod, "_sweep_orphaned_sol_processes", lambda *_a, **_k: 0) monkeypatch.setattr(mod.time, "sleep", lambda _seconds: None) monkeypatch.setattr(mod, "start_callosum_in_process", lambda: None) @@ -365,7 +366,7 @@ def test_supervisor_readiness_marker_requires_started_convey_accepting( "argv", ["supervisor", "0", "--no-daily", "--no-schedule"], ) - monkeypatch.setattr(mod, "run_pending_tasks", lambda *a, **k: (0, 0)) + monkeypatch.setattr(mod, "run_pending_tasks", lambda *a, **k: []) monkeypatch.setattr(mod, "_sweep_orphaned_sol_processes", lambda *_a, **_k: 0) monkeypatch.setattr(mod.time, "sleep", lambda _seconds: None) monkeypatch.setattr(mod, "start_callosum_in_process", lambda: None) @@ -429,7 +430,7 @@ def _run_supervisor_main_for_shutdown_knobs(tmp_path, monkeypatch, *, argv): monkeypatch.delenv("SOL_SUPERVISOR_SPAWNED", raising=False) monkeypatch.delenv("SOLSTONE_APP_SUPERVISED", raising=False) monkeypatch.setattr(sys, "argv", argv) - monkeypatch.setattr(mod, "run_pending_tasks", lambda *a, **k: (0, 0)) + monkeypatch.setattr(mod, "run_pending_tasks", lambda *a, **k: []) monkeypatch.setattr(mod, "_sweep_orphaned_sol_processes", lambda *_a, **_k: 0) monkeypatch.setattr(mod.time, "sleep", lambda _seconds: None) monkeypatch.setattr(mod, "start_callosum_in_process", lambda: None) @@ -3889,7 +3890,7 @@ def test_supervisor_singleton_lock_acquired(tmp_path, monkeypatch): # Skip maint discovery/subprocess runs — unrelated to lock acquisition and # slow enough on a fresh tmp_path to blow the 5s pytest-timeout under load. - monkeypatch.setattr(mod, "run_pending_tasks", lambda *a, **k: (0, 0)) + monkeypatch.setattr(mod, "run_pending_tasks", lambda *a, **k: []) monkeypatch.setattr(mod, "_sweep_orphaned_sol_processes", lambda *_a, **_k: 0) monkeypatch.setattr(mod.time, "sleep", lambda _seconds: None) monkeypatch.setattr(mod, "start_callosum_in_process", stop_after_lock) @@ -3909,6 +3910,43 @@ def test_supervisor_singleton_lock_acquired(tmp_path, monkeypatch): assert mod.is_supervisor_up() is True +def test_supervisor_blocks_before_callosum_on_blocking_maint_failure( + tmp_path, monkeypatch, capsys +): + mod = importlib.reload(importlib.import_module("solstone.think.supervisor")) + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + (tmp_path / "health").mkdir(parents=True, exist_ok=True) + monkeypatch.setattr(sys, "argv", ["supervisor"]) + task = MaintTask( + app="thinking", + name="001_migrate_provider_install_state", + script_path=Path("/dummy.py"), + retry_on_next_start=True, + blocks_supervisor_start=True, + ) + state_file = tmp_path / "maint" / "thinking" / f"{task.name}.jsonl" + result = MaintTaskResult( + task=task, + success=False, + exit_code=7, + state_file=state_file, + ) + monkeypatch.setattr(mod, "run_pending_tasks", lambda *a, **k: [result]) + monkeypatch.setattr(mod, "_sweep_orphaned_sol_processes", lambda *_a, **_k: 0) + start_mock = MagicMock() + monkeypatch.setattr(mod, "start_callosum_in_process", start_mock) + + with pytest.raises(SystemExit) as exc: + mod.main() + + assert exc.value.code == 1 + start_mock.assert_not_called() + captured = capsys.readouterr() + assert "thinking:001_migrate_provider_install_state" in captured.err + assert str(state_file) in captured.err + assert "retry-on-next-start" in captured.err + + def test_supervisor_singleton_lock_blocked(tmp_path, monkeypatch, capsys): import fcntl diff --git a/tests/test_supervisor_sync_gate.py b/tests/test_supervisor_sync_gate.py index bf7f1a9fd..e8af28309 100644 --- a/tests/test_supervisor_sync_gate.py +++ b/tests/test_supervisor_sync_gate.py @@ -53,7 +53,7 @@ def _load_supervisor(tmp_path, monkeypatch, argv=None): monkeypatch.delenv("SOL_SUPERVISOR_SPAWNED", raising=False) monkeypatch.setattr(sys, "argv", argv or ["supervisor"]) monkeypatch.setattr(mod.time, "sleep", lambda _seconds: None) - monkeypatch.setattr(mod, "run_pending_tasks", lambda *args, **kwargs: (0, 0)) + monkeypatch.setattr(mod, "run_pending_tasks", lambda *args, **kwargs: []) monkeypatch.setattr(mod, "_sweep_orphaned_sol_processes", lambda *_a, **_k: 0) return mod