← Files NGS Analysis WorkbenchARCHIVED FILE
mcp/ngs_workbench_mcp/workflows/nextflow.py
40.8 KB · Sep 30, 2026 · 23:20 UTC
"""Plan and execute Nextflow workflows, including curated nf-core collections."""
from __future__ import annotations
import json
import re
import uuid
from datetime import UTC, datetime
from pathlib import Path, PurePosixPath
from string import Template
from typing import Any
from ngs_workbench_daemon.hashing import sha256_bytes
from ngs_workbench_daemon.protocol import (
ExecutionRequest,
FileOperation,
current_controller_environment,
run_metadata_directory,
)
from pydantic import ValidationError
from .. import compute_targets, preparation, readiness_service, runs
from ..plans import (
InputSummary,
MonitoringPlan,
NextflowPlanRequest,
NextflowRunPlan,
PreparationSpec,
RemoteExecution,
RemoteStagedFile,
RunEffects,
controller_launch_argv,
normalize_display_name,
plan_checksum,
)
from ..readiness.models import (
ReadinessAssessment,
ReadinessEvidence,
RequirementSet,
RuntimeRequirement,
)
from ..report.artifacts import RESULT_ROOTS
from .remote import staged_source_files
from .source import WorkflowSource, observe_source
BINDING = "nextflow"
def _error(*errors: str, **details: Any) -> dict[str, Any]:
return {"ok": False, "errors": list(errors), **details}
def assess_readiness(
workflow_id: str,
runtime_snapshot_id: str | None = None,
*,
target_id: str = "local",
controller_candidate_id: str | None = None,
refresh: bool = False,
) -> ReadinessAssessment:
"""Observe the controller without assuming nf-core profiles or task environments."""
requirements = RequirementSet(
binding=BINDING,
pipeline=workflow_id,
requirements=[
RuntimeRequirement(
id="nextflow",
layer="workflow_controller",
capability="command",
value="nextflow",
source="workflow.engine",
probe_mode="version",
)
],
warnings=["workflow-specific task software and references are not independently verified"],
)
target = compute_targets.resolve_compute_target(target_id)
if target.controller_transport == "ssh" and target.executor != "local_process":
requirements.blockers.append(
"generic Nextflow does not yet support the selected remote task executor"
)
return readiness_service.assess(
requirements,
runtime_snapshot_id,
target_id=target_id,
controller_candidate_id=controller_candidate_id,
refresh=refresh,
)
def _remote_execution(
*,
target: compute_targets.ComputeTarget,
source: WorkflowSource,
run_dir: Path,
metadata_dir: Path,
) -> RemoteExecution:
if target.config_hash is None or target.host_access is None or target.workspace_root is None:
raise ValueError("SSH target is missing its approved connection or workspace")
remote_root = PurePosixPath(run_dir)
staged = staged_source_files(
source_root=Path(source.root),
local_run_dir=metadata_dir,
remote_run_dir=remote_root,
)
return RemoteExecution(
config_hash=target.config_hash,
host_access=target.host_access,
workspace_root=target.workspace_root,
run_dir=remote_root.as_posix(),
executor="local_process",
staged_files=staged,
)
def plan_run(
workflow_id: str,
run_dir: str,
workflow_source: WorkflowSource,
*,
workflow_parameters: dict[str, str | int | float | bool] | None = None,
params_file: str | None = None,
profile: str | None = None,
display_name: str | None = None,
runtime_snapshot_id: str | None = None,
controller_candidate_id: str | None = None,
target_id: str = "local",
preparation_spec: PreparationSpec | None = None,
run_id: str | None = None,
first_run_id: str | None = None,
attempt_number: int | None = 1,
workflow_version_id: str | None = None,
refresh_runtime: bool = False,
) -> dict[str, Any]:
"""Create one read-only source-bound plan without curated nf-core conventions."""
if not re.fullmatch(r"[a-z0-9][a-z0-9_-]{0,95}", workflow_id):
return _error("invalid workflow identifier")
try:
source = observe_source(workflow_source)
target = compute_targets.resolve_compute_target(target_id)
execution_dir = runs.resolve_run_directory(run_dir, target_id)
except (OSError, ValueError) as exc:
return _error(str(exc))
if source.engine != BINDING:
return _error("workflow engine does not match the Nextflow binding")
if target.controller_transport == "ssh" and target.executor != "local_process":
return _error("generic Nextflow does not yet support the selected remote task executor")
resolved_run_id = run_id or _new_run_id(workflow_id)
metadata_dir = run_metadata_directory(target_id, str(execution_dir))
staging_dir = metadata_dir if target.controller_transport == "ssh" else execution_dir
if (
target.controller_transport != "ssh"
and source.root is not None
and execution_dir.is_relative_to(Path(source.root))
):
return _error("workflow source must not contain its execution directory")
parameters = dict(workflow_parameters or {})
for name in parameters:
if not re.fullmatch(r"[a-zA-Z][a-zA-Z0-9_-]*", name):
return _error(f"unsupported Nextflow workflow parameter: {name}")
if profile is not None and (
not profile.strip() or any(token.strip().startswith("-") for token in profile.split(","))
):
return _error("profile tokens must be non-empty and must not start with '-'")
try:
resolved_params = preparation.resolve_input_path(params_file, target_id)
except ValueError as exc:
return _error(str(exc))
blockers: list[str] = []
params_digest = None
if resolved_params is not None:
try:
params_digest = preparation.inspect_input(resolved_params, preparation_spec, target_id)[
"sha256"
]
except (OSError, ValueError) as exc:
blockers.append(f"params_file could not be inspected: {exc}")
if (
target.controller_transport != "ssh"
and source.root is not None
and resolved_params is not None
and resolved_params.is_relative_to(Path(source.root))
):
return _error("run parameters must remain outside reusable workflow source")
try:
readiness = ReadinessAssessment.model_validate(
assess_readiness(
workflow_id,
runtime_snapshot_id,
target_id=target_id,
controller_candidate_id=controller_candidate_id,
refresh=refresh_runtime,
)
)
except ValueError as exc:
return _error(str(exc))
if readiness.binding is not None and readiness.binding != BINDING:
return _error("readiness does not match the selected workflow engine")
readiness = readiness.model_copy(update={"binding": BINDING})
blockers.extend(readiness.blockers)
blockers.extend(f"readiness unresolved: {item}" for item in readiness.unknowns)
try:
runs.check_run_directory(execution_dir, target_id)
except (OSError, ValueError) as exc:
blockers.append(str(exc))
remote = None
if target.controller_transport == "ssh":
if readiness.host is None or readiness.host.os.lower() != "linux":
blockers.append("SSH workflow execution requires an observed Linux controller host")
else:
try:
remote = _remote_execution(
target=target,
source=source,
run_dir=execution_dir,
metadata_dir=metadata_dir,
)
except (OSError, ValueError) as exc:
blockers.append(f"SSH workflow staging could not be resolved: {exc}")
try:
request = NextflowPlanRequest(
pipeline=workflow_id,
workflow_version_id=workflow_version_id,
parameter_contract="generic",
target=target.ref(),
remote=remote,
display_name=normalize_display_name(display_name, workflow_id),
# LEGACY_CLEANUP(2026-09-07): Revalidating pre-lineage approvals passes None;
# preserve absent fields/checksum until those saved plans are unsupported.
first_run_id=first_run_id or (resolved_run_id if attempt_number is not None else None),
attempt_number=attempt_number,
workflow=source.entrypoint,
run_dir=str(execution_dir),
profile=profile or "",
sample_sheet=None,
sample_sheet_sha256=None,
revision=None,
params_file=str(resolved_params) if resolved_params is not None else None,
params_file_sha256=params_digest,
runtime_snapshot_id=readiness.snapshot_id,
controller_candidate_id=(
readiness.selected_controller.candidate_id
if readiness.selected_controller is not None
else controller_candidate_id
),
workflow_source=source,
workflow_parameters=parameters or None,
run_id=resolved_run_id,
)
except ValidationError as exc:
return _error(*[error["msg"] for error in exc.errors()])
controller_root = execution_dir
argv = [
*controller_launch_argv(readiness, "nextflow"),
"run",
str(controller_root / "workflow" / source.entrypoint),
"-output-dir",
str(controller_root / "results"),
"-work-dir",
str(controller_root / "work"),
"-with-trace",
str(controller_root / "logs" / "nextflow.trace.txt"),
*(["-profile", profile] if profile else []),
*(
[
"-params-file",
str(
resolved_params
if remote is not None
else execution_dir / "config" / resolved_params.name
),
]
if resolved_params is not None
else []
),
*[
f"--{name}={str(value).lower() if isinstance(value, bool) else value}"
for name, value in sorted(parameters.items())
],
]
approved_plan = metadata_dir / "workflow" / "approved_plan.json"
local_writes = [
*(preparation_spec.writes if preparation_spec is not None and remote is None else []),
*runs.registry_effect_paths(),
str(staging_dir),
str(staging_dir / "workflow"),
str(approved_plan),
*(
[str(execution_dir / "config" / resolved_params.name)]
if resolved_params is not None and remote is None
else []
),
*(
[
str(execution_dir / "results"),
str(execution_dir / "logs" / "nextflow.trace.txt"),
str(execution_dir / "logs" / "nextflow.log"),
str(execution_dir / "work"),
str(execution_dir / ".nextflow"),
str(execution_dir / ".nextflow.log"),
]
if remote is None
else []
),
]
remote_writes = (
[
*(preparation_spec.writes if preparation_spec is not None else []),
str(controller_root / "workflow" / "approved_plan.json"),
str(controller_root / "workflow" / "controller.json"),
str(controller_root / "workflow" / "controller.exit"),
str(controller_root / "workflow" / "controller.exit.tmp"),
str(controller_root / "logs" / "nextflow.log"),
str(controller_root / "logs" / "nextflow.trace.txt"),
*[str(controller_root / path) for path in RESULT_ROOTS],
str(controller_root / "work"),
str(controller_root / ".nextflow"),
str(controller_root / ".nextflow.log"),
*[item.destination for item in remote.staged_files],
]
if remote is not None
else []
)
effects_root = execution_dir
draft = NextflowRunPlan(
runnable=not blockers and readiness.status == "ready",
request=request,
command_argv=argv,
readiness=readiness,
input_summary=InputSummary(
source="params_file" if resolved_params is not None else "missing",
path=str(resolved_params) if resolved_params is not None else None,
sha256=params_digest,
),
effects=RunEffects(
run_dir=str(effects_root),
output_dir=str(effects_root / "results"),
work_dir=str(effects_root / "work"),
launch_log=str(effects_root / "logs" / "nextflow.log"),
local_writes=local_writes,
remote_writes=remote_writes,
downloads=[
f"{item.url} ({item.bytes} bytes, {item.sha256})"
for item in preparation_spec.operations
if item.url is not None
]
if preparation_spec is not None
else [],
network_access=[
f"verified HTTPS download from {item.url}"
for item in preparation_spec.operations
if item.url is not None
]
if preparation_spec is not None
else [],
),
preparation=preparation_spec,
blockers=blockers,
warnings=readiness.warnings,
monitoring=MonitoringPlan(
terminal_statuses=["completed", "failed", "canceled", "orphaned"],
log_paths=[str(effects_root / "logs" / "nextflow.log")],
artifact_paths=[
str(effects_root / "logs" / "nextflow.trace.txt"),
str(effects_root / "results"),
],
),
)
return draft.model_dump(mode="json")
def approved_execution_request(
approved_plan: NextflowRunPlan, approved_checksum: str
) -> ExecutionRequest | dict[str, Any]:
"""Replay generic source identity and describe only already-approved file effects."""
request = approved_plan.request
if request.workflow_source is None or request.workflow_source.engine != BINDING:
return _error("approved Nextflow plan does not contain its exact workflow source")
existing = runs.find_registered_run(request.target.target_id, request.run_dir)
if existing is not None:
if existing.plan_checksum != approved_checksum:
return _error("run identity already belongs to a different approved plan")
plan = approved_plan
else:
planned = plan_run(
workflow_id=request.pipeline,
run_dir=request.run_dir,
workflow_source=request.workflow_source,
workflow_parameters=request.workflow_parameters,
params_file=request.params_file,
profile=request.profile or None,
display_name=request.display_name,
runtime_snapshot_id=request.runtime_snapshot_id,
controller_candidate_id=request.controller_candidate_id,
target_id=request.target.target_id,
preparation_spec=approved_plan.preparation,
run_id=request.run_id,
first_run_id=request.first_run_id,
attempt_number=request.attempt_number,
workflow_version_id=request.workflow_version_id,
refresh_runtime=True,
)
if not planned.get("ok"):
return planned
plan = NextflowRunPlan.model_validate(planned)
if plan_checksum(plan) != approved_checksum:
return _error("plan checksum does not match the current workflow, parameters, or runtime")
if not plan.runnable:
return _error(*(plan.blockers or plan.readiness.unknowns), status="blocked")
metadata_dir = run_metadata_directory(request.target.target_id, request.run_dir)
run_dir = metadata_dir if request.remote is not None else Path(plan.effects.run_dir)
source = request.workflow_source
workflow = FileOperation(
kind="copy_tree",
destination=str(run_dir / "workflow"),
source=source.root,
source_sha256=source.source_sha256,
)
files = [workflow]
if request.params_file is not None and request.remote is None:
path = Path(request.params_file)
files.append(
FileOperation(
kind="copy_file",
destination=str(run_dir / "config" / path.name),
source=str(path),
source_sha256=request.params_file_sha256,
)
)
files.append(
FileOperation(
kind="write_approved_plan",
destination=str(metadata_dir / "workflow" / "approved_plan.json"),
)
)
return ExecutionRequest(
binding=BINDING,
plan_checksum=approved_checksum,
plan=plan.model_dump(mode="json"),
files=files,
inputs={request.params_file: request.params_file_sha256}
if request.params_file is not None and request.params_file_sha256 is not None
else {},
metadata=runs.execution_metadata(plan),
environment=current_controller_environment(),
)
# Collection-specific conventions for curated nf-core workflows.
_STANDARD_PROFILE_TASK_ENVIRONMENTS = {
"docker": "docker",
"podman": "podman",
"apptainer": "apptainer",
"singularity": "singularity",
"conda": "conda|mamba|micromamba",
}
_STANDARD_NON_ENVIRONMENT_PROFILES = {"test"}
_SLURM_CONFIGURATION = Template(
"""process {
executor = 'slurm'
withName: '.*' {
$partition
$account
}
}
"""
)
def _remote_nextflow_configuration(
target: compute_targets.ComputeTarget,
) -> str | None:
if target.executor != "slurm":
return None
settings = target.executor_configuration
partition = settings.get("partition")
account = settings.get("account")
return _SLURM_CONFIGURATION.substitute(
partition=f"queue = {json.dumps(partition)}" if partition else "",
account=f"clusterOptions = {json.dumps(f'--account={account}')}" if account else "",
)
def resolve_runtime_requirements(
pipeline: str,
profile: str,
*,
workflow: str | None = None,
revision: str | None = None,
) -> RequirementSet:
"""Describe what one exact curated nf-core workflow requires."""
tokens = {token.strip().lower() for token in profile.split(",") if token.strip()}
requirements = [
RuntimeRequirement(
id="nextflow",
layer="workflow_controller",
capability="command",
value="nextflow",
source="request.binding",
detail="the nf-core collection executes workflows with Nextflow",
probe_mode="version",
)
]
unknowns: list[str] = []
warnings: list[str] = []
selected_task_environments = [
name for name in _STANDARD_PROFILE_TASK_ENVIRONMENTS if name in tokens
]
unresolved_profile_tokens = sorted(
tokens - _STANDARD_PROFILE_TASK_ENVIRONMENTS.keys() - _STANDARD_NON_ENVIRONMENT_PROFILES
)
if not selected_task_environments:
unknowns.append(
"effective task software environment is unresolved; the requested profile does not "
"expose whether Nextflow tasks use containers, Conda, environment modules, or host "
"tools"
)
elif len(selected_task_environments) > 1:
unknowns.append(
"multiple standard task-environment profiles were requested; resolve the effective "
"Nextflow configuration to determine what is enabled: "
f"{', '.join(selected_task_environments)}"
)
if unresolved_profile_tokens:
warnings.append(
"readiness preserves but does not independently interpret additional profile "
f"components: {', '.join(unresolved_profile_tokens)}; confirm their effects from the "
"selected workflow revision before approval"
)
for task_environment in selected_task_environments:
requirements.append(
RuntimeRequirement(
id=f"{task_environment}-command",
layer="task_environment",
capability="command_group" if task_environment == "conda" else "command",
value=_STANDARD_PROFILE_TASK_ENVIRONMENTS[task_environment],
source="request.profile",
detail=(
f"the standard {task_environment} profile enables {task_environment} "
"for task software"
),
probe_mode="version",
)
)
if task_environment == "docker":
requirements.append(
RuntimeRequirement(
id="docker-daemon",
layer="task_environment",
capability="daemon",
value="docker",
source="request.profile",
)
)
return RequirementSet(
binding="nextflow",
pipeline=pipeline,
requirements=requirements,
evidence=[
ReadinessEvidence(
kind="requested_profile",
source="request.profile",
detail=profile,
),
ReadinessEvidence(
kind="workflow_binding",
source=workflow or pipeline,
detail=f"revision {revision}" if revision else "revision is not pinned",
),
],
unknowns=unknowns,
warnings=warnings,
)
def assess_nfcore_readiness(
pipeline: str,
profile: str,
runtime_snapshot_id: str | None = None,
*,
target_id: str = "local",
workflow: str | None = None,
revision: str | None = None,
controller_candidate_id: str | None = None,
refresh: bool = False,
) -> ReadinessAssessment:
"""Compose nf-core requirements with generic runtime observation and evaluation."""
requirement_set = resolve_runtime_requirements(
pipeline,
profile,
workflow=workflow,
revision=revision,
)
return readiness_service.assess(
requirement_set,
runtime_snapshot_id,
target_id=target_id,
controller_candidate_id=controller_candidate_id,
refresh=refresh,
)
def _remote_content(destination: PurePosixPath, content: str) -> RemoteStagedFile:
encoded = content.encode("utf-8")
return RemoteStagedFile(
destination=destination.as_posix(),
sha256=sha256_bytes(encoded),
bytes=len(encoded),
content=content,
)
def _new_run_id(pipeline: str) -> str:
stamp = datetime.now(UTC).strftime("%Y%m%d-%H%M%S")
pipeline_slug = pipeline.replace("_", "-")[:62]
return f"nextflow-{pipeline_slug}-{stamp}-{uuid.uuid4().hex[:8]}"
def _input_summary(
sample_sheet: Path | None,
sample_sheet_sha256: str | None,
*,
test_profile: bool,
) -> InputSummary:
if sample_sheet is None:
return InputSummary(source="workflow_test_profile" if test_profile else "missing")
return InputSummary(
source="sample_sheet",
path=str(sample_sheet),
sha256=sample_sheet_sha256,
)
def build_nextflow_argv(
*,
nextflow_argv_prefix: list[str] | None = None,
workflow: str,
sample_sheet: Path | PurePosixPath | None,
run_dir: Path | PurePosixPath,
profile: str,
revision: str | None,
params_file: Path | PurePosixPath | None,
parameter_contract: str,
) -> list[str]:
"""Build the literal argv used to start Nextflow."""
generated_params = run_dir / "workflow" / "params.generated.json"
argv = [
*(nextflow_argv_prefix or ["nextflow"]),
"run",
workflow,
]
if params_file is not None or sample_sheet is not None:
argv.extend(["-params-file", str(params_file or generated_params)])
argv.extend(
[
"-work-dir",
str(run_dir / "work"),
"-with-report",
str(run_dir / "workflow" / "nextflow_report.html"),
"-with-timeline",
str(run_dir / "workflow" / "timeline.html"),
"-with-trace",
str(run_dir / "workflow" / "trace.txt"),
"-with-dag",
str(run_dir / "workflow" / "dag.html"),
]
)
if revision:
argv.extend(["-r", revision])
if profile:
argv.extend(["-profile", profile])
if parameter_contract == "nf-core" and params_file and sample_sheet:
argv.extend(
[
"--input",
str(sample_sheet),
"--outdir",
str(run_dir / "results"),
]
)
elif parameter_contract == "nf-core" and sample_sheet is None:
argv.extend(["--outdir", str(run_dir / "results")])
elif parameter_contract == "generic":
argv.extend(["-output-dir", str(run_dir / "results")])
return argv
def _plan_error(*errors: str, **details: Any) -> dict[str, Any]:
return {"ok": False, "errors": list(errors), **details}
def plan_remote_run(
pipeline: str,
run_dir: str,
profile: str,
display_name: str | None = None,
sample_sheet: str | None = None,
revision: str | None = None,
params_file: str | None = None,
run_id: str | None = None,
first_run_id: str | None = None,
attempt_number: int | None = 1,
runtime_snapshot_id: str | None = None,
controller_candidate_id: str | None = None,
preparation_spec: PreparationSpec | None = None,
*,
workflow: str | None = None,
workflow_title: str | None = None,
workflow_version_id: str | None = None,
parameter_contract: str = "nf-core",
target_id: str = "local",
refresh_runtime: bool = False,
) -> dict[str, Any]:
"""Build a typed engine-resolved remote Nextflow plan."""
workflow_title = workflow_title or pipeline
if workflow is None:
return _plan_error(f"workflow source is unavailable: {pipeline}")
try:
target = compute_targets.resolve_compute_target(target_id)
execution_dir = runs.resolve_run_directory(run_dir, target_id)
except ValueError as exc:
return _plan_error(str(exc))
profile_tokens_list = [token.strip() for token in profile.split(",") if token.strip()]
profile = ",".join(profile_tokens_list)
revision = revision.strip() if revision else None
if (parameter_contract == "nf-core" and not profile) or any(
token.startswith("-") for token in profile_tokens_list
):
return _plan_error("profile tokens must be non-empty and must not start with '-'")
if revision is not None and revision.startswith("-"):
return _plan_error("revision must not start with '-'")
try:
resolved_sample_sheet = preparation.resolve_input_path(sample_sheet, target_id)
resolved_params_file = preparation.resolve_input_path(params_file, target_id)
except ValueError as exc:
return _plan_error(str(exc))
blockers = []
profile_tokens = {token.strip().lower() for token in profile.split(",")}
if (
parameter_contract == "nf-core"
and resolved_sample_sheet is None
and "test" not in profile_tokens
):
blockers.append("sample_sheet is required unless profile includes 'test'")
input_hashes: dict[Path, str] = {}
for path in (resolved_sample_sheet, resolved_params_file):
if path is not None:
try:
input_hashes[path] = preparation.inspect_input(path, preparation_spec, target_id)[
"sha256"
]
except (OSError, ValueError) as exc:
blockers.append(f"input could not be inspected: {path}: {exc}")
sample_sheet_sha256 = (
input_hashes.get(resolved_sample_sheet) if resolved_sample_sheet is not None else None
)
params_file_sha256 = (
input_hashes.get(resolved_params_file) if resolved_params_file is not None else None
)
try:
readiness_result = (
assess_nfcore_readiness(
pipeline,
profile,
runtime_snapshot_id,
target_id=target.target_id,
workflow=workflow,
revision=revision,
controller_candidate_id=controller_candidate_id,
refresh=refresh_runtime,
)
if parameter_contract == "nf-core"
else assess_readiness(
pipeline,
runtime_snapshot_id,
target_id=target.target_id,
controller_candidate_id=controller_candidate_id,
refresh=refresh_runtime,
)
)
readiness = ReadinessAssessment.model_validate(readiness_result)
except ValueError as exc:
return _plan_error(str(exc))
if readiness.binding is not None and readiness.binding != BINDING:
return _plan_error("readiness does not match the selected workflow engine")
readiness = readiness.model_copy(update={"binding": BINDING})
blockers.extend(readiness.blockers)
blockers.extend(f"readiness unresolved: {item}" for item in readiness.unknowns)
resolved_run_id = run_id or _new_run_id(pipeline)
metadata_dir = run_metadata_directory(target_id, str(execution_dir))
try:
runs.check_run_directory(execution_dir, target_id)
except (OSError, ValueError) as exc:
blockers.append(str(exc))
remote = None
controller_run_dir: Path | PurePosixPath = execution_dir
if target.controller_transport == "ssh":
if readiness.host is None or readiness.host.os.lower() != "linux":
blockers.append("SSH workflow execution requires an observed Linux controller host")
if (
target.config_hash is None
or target.host_access is None
or target.workspace_root is None
):
return _plan_error("SSH target is missing its approved connection or workspace")
remote_run_dir = PurePosixPath(execution_dir)
staged_files = []
if resolved_params_file is None and resolved_sample_sheet is not None:
staged_files.append(
_remote_content(
remote_run_dir / "workflow" / "params.generated.json",
json.dumps(
{
"input": str(resolved_sample_sheet),
"outdir": str(remote_run_dir / "results"),
}
)
+ "\n",
)
)
if configuration := _remote_nextflow_configuration(target):
staged_files.append(
_remote_content(
remote_run_dir / "workflow" / "nextflow.config",
configuration,
)
)
remote = RemoteExecution(
config_hash=target.config_hash,
host_access=target.host_access,
workspace_root=target.workspace_root,
run_dir=remote_run_dir.as_posix(),
executor=target.executor,
executor_configuration=target.executor_configuration,
staged_files=staged_files,
)
controller_run_dir = remote_run_dir
controller_argv = controller_launch_argv(readiness, "nextflow")
if remote is not None and remote.executor == "slurm":
controller_argv.extend(["-c", str(controller_run_dir / "workflow" / "nextflow.config")])
argv = build_nextflow_argv(
nextflow_argv_prefix=controller_argv,
workflow=workflow,
sample_sheet=resolved_sample_sheet,
run_dir=controller_run_dir,
profile=profile,
revision=revision,
params_file=resolved_params_file,
parameter_contract=parameter_contract,
)
local_writes = [
*(preparation_spec.writes if preparation_spec is not None and remote is None else []),
*runs.registry_effect_paths(),
str(metadata_dir / "workflow" / "approved_plan.json"),
]
if remote is None:
local_writes.extend(
[
str(execution_dir / "logs" / "nextflow.log"),
str(execution_dir / "results"),
str(execution_dir / "work"),
]
)
if remote is None and resolved_params_file is None and resolved_sample_sheet is not None:
local_writes.append(str(execution_dir / "workflow" / "params.generated.json"))
warnings = list(readiness.warnings)
if not revision:
warnings.append("no nf-core revision is pinned")
try:
request = NextflowPlanRequest(
pipeline=pipeline,
workflow_version_id=workflow_version_id,
parameter_contract=parameter_contract,
target=target.ref(),
remote=remote,
display_name=normalize_display_name(display_name, workflow_title),
# LEGACY_CLEANUP(2026-09-07): Revalidating pre-lineage approvals passes None;
# preserve absent fields/checksum until those saved plans are unsupported.
first_run_id=first_run_id or (resolved_run_id if attempt_number is not None else None),
attempt_number=attempt_number,
workflow=workflow,
run_dir=str(execution_dir),
profile=profile,
sample_sheet=str(resolved_sample_sheet) if resolved_sample_sheet else None,
sample_sheet_sha256=sample_sheet_sha256,
revision=revision,
params_file=str(resolved_params_file) if resolved_params_file else None,
params_file_sha256=params_file_sha256,
runtime_snapshot_id=readiness.snapshot_id,
controller_candidate_id=(
readiness.selected_controller.candidate_id
if readiness.selected_controller is not None
else controller_candidate_id
),
run_id=resolved_run_id,
)
except ValidationError as exc:
return _plan_error(*[error["msg"] for error in exc.errors()])
downloads = [
f"workflow source for {workflow}"
+ (f" at revision {revision}" if revision else " at its default revision"),
]
network_access = ["Nextflow workflow source resolution"]
downloads.append("container images referenced by the workflow if absent from the local cache")
network_access.append("container registry access for uncached images")
if parameter_contract == "nf-core" and "test" in profile_tokens:
downloads.append(
"test-profile inputs referenced by the workflow if absent from the local cache"
)
network_access.append("remote test-data access declared by the workflow")
if preparation_spec is not None:
downloads.extend(
f"{item.url} ({item.bytes} bytes, {item.sha256})"
for item in preparation_spec.operations
if item.url is not None
)
network_access.extend(
f"verified HTTPS download from {item.url}"
for item in preparation_spec.operations
if item.url is not None
)
effects_root = execution_dir
draft = NextflowRunPlan(
runnable=not blockers and readiness.status == "ready",
request=request,
command_argv=argv,
readiness=readiness,
input_summary=_input_summary(
resolved_sample_sheet,
sample_sheet_sha256,
test_profile="test" in profile_tokens,
),
effects=RunEffects(
run_dir=str(effects_root),
output_dir=str(effects_root / "results"),
work_dir=str(effects_root / "work"),
launch_log=str(effects_root / "logs" / "nextflow.log"),
local_writes=local_writes,
remote_writes=(
[
*(preparation_spec.writes if preparation_spec is not None else []),
str(controller_run_dir / "workflow" / "approved_plan.json"),
str(controller_run_dir / "workflow" / "nextflow_report.html"),
str(controller_run_dir / "workflow" / "timeline.html"),
str(controller_run_dir / "workflow" / "trace.txt"),
str(controller_run_dir / "workflow" / "dag.html"),
str(controller_run_dir / "workflow" / "controller.json"),
str(controller_run_dir / "workflow" / "controller.exit"),
str(controller_run_dir / "workflow" / "controller.exit.tmp"),
str(controller_run_dir / "logs" / "nextflow.log"),
str(controller_run_dir / "results"),
str(controller_run_dir / "work"),
*[item.destination for item in remote.staged_files],
]
if remote is not None
else []
),
downloads=downloads,
network_access=network_access,
),
preparation=preparation_spec,
blockers=blockers,
warnings=warnings,
monitoring=MonitoringPlan(
terminal_statuses=[
"completed",
"failed",
"canceled",
"orphaned",
],
log_paths=[str(effects_root / "logs" / "nextflow.log")],
artifact_paths=[
str(effects_root / "results"),
str(effects_root / "workflow" / "nextflow_report.html"),
str(effects_root / "workflow" / "timeline.html"),
str(effects_root / "workflow" / "trace.txt"),
str(effects_root / "workflow" / "dag.html"),
],
),
)
return draft.model_dump(mode="json")
def approved_nfcore_execution_request(
approved_plan: NextflowRunPlan,
approved_checksum: str,
*,
refresh_runtime: bool = True,
) -> ExecutionRequest | dict[str, Any]:
"""Revalidate one immutable plan and describe its approved daemon-owned effects."""
approved = approved_plan.request
existing = runs.find_registered_run(approved.target.target_id, approved.run_dir)
if existing is not None:
if existing.plan_checksum != approved_checksum:
return _plan_error("run identity already belongs to a different approved plan")
plan = approved_plan
elif refresh_runtime:
planned = plan_remote_run(
pipeline=approved.pipeline,
run_dir=approved.run_dir,
profile=approved.profile,
display_name=approved.display_name,
sample_sheet=approved.sample_sheet,
revision=approved.revision,
params_file=approved.params_file,
run_id=approved.run_id,
first_run_id=approved.first_run_id,
attempt_number=approved.attempt_number,
runtime_snapshot_id=approved.runtime_snapshot_id,
controller_candidate_id=approved.controller_candidate_id,
preparation_spec=approved_plan.preparation,
workflow=approved.workflow,
workflow_title=approved.display_name,
workflow_version_id=approved.workflow_version_id,
parameter_contract=approved.parameter_contract or "nf-core",
target_id=approved.target.target_id,
refresh_runtime=True,
)
if not planned.get("ok"):
return planned
plan = NextflowRunPlan.model_validate(planned)
else:
plan = approved_plan
current_checksum = plan_checksum(plan)
if current_checksum != approved_checksum:
return _plan_error(
"plan checksum does not match the current inputs or runtime",
expected_checksum=current_checksum,
supplied_checksum=approved_checksum,
)
if not plan.runnable:
return _plan_error(
*(plan.blockers or plan.readiness.unknowns),
status="blocked",
readiness=plan.readiness.model_dump(mode="json"),
)
payload = plan.model_dump(mode="json")
request = plan.request
metadata_dir = run_metadata_directory(request.target.target_id, request.run_dir)
run_dir = metadata_dir if request.remote is not None else Path(plan.effects.run_dir)
workflow_dir = run_dir / "workflow"
files = []
if request.remote is None and request.params_file is None and request.sample_sheet is not None:
files.append(
FileOperation(
kind="write_json",
destination=str(workflow_dir / "params.generated.json"),
content={"input": request.sample_sheet, "outdir": plan.effects.output_dir},
)
)
files.append(
FileOperation(
kind="write_approved_plan",
destination=str(metadata_dir / "workflow" / "approved_plan.json"),
)
)
response = {"runtime_snapshot_id": request.runtime_snapshot_id}
return ExecutionRequest(
binding="nextflow",
plan_checksum=approved_checksum,
plan=payload,
files=files,
inputs={
path: checksum
for path, checksum in (
(request.sample_sheet, request.sample_sheet_sha256),
(request.params_file, request.params_file_sha256),
)
if path is not None and checksum is not None
},
metadata={"response": response, **runs.execution_metadata(plan)},
environment=current_controller_environment(),
)
SHA-256: f0e1108a123bb4905e9eb9d1ad546133998bf1fb649be4471b763aca851323f5