Merge pull request #17558 from mvdbeek/subworkflow_filtering

Filter out subworkflow invocations
This commit is contained in:
Marius van den Beek
2024-02-28 11:17:37 -06:00
committed by GitHub
8 changed files with 39 additions and 4 deletions
+1
View File
@@ -17888,6 +17888,7 @@ export interface operations {
instance?: boolean | null;
view?: string | null;
step_details?: boolean;
include_nested_invocations?: boolean;
};
/** @description The user ID that will be used to effectively make this API call. Only admins and designated users can make API calls on behalf of other users. */
header?: {
@@ -200,6 +200,8 @@ export default {
const extraParams = this.ownerGrid ? {} : { include_terminal: false };
if (this.storedWorkflowId) {
extraParams["workflow_id"] = this.storedWorkflowId;
} else {
extraParams["include_nested_invocations"] = false;
}
if (this.historyId) {
extraParams["history_id"] = this.historyId;
+13 -1
View File
@@ -69,6 +69,7 @@ from galaxy.model import (
Workflow,
WorkflowInvocation,
WorkflowInvocationStep,
WorkflowInvocationToSubworkflowInvocationAssociation,
)
from galaxy.model.base import (
ensure_object_added_to_session,
@@ -486,6 +487,7 @@ class WorkflowsManager(sharable.SharableModelManager, deletable.DeletableManager
offset=None,
sort_by=None,
sort_desc=None,
include_nested_invocations=True,
) -> Tuple[Query, int]:
"""Get invocations owned by the current user."""
@@ -503,6 +505,16 @@ class WorkflowsManager(sharable.SharableModelManager, deletable.DeletableManager
stmt = stmt.join(WorkflowInvocationStep).where(WorkflowInvocationStep.job_id == job_id)
if not include_terminal:
stmt = stmt.where(WorkflowInvocation.state.in_(WorkflowInvocation.non_terminal_states))
if not include_nested_invocations:
subquery = (
select(WorkflowInvocationToSubworkflowInvocationAssociation.id)
.where(
WorkflowInvocationToSubworkflowInvocationAssociation.subworkflow_invocation_id
== WorkflowInvocation.id
)
.exists()
)
stmt = stmt.where(~subquery)
total_matches = get_count(trans.sa_session, stmt)
@@ -2091,5 +2103,5 @@ def _get_invocation(session, eager, invocation_id):
def get_count(session, statement):
stmt = select(func.count()).select_from(statement)
stmt = select(func.count()).select_from(statement.subquery())
return session.scalar(stmt)
+1
View File
@@ -1452,6 +1452,7 @@ class InvocationIndexQueryPayload(Model):
lt=1000,
)
offset: Optional[int] = Field(default=0, description="Number of invocations to skip")
include_nested_invocations: bool = True
PageSortByEnum = Literal["create_time", "title", "update_time", "username"]
@@ -1293,6 +1293,7 @@ class FastAPIInvocations:
instance: InvocationsInstanceQueryParam = False,
view: SerializationViewQueryParam = None,
step_details: StepDetailQueryParam = False,
include_nested_invocations: bool = True,
trans: ProvidesUserContext = DependsOnTrans,
) -> List[WorkflowInvocationResponse]:
invocation_payload = InvocationIndexPayload(
@@ -1306,6 +1307,7 @@ class FastAPIInvocations:
limit=limit,
offset=offset,
instance=instance,
include_nested_invocations=include_nested_invocations,
)
serialization_params = InvocationSerializationParams(
view=view,
@@ -126,6 +126,7 @@ class InvocationsService(ServiceBase, ConsumesModelStores):
offset=invocation_payload.offset,
sort_by=invocation_payload.sort_by,
sort_desc=invocation_payload.sort_desc,
include_nested_invocations=invocation_payload.include_nested_invocations,
)
invocation_dict = self.serialize_workflow_invocations(invocations, serialization_params)
return invocation_dict, total_matches
+14
View File
@@ -7145,6 +7145,20 @@ input_c:
usage_details_response = self._get(f"workflows/{workflow_id}/usage/{invocation_id}")
self._assert_status_code_is(usage_details_response, 403)
def test_invocation_filtering_exclude_subworkflow(self):
with self.dataset_populator.test_history() as history_id:
self._run_workflow(
WORKFLOW_NESTED_SIMPLE,
test_data="""
outer_input:
value: 1.bed
type: File
""",
history_id=history_id,
)
assert len(self.workflow_populator.history_invocations(history_id)) == 2
assert len(self.workflow_populator.history_invocations(history_id, include_nested_invocations=False)) == 1
def test_workflow_publishing(self):
workflow_id = self.workflow_populator.simple_workflow("dummy")
response = self._show_workflow(workflow_id)
+5 -3
View File
@@ -1701,7 +1701,7 @@ class BaseWorkflowPopulator(BasePopulator):
return wait_on_state(workflow_state, desc="workflow invocation state", timeout=timeout, assert_ok=assert_ok)
def workflow_invocations(self, workflow_id: str) -> List[Dict[str, Any]]:
def workflow_invocations(self, workflow_id: str, include_nested_invocations=True) -> List[Dict[str, Any]]:
response = self._get(f"workflows/{workflow_id}/invocations")
api_asserts.assert_status_code_is(response, 200)
return response.json()
@@ -1711,8 +1711,10 @@ class BaseWorkflowPopulator(BasePopulator):
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})
def history_invocations(self, history_id: str, include_nested_invocations: bool = True) -> List[Dict[str, Any]]:
history_invocations_response = self._get(
"invocations", {"history_id": history_id, "include_nested_invocations": include_nested_invocations}
)
api_asserts.assert_status_code_is(history_invocations_response, 200)
return history_invocations_response.json()