Delete non-terminal jobs and subworkflow invocation when cancelling an invocation

This commit is contained in:
mvdbeek
2023-11-13 15:45:26 +01:00
parent c85298b48e
commit 64d608da3a
3 changed files with 85 additions and 4 deletions
+16 -4
View File
@@ -421,7 +421,7 @@ class WorkflowsManager(sharable.SharableModelManager, deletable.DeletableManager
target_format=target_format,
)
def cancel_invocation(self, trans, decoded_invocation_id):
def cancel_invocation(self, trans, decoded_invocation_id: int):
workflow_invocation = self.get_invocation(trans, decoded_invocation_id)
cancelled = workflow_invocation.cancel()
@@ -430,9 +430,21 @@ class WorkflowsManager(sharable.SharableModelManager, deletable.DeletableManager
trans.sa_session.add(workflow_invocation)
with transaction(trans.sa_session):
trans.sa_session.commit()
else:
# TODO: More specific exception?
raise exceptions.MessageException("Cannot cancel an inactive workflow invocation.")
for step in workflow_invocation.steps:
for job in step.jobs:
job.mark_deleted()
trans.sa_session.add(job)
if step.implicit_collection_jobs:
for icjja in step.implicit_collection_jobs.jobs:
icjja.job.mark_deleted()
trans.sa_session.add(icjja.job)
with transaction(trans.sa_session):
trans.sa_session.commit()
for invocation in workflow_invocation.subworkflow_invocations:
self.cancel_invocation(trans, invocation.subworkflow_invocation_id)
return workflow_invocation
+57
View File
@@ -4040,6 +4040,63 @@ input1:
message = invocation["messages"][0]
assert message["reason"] == "user_request"
@skip_without_tool("collection_creates_dynamic_nested")
def test_cancel_workflow_invocation_deletes_jobs(self):
with self.dataset_populator.test_history() as history_id:
summary = self._run_workflow(
"""
class: GalaxyWorkflow
inputs:
list_input:
type: collection
collection_type: list
steps:
first_step:
tool_id: cat_data_and_sleep
in:
input1: list_input
state:
sleep_time: 60
subworkflow_step:
run:
class: GalaxyWorkflow
inputs:
list_input:
type: collection
collection_type: list
steps:
intermediate_step:
tool_id: identifier_multiple
in:
input1: list_input
subworkflow:
in:
list_input: first_step/out_file1
test_data:
list_input:
collection_type: list
elements:
- identifier: 1
content: A
- identifier: 2
content: B
""",
history_id=history_id,
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)
invocation_before_cancellation = self.workflow_populator.get_invocation(summary.invocation_id)
assert invocation_before_cancellation["state"] == "scheduled"
subworkflow_invocation_id = invocation_before_cancellation["steps"][2]["subworkflow_invocation_id"]
self.workflow_populator.cancel_invocation(summary.invocation_id)
invocation_jobs = self.workflow_populator.get_invocation_jobs(summary.invocation_id)
for job in invocation_jobs:
assert job["state"] == "deleted" or job["state"] == "deleted_new"
subworkflow_invocation_jobs = self.workflow_populator.get_invocation_jobs(subworkflow_invocation_id)
for job in subworkflow_invocation_jobs:
assert job["state"] == "deleted" or job["state"] == "deleted_new"
def test_workflow_failed_output_not_found(self, history_id):
summary = self._run_workflow(
"""
+12
View File
@@ -1697,6 +1697,11 @@ class BaseWorkflowPopulator(BasePopulator):
api_asserts.assert_status_code_is(response, 200)
return response.json()
def cancel_invocation(self, invocation_id: str):
response = self._delete(f"invocations/{invocation_id}")
api_asserts.assert_status_code_is(response, 200)
return response.json()
def history_invocations(self, history_id: str) -> List[Dict[str, Any]]:
history_invocations_response = self._get("invocations", {"history_id": history_id})
api_asserts.assert_status_code_is(history_invocations_response, 200)
@@ -2079,6 +2084,13 @@ class BaseWorkflowPopulator(BasePopulator):
return workflow_request, history_id, workflow_id
def get_invocation_jobs(self, invocation_id: str) -> List[Dict[str, Any]]:
jobs_response = self._get("jobs", data={"invocation_id": invocation_id})
api_asserts.assert_status_code_is(jobs_response, 200)
jobs = jobs_response.json()
assert isinstance(jobs, list)
return jobs
def wait_for_invocation_and_jobs(
self, history_id: str, workflow_id: str, invocation_id: str, assert_ok: bool = True
) -> None: