diff --git a/lib/galaxy_test/api/test_exports.py b/lib/galaxy_test/api/test_exports.py index 88072e1517b..f120abe1b74 100644 --- a/lib/galaxy_test/api/test_exports.py +++ b/lib/galaxy_test/api/test_exports.py @@ -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, diff --git a/lib/galaxy_test/api/test_workflow_completion.py b/lib/galaxy_test/api/test_workflow_completion.py index 22457b3b8a8..7e4a40b5ed3 100644 --- a/lib/galaxy_test/api/test_workflow_completion.py +++ b/lib/galaxy_test/api/test_workflow_completion.py @@ -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 diff --git a/lib/galaxy_test/api/test_workflows.py b/lib/galaxy_test/api/test_workflows.py index 353b2ba5a60..8092dacf803 100644 --- a/lib/galaxy_test/api/test_workflows.py +++ b/lib/galaxy_test/api/test_workflows.py @@ -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: diff --git a/lib/galaxy_test/base/populators.py b/lib/galaxy_test/base/populators.py index 0f3ccb3878c..6137aaa6e57 100644 --- a/lib/galaxy_test/base/populators.py +++ b/lib/galaxy_test/base/populators.py @@ -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, diff --git a/lib/galaxy_test/driver/driver_util.py b/lib/galaxy_test/driver/driver_util.py index 37eaa5f73a6..5e368bf0e1a 100644 --- a/lib/galaxy_test/driver/driver_util.py +++ b/lib/galaxy_test/driver/driver_util.py @@ -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 diff --git a/test/integration/test_workflow_completion_hooks.py b/test/integration/test_workflow_completion_hooks.py index 6b81b1c1c2b..b9c85b7d6ef 100644 --- a/test/integration/test_workflow_completion_hooks.py +++ b/test/integration/test_workflow_completion_hooks.py @@ -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()