Reuse existing workflow fixtures and abstractions

This commit is contained in:
mvdbeek
2026-01-21 12:21:10 +01:00
parent 307c8e351e
commit c80435b349
6 changed files with 98 additions and 229 deletions
+2 -17
View File
@@ -12,24 +12,9 @@ from galaxy_test.base.populators import (
DatasetPopulator,
WorkflowPopulator,
)
from galaxy_test.base.workflow_fixtures import WORKFLOW_SIMPLE_CAT_TWICE
from ._framework import ApiTestCase
# Simple workflow for testing
SIMPLE_WORKFLOW = """
class: GalaxyWorkflow
name: Simple Export Test Workflow
inputs:
input1: data
outputs:
output1:
outputSource: cat/out_file1
steps:
cat:
tool_id: cat1
in:
input1: input1
"""
class TestExportsEndpoint(ApiTestCase, UsesCeleryTasks):
"""Tests for the /api/exports endpoint."""
@@ -97,7 +82,7 @@ class TestExportsEndpoint(ApiTestCase, UsesCeleryTasks):
"""Test that exports endpoint returns invocation exports via prepare_store_download."""
with self.dataset_populator.test_history() as history_id:
summary = self.workflow_populator.run_workflow(
SIMPLE_WORKFLOW,
WORKFLOW_SIMPLE_CAT_TWICE,
test_data={"input1": "hello world"},
history_id=history_id,
wait=True,
+32 -132
View File
@@ -6,52 +6,14 @@ These tests verify:
- Completion records are created with accurate job state summaries
"""
from typing import Optional
from galaxy_test.base.api import UsesCeleryTasks
from galaxy_test.base.populators import (
DatasetPopulator,
wait_on,
WorkflowPopulator,
)
from galaxy_test.base.workflow_fixtures import WORKFLOW_SIMPLE_CAT_TWICE
from ._framework import ApiTestCase
# Simple workflow for testing completion
SIMPLE_WORKFLOW = """
class: GalaxyWorkflow
name: Simple Completion Test Workflow
inputs:
input1: data
outputs:
output1:
outputSource: cat/out_file1
steps:
cat:
tool_id: cat1
in:
input1: input1
"""
# Workflow with multiple steps
MULTI_STEP_WORKFLOW = """
class: GalaxyWorkflow
name: Multi Step Workflow
inputs:
input1: data
outputs:
output1:
outputSource: cat2/out_file1
steps:
cat1:
tool_id: cat1
in:
input1: input1
cat2:
tool_id: cat1
in:
input1: cat1/out_file1
"""
class TestWorkflowCompletionEndpoint(ApiTestCase, UsesCeleryTasks):
"""Tests for the workflow completion API endpoint."""
@@ -68,7 +30,7 @@ class TestWorkflowCompletionEndpoint(ApiTestCase, UsesCeleryTasks):
"""Test that the completion endpoint exists and returns proper response."""
with self.dataset_populator.test_history() as history_id:
summary = self.workflow_populator.run_workflow(
SIMPLE_WORKFLOW,
WORKFLOW_SIMPLE_CAT_TWICE,
test_data={"input1": "hello world"},
history_id=history_id,
wait=True,
@@ -76,48 +38,34 @@ class TestWorkflowCompletionEndpoint(ApiTestCase, UsesCeleryTasks):
)
# The endpoint should exist and return 200
response = self._get(f"invocations/{summary.invocation_id}/completion")
assert response.status_code == 200
completion = self.workflow_populator.get_invocation_completion(summary.invocation_id)
# Completion may or may not exist yet depending on timing
if completion is not None:
assert "completion_time" in completion
assert "all_jobs_ok" in completion
def test_completion_endpoint_returns_none_before_completion(self):
"""Test that completion endpoint returns None for incomplete invocation."""
with self.dataset_populator.test_history() as history_id:
summary = self.workflow_populator.run_workflow(
SIMPLE_WORKFLOW,
WORKFLOW_SIMPLE_CAT_TWICE,
test_data={"input1": "hello world"},
history_id=history_id,
wait=False, # Don't wait for completion
)
# Check immediately - might be None or already complete depending on timing
response = self._get(f"invocations/{summary.invocation_id}/completion")
assert response.status_code == 200
completion = self.workflow_populator.get_invocation_completion(summary.invocation_id)
# Either None (not complete) or has completion data
result = response.json()
if result is not None:
assert "completion_time" in result
assert "all_jobs_ok" in result
def test_workflow_reaches_scheduled_state(self):
"""Test that workflow reaches SCHEDULED state after jobs complete."""
with self.dataset_populator.test_history() as history_id:
summary = self.workflow_populator.run_workflow(
SIMPLE_WORKFLOW,
test_data={"input1": "hello world"},
history_id=history_id,
wait=True,
assert_ok=True,
)
# After jobs complete, invocation should be at least SCHEDULED
invocation = self._get_invocation(summary.invocation_id)
assert invocation["state"] in ["scheduled", "completed"]
if completion is not None:
assert "completion_time" in completion
assert "all_jobs_ok" in completion
def test_monitor_transitions_to_completed(self):
"""Test that background monitor transitions invocation to COMPLETED state."""
with self.dataset_populator.test_history() as history_id:
summary = self.workflow_populator.run_workflow(
SIMPLE_WORKFLOW,
WORKFLOW_SIMPLE_CAT_TWICE,
test_data={"input1": "hello world"},
history_id=history_id,
wait=True,
@@ -125,72 +73,46 @@ class TestWorkflowCompletionEndpoint(ApiTestCase, UsesCeleryTasks):
)
# Wait for monitor to detect completion and transition state
def check_completed():
invocation = self._get_invocation(summary.invocation_id)
return invocation if invocation["state"] == "completed" else None
invocation = wait_on(check_completed, "invocation to reach completed state", timeout=30)
invocation = self.workflow_populator.wait_for_invocation_and_completion(summary.invocation_id, timeout=30)
# Verify the monitor transitioned to completed state
assert invocation["state"] == "completed"
# Verify completion record exists
response = self._get(f"invocations/{summary.invocation_id}/completion")
assert response.status_code == 200
completion = response.json()
completion = self.workflow_populator.get_invocation_completion(summary.invocation_id)
assert completion is not None, "Expected completion record to exist"
assert completion["all_jobs_ok"] is True
def test_multi_step_workflow_jobs_complete(self):
"""Test that a multi-step workflow jobs all complete."""
with self.dataset_populator.test_history() as history_id:
summary = self.workflow_populator.run_workflow(
MULTI_STEP_WORKFLOW,
test_data={"input1": "hello world"},
history_id=history_id,
wait=True,
assert_ok=True,
)
# Verify workflow reached scheduled state
invocation = self._get_invocation(summary.invocation_id)
assert invocation["state"] in ["scheduled", "completed"]
# Verify we have the expected number of tool steps
steps = invocation["steps"]
tool_steps = [s for s in steps if s.get("job_id")]
assert len(tool_steps) >= 2
def test_completion_response_structure(self):
"""Test that completion response has correct structure when present."""
"""Test that completion response has correct structure."""
with self.dataset_populator.test_history() as history_id:
summary = self.workflow_populator.run_workflow(
SIMPLE_WORKFLOW,
WORKFLOW_SIMPLE_CAT_TWICE,
test_data={"input1": "hello world"},
history_id=history_id,
wait=True,
assert_ok=True,
)
response = self._get(f"invocations/{summary.invocation_id}/completion")
assert response.status_code == 200
result = response.json()
# Wait for completion
self.workflow_populator.wait_for_invocation_and_completion(summary.invocation_id, timeout=30)
# If completion exists, verify structure
if result is not None:
assert "completion_time" in result
assert "job_state_summary" in result
assert "all_jobs_ok" in result
assert "hooks_executed" in result
assert isinstance(result["job_state_summary"], dict)
assert isinstance(result["hooks_executed"], list)
# Verify completion structure
completion = self.workflow_populator.get_invocation_completion(summary.invocation_id)
assert completion is not None, "Expected completion record to exist after waiting"
assert "completion_time" in completion
assert "job_state_summary" in completion
assert "all_jobs_ok" in completion
assert "hooks_executed" in completion
assert isinstance(completion["job_state_summary"], dict)
assert isinstance(completion["hooks_executed"], list)
def test_on_complete_actions_stored(self):
"""Test that on_complete actions are stored when invoking workflow."""
with self.dataset_populator.test_history() as history_id:
# Run workflow with on_complete actions (new format with action objects)
summary = self.workflow_populator.run_workflow(
SIMPLE_WORKFLOW,
WORKFLOW_SIMPLE_CAT_TWICE,
test_data={"input1": "hello world"},
history_id=history_id,
wait=True,
@@ -199,11 +121,7 @@ class TestWorkflowCompletionEndpoint(ApiTestCase, UsesCeleryTasks):
)
# Wait for monitor to process (completion transition)
def check_completed():
invocation = self._get_invocation(summary.invocation_id)
return invocation if invocation["state"] == "completed" else None
invocation = wait_on(check_completed, "invocation to reach completed state", timeout=30)
invocation = self.workflow_populator.wait_for_invocation_and_completion(summary.invocation_id, timeout=30)
# Verify state transitioned to completed
assert invocation["state"] == "completed"
@@ -215,7 +133,7 @@ class TestWorkflowCompletionEndpoint(ApiTestCase, UsesCeleryTasks):
# Note: The export won't actually work without a configured file source,
# but the API should accept the configuration and the workflow should complete
summary = self.workflow_populator.run_workflow(
SIMPLE_WORKFLOW,
WORKFLOW_SIMPLE_CAT_TWICE,
test_data={"input1": "hello world"},
history_id=history_id,
wait=True,
@@ -235,25 +153,7 @@ class TestWorkflowCompletionEndpoint(ApiTestCase, UsesCeleryTasks):
# Wait for monitor to process (completion transition)
# (export hook will be skipped since no file source is configured)
def check_completed():
invocation = self._get_invocation(summary.invocation_id)
return invocation if invocation["state"] == "completed" else None
invocation = wait_on(check_completed, "invocation to reach completed state", timeout=30)
invocation = self.workflow_populator.wait_for_invocation_and_completion(summary.invocation_id, timeout=30)
# Verify state transitioned to completed
assert invocation["state"] == "completed"
# Helper methods
def _get_invocation(self, invocation_id: str) -> dict:
"""Get invocation details."""
response = self._get(f"invocations/{invocation_id}")
response.raise_for_status()
return response.json()
def _get_completion_details(self, invocation_id: str) -> Optional[dict]:
"""Get completion details for an invocation."""
response = self._get(f"invocations/{invocation_id}/completion")
response.raise_for_status()
result = response.json()
return result if result else None
+2 -2
View File
@@ -2105,8 +2105,8 @@ steps:
invocation_id = self.workflow_populator.invoke_workflow_and_wait(workflow_id, request=workflow_request).json()[
"id"
]
invocation = self._invocation_details(workflow_id, invocation_id)
assert invocation["state"] in ("scheduled", "completed"), invocation
invocation = self.workflow_populator.wait_for_invocation_and_completion(invocation_id)
assert invocation["state"] == "completed", invocation
@skip_without_tool("collection_creates_pair")
def test_workflow_run_output_collections(self) -> None:
+43
View File
@@ -2662,6 +2662,12 @@ class BaseWorkflowPopulator(BasePopulator):
def wait_for_invocation_and_jobs(
self, history_id: str, workflow_id: str, invocation_id: str, assert_ok: bool = True
) -> None:
"""Wait for invocation to be scheduled and all jobs to complete.
.. deprecated::
Use :meth:`wait_for_invocation_and_completion` for new tests,
which waits for the invocation to reach the 'completed' state.
"""
state = self.wait_for_invocation(workflow_id, invocation_id, assert_ok=assert_ok)
if assert_ok:
assert state in ("scheduled", "completed"), state
@@ -2669,6 +2675,43 @@ class BaseWorkflowPopulator(BasePopulator):
self.dataset_populator.wait_for_history_jobs(history_id, assert_ok=assert_ok)
time.sleep(0.5)
def get_invocation_completion(self, invocation_id: str) -> Optional[dict[str, Any]]:
"""Get completion record for an invocation.
Returns the completion record if it exists, or None if the invocation
has not yet completed.
"""
response = self._get(f"invocations/{invocation_id}/completion")
api_asserts.assert_status_code_is(response, 200)
result = response.json()
return result if result else None
def wait_for_invocation_and_completion(
self,
invocation_id: str,
timeout: timeout_type = DEFAULT_TIMEOUT,
assert_ok: bool = True,
) -> dict[str, Any]:
"""Wait for invocation to reach 'completed' state and return invocation details.
This waits for the workflow completion monitor to process the invocation
and transition it to the 'completed' state. Use this instead of
wait_for_invocation_and_jobs for new tests.
Returns the invocation dict with state='completed'.
"""
def check_completed():
invocation = self.get_invocation(invocation_id)
if invocation["state"] == "completed":
return invocation
elif assert_ok and invocation["state"] in ("failed", "cancelled"):
raise AssertionError(f"Invocation reached terminal state: {invocation['state']}")
return None
invocation = wait_on(check_completed, "invocation to reach completed state", timeout=timeout)
return invocation
def index(
self,
show_shared: Optional[bool] = None,
+1
View File
@@ -239,6 +239,7 @@ def setup_galaxy_config(
job_handler_monitor_sleep=0.2,
job_runner_monitor_sleep=0.2,
workflow_monitor_sleep=0.2,
workflow_completion_monitor_sleep=1.0,
)
if default_shed_tool_data_table_config:
config["shed_tool_data_table_config"] = default_shed_tool_data_table_config
@@ -6,8 +6,8 @@ workflow invocations to a configured file source.
"""
import os
import time
import zipfile
from tempfile import mkdtemp
from galaxy_test.base.api import UsesCeleryTasks
from galaxy_test.base.populators import (
@@ -15,6 +15,7 @@ from galaxy_test.base.populators import (
wait_on,
WorkflowPopulator,
)
from galaxy_test.base.workflow_fixtures import WORKFLOW_SIMPLE_CAT_TWICE
from galaxy_test.driver.integration_util import IntegrationTestCase
@@ -43,9 +44,8 @@ class TestWorkflowCompletionExportHook(IntegrationTestCase, UsesCeleryTasks):
super().handle_galaxy_config_kwds(config)
UsesCeleryTasks.handle_galaxy_config_kwds(config)
# Set up temp directory for file source
temp_dir = os.path.realpath(mkdtemp())
cls._test_driver.temp_directories.append(temp_dir)
# Set up temp directory for file source using test driver's mkdtemp
temp_dir = os.path.realpath(cls._test_driver.mkdtemp())
cls.root_dir = os.path.join(temp_dir, "root")
os.makedirs(cls.root_dir, exist_ok=True)
@@ -61,9 +61,6 @@ class TestWorkflowCompletionExportHook(IntegrationTestCase, UsesCeleryTasks):
config["library_import_dir"] = None
config["user_library_import_dir"] = None
# Speed up monitor for faster tests
config["workflow_completion_monitor_sleep"] = 1.0
def setUp(self):
super().setUp()
self.dataset_populator = DatasetPopulator(self.galaxy_interactor)
@@ -71,31 +68,13 @@ class TestWorkflowCompletionExportHook(IntegrationTestCase, UsesCeleryTasks):
def test_export_to_file_source_on_completion(self):
"""Test that export_to_file_source hook exports invocation when workflow completes."""
# Define a simple workflow
workflow = """
class: GalaxyWorkflow
name: Completion Export Test Workflow
inputs:
input_data:
type: data
steps:
cat_step:
tool_id: cat
in:
input1: input_data
outputs:
output_data:
outputSource: cat_step/out_file1
"""
# Define the export filename
export_filename = "test_export.rocrate.zip"
target_uri = f"gxfiles://completion_export_test/{export_filename}"
with self.dataset_populator.test_history() as history_id:
# Run workflow with on_complete export action
summary = self.workflow_populator.run_workflow(
workflow,
test_data={"input_data": "hello world"},
WORKFLOW_SIMPLE_CAT_TWICE,
test_data={"input1": "hello world"},
history_id=history_id,
wait=True,
assert_ok=True,
@@ -112,12 +91,8 @@ outputs:
},
)
# Wait for invocation to reach completed state (monitor has processed it)
def check_completed():
invocation = self._get_invocation(summary.invocation_id)
return invocation if invocation["state"] == "completed" else None
invocation = wait_on(check_completed, "invocation to reach completed state", timeout=60)
# Wait for invocation to reach completed state
invocation = self.workflow_populator.wait_for_invocation_and_completion(summary.invocation_id, timeout=60)
assert invocation["state"] == "completed"
# Wait for the export file to appear (hook execution is async via Celery)
@@ -135,7 +110,6 @@ outputs:
# Verify the zip contains expected rocrate metadata
with zipfile.ZipFile(export_path, "r") as zf:
names = zf.namelist()
# RO-Crate should have ro-crate-metadata.json
assert "ro-crate-metadata.json" in names, f"Missing ro-crate-metadata.json in {names}"
def test_export_with_multiple_outputs(self):
@@ -144,17 +118,16 @@ outputs:
class: GalaxyWorkflow
name: Multi-Output Completion Export Test
inputs:
input_data:
type: data
input1: data
steps:
cat1:
tool_id: cat
tool_id: cat1
in:
input1: input_data
input1: input1
cat2:
tool_id: cat
tool_id: cat1
in:
input1: input_data
input1: input1
outputs:
output1:
outputSource: cat1/out_file1
@@ -167,7 +140,7 @@ outputs:
with self.dataset_populator.test_history() as history_id:
summary = self.workflow_populator.run_workflow(
workflow,
test_data={"input_data": "test data content"},
test_data={"input1": "test data content"},
history_id=history_id,
wait=True,
assert_ok=True,
@@ -185,11 +158,7 @@ outputs:
)
# Wait for completion
def check_completed():
invocation = self._get_invocation(summary.invocation_id)
return invocation if invocation["state"] == "completed" else None
wait_on(check_completed, "invocation to complete", timeout=60)
self.workflow_populator.wait_for_invocation_and_completion(summary.invocation_id, timeout=60)
# Wait for export
export_path = os.path.join(self.root_dir, export_filename)
@@ -210,52 +179,23 @@ outputs:
def test_no_export_without_on_complete(self):
"""Test that no export happens when on_complete is not specified."""
workflow = """
class: GalaxyWorkflow
name: No Export Test
inputs:
input_data:
type: data
steps:
cat_step:
tool_id: cat
in:
input1: input_data
outputs:
output_data:
outputSource: cat_step/out_file1
"""
# Use a unique filename that shouldn't be created
export_filename = "should_not_exist.rocrate.zip"
export_path = os.path.join(self.root_dir, export_filename)
with self.dataset_populator.test_history() as history_id:
summary = self.workflow_populator.run_workflow(
workflow,
test_data={"input_data": "hello world"},
WORKFLOW_SIMPLE_CAT_TWICE,
test_data={"input1": "hello world"},
history_id=history_id,
wait=True,
assert_ok=True,
# No on_complete specified
)
# Wait for completion
def check_completed():
invocation = self._get_invocation(summary.invocation_id)
return invocation if invocation["state"] == "completed" else None
wait_on(check_completed, "invocation to complete", timeout=60)
self.workflow_populator.wait_for_invocation_and_completion(summary.invocation_id, timeout=60)
# Give a moment for any erroneous export to happen
import time
time.sleep(2)
# Verify no export file was created
assert not os.path.exists(export_path), f"Export file should not exist: {export_path}"
def _get_invocation(self, invocation_id: str) -> dict:
"""Get invocation details."""
response = self._get(f"invocations/{invocation_id}")
response.raise_for_status()
return response.json()