← Files NGS Analysis WorkbenchARCHIVED FILE
tests/test_snakemake_mcp.py
40.1 KB · Sep 30, 2026 · 23:20 UTC
from __future__ import annotations
import hashlib
import json
import os
import re
import shlex
import sys
import tempfile
import time
import unittest
from pathlib import Path
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 client as daemon_client # noqa: E402
from ngs_workbench_mcp import ( # noqa: E402
app,
approved_execution,
plan_registry,
preparation,
runs,
)
from ngs_workbench_mcp.plans import SnakemakeRunPlan, plan_checksum # noqa: E402
from ngs_workbench_mcp.workflows import catalog_store, defaults, snakemake # noqa: E402
from ngs_workbench_mcp.workflows.source import WorkflowSource # noqa: E402
from test_nextflow_mcp import target_input_transport # noqa: E402
READINESS = {
"ok": True,
"status": "ready",
"scope": "executable_presence_only",
"commands": [{"name": "snakemake", "path": "/runtime/snakemake"}],
"blockers": [],
}
def _plan_checksum(payload: dict[str, object]) -> str:
return plan_checksum(SnakemakeRunPlan.model_validate(payload))
def _observed_readiness(observation: int) -> dict[str, object]:
identity = f"{observation:032x}"
snapshot_id = f"runtime-{identity}"
return {
**READINESS,
"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",
}
],
}
def _config(workspace: Path, name: str = "config.json") -> Path:
path = workspace / name
path.write_text(
json.dumps({"samples": {"sample": {"r1": "/data/r1.fastq.gz"}}}),
encoding="utf-8",
)
return path
def _argv(run_dir: Path, config_name: str, cores: int) -> list[str]:
return [
"/runtime/snakemake",
"--snakefile",
str(run_dir / "workflow" / "Snakefile"),
"--configfile",
str(run_dir / "config" / config_name),
"--directory",
str(run_dir / "results"),
"--cores",
str(cores),
"--printshellcmds",
]
def _external_workflow(workspace: Path) -> WorkflowSource:
root = workspace / "reusable-workflow"
root.mkdir()
(root / "Snakefile").write_text("rule all:\n input: []\n", encoding="utf-8")
return WorkflowSource(engine="snakemake", root=str(root), entrypoint="Snakefile")
def _bundled_workflow() -> WorkflowSource:
return WorkflowSource(
engine="snakemake",
root=str(PLUGIN_ROOT / "workflows" / "fastq_qc"),
entrypoint="workflow/Snakefile",
)
class SnakemakeBindingTests(unittest.TestCase):
def test_bundled_workflow_uses_its_packaged_default_configuration(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
with (
mock.patch.dict(
os.environ,
{"NGS_ANALYSIS_WORKBENCH_STATE_DIR": str(workspace / "state")},
),
mock.patch.object(
snakemake, "assess_readiness", return_value=READINESS
) as assess_readiness,
):
defaults.bootstrap_default_workflows()
planned = app.plan_snakemake("oai_fastq_qc", str(workspace / "run"))
default_config = PLUGIN_ROOT / "workflows" / "fastq_qc" / "config" / "config.json"
self.assertEqual(assess_readiness.call_args_list[0].args[1], default_config)
self.assertEqual(planned["request"]["config_file"], str(default_config))
self.assertEqual(planned["request"]["display_name"], "Bundled FASTQ QC")
def test_bundled_multiqc_commands_support_older_versions_without_network_checks(self) -> None:
workflow_files = [
PLUGIN_ROOT / "workflows" / "fastq_qc" / "workflow" / "Snakefile",
PLUGIN_ROOT / "workflows" / "fastq_qc" / "workflow" / "rules" / "trimmed_qc.smk",
PLUGIN_ROOT / "workflows" / "bulk_rnaseq_counts_qc" / "workflow" / "Snakefile",
]
multiqc_commands = [
shlex.split(arguments)
for path in workflow_files
for arguments in re.findall(
r"\{MULTIQC:q\}([^\"\n]+)", path.read_text(encoding="utf-8")
)
]
self.assertEqual(len(multiqc_commands), 4)
for arguments in multiqc_commands:
with self.subTest(arguments=arguments):
self.assertNotIn("--no-version-check", arguments)
self.assertIn("--no-megaqc-upload", arguments)
self.assertEqual(
arguments[arguments.index("--cl-config") + 1],
"no_version_check: true",
)
def test_globally_unique_workflow_identifier_selects_snakemake(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
config = _config(workspace)
with (
mock.patch.dict(
os.environ,
{"NGS_ANALYSIS_WORKBENCH_STATE_DIR": str(workspace / "state")},
),
mock.patch.object(snakemake, "assess_readiness", return_value=READINESS),
):
defaults.bootstrap_default_workflows()
planned = app.plan_snakemake(
"oai_fastq_qc", str(workspace / "run"), config_file=str(config)
)
self.assertEqual(planned["binding"], "snakemake")
self.assertEqual(planned["request"]["pipeline"], "oai_fastq_qc")
self.assertEqual(planned["request"]["config_file"], str(config.resolve()))
def test_readiness_for_another_engine_cannot_authorize_a_snakemake_run(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
config = _config(workspace)
with mock.patch.object(
snakemake,
"assess_readiness",
return_value={**READINESS, "binding": "nextflow"},
):
planned = snakemake.plan_run(
"fastq_qc",
str(workspace / "run"),
str(config),
workflow_source=_bundled_workflow(),
)
self.assertFalse(planned["ok"])
self.assertFalse((workspace / "run").exists())
def test_plan_preserves_unknown_readiness_and_is_not_runnable(self) -> None:
unknown = {
**READINESS,
"ok": False,
"status": "unknown",
"unknowns": ["workflow dependencies are not machine-resolvable"],
}
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
config = _config(workspace)
with mock.patch.object(
snakemake,
"assess_readiness",
return_value=unknown,
):
result = snakemake.plan_run(
pipeline="fastq_qc",
run_dir=str(workspace / "run"),
config_file=str(config),
run_id="snakemake-fastq-qc-unknown",
workflow_source=_bundled_workflow(),
)
wrote_runs = (workspace / "run").exists()
self.assertTrue(result["ok"])
self.assertFalse(result["runnable"])
self.assertEqual(result["readiness"]["status"], "unknown")
self.assertIn("readiness unresolved", result["blockers"][0])
self.assertFalse(wrote_runs)
def test_plan_is_read_only_and_binds_config_workflow_and_command(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
config = _config(workspace)
run_id = "snakemake-fastq-qc-plan"
run_dir = workspace / "run"
with mock.patch.object(
snakemake,
"assess_readiness",
return_value=READINESS,
) as readiness:
payload = snakemake.plan_run(
pipeline="fastq_qc",
run_dir=str(workspace / "run"),
config_file=str(config),
cores=3,
run_id=run_id,
workflow_source=_bundled_workflow(),
)
plan = SnakemakeRunPlan.model_validate(payload)
canonical_workspace = workspace.resolve()
canonical_config = config.resolve()
readiness.assert_called_once_with(
"fastq_qc",
canonical_config,
None,
None,
target_id="local",
controller_candidate_id=None,
refresh=False,
)
self.assertTrue(plan.runnable)
self.assertEqual(plan.request.pipeline, "fastq_qc")
self.assertEqual(plan.request.run_dir, str(canonical_workspace / "run"))
self.assertEqual(plan.request.config_file, str(canonical_config))
self.assertEqual(
plan.request.config_sha256,
f"sha256:{hashlib.sha256(config.read_bytes()).hexdigest()}",
)
self.assertEqual(
plan.request.workflow_source,
_bundled_workflow().model_copy(
update={
"root": str((PLUGIN_ROOT / "workflows" / "fastq_qc").resolve()),
"source_sha256": plan.request.workflow_sha256,
}
),
)
self.assertEqual(
plan.command_argv[2],
str(canonical_workspace / "run" / "workflow" / "workflow" / "Snakefile"),
)
canonical_run_dir = canonical_workspace / run_dir.relative_to(workspace)
self.assertNotIn("--log-handler-script", plan.command_argv)
self.assertEqual(plan.effects.output_dir, str(canonical_run_dir / "results"))
self.assertEqual(plan.effects.work_dir, plan.effects.output_dir)
self.assertIn(str(canonical_run_dir / "results"), plan.effects.local_writes)
self.assertIn(
str(canonical_run_dir / "results" / ".snakemake"), plan.effects.local_writes
)
self.assertIn(
str(canonical_run_dir / "results" / ".snakemake" / "log"),
plan.monitoring.log_paths,
)
self.assertEqual(
{path.relative_to(workspace) for path in workspace.rglob("*") if path.is_file()},
{Path(config.name)},
)
def test_plan_binds_prospective_config_and_downloads_in_one_identity(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
destination = workspace / "demo"
config = destination / "fastq-qc-config.json"
config_text = json.dumps(
{"samples": {"sample": {"r1": str(destination / "inputs" / "reads.fastq.gz")}}}
)
preparation_spec = preparation.normalize_preparation(
{
"destination_dir": str(destination),
"downloads": [
{
"relative_path": "inputs/reads.fastq.gz",
"url": (
"https://raw.githubusercontent.com/example/data/"
"abc123/reads.fastq.gz"
),
"bytes": 12,
"sha256": "1" * 64,
}
],
"generated_files": [
{
"relative_path": config.name,
"media_type": "application/json",
"content": config_text,
}
],
}
)
with mock.patch.object(
snakemake,
"assess_readiness",
return_value=READINESS,
) as readiness:
payload = snakemake.plan_run(
pipeline="fastq_qc",
run_dir=str(workspace / "run"),
config_file=str(config),
run_id="snakemake-fastq-qc-compound-plan",
preparation_spec=preparation_spec,
workflow_source=_bundled_workflow(),
)
plan = SnakemakeRunPlan.model_validate(payload)
self.assertTrue(plan.runnable)
self.assertFalse(destination.exists())
self.assertEqual(plan.preparation, preparation_spec)
self.assertEqual(plan.request.config_sha256, preparation_spec.operations[-1].sha256)
self.assertIn(preparation_spec.writes[0], plan.effects.local_writes)
readiness.assert_called_once_with(
"fastq_qc",
config.resolve(),
None,
config_text,
target_id="local",
controller_candidate_id=None,
refresh=False,
)
def test_plan_binds_the_selected_managed_environment_invocation(self) -> None:
candidate_id = f"controller-{'6' * 32}"
managed_readiness = {
**READINESS,
"selected_controller": {
"candidate_id": candidate_id,
"controller": "snakemake",
"source": "managed_environment",
"manager": "micromamba",
"executable_path": "/env/bin/snakemake",
"environment_path": "/env",
"platform_matches_host": True,
"launch_argv_prefix": [
"/runtime/micromamba",
"run",
"--prefix",
"/env",
"snakemake",
],
"recommended": True,
"recommendation_reason": "managed environment",
},
}
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
config = _config(workspace)
with mock.patch.object(
snakemake,
"assess_readiness",
return_value=managed_readiness,
):
payload = snakemake.plan_run(
pipeline="fastq_qc",
run_dir=str(workspace / "run"),
config_file=str(config),
controller_candidate_id=candidate_id,
run_id="snakemake-fastq-qc-managed",
workflow_source=_bundled_workflow(),
)
plan = SnakemakeRunPlan.model_validate(payload)
self.assertTrue(plan.runnable)
self.assertEqual(plan.request.controller_candidate_id, candidate_id)
self.assertEqual(
plan.command_argv[:5],
[
"/runtime/micromamba",
"run",
"--prefix",
"/env",
"snakemake",
],
)
def test_planner_accepts_a_catalog_resolved_source(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
root = Path(temporary)
workflow_root = root / "workflows" / "custom"
workflow_dir = workflow_root / "workflow"
workflow_dir.mkdir(parents=True)
(workflow_dir / "Snakefile").write_text(
"rule all:\n input: []\n",
encoding="utf-8",
)
workspace = root / "workspace"
workspace.mkdir()
config = _config(workspace, "custom.yaml")
source = WorkflowSource(
engine="snakemake",
root=str(workflow_root),
entrypoint="workflow/Snakefile",
)
with mock.patch.object(
snakemake,
"assess_readiness",
return_value=READINESS,
):
payload = snakemake.plan_run(
pipeline="custom",
run_dir=str(workspace / "run"),
config_file=str(config),
run_id="snakemake-custom-plan",
workflow_source=source,
)
plan = SnakemakeRunPlan.model_validate(payload)
self.assertTrue(plan.runnable)
self.assertEqual(plan.request.workflow, "workflow/Snakefile")
self.assertEqual(
Path(plan.command_argv[2]),
workspace.resolve() / "run" / "workflow" / "workflow" / "Snakefile",
)
def test_config_drift_rejects_start_without_creating_a_run(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
config = _config(workspace)
arguments = {
"pipeline": "scrnaseq",
"run_dir": str(workspace / "run"),
"config_file": str(config),
"run_id": "snakemake-scrnaseq-drift",
"workflow_source": _bundled_workflow(),
}
with mock.patch.object(
snakemake,
"assess_readiness",
return_value=READINESS,
):
approved_plan = snakemake.plan_run(**arguments)
approved_checksum = _plan_checksum(approved_plan)
config.write_text(json.dumps({"changed": True}), encoding="utf-8")
result = snakemake.approved_execution_request(
SnakemakeRunPlan.model_validate(approved_plan), approved_checksum
)
self.assertFalse(result["ok"])
self.assertEqual(result["supplied_checksum"], approved_checksum)
self.assertNotEqual(result["expected_checksum"], approved_checksum)
self.assertFalse((workspace / "run").exists())
def test_external_workflow_runs_without_being_bundled(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
root = Path(temporary)
workspace = root / "workspace"
workspace.mkdir()
source = _external_workflow(workspace)
config = _config(workspace)
with mock.patch.dict(
os.environ, {"NGS_ANALYSIS_WORKBENCH_STATE_DIR": str(root / "state")}
):
assert source.root is not None
catalog_store.save_workflow(
"custom_workflow",
"Custom workflow",
"snakemake",
catalog_store.LocalWorkflowSource(
kind="local", root=source.root, entrypoint=source.entrypoint
),
)
try:
with (
mock.patch.object(snakemake, "assess_readiness", return_value=READINESS),
mock.patch.object(
snakemake,
"build_snakemake_argv",
return_value=[sys.executable, "-c", "raise SystemExit(0)"],
),
):
registered = app.plan_snakemake(
workflow_id="custom_workflow",
run_dir=str(workspace / "run"),
config_file=str(config),
)
plan = registered
self.assertTrue(plan["runnable"])
self.assertFalse((workspace / "run").exists())
started = approved_execution.execute_registered_plan(
str(registered["plan_name"]),
str(registered["plan_id"]),
str(registered["plan_checksum"]),
)
deadline = time.monotonic() + 10
while time.monotonic() < deadline:
observation = app.observe_ngs_run(str(started["registry_run_id"]))
if observation["status"] in {"completed", "failed"}:
break
time.sleep(0.025)
result = app.get_ngs_run(str(started["registry_run_id"]))
self.assertEqual(result["status"], "completed", result)
staged_snakefile = Path(str(result["run_dir"])) / "workflow" / "Snakefile"
self.assertTrue(staged_snakefile.is_file())
self.assertEqual(
result["workflow_source"]["source_sha256"],
plan["request"]["workflow_source"]["source_sha256"],
)
self.assertEqual(result["workflow_source"]["kind"], "local")
finally:
process = daemon_client._SPAWNED_DAEMON
if process is not None and process.poll() is None:
process.terminate()
process.wait(timeout=5)
def test_source_cannot_contain_its_execution_directory(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary).resolve()
(workspace / "Snakefile").write_text("rule all:\n input: []\n", encoding="utf-8")
config = _config(workspace)
source = WorkflowSource(engine="snakemake", root=str(workspace), entrypoint="Snakefile")
planned = snakemake.plan_run(
"overlapping", str(workspace / "run"), str(config), workflow_source=source
)
self.assertFalse(planned["ok"])
self.assertFalse((workspace / "run").exists())
def test_long_workflow_identifier_preserves_a_valid_run_identity(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary).resolve()
source = _external_workflow(workspace)
config = _config(workspace)
workflow_id = "a" * 96
with mock.patch.object(snakemake, "assess_readiness", return_value=READINESS):
planned = snakemake.plan_run(
workflow_id, str(workspace / "run"), str(config), workflow_source=source
)
self.assertTrue(planned["runnable"], planned)
self.assertEqual(planned["request"]["pipeline"], workflow_id)
self.assertLessEqual(len(planned["request"]["run_id"]), 96)
self.assertFalse((workspace / "run").exists())
def test_external_workflow_without_default_profile_can_be_planned(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
source = _external_workflow(workspace)
config = _config(workspace)
with (
mock.patch.dict(
os.environ,
{"NGS_ANALYSIS_WORKBENCH_STATE_DIR": str(workspace / "private-state")},
),
mock.patch.object(
snakemake.readiness_service,
"assess",
side_effect=lambda requirements, *_args, **_kwargs: {
**READINESS,
"warnings": requirements.warnings,
},
),
):
assert source.root is not None
catalog_store.save_workflow(
"ordinary_workflow",
"Ordinary workflow",
"snakemake",
catalog_store.LocalWorkflowSource(
kind="local", root=source.root, entrypoint=source.entrypoint
),
)
readiness = app.check_snakemake_readiness(
workflow_id="ordinary_workflow",
run_dir=str(workspace / "run"),
config_file=str(config),
)
result = snakemake.plan_run(
"ordinary_workflow",
str(workspace / "run"),
str(config),
workflow_source=source,
)
self.assertTrue(result["runnable"])
self.assertEqual(readiness.warnings, [])
self.assertEqual(readiness.warnings, result["readiness"]["warnings"])
self.assertFalse((workspace / "run").exists())
def test_external_task_without_default_profile_preserves_readiness_limitation(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
source = _external_workflow(workspace)
root = Path(str(source.root))
(root / "Snakefile").write_text(
"rule run:\n output: 'result.txt'\n shell: 'missing_tool > {output}'\n",
encoding="utf-8",
)
config = _config(workspace)
with mock.patch.object(
snakemake.readiness_service,
"assess",
side_effect=lambda requirements, *_args, **_kwargs: {
**READINESS,
"warnings": requirements.warnings,
},
):
result = snakemake.plan_run(
"ordinary_workflow", str(workspace / "run"), str(config), workflow_source=source
)
self.assertTrue(result["runnable"])
self.assertEqual(result["readiness"]["warnings"], [])
def test_external_source_drift_rejects_execution_without_creating_run(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
source = _external_workflow(workspace)
config = _config(workspace)
with mock.patch.object(snakemake, "assess_readiness", return_value=READINESS):
plan = snakemake.plan_run(
"custom_workflow",
str(workspace / "run"),
str(config),
run_id="snakemake-custom-source-drift",
workflow_source=source,
)
(Path(str(source.root)) / "Snakefile").write_text(
"rule changed:\n input: []\n", encoding="utf-8"
)
result = snakemake.approved_execution_request(
SnakemakeRunPlan.model_validate(plan), _plan_checksum(plan)
)
self.assertFalse(result["ok"])
self.assertFalse((workspace / "run").exists())
def test_run_configuration_does_not_change_reusable_workflow_identity(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
source = _external_workflow(workspace)
first_config = _config(workspace, "first.json")
second_config = _config(workspace, "second.json")
second_config.write_text(json.dumps({"different_sample": True}), encoding="utf-8")
with mock.patch.object(snakemake, "assess_readiness", return_value=READINESS):
first = snakemake.plan_run(
"custom_workflow",
str(workspace / "run"),
str(first_config),
workflow_source=source,
)
second = snakemake.plan_run(
"custom_workflow",
str(workspace / "run"),
str(second_config),
workflow_source=source,
)
self.assertEqual(
first["request"]["workflow_source"]["source_sha256"],
second["request"]["workflow_source"]["source_sha256"],
)
self.assertNotEqual(
first["request"]["config_sha256"], second["request"]["config_sha256"]
)
def test_run_configuration_inside_reusable_source_is_bound_by_plan(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
workspace = Path(temporary)
source = _external_workflow(workspace)
config = _config(Path(str(source.root)))
result = snakemake.plan_run(
"custom_workflow", str(workspace / "run"), str(config), workflow_source=source
)
self.assertTrue(result["ok"], result)
self.assertEqual(result["request"]["config_file"], str(config.resolve()))
self.assertIsNotNone(result["request"]["config_sha256"])
self.assertFalse((workspace / "run").exists())
def test_approved_workflow_stages_configuration_and_completes(self) -> None:
with tempfile.TemporaryDirectory() as temporary:
root = Path(temporary)
workspace = root / "workspace"
workspace.mkdir()
config = _config(workspace)
controller = root / "snakemake"
controller.write_text(
f"#!{sys.executable}\n"
"import pathlib, sys\n"
"output = pathlib.Path(sys.argv[sys.argv.index('--directory') + 1])\n"
"report = output / 'multiqc' / 'report.html'\n"
"report.parent.mkdir(parents=True, exist_ok=True)\n"
"report.write_text('<html>observed</html>')\n"
"(output / 'artifact_index.json').write_text("
'\'{"artifacts":[{"path":"multiqc/report.html"}]}\')\n'
"native = output / '.snakemake' / 'log' / 'run.snakemake.log'\n"
"native.parent.mkdir(parents=True, exist_ok=True)\n"
"native.write_text('[Thu Aug 6 01:00:00 2026]\\nlocalrule report:\\n"
" jobid: 1\\nFinished jobid: 1 (Rule: report)\\n')\n",
encoding="utf-8",
)
controller.chmod(0o755)
readiness = {
**READINESS,
"commands": [{"name": "snakemake", "path": str(controller)}],
}
with mock.patch.dict(
os.environ, {"NGS_ANALYSIS_WORKBENCH_STATE_DIR": str(root / "state")}
):
try:
with mock.patch.object(snakemake, "assess_readiness", return_value=readiness):
plan = snakemake.plan_run(
"fastq_qc",
str(workspace / "run"),
str(config),
run_id="snakemake-approved-workflow",
workflow_source=_bundled_workflow(),
)
registered = plan_registry.register_plan("snakemake", plan)
started = approved_execution.execute_registered_plan(
str(registered["plan_name"]),
str(registered["plan_id"]),
str(registered["plan_checksum"]),
)
deadline = time.monotonic() + 10
while time.monotonic() < deadline:
observation = app.observe_ngs_run(str(started["registry_run_id"]))
if observation["status"] in {"completed", "failed"}:
break
time.sleep(0.025)
result = app.get_ngs_run(str(started["registry_run_id"]))
self.assertEqual(result["status"], "completed", result)
staged = Path(str(result["run_dir"])) / "config" / config.name
self.assertEqual(staged.read_bytes(), config.read_bytes())
observation = app.observe_ngs_run(str(started["registry_run_id"]))
report = runs.get_registry_run_report(str(started["registry_run_id"]))["report"]
self.assertTrue(observation["execution"]["evidence"]["available"])
self.assertIn(
"results/multiqc/report.html",
{entry["path"] for entry in report["entries"]},
)
self.assertFalse(
any(".snakemake" in entry["path"] for entry in report["entries"])
)
self.assertEqual(result["plan_checksum"], registered["plan_checksum"])
finally:
process = daemon_client._SPAWNED_DAEMON
if process is not None and process.poll() is None:
process.terminate()
process.wait(timeout=5)
class RemoteSnakemakePlanTests(DaemonTestCase):
def setUp(self) -> None:
super().setUp()
self.enterContext(target_input_transport())
def test_snakemake_does_not_claim_slurm_execution(self) -> None:
self._configure_ssh_target()
config = _config(self.workspace)
with mock.patch.object(
snakemake.readiness_service,
"assess",
side_effect=lambda requirements, *_args, **_kwargs: {
**READINESS,
"ok": not requirements.blockers,
"status": "blocked" if requirements.blockers else "ready",
"blockers": requirements.blockers,
},
):
defaults.bootstrap_default_workflows()
readiness = app.check_snakemake_readiness(
"oai_fastq_qc",
str(self.workspace / "run"),
config_file=str(config),
target_id="lab",
)
planned = snakemake.plan_run(
"fastq_qc",
str(self.workspace / "run"),
str(config),
target_id="lab",
workflow_source=_bundled_workflow(),
)
self.assertEqual(readiness.status, "blocked")
self.assertFalse(planned["ok"])
self.assertFalse((self.workspace / "run").exists())
def test_remote_plan_uses_target_config_and_stages_only_workflow(self) -> None:
target = self._configure_ssh_target(executor="local_process", partition=None)
source = _external_workflow(self.workspace)
config = self.workspace / "run.json"
config.write_text(
json.dumps(
{
"input": "/data/reads.fastq",
"sample_id": "S1.v2",
"reference_uri": "s3://references/GRCh38.fa",
}
),
encoding="utf-8",
)
readiness = {**READINESS, "host": {"os": "linux", "arch": "arm64"}}
with mock.patch.object(snakemake, "assess_readiness", return_value=readiness):
planned = snakemake.plan_run(
"custom_workflow",
str(self.workspace / "run"),
str(config),
workflow_source=source,
run_id="snakemake-custom-ssh",
target_id="lab",
)
self.assertTrue(planned["runnable"], planned)
remote = planned["request"]["remote"]
staged = {item["destination"]: item for item in remote["staged_files"]}
self.assertEqual(remote["config_hash"], target["config_hash"])
self.assertEqual(planned["effects"]["run_dir"], remote["run_dir"])
self.assertEqual(planned["effects"]["output_dir"], f"{remote['run_dir']}/results")
self.assertEqual(
planned["effects"]["launch_log"], f"{remote['run_dir']}/logs/snakemake.log"
)
self.assertTrue(
all(
path.startswith(f"{remote['run_dir']}/")
for path in planned["monitoring"]["log_paths"]
)
)
self.assertEqual(set(staged), {f"{remote['run_dir']}/workflow/Snakefile"})
config_argument = planned["command_argv"].index("--configfile")
self.assertEqual(planned["command_argv"][config_argument + 1], str(config))
self.assertIn(f"{remote['run_dir']}/workflow/Snakefile", planned["command_argv"])
directory = planned["command_argv"].index("--directory")
self.assertEqual(planned["command_argv"][directory + 1], f"{remote['run_dir']}/results")
self.assertTrue(set(staged).issubset(planned["effects"]["remote_writes"]))
self.assertFalse((self.workspace / "run").exists())
registered = plan_registry.register_plan("snakemake", planned)
with (
mock.patch.object(snakemake, "assess_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(
dispatch.call_args.args[0].inputs,
{str(config): f"sha256:{hashlib.sha256(config.read_bytes()).hexdigest()}"},
)
self.assertEqual(
{item.kind for item in dispatch.call_args.args[0].files},
{"copy_tree", "write_approved_plan"},
)
self.assertEqual(accepted["run_dir"], remote["run_dir"])
metadata_root = Path(dispatch.call_args.args[0].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 dispatch.call_args.args[0].files
)
)
def test_remote_plan_blocks_missing_target_config(self) -> None:
self._configure_ssh_target(executor="local_process", partition=None)
source = _external_workflow(self.workspace)
readiness = {**READINESS, "host": {"os": "linux", "arch": "arm64"}}
with mock.patch.object(snakemake, "assess_readiness", return_value=readiness):
planned = snakemake.plan_run(
"custom_workflow",
str(self.workspace / "run"),
str(self.workspace / "missing.json"),
target_id="lab",
workflow_source=source,
)
self.assertFalse(planned["runnable"], planned)
self.assertFalse((self.workspace / "run").exists())
def test_remote_plan_prepares_inputs_directly_on_target(self) -> None:
self._configure_ssh_target(executor="local_process", partition=None)
source = _external_workflow(self.workspace)
generated = self.workspace / "prepared" / "reads.fastq"
config = generated.parent / "remote.json"
content = json.dumps({"input": str(generated)})
prepared = preparation.normalize_preparation(
{
"destination_dir": str(generated.parent),
"generated_files": [
{"relative_path": generated.name, "content": "ACGT\n"},
{"relative_path": config.name, "content": content},
],
},
target_id="lab",
)
readiness = {**READINESS, "host": {"os": "linux", "arch": "arm64"}}
with mock.patch.object(snakemake, "assess_readiness", return_value=readiness):
planned = snakemake.plan_run(
"custom_workflow",
str(self.workspace / "run"),
str(config),
preparation_spec=prepared,
target_id="lab",
workflow_source=source,
)
self.assertTrue(planned["runnable"], planned)
self.assertFalse(generated.exists())
self.assertFalse(config.exists())
self.assertEqual(
planned["request"]["config_sha256"],
f"sha256:{hashlib.sha256(content.encode()).hexdigest()}",
)
self.assertTrue(set(prepared.writes).issubset(planned["effects"]["remote_writes"]))
self.assertTrue(set(prepared.writes).isdisjoint(planned["effects"]["local_writes"]))
self.assertEqual(
[
Path(item["destination"]).name
for item in planned["request"]["remote"]["staged_files"]
],
["Snakefile"],
)
if __name__ == "__main__":
unittest.main()
SHA-256: 9756d3e527fe28c149a428682f275761a4a6f0e7e398c459b0891c591f273654