Merge pull request #13767 from jmchilton/jobs_index_fastapi

[WIP] Migrate jobs index to FastAPI controller.
This commit is contained in:
Sergey Golitsynskiy
2022-04-20 08:33:13 -04:00
committed by GitHub
+152 -87
View File
@@ -5,8 +5,15 @@ API operations on a jobs.
"""
import logging
import typing
from enum import Enum
from typing import (
Any,
Dict,
List,
Optional,
)
from fastapi import Query
from sqlalchemy import or_
from galaxy import (
@@ -28,6 +35,7 @@ from galaxy.managers.jobs import (
view_show_job,
)
from galaxy.schema.fields import EncodedDatabaseIdField
from galaxy.util import listify
from galaxy.web import (
expose_api,
expose_api_anonymous,
@@ -47,6 +55,99 @@ log = logging.getLogger(__name__)
router = Router(tags=["jobs"])
class JobIndexViewEnum(str, Enum):
collection = "collection"
admin_job_list = "admin_job_list"
class JobIndexSortByEnum(str, Enum):
create_time = "create_time"
update_time = "update_time"
StateQueryParam: Optional[str] = Query(
default=None,
title="States",
description="Comma-separated list of states to filter job query on. If unspecified, jobs of any state may be returned.",
)
UserDetailsQueryParam: bool = Query(
default=False,
title="Include user details",
description="If true, and requestor is an admin, will return external job id and user email.",
)
UserIdQueryParam: Optional[str] = Query(
default=None,
title="User ID",
description="an encoded user id to restrict query to, must be own id if not admin user",
)
ViewQueryParam: JobIndexViewEnum = Query(
default="collection",
title="View",
description="Determines columns to return. Defaults to 'collection'.",
)
ToolIdQueryParam: Optional[str] = Query(
default=None,
title="Tool ID(s)",
description="Limit listing of jobs to those that match one of the included tool_ids. If none, all are returned",
)
ToolIdLikeQueryParam: Optional[str] = Query(
default=None,
title="Tool ID Pattern(s)",
description="Limit listing of jobs to those that match one of the included tool ID sql-like patterns. If none, all are returned",
)
DateRangeMinQueryParam: Optional[str] = Query(
default=None,
title="Date Range Minimum",
description="Limit listing of jobs to those that are updated after specified date (e.g. '2014-01-01')",
)
DateRangeMaxQueryParam: Optional[str] = Query(
default=None,
title="Date Range Maximum",
description="Limit listing of jobs to those that are updated before specified date (e.g. '2014-01-01')",
)
HistoryIdQueryParam: Optional[str] = Query(
default=None,
title="History ID",
description="Limit listing of jobs to those that match the history_id. If none, jobs from any history may be returned.",
)
WorkflowIdQueryParam: Optional[str] = Query(
default=None,
title="Workflow ID",
description="Limit listing of jobs to those that match the specified workflow ID. If none, jobs from any workflow (or from no workflows) may be returned.",
)
InvocationIdQueryParam: Optional[str] = Query(
default=None,
title="Invocation ID",
description="Limit listing of jobs to those that match the specified workflow invocation ID. If none, jobs from any workflow invocation (or from no workflows) may be returned.",
)
SortByQueryParam: JobIndexSortByEnum = Query(
default=JobIndexSortByEnum.update_time,
title="Sort By",
description="Sort results by specified field.",
)
LimitQueryParam: int = Query(default=500, title="Limit", description="Maximum number of jobs to return.")
OffsetQueryParam: int = Query(
default=0,
title="Offset",
description="Return jobs starting from this specified position. For example, if ``limit`` is set to 100 and ``offset`` to 200, jobs 200-299 will be returned.",
)
@router.cbv
class FastAPIJobs:
job_manager: JobManager = depends(JobManager)
@@ -58,8 +159,8 @@ class FastAPIJobs:
self,
id: EncodedDatabaseIdField,
trans: ProvidesUserContext = DependsOnTrans,
full: typing.Optional[bool] = False,
) -> typing.Dict:
full: Optional[bool] = False,
) -> Dict[str, Any]:
"""
Return dictionary containing description of job data
@@ -71,77 +172,44 @@ class FastAPIJobs:
job = self.job_manager.get_accessible_job(trans, id)
return view_show_job(trans, job, bool(full))
@router.get("/api/jobs")
def index(
self,
trans: ProvidesUserContext = DependsOnTrans,
state: Optional[str] = StateQueryParam,
user_details: bool = UserDetailsQueryParam,
user_id: Optional[str] = UserIdQueryParam,
view: JobIndexViewEnum = ViewQueryParam,
tool_id: Optional[str] = ToolIdQueryParam,
tool_id_like: Optional[str] = ToolIdLikeQueryParam,
date_range_min: Optional[str] = DateRangeMinQueryParam,
date_range_max: Optional[str] = DateRangeMaxQueryParam,
history_id: Optional[str] = HistoryIdQueryParam,
workflow_id: Optional[str] = WorkflowIdQueryParam,
invocation_id: Optional[str] = InvocationIdQueryParam,
order_by: JobIndexSortByEnum = SortByQueryParam,
limit: int = LimitQueryParam,
offset: int = OffsetQueryParam,
) -> List[Dict[str, Any]]:
security = trans.security
class JobController(BaseGalaxyAPIController, UsesVisualizationMixin):
job_manager = depends(JobManager)
job_search = depends(JobSearch)
hda_manager = depends(hdas.HDAManager)
def optional_list(input: Optional[str]) -> Optional[List[str]]:
if input is None:
return None
else:
return listify(input)
@expose_api
def index(self, trans: ProvidesUserContext, limit=500, offset=0, **kwd):
"""
GET /api/jobs
return jobs for current user
if user is admin and user_details is True, then
return jobs for all galaxy users based on filtering - this is an extended service
:type state: string or list
:param state: limit listing of jobs to those that match one of the included states. If none, all are returned.
:type tool_id: string or list
:param tool_id: limit listing of jobs to those that match one of the included tool_ids. If none, all are returned.
:type user_details: boolean
:param user_details: if true, and requestor is an admin, will return external job id and user email.
:type user_id: str
:param user_id: an encoded user id to restrict query to, must be own id if not admin user
:type limit: int
:param limit: Maximum number of jobs to return.
:type offset: int
:param offset: Return jobs starting from this specified position.
For example, if ``limit`` is set to 100 and ``offset`` to 200,
jobs 200-299 will be returned.
:type date_range_min: string '2014-01-01'
:param date_range_min: limit the listing of jobs to those updated on or after requested date
:type date_range_max: string '2014-12-31'
:param date_range_max: limit the listing of jobs to those updated on or before requested date
:type history_id: string
:param history_id: limit listing of jobs to those that match the history_id. If none, all are returned.
:type workflow_id: string
:param workflow_id: limit listing of jobs to those that match the workflow_id. If none, all are returned.
:type invocation_id: string
:param invocation_id: limit listing of jobs to those that match the invocation_id. If none, all are returned.
:type view: string
:param view: Determines columns to return. Defaults to 'collection'.
:rtype: list
:returns: list of dictionaries containing summary job information
"""
state = kwd.get("state", None)
states = optional_list(state)
tool_ids = optional_list(tool_id)
tool_id_likes = optional_list(tool_id_like)
is_admin = trans.user_is_admin
user_details = kwd.get("user_details", False)
user_id = kwd.get("user_id", None)
view = kwd.get("view", "collection")
if view not in ("collection", "admin_job_list"):
raise exceptions.RequestParameterInvalidException(f"view parameter '{view} is invalid")
if view == "admin_job_list" and not is_admin:
if view == JobIndexViewEnum.admin_job_list and not is_admin:
raise exceptions.AdminRequiredException("Only admins can use the admin_job_list view")
if user_id:
decoded_user_id = self.decode_id(user_id)
decoded_user_id = security.decode_id(user_id)
else:
decoded_user_id = None
if is_admin:
if decoded_user_id is not None:
query = trans.sa_session.query(model.Job).filter(model.Job.user_id == decoded_user_id)
@@ -163,27 +231,18 @@ class JobController(BaseGalaxyAPIController, UsesVisualizationMixin):
query = query.filter(or_(*t))
return query
query = build_and_apply_filters(query, state, lambda s: model.Job.state == s)
query = build_and_apply_filters(query, states, lambda s: model.Job.state == s)
query = build_and_apply_filters(query, tool_ids, lambda t: model.Job.tool_id == t)
query = build_and_apply_filters(query, tool_id_likes, lambda t: model.Job.tool_id.like(t))
query = build_and_apply_filters(query, date_range_min, lambda dmin: model.Job.update_time >= dmin)
query = build_and_apply_filters(query, date_range_max, lambda dmax: model.Job.update_time <= dmax)
query = build_and_apply_filters(query, kwd.get("tool_id", None), lambda t: model.Job.tool_id == t)
query = build_and_apply_filters(query, kwd.get("tool_id_like", None), lambda t: model.Job.tool_id.like(t))
query = build_and_apply_filters(
query, kwd.get("date_range_min", None), lambda dmin: model.Job.update_time >= dmin
)
query = build_and_apply_filters(
query, kwd.get("date_range_max", None), lambda dmax: model.Job.update_time <= dmax
)
history_id = kwd.get("history_id", None)
workflow_id = kwd.get("workflow_id", None)
invocation_id = kwd.get("invocation_id", None)
if history_id is not None:
decoded_history_id = self.decode_id(history_id)
decoded_history_id = security.decode_id(history_id)
query = query.filter(model.Job.history_id == decoded_history_id)
if workflow_id or invocation_id:
if workflow_id is not None:
decoded_workflow_id = self.decode_id(workflow_id)
decoded_workflow_id = security.decode_id(workflow_id)
wfi_step = (
trans.sa_session.query(model.WorkflowInvocationStep)
.join(model.WorkflowInvocation)
@@ -194,7 +253,7 @@ class JobController(BaseGalaxyAPIController, UsesVisualizationMixin):
.subquery()
)
elif invocation_id is not None:
decoded_invocation_id = self.decode_id(invocation_id)
decoded_invocation_id = security.decode_id(invocation_id)
wfi_step = (
trans.sa_session.query(model.WorkflowInvocationStep)
.filter(model.WorkflowInvocationStep.workflow_invocation_id == decoded_invocation_id)
@@ -208,7 +267,7 @@ class JobController(BaseGalaxyAPIController, UsesVisualizationMixin):
)
query = query1.union(query2)
if kwd.get("order_by") == "create_time":
if order_by == JobIndexSortByEnum.create_time:
order_by = model.Job.create_time.desc()
else:
order_by = model.Job.update_time.desc()
@@ -220,8 +279,8 @@ class JobController(BaseGalaxyAPIController, UsesVisualizationMixin):
out = []
for job in query.yield_per(model.YIELD_PER_ROWS):
job_dict = job.to_dict(view, system_details=is_admin)
j = self.encode_all_ids(trans, job_dict, True)
if view == "admin_job_list":
j = security.encode_all_ids(job_dict, True)
if view == JobIndexViewEnum.admin_job_list:
j["decoded_job_id"] = job.id
if user_details:
j["user_email"] = job.get_user_email()
@@ -229,6 +288,12 @@ class JobController(BaseGalaxyAPIController, UsesVisualizationMixin):
return out
class JobController(BaseGalaxyAPIController, UsesVisualizationMixin):
job_manager = depends(JobManager)
job_search = depends(JobSearch)
hda_manager = depends(hdas.HDAManager)
@expose_api
def common_problems(self, trans: ProvidesUserContext, id, **kwd):
"""