Use wait_for_invocation_and_completion instead of asking for state to be scheduled or completed

This commit is contained in:
mvdbeek
2026-01-21 12:21:11 +01:00
parent c1d75c5bcf
commit f611d7357b
3 changed files with 19 additions and 38 deletions
@@ -73,11 +73,7 @@ class TestWorkflowCompletionEndpoint(ApiTestCase, UsesCeleryTasks):
)
# Wait for monitor to detect completion and transition state
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"
self.workflow_populator.wait_for_invocation_and_completion(summary.invocation_id, timeout=30)
# Verify completion record exists
completion = self.workflow_populator.get_invocation_completion(summary.invocation_id)
assert completion is not None, "Expected completion record to exist"
@@ -153,7 +149,4 @@ class TestWorkflowCompletionEndpoint(ApiTestCase, UsesCeleryTasks):
# Wait for monitor to process (completion transition)
# (export hook will be skipped since no file source is configured)
invocation = self.workflow_populator.wait_for_invocation_and_completion(summary.invocation_id, timeout=30)
# Verify state transitioned to completed
assert invocation["state"] == "completed"
self.workflow_populator.wait_for_invocation_and_completion(summary.invocation_id, timeout=30)
+16 -27
View File
@@ -1742,8 +1742,7 @@ 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
self.workflow_populator.wait_for_invocation_and_completion(invocation_id)
invocation_jobs = self.workflow_populator.get_invocation_jobs(invocation_id)
for job in invocation_jobs:
assert job["state"] == "ok"
@@ -1826,7 +1825,7 @@ steps:
def test_run_workflow_with_url_collection(self):
with self.dataset_populator.test_history() as history_id:
invocation = self._run_multi_data_workflow(history_id)
assert invocation["state"] in ("scheduled", "completed"), invocation
assert invocation["state"] == "completed", invocation
invocation_jobs = self.workflow_populator.get_invocation_jobs(invocation["id"])
assert len(invocation_jobs) == 1
job = invocation_jobs[0]
@@ -1891,7 +1890,9 @@ steps:
invocation_id = self.workflow_populator.invoke_workflow_and_wait(
workflow_id, request=workflow_request, assert_ok=not invalid_hash
).json()["id"]
return self._invocation_details(workflow_id, invocation_id)
if invalid_hash:
return self._invocation_details(workflow_id, invocation_id)
return self.workflow_populator.wait_for_invocation_and_completion(invocation_id)
@skip_without_tool("collection_paired_default")
def test_run_workflow_with_url_paired_collection(self):
@@ -1952,8 +1953,7 @@ 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["inputs"]["0"]["src"] == "hdca"
input_hdca = self.dataset_populator.get_history_collection_details(
history_id=history_id, content_id=invocation["inputs"]["0"]["id"]
@@ -2105,8 +2105,7 @@ steps:
invocation_id = self.workflow_populator.invoke_workflow_and_wait(workflow_id, request=workflow_request).json()[
"id"
]
invocation = self.workflow_populator.wait_for_invocation_and_completion(invocation_id)
assert invocation["state"] == "completed", invocation
self.workflow_populator.wait_for_invocation_and_completion(invocation_id)
@skip_without_tool("collection_creates_pair")
def test_workflow_run_output_collections(self) -> None:
@@ -3430,12 +3429,9 @@ test_data:
# Review the paused steps to allow the workflow to continue.
self.__review_paused_steps(uploaded_workflow_id, invocation_id, order_index=1, action=True)
# Wait for the workflow to finish scheduling and ensure both the invocation
# Wait for the workflow to finish and ensure both the invocation
# and the history are in valid states.
invocation_scheduled = self._wait_for_invocation_state(
uploaded_workflow_id, invocation_id, ("scheduled", "completed")
)
assert invocation_scheduled, "Workflow state is not scheduled or completed..."
self.workflow_populator.wait_for_invocation_and_completion(invocation_id)
self.dataset_populator.wait_for_history(history_id, assert_ok=True)
content = self.dataset_populator.get_history_dataset_content(history_id)
@@ -5102,12 +5098,9 @@ input1:
# Review the paused steps to allow the workflow to continue.
self.__review_paused_steps(uploaded_workflow_id, invocation_id, order_index=2, action=True)
# Wait for the workflow to finish scheduling and ensure both the invocation
# Wait for the workflow to finish and ensure both the invocation
# and the history are in valid states.
invocation_scheduled = self._wait_for_invocation_state(
uploaded_workflow_id, invocation_id, ("scheduled", "completed")
)
assert invocation_scheduled, "Workflow state is not scheduled or completed..."
self.workflow_populator.wait_for_invocation_and_completion(invocation_id)
self.dataset_populator.wait_for_history(history_id, assert_ok=True)
@skip_without_tool("cat")
@@ -5165,9 +5158,7 @@ input1:
self._assert_invocation_non_terminal(uploaded_workflow_id, invocation_id)
self.__review_paused_steps(uploaded_workflow_id, invocation_id, order_index=4, action=True)
self.workflow_populator.wait_for_invocation_and_jobs(history_id, uploaded_workflow_id, invocation_id)
invocation = self._invocation_details(uploaded_workflow_id, invocation_id)
assert invocation["state"] in ("scheduled", "completed")
self.workflow_populator.wait_for_invocation_and_completion(invocation_id)
assert "reviewed\n1\nreviewed\n4\n" == self.dataset_populator.get_history_dataset_content(history_id)
@skip_without_tool("cat")
@@ -5245,9 +5236,8 @@ test_data:
wait=False,
)
# wait_for_invocation just waits until scheduling complete, not jobs or subworkflow invocations
self.workflow_populator.wait_for_invocation("null", summary.invocation_id, assert_ok=True)
self.workflow_populator.wait_for_invocation(summary.history_id, summary.invocation_id, assert_ok=True)
invocation_before_cancellation = self.workflow_populator.get_invocation(summary.invocation_id)
assert invocation_before_cancellation["state"] in ("scheduled", "completed")
subworkflow_invocation_id = invocation_before_cancellation["steps"][2]["subworkflow_invocation_id"]
self.workflow_populator.cancel_invocation(summary.invocation_id)
self.workflow_populator.wait_for_invocation_and_jobs(
@@ -5318,8 +5308,7 @@ outputs:
assert_ok=False,
wait=True,
)
invocation = self.workflow_populator.get_invocation(summary.invocation_id)
assert invocation["state"] in ("scheduled", "completed")
invocation = self.workflow_populator.wait_for_invocation_and_completion(summary.invocation_id)
assert len(invocation["messages"]) == 1
message = invocation["messages"][0]
assert message["reason"] == "workflow_output_not_found"
@@ -5984,8 +5973,8 @@ test_data:
summary = self._run_workflow(workflow, history_id=history_id, wait=True, assert_ok=True)
# Verify parent workflow executed successfully
self.workflow_populator.wait_for_invocation_and_completion(summary.invocation_id)
parent_invocation = self.workflow_populator.get_invocation(summary.invocation_id, step_details=True)
assert parent_invocation["state"] in ("scheduled", "completed")
# Find the subworkflow step and get its invocation
subworkflow_step = None
@@ -6002,7 +5991,7 @@ test_data:
# The subworkflow should have succeeded
assert (
subworkflow_invocation["state"] == "scheduled"
subworkflow_invocation["state"] == "completed"
), f"Expected subworkflow to succeed, got state: {subworkflow_invocation['state']}"
# Should not have error messages (previously failed with "when_not_boolean")
@@ -92,8 +92,7 @@ class TestWorkflowCompletionExportHook(IntegrationTestCase, UsesCeleryTasks):
)
# 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"
self.workflow_populator.wait_for_invocation_and_completion(summary.invocation_id, timeout=60)
# Wait for the export file to appear (hook execution is async via Celery)
export_path = os.path.join(self.root_dir, export_filename)