Source code for spacr.cli_repro

"""
``spacr repro`` — replay a pipeline run from its recorded journal.

Every pipeline invocation writes a run journal (see
:mod:`spacr.run_journal`) containing the exact settings + environment
that produced a result. This CLI re-runs those settings, opens a
FRESH journal folder, and reports whether the outcome matches.

Usage::

    spacr repro ~/.spacr/runs/2026-07-23_143507_ab12cd34__mask
    spacr repro ~/.spacr/runs/2026-07-23_143507_ab12cd34__mask --dry
    spacr repro ~/.spacr/runs/2026-07-23_143507_ab12cd34__mask --show

* ``--dry``  prints the resolved settings + which pipeline entry
  will run; doesn't invoke it.
* ``--show`` prints the manifest + settings; doesn't invoke it.
* ``--export snakemake|nextflow --out DIR`` writes the run as a
  workflow that runs the same module once per plate with ``spacr-run``,
  locally, on a cluster, or inside the spaCR Docker/Apptainer image.

Exit codes:
  0  — replay ran to completion (regardless of scientific outcome)
  1  — replay raised (journal captures the traceback for triage)
  2  — bad input (missing run folder / unresolvable app_key)
"""
from __future__ import annotations

import argparse
import json
import os
import re
import sys
from pathlib import Path
from typing import Any, Dict, List, Optional

from .run_journal import load_run_settings, open_run, runs_root


def _print_manifest(run_dir: Path) -> None:
    """Print the recorded run summary and optional model hashes.

    :param run_dir: Run-journal directory containing ``manifest.json``.
    :returns: ``None``; the formatted summary is written to standard output.
    """
    m = json.loads((run_dir / "manifest.json").read_text())
    print(f"run:       {run_dir.name}")
    print(f"app:       {m.get('app_key')}")
    print(f"status:    {m.get('status')}")
    print(f"start:     {m.get('start_utc')}")
    print(f"elapsed:   {m.get('elapsed_s')}s")
    print(f"n_settings:{m.get('n_settings')}")
    env = m.get("env", {})
    print(f"spacr:     {env.get('spacr')} (git {env.get('spacr_git')})")
    print(f"python:    {env.get('python')}  torch {env.get('torch')}  "
          f"cellpose {env.get('cellpose')}")
    if m.get("model_hashes"):
        print("models:")
        for k, v in m["model_hashes"].items():
            print(f"  {k}: {v}")


def _print_settings(settings: dict) -> None:
    """Print settings as deterministic, key-sorted rows.

    :param settings: Setting names and values to display.
    :returns: ``None``; the rows are written to standard output.
    """
    print("settings:")
    for k, v in sorted(settings.items()):
        print(f"  {k:32s} = {v!r}")


def _resolve_pipeline(app_key: str):
    """Return the pipeline callable for ``app_key`` or ``None``."""
    try:
        from .qt.bridge import resolve_pipeline_entry
        return resolve_pipeline_entry(app_key)
    except Exception:
        return None


_WORKFLOW_IMAGE = "ghcr.io/einarolafsson/spacr"
_WORKFLOW_ENGINES = ("snakemake", "nextflow")
_WORKFLOW_FILE_OUTPUTS = {
    "convert": ("db_path", "checkpoint_path"),
    "align": ("db_path",),
}


def _workflow_image(manifest: Dict[str, Any]) -> str:
    """Return the CPU container image for the spaCR version that ran.

    A released version (digits and dots only) pins its own tag; any other
    version falls back to the rolling ``:cpu`` tag.
    """
    version = str((manifest.get("env") or {}).get("spacr") or "")
    if re.fullmatch(r"\d+(\.\d+)+", version):
        return f"{_WORKFLOW_IMAGE}:{version}-cpu"
    return f"{_WORKFLOW_IMAGE}:cpu"


def _workflow_plates(settings: Dict[str, Any],
                     plates: Optional[List[str]] = None, *,
                     module: str = "") -> Dict[str, Dict[str, Any]]:
    """Split recorded settings into one settings dict per plate.

    Each entry of ``src`` (a path or a list of paths) becomes its own job
    with ``src`` set to that one plate. ``plates`` replaces the recorded
    plates. Settings with no ``src`` stay one job named ``run``. When there
    is more than one job, explicit ``dst`` and ``dst_root`` output folders
    gain the unique plate id as a subfolder, so jobs do not overwrite each
    other's outputs. Single-job destinations and unset folders stay unchanged.
    Convert and Align database outputs, and new Convert checkpoints, are
    separated too; the same setting names in other modules are untouched.
    Convert map names that escape the destination folder are separated as
    explicit output paths; relative names inside it retain their spelling.

    :param settings: recorded settings; the input dictionary is not changed.
    :param plates: optional replacement source folders.
    :param module: canonical CLI module key, identifying output-file settings.
    :returns: ``{plate_id: settings}``, ids unique and safe as file names.
    :raises ValueError: when multiple Convert jobs would share an explicit
        checkpoint requested as resume input.
    """
    src = settings.get("src")
    if plates:
        sources = [str(p) for p in plates]
    elif isinstance(src, (list, tuple)):
        sources = [str(p) for p in src if str(p).strip()]
    elif isinstance(src, str) and src.strip():
        sources = [src]
    else:
        sources = []
    if not sources:
        return {"run": dict(settings)}
    if (len(sources) > 1 and module == "convert"
            and settings.get("resume") and settings.get("checkpoint_path")):
        raise ValueError(
            "Multiple Convert jobs cannot reuse one explicit resume checkpoint. "
            "Export one plate, or clear checkpoint_path to use each plate's "
            "destination checkpoint.")
    jobs: Dict[str, Dict[str, Any]] = {}
    for source in sources:
        stem = re.sub(r"[^A-Za-z0-9_.-]+", "_",
                      Path(source.rstrip("/\\")).name).strip("._") or "plate"
        plate_id, n = stem, 2
        while plate_id in jobs:
            plate_id, n = f"{stem}_{n}", n + 1
        jobs[plate_id] = {**settings, "src": source}
        if len(sources) > 1:
            for key in ("dst", "dst_root"):
                destination = settings.get(key)
                if isinstance(destination, str) and destination:
                    jobs[plate_id][key] = str(Path(destination) / plate_id)
            file_outputs = {key: settings.get(key)
                            for key in _WORKFLOW_FILE_OUTPUTS.get(module, ())}
            map_name = settings.get("map_name")
            if module == "convert" and isinstance(map_name, str) and map_name:
                map_path = Path(os.path.normpath(map_name))
                if map_path.is_absolute() or map_path.parts[:1] == ("..",):
                    original_root = settings.get("dst") or (
                        os.path.normpath(os.path.abspath(source)) + "_yokogawa")
                    file_outputs["map_name"] = os.path.normpath(os.path.join(
                        os.path.abspath(str(original_root)), map_name))
            for key, destination in file_outputs.items():
                if not isinstance(destination, str) or not destination:
                    continue
                path = Path(os.path.normpath(destination))
                roots = [(Path(os.path.abspath(settings[root]) if key == "map_name"
                               else os.path.normpath(settings[root])), root)
                         for root in ("dst", "dst_root")
                         if isinstance(settings.get(root), str) and settings[root]]
                for root, root_key in sorted(roots, key=lambda row: len(row[0].parts),
                                             reverse=True):
                    try:
                        relative = path.relative_to(root)
                    except ValueError:
                        continue
                    jobs[plate_id][key] = str(Path(jobs[plate_id][root_key]) / relative)
                    break
                else:
                    jobs[plate_id][key] = str(path.parent / plate_id / path.name)
                if key == "map_name":
                    jobs[plate_id][key] = os.path.abspath(jobs[plate_id][key])
    return jobs


def _snakefile(module: str, image: str, run_name: str) -> str:
    """Return the Snakefile text that runs ``module`` once per settings file."""
    return f"""# spaCR workflow exported from run {run_name}.
#
# One job per plate: every settings/<plate>.json is one `spacr-run {module}`.
# Add a plate by copying a settings file and changing its "src".
#
#   snakemake --cores 4                          # this machine
#   snakemake --cores 4 --use-apptainer          # inside {image}
#   snakemake --executor slurm --jobs 20 --use-apptainer   # a cluster
#
# Inside a container the data folders must be visible at the same paths,
# e.g. --apptainer-args "--bind /data". config.yaml holds the image and the
# spacr-run command.

configfile: "config.yaml"

PLATES = glob_wildcards("settings/{{plate}}.json").plate


rule all:
    input:
        expand("done/{{plate}}.ok", plate=PLATES)


rule spacr_{module}:
    input:
        "settings/{{plate}}.json"
    output:
        touch("done/{{plate}}.ok")
    log:
        "logs/{{plate}}.log"
    container:
        config["image"]
    params:
        spacr_run=config["spacr_run"]
    threads: config.get("threads", 1)
    shell:
        "{{params.spacr_run}} {module} --settings {{input}} > {{log}} 2>&1"
"""


def _nextflow_main(module: str, run_name: str) -> str:
    """Return the Nextflow DSL2 ``main.nf`` that runs ``module`` per plate."""
    return f"""#!/usr/bin/env nextflow
// spaCR workflow exported from run {run_name}.
//
// One task per plate: every settings/<plate>.json is one `spacr-run {module}`.
// Add a plate by copying a settings file and changing its "src".
//
//   nextflow run main.nf                              // this machine
//   nextflow run main.nf -profile apptainer           // inside the image
//   nextflow run main.nf -profile slurm,apptainer     // a cluster
//
// Inside a container the data folders must be visible at the same paths:
// add e.g. -v /data:/data to docker.runOptions in nextflow.config.

nextflow.enable.dsl = 2

process SPACR_{module.upper()} {{
    tag "${{plate}}"
    publishDir "${{params.outdir}}/logs", mode: 'copy', pattern: '*.log'

    input:
    tuple val(plate), path(settings)

    output:
    tuple val(plate), path("${{plate}}.log")

    script:
    \"\"\"
    ${{params.spacr_run}} {module} --settings ${{settings}} > ${{plate}}.log 2>&1
    \"\"\"
}}

workflow {{
    channel
        .fromPath("${{params.settings_dir}}/*.json")
        .map {{ f -> tuple(f.baseName, f) }}
        | SPACR_{module.upper()}
}}
"""


def _nextflow_config(image: str) -> str:
    """Return ``nextflow.config`` with local, container and SLURM profiles."""
    return f"""params {{
    settings_dir = "${{projectDir}}/settings"
    outdir       = "${{projectDir}}/results"
    spacr_run    = "spacr-run"
    image        = "{image}"
}}

process {{
    cpus      = 1
    container = params.image
}}

profiles {{
    docker {{
        docker.enabled    = true
        docker.runOptions = '-u $(id -u):$(id -g)'
    }}
    apptainer {{
        apptainer.enabled    = true
        apptainer.autoMounts = true
    }}
    slurm {{
        process.executor = 'slurm'
    }}
}}
"""


def _export_workflow(run_dir: Any, out_dir: Any, engine: str = "snakemake",
                     plates: Optional[List[str]] = None,
                     image: Optional[str] = None,
                     spacr_run: str = "spacr-run") -> Path:
    """Write a recorded run as a Snakemake or Nextflow workflow.

    The workflow runs the run's module once per plate with ``spacr-run`` and
    the recorded settings, one ``settings/<plate>.json`` each. Multiple jobs
    receive unique subfolders of explicit ``dst`` or ``dst_root`` folders.
    Convert/Align databases, fresh Convert checkpoints and Convert map paths
    escaping their destination are separated too. Input paths, other settings
    and single-job destinations are preserved.

    :param run_dir: run-journal folder (or its name under the runs root).
    :param out_dir: folder to write the workflow into; created if missing.
    :param engine: ``"snakemake"`` or ``"nextflow"``.
    :param plates: plate folders to run instead of the recorded ``src``.
    :param image: container image; defaults to the CPU image of the
        spaCR version that ran.
    :param spacr_run: command that starts ``spacr-run`` on the nodes.
    :returns: the workflow's main file (``Snakefile`` or ``main.nf``).
    :raises ValueError: for an unknown engine, a folder that is not a run,
        a module that cannot run headless, or an ambiguous multi-plate
        explicit Convert resume checkpoint.
    """
    from .cli import resolve_module

    if engine not in _WORKFLOW_ENGINES:
        raise ValueError(f"engine must be one of {_WORKFLOW_ENGINES}, "
                         f"not {engine!r}")
    run_dir = Path(run_dir)
    if not run_dir.exists() and (runs_root() / run_dir.name).exists():
        run_dir = runs_root() / run_dir.name
    manifest_path = run_dir / "manifest.json"
    if not manifest_path.is_file():
        raise ValueError(f"{run_dir} is not a run folder (no manifest.json)")
    manifest = json.loads(manifest_path.read_text(encoding="utf-8"))
    module = resolve_module(str(manifest.get("app_key") or ""))
    if module is None:
        raise ValueError(f"module {manifest.get('app_key')!r} cannot run "
                         f"headless, so it cannot be exported")
    settings = load_run_settings(run_dir)
    image = image or _workflow_image(manifest)
    jobs = _workflow_plates(settings, plates, module=module.key)
    out = Path(out_dir)
    (out / "settings").mkdir(parents=True, exist_ok=True)
    for plate_id, plate_settings in jobs.items():
        (out / "settings" / f"{plate_id}.json").write_text(
            json.dumps(plate_settings, indent=2, default=str), encoding="utf-8")
    (out / "run_manifest.json").write_text(
        json.dumps(manifest, indent=2, default=str), encoding="utf-8")
    if engine == "snakemake":
        (out / "config.yaml").write_text(
            f"image: {json.dumps('docker://' + image)}\n"
            f"spacr_run: {json.dumps(spacr_run)}\nthreads: 1\n",
            encoding="utf-8")
        main = out / "Snakefile"
        main.write_text(_snakefile(module.key, image, run_dir.name),
                        encoding="utf-8")
    else:
        config = _nextflow_config(image)
        if spacr_run != "spacr-run":
            config = config.replace('spacr_run    = "spacr-run"',
                                    f"spacr_run    = {json.dumps(spacr_run)}")
        (out / "nextflow.config").write_text(config, encoding="utf-8")
        main = out / "main.nf"
        main.write_text(_nextflow_main(module.key, run_dir.name),
                        encoding="utf-8")
    return main


[docs] def main(argv=None) -> int: """CLI entry point wired as the ``spacr-repro`` console script. :param argv: optional argv list; defaults to ``sys.argv[1:]``. :returns: process exit code. """ p = argparse.ArgumentParser( prog="spacr repro", description="Replay a spaCR pipeline run from its recorded " "journal folder.", ) p.add_argument("run_dir", help="Path to a folder under ~/.spacr/runs/, or the " "folder's basename.") p.add_argument("--dry", action="store_true", help="Print resolved settings + app; don't run.") p.add_argument("--show", action="store_true", help="Print manifest + settings; don't run.") p.add_argument("--export", choices=_WORKFLOW_ENGINES, help="Write the run as a Snakemake or Nextflow workflow " "into --out instead of running it.") p.add_argument("--out", metavar="DIR", help="Folder for --export.") p.add_argument("--plates", nargs="+", metavar="FOLDER", help="With --export: plate folders to run instead of " "the recorded src.") p.add_argument("--image", metavar="IMAGE", help="With --export: container image; default is the " "CPU image of the spaCR version that ran.") args = p.parse_args(argv) run_dir = Path(args.run_dir) if not run_dir.exists(): candidate = runs_root() / args.run_dir if candidate.exists(): run_dir = candidate else: print(f"error: no such run folder: {args.run_dir}", file=sys.stderr) return 2 manifest_path = run_dir / "manifest.json" if not manifest_path.exists(): print(f"error: {run_dir} is not a valid run folder " f"(no manifest.json)", file=sys.stderr) return 2 if args.export: if not args.out: print("error: --export needs --out DIR", file=sys.stderr) return 2 try: main_file = _export_workflow(run_dir, args.out, args.export, plates=args.plates, image=args.image) except ValueError as e: print(f"error: {e}", file=sys.stderr) return 2 print(f"wrote {main_file}") return 0 manifest = json.loads(manifest_path.read_text()) app_key = manifest.get("app_key") settings = load_run_settings(run_dir) if args.show: _print_manifest(run_dir) print() _print_settings(settings) return 0 entry = _resolve_pipeline(app_key) if entry is None: print(f"error: no pipeline entry for app_key={app_key!r}", file=sys.stderr) return 2 if args.dry: print(f"would run: {entry.__module__}.{entry.__name__}(settings)") _print_settings(settings) return 0 from .figure_font import _open_sans_is_the_default print(f"replaying {app_key} — this opens a NEW run journal folder.") with _open_sans_is_the_default(), open_run(app_key, settings) as run: try: entry(settings) run.set_status("success") except Exception as e: run.set_status("failed") print(f"replay raised: {type(e).__name__}: {e}", file=sys.stderr) return 1 print(f"done → {run.dir}") return 0
if __name__ == "__main__": raise SystemExit(main())