← Files NGS Analysis WorkbenchARCHIVED FILE
tests/test_nextflow_mcp.py
44.6 KB · Sep 30, 2026 · 23:20 UTC
from __future__ import annotations
import hashlib
import json
import os
import subprocess
import sys
import tempfile
import time
import unittest
from collections.abc import Iterator
from contextlib import contextmanager
from pathlib import Path
from typing import Any
from unittest import mock
PLUGIN_ROOT = Path(__file__).resolve().parents[1]
MCP_ROOT = PLUGIN_ROOT / "mcp"
if str(MCP_ROOT) not in sys.path:
sys.path.insert(0, str(MCP_ROOT))
from daemon_test_case import DaemonTestCase # noqa: E402
from ngs_workbench_daemon import runs as daemon_runs # noqa: E402
from ngs_workbench_daemon import ssh # noqa: E402
from ngs_workbench_mcp import app, approved_execution, plan_registry, runs # noqa: E402
from ngs_workbench_mcp.plans import ( # noqa: E402
NextflowRunPlan,
PreparationOperation,
PreparationSpec,
plan_checksum,
)
from ngs_workbench_mcp.workflows import catalog_store, defaults, nextflow # noqa: E402
from ngs_workbench_mcp.workflows.source import WorkflowSource # noqa: E402
_plan_catalog_workflow = nextflow.plan_remote_run
def _plan_with_catalog_source(*args: Any, **kwargs: Any) -> dict[str, Any]:
kwargs.setdefault("workflow", "nf-core/demo")
return _plan_catalog_workflow(*args, **kwargs)
nextflow.plan_remote_run = _plan_with_catalog_source
def _plan_checksum(payload: dict[str, object]) -> str:
return plan_checksum(NextflowRunPlan.model_validate(payload))
@contextmanager
def target_input_transport() -> Iterator[None]:
run = subprocess.run
def execute(_remote: dict[str, Any], source: str, _operation: str, timeout: int = 30) -> Any:
result = run(
[sys.executable, "-"], input=source, capture_output=True, text=True, timeout=timeout
)
if result.returncode:
raise ValueError(result.stderr)
return json.loads(result.stdout)
with mock.patch.object(ssh, "_run_program", side_effect=execute):
yield
def ready(
nextflow: str = "/runtime/nextflow",
*,
observation: int | None = None,
docker_platform: tuple[str, str] = ("linux", "amd64"),
) -> dict[str, object]:
result: dict[str, object] = {
"ok": True,
"status": "ready",
"scope": "executable_presence_only",
"commands": [
{"name": "nextflow", "path": nextflow},
{"name": "docker", "path": "/runtime/docker"},
],
"blockers": [],
"docker": {
"path": "/runtime/docker",
"daemon_reachable": True,
"server_os": docker_platform[0],
"server_arch": docker_platform[1],
},
}
if observation is not None:
identity = f"{observation:032x}"
snapshot_id = f"runtime-{identity}"
result.update(
{
"readiness_id": f"readiness-{identity}",
"snapshot_id": snapshot_id,
"observed_at": f"2026-08-06T00:0{observation}:00Z",
"expires_at": f"2026-08-06T00:0{observation + 5}:00Z",
"evidence": [
{
"kind": "runtime_snapshot",
"source": snapshot_id,
"detail": "unchanged local runtime facts",
}
],
}
)
return result
class CuratedNextflowPlanTests(unittest.TestCase):
def test_nextflow_planner_discovers_samplesheet_from_native_parameters_file(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
sample_sheet = workspace / "samples.csv"
sample_sheet.write_text(
"sample,fastq_1,fastq_2\nSAMPLE,/data/R1.fastq.gz,/data/R2.fastq.gz\n",
encoding="utf-8",
)
params = workspace / "params.json"
params.write_text(
json.dumps({"input": str(sample_sheet), "genome": "GRCh38"}),
encoding="utf-8",
)
with (
mock.patch.dict(
os.environ,
{"NGS_ANALYSIS_WORKBENCH_STATE_DIR": str(workspace / "state")},
),
mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=ready()),
):
defaults.bootstrap_default_workflows()
planned = app.plan_nextflow(
"rnaseq",
str(workspace / "run"),
params_file=str(params),
profile="docker",
)
self.assertTrue(planned["runnable"], planned)
self.assertEqual(planned["request"]["sample_sheet"], str(sample_sheet.resolve()))
self.assertEqual(planned["request"]["params_file"], str(params.resolve()))
self.assertIn(str(sample_sheet.resolve()), planned["command_argv"])
with mock.patch.object(
approved_execution.daemon_client, "execute", return_value={"ok": True}
) as dispatch:
result = approved_execution.execute_registered_plan(
planned["plan_name"], planned["plan_id"], planned["plan_checksum"]
)
self.assertTrue(result["ok"])
self.assertEqual(dispatch.call_args.args[0].binding, "nextflow")
def test_nextflow_parameters_remain_workflow_owned_without_trim_overrides(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
params = workspace / "params.json"
params.write_text(json.dumps({"skip_trim": False}), encoding="utf-8")
with (
mock.patch.dict(
os.environ,
{"NGS_ANALYSIS_WORKBENCH_STATE_DIR": str(workspace / "state")},
),
mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=ready()),
):
defaults.bootstrap_default_workflows()
planned = app.plan_nextflow(
"fastq_qc",
str(workspace / "run"),
params_file=str(params),
profile="test,docker",
)
preserved_parameters = json.loads(params.read_text(encoding="utf-8"))
self.assertNotIn("trim", planned["request"])
self.assertNotIn("--skip_trim", planned["command_argv"])
self.assertEqual(preserved_parameters, {"skip_trim": False})
def test_nextflow_discovers_approved_generated_parameters_without_writing_them(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
directory = workspace / "prepared"
params = directory / "params.json"
sample_sheet = directory / "samples.csv"
preparation = {
"destination_dir": str(directory),
"generated_files": [
{
"relative_path": params.name,
"content": json.dumps({"input": str(sample_sheet)}),
},
{
"relative_path": sample_sheet.name,
"content": "sample,fastq_1,fastq_2\nSAMPLE,/data/R1.fastq.gz,\n",
},
],
}
with (
mock.patch.dict(
os.environ,
{"NGS_ANALYSIS_WORKBENCH_STATE_DIR": str(workspace / "state")},
),
mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=ready()),
):
defaults.bootstrap_default_workflows()
planned = app.plan_nextflow(
"rnaseq",
str(workspace / "run"),
params_file=str(params),
profile="docker",
preparation=preparation,
)
self.assertTrue(planned["runnable"], planned)
self.assertTrue(planned["request"]["params_file_sha256"].startswith("sha256:"))
self.assertTrue(planned["request"]["sample_sheet_sha256"].startswith("sha256:"))
self.assertFalse(params.exists())
self.assertFalse(sample_sheet.exists())
def test_curated_workflow_plans_and_executes_under_its_nextflow_engine(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
with mock.patch.dict(
os.environ,
{"NGS_ANALYSIS_WORKBENCH_STATE_DIR": str(Path(temporary) / "state")},
):
defaults.bootstrap_default_workflows()
with mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=ready()):
registered = app.plan_nextflow(
"fastq_qc",
str(Path(temporary) / "run"),
profile="test,docker",
)
self.assertEqual(registered["binding"], "nextflow")
self.assertEqual(registered["readiness"]["binding"], "nextflow")
self.assertEqual(registered["request"]["revision"], "1.2.0")
with (
mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=ready()),
mock.patch.object(
approved_execution.daemon_client,
"execute",
return_value={"ok": True},
) as dispatch,
):
result = approved_execution.execute_registered_plan(
registered["plan_name"],
registered["plan_id"],
registered["plan_checksum"],
)
self.assertTrue(result["ok"])
self.assertEqual(result["binding"], "nextflow")
self.assertEqual(dispatch.call_args.args[0].binding, "nextflow")
def test_plan_binds_the_selected_managed_environment_invocation(self) -> None:
candidate_id = f"controller-{'7' * 32}"
managed_readiness = {
**ready(),
"selected_controller": {
"candidate_id": candidate_id,
"controller": "nextflow",
"source": "managed_environment",
"manager": "conda",
"executable_path": "/env/bin/nextflow",
"environment_path": "/env",
"platform_matches_host": True,
"launch_argv_prefix": [
"/runtime/conda",
"run",
"--prefix",
"/env",
"nextflow",
],
"recommended": True,
"recommendation_reason": "managed environment",
},
}
with tempfile.TemporaryDirectory() as temporary:
with mock.patch.object(
nextflow,
"assess_nfcore_readiness",
return_value=managed_readiness,
):
plan = nextflow.plan_remote_run(
"fastq_qc",
str(Path(temporary) / "run"),
"test,docker",
controller_candidate_id=candidate_id,
run_id="nfcore-fastq-qc-managed",
)
self.assertTrue(plan["runnable"])
self.assertEqual(plan["request"]["controller_candidate_id"], candidate_id)
self.assertEqual(
plan["command_argv"][:6],
[
"/runtime/conda",
"run",
"--prefix",
"/env",
"nextflow",
"run",
],
)
def test_plan_preserves_unknown_readiness_and_is_not_runnable(self) -> None:
unknown = {
**ready(),
"ok": False,
"status": "unknown",
"unknowns": ["effective task software environment is unresolved"],
}
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
with mock.patch.object(
nextflow,
"assess_nfcore_readiness",
return_value=unknown,
):
result = nextflow.plan_remote_run(
"fastq_qc",
str(workspace / "run"),
"mylab",
run_id="nfcore-fastq-qc-unknown",
)
wrote_runs = (workspace / "run").exists()
self.assertTrue(result["ok"])
self.assertFalse(result["runnable"])
self.assertEqual(result["readiness"]["status"], "unknown")
self.assertTrue(any("readiness unresolved" in item for item in result["blockers"]))
self.assertFalse(wrote_runs)
def test_plan_rejects_an_invalid_explicit_run_id(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
with mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=ready()):
result = nextflow.plan_remote_run(
"fastq_qc",
str(workspace / "run"),
"test,docker",
run_id="../../escape",
)
self.assertFalse(result["ok"])
self.assertFalse((workspace / "run").exists())
def test_plan_is_read_only_and_checksum_is_stable_for_a_run_id(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
with mock.patch.object(
nextflow,
"assess_nfcore_readiness",
return_value=ready(docker_platform=("linux", "arm64")),
):
first = nextflow.plan_remote_run(
"fastq_qc",
str(workspace / "run"),
"test,docker,arm64",
revision="1.2.0",
run_id="nfcore-fastq-qc-fixed",
)
second = nextflow.plan_remote_run(
"fastq_qc",
str(workspace / "run"),
"test,docker,arm64",
revision="1.2.0",
run_id="nfcore-fastq-qc-fixed",
)
self.assertTrue(first["ok"])
self.assertTrue(first["runnable"])
self.assertEqual(_plan_checksum(first), _plan_checksum(second))
self.assertEqual(first["request"]["profile"], "test,docker,arm64")
self.assertFalse((workspace / "run").exists())
def test_plan_binds_agent_name_and_opaque_samplesheet(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
sample_sheet = workspace / "samplesheet.csv"
sample_sheet.write_text(
"sample,fastq_1,fastq_2\nS1,/data/S1_R1.fastq.gz,/data/S1_R2.fastq.gz\n"
"S2,/data/S2_R1.fastq.gz,/data/S2_R2.fastq.gz\n",
encoding="utf-8",
)
with mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=ready()):
plan = nextflow.plan_remote_run(
"rnaseq",
str(workspace / "run"),
"docker",
display_name=" PBMC pilot RNA-seq ",
sample_sheet=str(sample_sheet),
run_id="nfcore-rnaseq-named",
)
self.assertEqual(plan["request"]["display_name"], "PBMC pilot RNA-seq")
self.assertEqual(plan["input_summary"]["source"], "sample_sheet")
self.assertEqual(plan["input_summary"]["path"], str(sample_sheet.resolve()))
self.assertIsNone(plan["input_summary"]["sample_count"])
self.assertIsNone(plan["input_summary"]["record_count"])
self.assertIsNone(plan["input_summary"]["read_layout"])
def test_plan_rejects_an_unregistered_target_without_checking_local_readiness(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
with mock.patch.object(nextflow, "assess_nfcore_readiness") as readiness:
result = nextflow.plan_remote_run(
"fastq_qc",
str(Path(temporary) / "run"),
"test,docker",
target_id="lab-slurm",
run_id="nfcore-fastq-qc-remote",
)
self.assertFalse(result["ok"])
self.assertIn("compute target is not registered", result["errors"][0])
readiness.assert_not_called()
def test_missing_sample_sheet_is_not_reported_as_test_profile_input(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
with mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=ready()):
blocked = nextflow.plan_remote_run(
"rnaseq",
str(Path(temporary) / "run"),
"contest,docker",
run_id="nfcore-rnaseq-missing-input",
)
test_plan = nextflow.plan_remote_run(
"fastq_qc",
str(Path(temporary) / "run"),
"test,docker",
run_id="nfcore-fastq-qc-test-input",
)
self.assertFalse(blocked["runnable"])
self.assertEqual(blocked["input_summary"]["source"], "missing")
self.assertEqual(test_plan["input_summary"]["source"], "workflow_test_profile")
def test_approved_request_rejects_docker_platform_drift_before_writing(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
with mock.patch.object(
nextflow,
"assess_nfcore_readiness",
return_value=ready(docker_platform=("linux", "arm64")),
):
plan = nextflow.plan_remote_run(
"fastq_qc",
str(workspace / "run"),
"test,docker,arm64",
revision="1.2.0",
run_id="nfcore-fastq-qc-drift",
)
with mock.patch.object(
nextflow,
"assess_nfcore_readiness",
return_value=ready(docker_platform=("linux", "amd64")),
):
result = nextflow.approved_nfcore_execution_request(
NextflowRunPlan.model_validate(plan), _plan_checksum(plan)
)
self.assertFalse(result["ok"])
self.assertIn("checksum", result["errors"][0])
self.assertFalse((workspace / "run").exists())
def test_approved_request_accepts_refreshed_readiness_when_facts_are_unchanged(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
with mock.patch.object(
nextflow,
"assess_nfcore_readiness",
side_effect=[ready(observation=1), ready(observation=2)],
):
plan = nextflow.plan_remote_run(
"fastq_qc",
str(workspace / "run"),
"test,docker",
revision="1.2.0",
run_id="nfcore-fastq-qc-refreshed",
)
request = nextflow.approved_nfcore_execution_request(
NextflowRunPlan.model_validate(plan), _plan_checksum(plan)
)
self.assertEqual(request.plan_checksum, _plan_checksum(plan))
self.assertFalse((workspace / "run").exists())
def test_approved_request_rejects_params_file_content_drift_before_writing(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
params_file = workspace / "params.json"
params_file.write_text('{"aligner":"star_salmon"}\n', encoding="utf-8")
with mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=ready()):
plan = nextflow.plan_remote_run(
"rnaseq",
str(workspace / "run"),
"test,docker",
params_file=str(params_file),
revision="3.26.0",
run_id="nfcore-rnaseq-params-drift",
)
params_file.write_text('{"aligner":"hisat2"}\n', encoding="utf-8")
result = nextflow.approved_nfcore_execution_request(
NextflowRunPlan.model_validate(plan), _plan_checksum(plan)
)
self.assertFalse(result["ok"])
self.assertFalse((workspace / "run").exists())
class RemoteNextflowPlanTests(DaemonTestCase):
def setUp(self) -> None:
super().setUp()
defaults.bootstrap_default_workflows()
self.enterContext(target_input_transport())
self.target_files = self.root / "target-files"
self.target_files.mkdir()
def test_slurm_plan_binds_target_configuration_and_exact_controller_argv(self) -> None:
target = self._configure_ssh_target()
readiness = {
**ready(nextflow="/remote/bin/nextflow"),
"host": {"os": "linux", "arch": "arm64"},
}
with mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=readiness):
planned = nextflow.plan_remote_run(
"fastq_qc",
str(self.workspace / "run"),
"test,docker",
run_id="nfcore-fastq-qc-ssh-plan",
target_id="lab",
)
self.assertTrue(planned["runnable"])
remote = planned["request"]["remote"]
remote_root = str(self.workspace / "run")
self.assertEqual(remote["config_hash"], target["config_hash"])
self.assertEqual(remote["host_access"], target["host_access"])
self.assertEqual(remote["executor_configuration"], {"partition": "cpu"})
self.assertEqual(remote["run_dir"], remote_root)
self.assertEqual(planned["effects"]["run_dir"], remote_root)
for field, relative in (
("output_dir", "results"),
("work_dir", "work"),
("launch_log", "logs/nextflow.log"),
):
self.assertEqual(planned["effects"][field], f"{remote_root}/{relative}")
self.assertTrue(
all(
path.startswith(f"{remote_root}/")
for field in ("log_paths", "artifact_paths")
for path in planned["monitoring"][field]
)
)
self.assertIn(f"{remote_root}/results", planned["command_argv"])
self.assertIn(f"{remote_root}/workflow/trace.txt", planned["command_argv"])
self.assertIn(f"{remote_root}/workflow/trace.txt", planned["effects"]["remote_writes"])
self.assertIn(f"{remote_root}/results", planned["effects"]["remote_writes"])
configuration = next(
item
for item in remote["staged_files"]
if item["destination"].endswith("nextflow.config")
)
self.assertIn("executor = 'slurm'", configuration["content"])
self.assertIn('queue = "cpu"', configuration["content"])
self.assertNotIn("memory = '1536 MB'", configuration["content"])
self.assertIn(configuration["destination"], planned["command_argv"])
self.assertFalse((self.workspace / "run").exists())
registered = plan_registry.register_plan("nextflow", planned)
with (
mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=readiness),
mock.patch.object(
approved_execution.daemon_client, "execute", return_value={"ok": True}
) as dispatch,
):
accepted = approved_execution.execute_registered_plan(
str(registered["plan_name"]),
str(registered["plan_id"]),
str(registered["plan_checksum"]),
)
self.assertTrue(accepted["ok"])
dispatch.assert_called_once()
self.assertEqual(accepted["run_dir"], remote_root)
request = dispatch.call_args.args[0]
metadata_root = Path(request.files[-1].destination).parents[1]
self.assertTrue(metadata_root.is_relative_to(self.root / "state/runs"))
self.assertTrue(
all(Path(item.destination).is_relative_to(metadata_root) for item in request.files)
)
self.assertEqual(daemon_runs._validate_request(request)[0], metadata_root)
changed = self._configure_ssh_target(partition="gpu")
self.assertNotEqual(changed["config_hash"], target["config_hash"])
with mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=readiness):
updated = nextflow.plan_remote_run(
"fastq_qc",
str(self.workspace / "run"),
"test,docker",
run_id="nfcore-fastq-qc-ssh-plan",
target_id="lab",
)
self.assertNotEqual(_plan_checksum(updated), _plan_checksum(planned))
def test_remote_inputs_remain_in_place_and_checksum_bound(self) -> None:
self._configure_ssh_target(executor="local_process", partition=None)
sample_sheet = self.target_files / "samples.csv"
sample_content = "sample,read_one,group\nS1,/data/S1_R1.fastq.gz,control\n"
sample_sheet.write_text(sample_content, encoding="utf-8")
params = self.target_files / "params.json"
content = json.dumps(
{
"input": str(sample_sheet),
"reference": "/data/genome",
"nested": {"label": "sample.v2"},
}
)
params.write_text(content, encoding="utf-8")
readiness = {**ready(), "host": {"os": "linux", "arch": "arm64"}}
with mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=readiness):
planned = app.plan_nextflow(
"fastq_qc",
str(self.workspace / "run"),
profile="test,docker",
params_file=str(params),
target_id="lab",
)
self.assertTrue(planned["runnable"], planned)
self.assertIn(str(params), planned["command_argv"])
self.assertIn(str(sample_sheet), planned["command_argv"])
self.assertEqual(planned["request"]["remote"]["staged_files"], [])
self.assertEqual(params.read_text(), content)
request = nextflow.approved_nfcore_execution_request(
NextflowRunPlan.model_validate(planned), planned["plan_checksum"]
)
self.assertEqual(
request.inputs[str(sample_sheet)],
f"sha256:{hashlib.sha256(sample_content.encode()).hexdigest()}",
)
sample_sheet.write_text(sample_content.replace("S1_R1", "changed_R1"), encoding="utf-8")
rejected = nextflow.approved_nfcore_execution_request(
NextflowRunPlan.model_validate(planned), planned["plan_checksum"]
)
self.assertFalse(rejected["ok"])
with self.assertRaises(ValueError):
app.plan_nextflow(
"fastq_qc",
str(self.workspace / "run"),
profile="test,docker",
params_file="params.json",
target_id="lab",
)
future = self.target_files / "downloads" / "samples.csv"
prepared = PreparationSpec(
destination_dir=str(future.parent),
operations=[
PreparationOperation(
operation="download_verified_file",
relative_path=future.name,
path=str(future),
bytes=0,
sha256=f"sha256:{'0' * 64}",
)
],
download_bytes=0,
writes=[str(future)],
)
with mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=readiness):
prepared_plan = nextflow.plan_remote_run(
"fastq_qc",
str(self.workspace / "run"),
"test,docker",
sample_sheet=str(future),
preparation_spec=prepared,
run_id="nfcore-fastq-qc-ssh-download",
target_id="lab",
)
self.assertTrue(prepared_plan["runnable"], prepared_plan)
self.assertIn(str(future), prepared_plan["effects"]["remote_writes"])
self.assertNotIn(str(future), prepared_plan["effects"]["local_writes"])
generated = prepared_plan["request"]["remote"]["staged_files"][0]
self.assertEqual(json.loads(generated["content"])["input"], str(future))
self.assertFalse(future.exists())
self.assertFalse((self.workspace / "run").exists())
def test_unknown_profiles_remain_workflow_owned_without_generated_configuration(self) -> None:
self._configure_ssh_target(executor="local_process", partition=None)
requirements = nextflow.resolve_runtime_requirements("fastq_qc", "test,host")
self.assertEqual({item.id for item in requirements.requirements}, {"nextflow"})
self.assertTrue(requirements.unknowns)
readiness = {
**ready(nextflow=sys.executable),
"host": {"os": "linux", "arch": "arm64"},
}
with mock.patch.object(nextflow, "assess_nfcore_readiness", return_value=readiness):
planned = nextflow.plan_remote_run(
"fastq_qc",
str(self.workspace / "run"),
"test,docker,HOST",
run_id="nfcore-fastq-qc-ssh-native-profile",
target_id="lab",
)
remote = planned["request"]["remote"]
self.assertFalse(
any(item["destination"].endswith("/nextflow.config") for item in remote["staged_files"])
)
self.assertEqual(planned["command_argv"][:2], [sys.executable, "run"])
self.assertIn("test,docker,HOST", planned["command_argv"])
self.assertTrue(any("container" in item for item in planned["effects"]["downloads"]))
self.assertTrue(planned["runnable"])
with mock.patch.object(
nextflow,
"assess_nfcore_readiness",
return_value={**readiness, "host": {"os": "darwin", "arch": "arm64"}},
):
unsupported = nextflow.plan_remote_run(
"fastq_qc", str(self.workspace / "run"), "test,docker", target_id="lab"
)
self.assertFalse(unsupported["runnable"])
class GenericNextflowTests(DaemonTestCase):
def setUp(self) -> None:
super().setUp()
self.controller = self.root / "nextflow"
self.controller.write_text(
f"#!{sys.executable}\n"
"import json, pathlib, sys\n"
"assert sys.argv[1] == 'run'\n"
"assert pathlib.Path(sys.argv[2]).is_file()\n"
"trace = pathlib.Path(sys.argv[sys.argv.index('-with-trace') + 1])\n"
"trace.write_text('task_id\\tname\\tstatus\\n1\\tREPORT\\tCOMPLETED\\n')\n"
"params_path = pathlib.Path(sys.argv[sys.argv.index('-params-file') + 1]) "
"if '-params-file' in sys.argv else None\n"
"params = json.loads(params_path.read_text()) if params_path else {}\n"
"output = pathlib.Path(sys.argv[sys.argv.index('-output-dir') + 1])\n"
"result = output / 'output' / 'message.txt'\n"
"result.parent.mkdir(parents=True, exist_ok=True)\n"
"result.write_text(params.get('message', 'default'))\n",
encoding="utf-8",
)
self.controller.chmod(0o755)
root = self.workspace / "source"
root.mkdir()
(root / "main.nf").write_text("workflow { log.info 'example' }\n", encoding="utf-8")
self.source = WorkflowSource(engine="nextflow", root=str(root), entrypoint="main.nf")
self._save_source("custom_nextflow")
self.readiness = {
"ok": True,
"status": "ready",
"scope": "executable_presence_only",
"commands": [{"name": "nextflow", "path": str(self.controller)}],
"blockers": [],
}
def _save_source(self, workflow_id: str, source: WorkflowSource | None = None) -> None:
selected = source or self.source
assert selected.root is not None
catalog_store.save_workflow(
workflow_id,
workflow_id,
"nextflow",
catalog_store.LocalWorkflowSource(
kind="local", root=selected.root, entrypoint=selected.entrypoint
),
)
def _execute(self, plan: dict[str, object]) -> dict[str, object]:
started = approved_execution.execute_registered_plan(
str(plan["plan_name"]), str(plan["plan_id"]), str(plan["plan_checksum"])
)
self.assertTrue(started["ok"], started)
deadline = time.monotonic() + 10
while time.monotonic() < deadline:
observation = app.observe_ngs_run(str(started["registry_run_id"]))
if observation["status"] in {"completed", "failed"}:
self.assertEqual(observation["status"], "completed", observation)
return app.get_ngs_run(str(started["registry_run_id"]))
time.sleep(0.025)
self.fail("approved Nextflow workflow did not finish")
def test_generic_remote_parameters_remain_opaque(self) -> None:
catalog_store.save_workflow(
"community_remote",
"Community remote",
"nextflow",
catalog_store.RemoteWorkflowSource(
kind="remote", workflow="community/example", revision="1.0.0"
),
)
params = self.workspace / "remote-params.json"
params.write_text(json.dumps({"input": ["a", "b"]}), encoding="utf-8")
with mock.patch.object(nextflow, "assess_readiness", return_value=self.readiness):
planned = app.plan_nextflow(
"community_remote", str(self.workspace / "remote-run"), params_file=str(params)
)
self.assertTrue(planned["ok"], planned)
self.assertEqual(planned["request"]["parameter_contract"], "generic")
self.assertEqual(planned["request"]["workflow"], "community/example")
self.assertEqual(planned["request"]["revision"], "1.0.0")
self.assertIn("-output-dir", planned["command_argv"])
self.assertNotIn("--outdir", planned["command_argv"])
def test_generic_source_executes(self) -> None:
(self.workspace / "first").write_text("incidental filename collision", encoding="utf-8")
params = self.workspace / "params.json"
params.write_text(json.dumps({"message": "first"}), encoding="utf-8")
with mock.patch.object(nextflow, "assess_readiness", return_value=self.readiness):
original = app.plan_nextflow(
"custom_nextflow",
str(self.workspace / "run"),
params_file=str(params),
)
self.assertEqual(original["binding"], "nextflow")
self.assertEqual(original["request"]["params_file"], str(params))
self.assertFalse((self.workspace / "run").exists())
self.assertFalse(
{"-profile", "--input", "--outdir", "--skip_trim"}.intersection(
original["command_argv"]
)
)
output_index = original["command_argv"].index("-output-dir")
self.assertEqual(
original["command_argv"][output_index + 1], original["effects"]["output_dir"]
)
self.assertEqual(
original["effects"]["output_dir"],
str(Path(original["effects"]["run_dir"]) / "results"),
)
self.assertIn(original["effects"]["output_dir"], original["effects"]["local_writes"])
self.assertIn(
original["effects"]["output_dir"], original["monitoring"]["artifact_paths"]
)
first = self._execute(original)
observation = app.observe_ngs_run(str(first["registry_run_id"]))
report = runs.get_registry_run_report(str(first["registry_run_id"]))["report"]
self.assertTrue(observation["execution"]["evidence"]["available"])
self.assertEqual(
(Path(first["run_dir"]) / "results" / "output" / "message.txt").read_text(),
"first",
)
self.assertIn(
"results/output/message.txt",
{entry["path"] for entry in report["entries"]},
)
repeated = approved_execution.execute_registered_plan(
original["plan_name"], original["plan_id"], original["plan_checksum"]
)
self.assertEqual(repeated["run_id"], first["run_id"])
def test_saved_snapshot_is_separate_from_its_execution_directory(self) -> None:
(self.workspace / "main.nf").write_text("workflow {}\n", encoding="utf-8")
source = WorkflowSource(engine="nextflow", root=str(self.workspace), entrypoint="main.nf")
self._save_source("overlapping", source)
planned = app.plan_nextflow("overlapping", str(self.workspace / "run"))
self.assertTrue(planned["ok"], planned)
self.assertNotEqual(planned["request"]["workflow_source"]["root"], str(self.workspace))
self.assertFalse((self.workspace / "run").exists())
def test_long_workflow_identifier_and_opaque_parameters_remain_executable(self) -> None:
params = self.workspace / "params.json"
workflow_parameters = {
"message": "workflow-owned",
"custom_options": {"strategy": "paired"},
}
params.write_text(json.dumps(workflow_parameters), encoding="utf-8")
workflow_id = "a" * 96
self._save_source(workflow_id)
with mock.patch.object(nextflow, "assess_readiness", return_value=self.readiness):
planned = app.plan_nextflow(
workflow_id,
str(self.workspace / "run"),
params_file=str(params),
)
self.assertTrue(planned["runnable"], planned)
self.assertLessEqual(len(planned["request"]["run_id"]), 96)
self.assertEqual(planned["request"]["params_file"], str(params))
result = self._execute(planned)
self.assertEqual(result["status"], "completed")
self.assertEqual(
json.loads((Path(result["run_dir"]) / "config" / params.name).read_text()),
workflow_parameters,
)
self.assertEqual(
(Path(result["run_dir"]) / "results" / "output" / "message.txt").read_text(),
"workflow-owned",
)
def test_parameters_changing_after_approval_cannot_start_a_run(self) -> None:
params = self.workspace / "params.json"
params.write_text('{"message":"original"}', encoding="utf-8")
with mock.patch.object(nextflow, "assess_readiness", return_value=self.readiness):
plan = app.plan_nextflow(
"custom_nextflow",
str(self.workspace / "run"),
params_file=str(params),
)
params.write_text('{"message":"changed"}', encoding="utf-8")
result = approved_execution.execute_registered_plan(
plan["plan_name"], plan["plan_id"], plan["plan_checksum"]
)
self.assertFalse(result["ok"])
self.assertFalse((self.workspace / "run").exists())
def test_readiness_for_another_engine_cannot_authorize_a_nextflow_run(self) -> None:
readiness = {**self.readiness, "binding": "snakemake"}
with mock.patch.object(nextflow, "assess_readiness", return_value=readiness):
planned = app.plan_nextflow("custom_nextflow", str(self.workspace / "run"))
self.assertFalse(planned["ok"])
self.assertFalse((self.workspace / "run").exists())
def test_remote_nextflow_deploys_source_and_uses_native_parameters(self) -> None:
self._configure_ssh_target(executor="local_process", partition=None)
params = self.root / "remote-params.yaml"
content = "sample: sample.v2\ninput: /data/samples.csv\nreference: /data/genome\n"
params.write_text(content, encoding="utf-8")
readiness = {**self.readiness, "host": {"os": "linux", "arch": "arm64"}}
with (
target_input_transport(),
mock.patch.object(nextflow, "assess_readiness", return_value=readiness),
):
self._save_source("remote_nextflow")
planned = app.plan_nextflow(
"remote_nextflow",
str(self.workspace / "run"),
params_file=str(params),
target_id="lab",
)
request = nextflow.approved_execution_request(
NextflowRunPlan.model_validate(planned), planned["plan_checksum"]
)
self.assertTrue(planned["runnable"], planned)
remote = planned["request"]["remote"]
self.assertEqual(
[item["destination"] for item in remote["staged_files"]],
[f"{remote['run_dir']}/workflow/main.nf"],
)
self.assertIn(str(params), planned["command_argv"])
self.assertIn(f"{remote['run_dir']}/workflow/main.nf", planned["command_argv"])
self.assertEqual(
request.inputs, {str(params): f"sha256:{hashlib.sha256(content.encode()).hexdigest()}"}
)
self.assertEqual(params.read_text(), content)
self.assertFalse((self.workspace / "run").exists())
def test_generic_nextflow_does_not_claim_unimplemented_slurm_execution(self) -> None:
self._configure_ssh_target()
with (
target_input_transport(),
mock.patch.object(
nextflow.readiness_service,
"assess",
side_effect=lambda requirements, *_args, **_kwargs: {
**self.readiness,
"ok": not requirements.blockers,
"status": "blocked" if requirements.blockers else "ready",
"blockers": requirements.blockers,
},
),
):
readiness = app.check_nextflow_readiness(
"custom_nextflow",
str(self.workspace / "run"),
target_id="lab",
)
planned = app.plan_nextflow("custom_nextflow", str(self.workspace / "run"), target_id="lab")
self.assertEqual(readiness.status, "blocked")
self.assertFalse(planned["ok"])
self.assertFalse((self.workspace / "run").exists())
def test_remote_parameters_can_be_prepared_after_approval(self) -> None:
self._configure_ssh_target(executor="local_process", partition=None)
destination = self.root / "target-files" / "params.json"
preparation = {
"destination_dir": str(destination.parent),
"downloads": [
{
"relative_path": destination.name,
"url": "https://raw.githubusercontent.com/example/project/main/params.json",
"bytes": 2,
"sha256": "0" * 64,
}
],
}
readiness = {**self.readiness, "host": {"os": "linux", "arch": "arm64"}}
with (
target_input_transport(),
mock.patch.object(nextflow, "assess_readiness", return_value=readiness),
):
planned = app.plan_nextflow(
"custom_nextflow",
str(self.workspace / "run"),
target_id="lab",
params_file=str(destination),
preparation=preparation,
)
self.assertTrue(planned["runnable"], planned)
self.assertIn(str(destination), planned["effects"]["remote_writes"])
self.assertNotIn(str(destination), planned["effects"]["local_writes"])
self.assertFalse(destination.exists())
self.assertFalse((self.workspace / "run").exists())
def test_preparation_downloads_are_visible_before_execution_approval(self) -> None:
download = "https://raw.githubusercontent.com/example/project/main/reads.fastq"
preparation = {
"destination_dir": str(self.workspace / "prepared"),
"downloads": [
{
"relative_path": "reads.fastq",
"url": download,
"bytes": 4,
"sha256": "0" * 64,
}
],
}
with mock.patch.object(nextflow, "assess_readiness", return_value=self.readiness):
planned = app.plan_nextflow(
"custom_nextflow",
str(self.workspace / "run"),
preparation=preparation,
)
self.assertTrue(any(download in item for item in planned["effects"]["downloads"]))
self.assertTrue(any(download in item for item in planned["effects"]["network_access"]))
self.assertFalse((self.workspace / "prepared").exists())
if __name__ == "__main__":
unittest.main()
SHA-256: d3cc927f3aea00047a48bcc8a16aba9126557a22a32c032e67b697ebe124c11e