"""``muse maintenance`` — scheduled store maintenance orchestration. Orchestrates multiple store-health tasks under a single command with a shared schedule configuration and persistent run-history. Subcommands ----------- ``run [--task ]... [--all] [--dry-run] [--json]`` Execute one or more maintenance tasks. Available tasks: ``gc`` Remove unreachable objects, commits, and snapshots from the store (wraps :func:`muse.core.gc.run_gc` with ``full=True``). ``verify-objects`` Rehash every object in the store and report any whose content does not match their filename. Non-destructive; always safe to run. Without ``--task``, runs the tasks listed in the persisted schedule config (default: ``gc`` only). ``--all`` runs every known task. Unless ``--dry-run`` is set, each task's completion time is written back to ``.muse/maintenance.json`` so ``status`` can report freshness. ``status [--json]`` Show the schedule configuration and the last-run timestamp for each task. ``schedule [--period-hours N] [--enable] [--disable]`` Write or update the schedule configuration in ``.muse/maintenance.json``. ``--period-hours`` sets the suggested inter-run interval (informational — Muse does not run background daemons; callers use this value to decide whether to trigger ``run``). State ----- All state is stored in ``.muse/maintenance.json``:: { "enabled": true, "period_hours": 24, "tasks": ["gc", "verify-objects"], "last_run": { "gc": "2026-04-14T12:00:00+00:00", "verify-objects": "2026-04-14T12:00:00+00:00" } } Exit codes:: 0 — all tasks completed (even if verify-objects found issues) 1 — a task argument was invalid 2 — usage error """ from __future__ import annotations import argparse import datetime import json as _json import logging import sys import time from typing import Any from muse.core.errors import ExitCode from muse.core.gc import _DEFAULT_GRACE_PERIOD_SECONDS, run_gc from muse.core.repo import require_repo logger = logging.getLogger(__name__) # --------------------------------------------------------------------------- # Constants # --------------------------------------------------------------------------- _ALL_TASKS: list[str] = ["gc", "verify-objects"] _DEFAULT_TASKS: list[str] = ["gc"] _DEFAULT_PERIOD_HOURS: int = 24 _CONFIG_FILE = "maintenance.json" # --------------------------------------------------------------------------- # Config I/O # --------------------------------------------------------------------------- def _config_path(root) -> "pathlib.Path": import pathlib return root / ".muse" / _CONFIG_FILE def _read_config(root) -> dict[str, Any]: path = _config_path(root) if not path.exists(): return { "enabled": True, "period_hours": _DEFAULT_PERIOD_HOURS, "tasks": list(_DEFAULT_TASKS), "last_run": {}, } try: return _json.loads(path.read_text(encoding="utf-8")) except (OSError, _json.JSONDecodeError): return { "enabled": True, "period_hours": _DEFAULT_PERIOD_HOURS, "tasks": list(_DEFAULT_TASKS), "last_run": {}, } def _write_config(root, cfg: dict[str, Any]) -> None: _config_path(root).write_text(_json.dumps(cfg, indent=2), encoding="utf-8") def _now_iso() -> str: return datetime.datetime.now(datetime.timezone.utc).isoformat() # --------------------------------------------------------------------------- # Task implementations # --------------------------------------------------------------------------- def _run_gc(root, *, dry_run: bool = False) -> dict[str, Any]: """Run garbage collection; return a result summary dict.""" t0 = time.monotonic() result = run_gc( root, dry_run=dry_run, grace_period_seconds=0 if dry_run else _DEFAULT_GRACE_PERIOD_SECONDS, full=True, ) return { "collected_count": result.collected_count, "collected_bytes": result.collected_bytes, "reachable_count": result.reachable_count, "commits_collected": result.commits_collected, "snapshots_collected": result.snapshots_collected, "dry_run": dry_run, "elapsed_seconds": time.monotonic() - t0, } def _run_verify_objects(root) -> dict[str, Any]: """Rehash every object; return summary with checked/failed counts.""" import hashlib import pathlib t0 = time.monotonic() objects_dir = root / ".muse" / "objects" checked = 0 failed = 0 failures: list[str] = [] if not objects_dir.exists(): return { "checked": 0, "failed": 0, "failures": [], "elapsed_seconds": time.monotonic() - t0, } for prefix_dir in sorted(objects_dir.iterdir()): if not prefix_dir.is_dir() or len(prefix_dir.name) != 2: continue for obj_file in sorted(prefix_dir.iterdir()): if not obj_file.is_file(): continue expected_id = prefix_dir.name + obj_file.name if len(expected_id) != 64: continue checked += 1 try: h = hashlib.sha256() with obj_file.open("rb") as fh: for chunk in iter(lambda: fh.read(65536), b""): h.update(chunk) actual = h.hexdigest() if actual != expected_id: failed += 1 failures.append(expected_id) except OSError: failed += 1 failures.append(expected_id) return { "checked": checked, "failed": failed, "failures": failures[:20], # cap to avoid huge JSON "elapsed_seconds": time.monotonic() - t0, } # --------------------------------------------------------------------------- # Subcommand handlers # --------------------------------------------------------------------------- _TASK_RUNNERS = { "gc": _run_gc, "verify-objects": _run_verify_objects, } def _cmd_run(args: argparse.Namespace, root) -> None: dry_run: bool = args.dry_run output_json: bool = args.output_json # Validate requested tasks. if args.all: tasks_to_run = list(_ALL_TASKS) elif args.task: tasks_to_run = list(args.task) unknown = [t for t in tasks_to_run if t not in _TASK_RUNNERS] if unknown: print( f"❌ Unknown task(s): {', '.join(unknown)}. " f"Available: {', '.join(_ALL_TASKS)}", file=sys.stderr, ) raise SystemExit(ExitCode.USER_ERROR) else: cfg = _read_config(root) tasks_to_run = cfg.get("tasks") or list(_DEFAULT_TASKS) t0_total = time.monotonic() results: dict[str, Any] = {} for task in tasks_to_run: runner_fn = _TASK_RUNNERS[task] if task == "gc": results[task] = runner_fn(root, dry_run=dry_run) else: results[task] = runner_fn(root) elapsed = time.monotonic() - t0_total # Persist timestamps (skipped in dry-run). if not dry_run: cfg = _read_config(root) now = _now_iso() if "last_run" not in cfg: cfg["last_run"] = {} for task in tasks_to_run: cfg["last_run"][task] = now _write_config(root, cfg) payload: dict[str, Any] = { "tasks_run": tasks_to_run, "results": results, "dry_run": dry_run, "elapsed_seconds": round(elapsed, 4), } if output_json: print(_json.dumps(payload)) else: prefix = "[dry-run] " if dry_run else "" for task in tasks_to_run: r = results[task] if task == "gc": action = "Would remove" if dry_run else "Removed" print( f"{prefix}{action} {r['collected_count']} unreachable object(s) " f"({r['collected_bytes']} bytes); " f"{r['reachable_count']} reachable" ) elif task == "verify-objects": status = "✅" if r["failed"] == 0 else "❌" print( f"{prefix}{status} verify-objects: " f"{r['checked']} checked, {r['failed']} failed" ) print( f"{prefix}Done in {elapsed:.3f}s" + (" (no changes written)" if dry_run else "") ) def _cmd_status(args: argparse.Namespace, root) -> None: cfg = _read_config(root) output_json: bool = args.output_json enabled = cfg.get("enabled", True) period_hours = cfg.get("period_hours", _DEFAULT_PERIOD_HOURS) last_run: dict[str, str] = cfg.get("last_run", {}) if output_json: print(_json.dumps({ "enabled": enabled, "period_hours": period_hours, "tasks": cfg.get("tasks", list(_DEFAULT_TASKS)), "last_run": last_run, })) return status_label = "enabled" if enabled else "disabled" print(f"Maintenance: {status_label} (period: {period_hours}h)") if not last_run: print(" Never run.") else: for task, ts in sorted(last_run.items()): print(f" {task:<20} last run: {ts}") def _cmd_schedule(args: argparse.Namespace, root) -> None: period_hours: int | None = args.period_hours enable: bool = args.enable disable: bool = args.disable if period_hours is not None and period_hours < 0: print("❌ --period-hours must be ≥ 0", file=sys.stderr) raise SystemExit(ExitCode.USER_ERROR) cfg = _read_config(root) if period_hours is not None: cfg["period_hours"] = period_hours if enable: cfg["enabled"] = True if disable: cfg["enabled"] = False _write_config(root, cfg) status = "enabled" if cfg.get("enabled", True) else "disabled" ph = cfg.get("period_hours", _DEFAULT_PERIOD_HOURS) print(f"Maintenance schedule updated: {status}, period: {ph}h") # --------------------------------------------------------------------------- # Registration # --------------------------------------------------------------------------- def register( subparsers: "argparse._SubParsersAction[argparse.ArgumentParser]", ) -> None: """Register the ``muse maintenance`` subcommand.""" parser = subparsers.add_parser( "maintenance", help="Scheduled store maintenance orchestration.", description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter, ) sub = parser.add_subparsers(dest="maint_command", metavar="SUBCOMMAND") sub.required = True # run p_run = sub.add_parser( "run", help="Run maintenance tasks.", ) p_run.add_argument( "--task", action="append", metavar="TASK", choices=_ALL_TASKS, dest="task", help=f"Task to run (repeatable). Choices: {', '.join(_ALL_TASKS)}", ) p_run.add_argument( "--all", action="store_true", default=False, help="Run all maintenance tasks.", ) p_run.add_argument( "--dry-run", action="store_true", default=False, dest="dry_run", help="Preview without making changes or updating timestamps.", ) p_run.add_argument( "--json", action="store_true", default=False, dest="output_json", help="Emit a JSON result object.", ) p_run.set_defaults(maint_func=_cmd_run) # status p_status = sub.add_parser("status", help="Show schedule and last-run info.") p_status.add_argument( "--json", action="store_true", default=False, dest="output_json", ) p_status.set_defaults(maint_func=_cmd_status) # schedule p_sched = sub.add_parser("schedule", help="Configure maintenance schedule.") p_sched.add_argument( "--period-hours", type=int, default=None, metavar="N", dest="period_hours", help="Suggested interval between runs (hours).", ) p_sched.add_argument( "--enable", action="store_true", default=False, help="Mark maintenance as enabled.", ) p_sched.add_argument( "--disable", action="store_true", default=False, help="Mark maintenance as disabled.", ) p_sched.set_defaults(maint_func=_cmd_schedule) parser.set_defaults(func=run) # --------------------------------------------------------------------------- # Entry point # --------------------------------------------------------------------------- def run(args: argparse.Namespace) -> None: root = require_repo() args.maint_func(args, root)