Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
108 changes: 79 additions & 29 deletions planemo/engine/interface.py
Original file line number Diff line number Diff line change
@@ -1,14 +1,20 @@
"""Module contianing the :class:`Engine` abstraction."""

import abc
import copy
import json
import os
import tempfile
from contextlib import contextmanager
from typing import (
Any,
Callable,
Dict,
Iterator,
List,
Optional,
)
from urllib.parse import urlparse

import click

Expand All @@ -17,10 +23,76 @@
from planemo.runnable import (
cases,
RunnableType,
TestCase,
)
from planemo.test.results import StructuredData


def _absolute_test_data_path(path: Any, tests_directory: str) -> Any:
if not isinstance(path, str) or os.path.isabs(path) or urlparse(path).scheme or path.startswith("#"):
return path
return os.path.abspath(os.path.join(tests_directory, path))


def _absolutize_composite_data(file_value: Dict[str, Any], tests_directory: str) -> None:
composite_data = file_value.get("composite_data") or []
for composite_index, composite_item in enumerate(composite_data):
if isinstance(composite_item, dict):
for path_key in ("path", "location"):
if path_key in composite_item:
composite_item[path_key] = _absolute_test_data_path(composite_item[path_key], tests_directory)
elif isinstance(composite_item, str):
composite_data[composite_index] = _absolute_test_data_path(composite_item, tests_directory)


def _absolutize_nested_job_paths(value: Any, tests_directory: str) -> None:
if isinstance(value, list):
for item in value:
_absolutize_nested_job_paths(item, tests_directory)
elif isinstance(value, dict):
if value.get("class") in ("File", "Directory"):
for path_key in ("path", "location"):
if path_key in value:
value[path_key] = _absolute_test_data_path(value[path_key], tests_directory)
_absolutize_composite_data(value, tests_directory)

for item in value.values():
_absolutize_nested_job_paths(item, tests_directory)


def _absolutize_job_paths(job: Dict[str, Any], tests_directory: str) -> Dict[str, Any]:
"""Copy a test job and resolve its local data paths against the test directory."""
prepared_job = copy.deepcopy(job)
_absolutize_nested_job_paths(prepared_job, tests_directory)
return prepared_job


@contextmanager
def materialized_job_paths(test_cases: List[TestCase]) -> Iterator[List[str]]:
"""Yield a job file path for each test case, for the duration of the context.

Jobs defined inline in a test definition get written to a temporary directory
rather than beside the definition, so the source directory is left untouched and
need not be writable. Their relative data paths are resolved against the test
definition's directory first so they survive the move - the test case's own job
is not modified. Test cases that already point at a job file are passed through
untouched.
"""
with tempfile.TemporaryDirectory(prefix="planemo-test-jobs-") as job_directory:
job_paths = []
for index, test_case in enumerate(test_cases):
if test_case.job_path is not None:
job_paths.append(test_case.job_path)
continue
# a test case defines exactly one of job_path and job
assert test_case.job is not None
job_path = os.path.join(job_directory, f"job-{index}.json")
with open(job_path, "w") as f:
json.dump(_absolutize_job_paths(test_case.job, test_case.tests_directory), f)
job_paths.append(job_path)
yield job_paths


class Engine(metaclass=abc.ABCMeta):
"""Abstract description of an external process for running tools or workflows."""

Expand Down Expand Up @@ -136,38 +208,16 @@ def _collect_test_results(self, test_cases, test_timeout):

def _run_test_cases(self, test_cases, test_timeout):
runnables = [test_case.runnable for test_case in test_cases]
job_paths = []
tmp_paths = []
output_collectors = []
for test_case in test_cases:
if test_case.job_path is None:
job = test_case.job
with tempfile.NamedTemporaryFile(
dir=test_case.tests_directory,
suffix=".json",
prefix="plnmotmptestjob",
delete=False,
mode="w+",
) as f:
tmp_path = f.name
job_path = tmp_path
tmp_paths.append(tmp_path)
json.dump(job, f)
job_paths.append(job_path)
else:
job_paths.append(test_case.job_path)
output_collectors.append(
lambda run_response, test_case=test_case: test_case.structured_test_data(run_response)
)
try:
run_responses = self._run(runnables, job_paths, output_collectors, test_timeout=test_timeout)
finally:
for tmp_path in tmp_paths:
os.remove(tmp_path)
return run_responses
output_collectors = [
lambda run_response, test_case=test_case: test_case.structured_test_data(run_response)
for test_case in test_cases
]
with materialized_job_paths(test_cases) as job_paths:
return self._run(runnables, job_paths, output_collectors, test_timeout=test_timeout)


__all__ = (
"Engine",
"BaseEngine",
"materialized_job_paths",
)
109 changes: 109 additions & 0 deletions tests/test_engines.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,19 @@
"""Unit tests for engines and runnables."""

import copy
import json
import os

import pytest

from planemo.engine import engine_context
from planemo.engine.galaxy import log_service_logs_on_failure
from planemo.engine.interface import (
_absolutize_job_paths,
materialized_job_paths,
)
from planemo.runnable import (
cases,
for_path,
get_outputs,
RunnableType,
Expand All @@ -21,6 +30,14 @@
A_GALAXY_GA_WORKFLOW = os.path.join(TEST_DATA_DIR, "test_workflow_1.ga")
A_GALAXY_YAML_WORKFLOW = os.path.join(TEST_DATA_DIR, "wf1.gxwf.yml")

# workflows whose test definitions embed the job rather than pointing at a job file
A_COLLECTION_INPUT_WORKFLOW = os.path.join(TEST_DATA_DIR, "wf5-collection-input.gxwf.yml")
A_COMPOSITE_INPUT_WORKFLOW = os.path.join(TEST_DATA_DIR, "wf6-composite-inputs.gxwf.yml")
A_NESTED_COLLECTION_WORKFLOW = os.path.join(TEST_DATA_DIR, "wf8-collection-nested-input.gxwf.yml")
A_REMOTE_INPUT_WORKFLOW = os.path.join(TEST_DATA_DIR, "wf13_tool_shed_repository_gxformat2.yml")
# ... and one that points at tests/data/wf2-job.yml
A_FILE_JOB_WORKFLOW = os.path.join(TEST_DATA_DIR, "wf2.ga")

CAN_HANDLE = {
"galaxy": {
A_CWL_TOOL: True,
Expand Down Expand Up @@ -92,3 +109,95 @@ def test_service_logs_logged_when_no_result_registered():
ctx = _RecordingContext()
log_service_logs_on_failure(ctx, _ConfigWithServiceLogs(), [])
assert len(ctx.messages) == 1


def _materialized_jobs(workflow_path):
"""Run a real workflow's test cases through job materialization, return the job dicts."""
jobs = []
with materialized_job_paths(cases(for_path(workflow_path))) as job_paths:
for job_path in job_paths:
with open(job_path) as f:
jobs.append(json.load(f))
return jobs


def test_inline_jobs_are_not_written_to_the_test_directory():
test_cases = cases(for_path(A_COLLECTION_INPUT_WORKFLOW))
before = sorted(os.listdir(TEST_DATA_DIR))
with materialized_job_paths(test_cases) as job_paths:
assert len(job_paths) == len(test_cases)
for job_path in job_paths:
assert os.path.exists(job_path)
assert os.path.dirname(job_path) != TEST_DATA_DIR
assert sorted(os.listdir(TEST_DATA_DIR)) == before
for job_path in job_paths:
assert not os.path.exists(job_path)


def test_inline_job_directory_removed_when_the_engine_raises():
with pytest.raises(RuntimeError, match="engine failed"):
with materialized_job_paths(cases(for_path(A_COLLECTION_INPUT_WORKFLOW))) as job_paths:
recorded = list(job_paths)
raise RuntimeError("engine failed")
for job_path in recorded:
assert not os.path.exists(job_path)


def test_inline_job_relative_paths_resolve_to_real_test_data():
collection_job, cwl_style_job = _materialized_jobs(A_COLLECTION_INPUT_WORKFLOW)
hello = os.path.join(TEST_DATA_DIR, "hello.txt")
assert collection_job["input1"]["elements"][0]["path"] == hello
assert cwl_style_job["input1"][0]["path"] == hello
assert os.path.exists(hello)


def test_nested_collection_paths_resolve_to_real_test_data():
(job,) = _materialized_jobs(A_NESTED_COLLECTION_WORKFLOW)
pair = job["input1"]["elements"][0]["elements"]
assert [element["path"] for element in pair] == [os.path.join(TEST_DATA_DIR, "hello.txt")] * 2


def test_composite_data_paths_resolve_to_real_test_data():
(job,) = _materialized_jobs(A_COMPOSITE_INPUT_WORKFLOW)
paths = [item["path"] for item in job["input1"]["composite_data"]]
assert paths == [
os.path.join(TEST_DATA_DIR, "Example_Continuous.imzML"),
os.path.join(TEST_DATA_DIR, "Example_Continuous.ibd"),
]
assert all(os.path.exists(path) for path in paths)


def test_remote_inputs_are_left_alone():
(job,) = _materialized_jobs(A_REMOTE_INPUT_WORKFLOW)
pair = job["pe-fastq"]["elements"][0]["elements"]
assert [element["location"] for element in pair] == [
"https://github.com/GoekeLab/bioinformatics-workflows/raw/master/test_data/reads_1.fq.gz",
"https://github.com/GoekeLab/bioinformatics-workflows/raw/master/test_data/reads_2.fq.gz",
]


def test_test_case_job_is_left_unmodified():
test_cases = cases(for_path(A_COLLECTION_INPUT_WORKFLOW))
jobs_before = copy.deepcopy([test_case.job for test_case in test_cases])
with materialized_job_paths(test_cases):
pass
assert [test_case.job for test_case in test_cases] == jobs_before


def test_job_files_are_used_in_place_and_survive():
test_cases = cases(for_path(A_FILE_JOB_WORKFLOW))
job_file = os.path.join(TEST_DATA_DIR, "wf2-job.yml")
with materialized_job_paths(test_cases) as job_paths:
assert job_paths == [job_file]
assert os.path.exists(job_file)


def test_paths_that_are_not_local_test_data_are_left_alone():
job = {
"absolute": {"class": "File", "path": "/already/absolute.txt"},
# planemo test definitions do use URLs under "path" - see tests/data/cat_tool_url_job.json
"url_under_path": {"class": "File", "path": "https://example.org/input.txt"},
"cwl_reference": {"class": "File", "path": "#main/step/out"},
"not_a_file": {"path": "this-is-just-a-parameter-value"},
}
assert _absolutize_job_paths(job, TEST_DATA_DIR) == job
Loading