mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge remote-tracking branch 'upstream/release_26.0' into dev
This commit is contained in:
@@ -10,7 +10,7 @@ import ToolsListTable from "@/components/ToolsList/ToolsListTable.vue";
|
||||
const toolStore = useToolStore();
|
||||
const { loading } = storeToRefs(toolStore);
|
||||
|
||||
const whooshQuery = computed(() => createWhooshQuery({ section: "Get Data" }));
|
||||
const whooshQuery = computed(() => createWhooshQuery({ section: '"Get Data"' }));
|
||||
|
||||
const toolsInGetDataSection = computed(() => Object.values(toolStore.getToolsById(whooshQuery.value)));
|
||||
|
||||
|
||||
@@ -1857,8 +1857,8 @@ class MinimalJobWrapper(HasResourceParameters):
|
||||
|
||||
if object_store_id is None:
|
||||
object_store_id = job.preferred_object_store_id
|
||||
if object_store_id is None and job.workflow_invocation_step:
|
||||
workflow_invocation_step = job.workflow_invocation_step
|
||||
workflow_invocation_step = job.effective_workflow_invocation_step
|
||||
if object_store_id is None and workflow_invocation_step:
|
||||
invocation_object_stores = workflow_invocation_step.preferred_object_stores
|
||||
if invocation_object_stores.is_split_configuration:
|
||||
# Redo for subworkflows...
|
||||
|
||||
@@ -1310,8 +1310,8 @@ def map_tool_to_destination(
|
||||
raise JobMappingException(e)
|
||||
|
||||
# Get all inputs from tool and databases
|
||||
inp_data: dict[str, DatasetInstance] = {da.name: da.dataset for da in job.input_datasets}
|
||||
inp_data.update([(da.name, da.dataset) for da in job.input_library_datasets])
|
||||
inp_data: dict[str, DatasetInstance] = {da.name: da.dataset for da in job.input_datasets if da.dataset}
|
||||
inp_data.update([(da.name, da.dataset) for da in job.input_library_datasets if da.dataset])
|
||||
|
||||
if config is not None and str(tool.old_id) in config["tools"]:
|
||||
if "rules" in config["tools"][str(tool.old_id)]:
|
||||
|
||||
+25
-20
@@ -558,7 +558,7 @@ class JobSearch:
|
||||
# and the ids that have been used in the job that has already been run in `used_ids`.
|
||||
requested_ids = []
|
||||
data_types = []
|
||||
used_ids: list[Label[int]] = []
|
||||
used_ids: list[Label[int] | Label[int | None]] = []
|
||||
for k, input_list in input_data.items():
|
||||
# k will be matched against the JobParameter.name column. This can be prefixed depending on whether
|
||||
# the input is in a repeat, or not (section and conditional)
|
||||
@@ -672,21 +672,18 @@ class JobSearch:
|
||||
history_id: Union[int, None],
|
||||
) -> "Select[tuple[int]]":
|
||||
"""Build subquery that selects a job with correct job parameters."""
|
||||
job_ids_materialized_cte = stmt.cte("job_ids_cte")
|
||||
outer_select_columns = [job_ids_materialized_cte.c[col.name] for col in stmt.selected_columns]
|
||||
stmt = select(*outer_select_columns).select_from(job_ids_materialized_cte)
|
||||
stmt = (
|
||||
stmt.join(model.Job, model.Job.id == job_ids_materialized_cte.c.job_id)
|
||||
.join(model.History, model.Job.history_id == model.History.id)
|
||||
.where(
|
||||
and_(
|
||||
model.Job.tool_id == tool_id,
|
||||
or_(
|
||||
model.Job.user_id == user_id,
|
||||
model.History.published == true(),
|
||||
),
|
||||
model.Job.copied_from_job_id.is_(None), # Always pick original job
|
||||
)
|
||||
# Apply job-level filters BEFORE the CTE so they are included in the
|
||||
# materialized result. This lets PostgreSQL use selective indexes
|
||||
# (e.g. on tool_id) inside the CTE instead of scanning the entire job
|
||||
# table first and filtering afterwards.
|
||||
stmt = stmt.join(model.History, model.Job.history_id == model.History.id).where(
|
||||
and_(
|
||||
model.Job.tool_id == tool_id,
|
||||
or_(
|
||||
model.Job.user_id == user_id,
|
||||
model.History.published == true(),
|
||||
),
|
||||
model.Job.copied_from_job_id.is_(None), # Always pick original job
|
||||
)
|
||||
)
|
||||
if tool_version:
|
||||
@@ -712,6 +709,14 @@ class JobSearch:
|
||||
job_states = {Job.states.SKIPPED}
|
||||
stmt = stmt.where(Job.state.in_(job_states))
|
||||
|
||||
# Wrap in a CTE to materialize the filtered job IDs. This prevents
|
||||
# the planner from choosing poor join orders for the subsequent
|
||||
# job_parameter joins (important when tool_id is not highly selective).
|
||||
job_ids_materialized_cte = stmt.cte("job_ids_cte")
|
||||
outer_select_columns = [job_ids_materialized_cte.c[col.name] for col in stmt.selected_columns]
|
||||
stmt = select(*outer_select_columns).select_from(job_ids_materialized_cte)
|
||||
stmt = stmt.join(model.Job, model.Job.id == job_ids_materialized_cte.c.job_id)
|
||||
|
||||
for k, v in wildcard_param_dump.items():
|
||||
if v == {"__class__": "RuntimeValue"}:
|
||||
# TODO: verify this is always None. e.g. run with runtime input input
|
||||
@@ -793,7 +798,7 @@ class JobSearch:
|
||||
self,
|
||||
stmt: "Select[tuple[int]]",
|
||||
data_conditions: list["ColumnElement[bool]"],
|
||||
used_ids: list["Label[int]"],
|
||||
used_ids: list["Label[int] | Label[int | None]"],
|
||||
k,
|
||||
v,
|
||||
identifier,
|
||||
@@ -860,7 +865,7 @@ class JobSearch:
|
||||
self,
|
||||
stmt: "Select[tuple[int]]",
|
||||
data_conditions: list["ColumnElement[bool]"],
|
||||
used_ids: list["Label[int]"],
|
||||
used_ids: list["Label[int] | Label[int | None]"],
|
||||
k,
|
||||
v,
|
||||
value_index: int,
|
||||
@@ -883,7 +888,7 @@ class JobSearch:
|
||||
self,
|
||||
stmt: "Select[tuple[int]]",
|
||||
data_conditions: list["ColumnElement[bool]"],
|
||||
used_ids: list["Label[int]"],
|
||||
used_ids: list["Label[int] | Label[int | None]"],
|
||||
k,
|
||||
v,
|
||||
user_id: int,
|
||||
@@ -1131,7 +1136,7 @@ class JobSearch:
|
||||
self,
|
||||
stmt: "Select[tuple[int]]",
|
||||
data_conditions: list["ColumnElement[bool]"],
|
||||
used_ids: list["Label[int]"],
|
||||
used_ids: list["Label[int] | Label[int | None]"],
|
||||
k,
|
||||
v,
|
||||
user_id: int,
|
||||
|
||||
@@ -1668,7 +1668,7 @@ class Job(Base, JobLike, UsesCreateAndUpdateTime, Dictifiable, Serializable):
|
||||
interactivetool_entry_points: Mapped[list["InteractiveToolEntryPoint"]] = relationship(
|
||||
back_populates="job", uselist=True
|
||||
)
|
||||
implicit_collection_jobs_association: Mapped["ImplicitCollectionJobsJobAssociation"] = relationship(
|
||||
implicit_collection_jobs_association: Mapped[Optional["ImplicitCollectionJobsJobAssociation"]] = relationship(
|
||||
back_populates="job", uselist=False
|
||||
)
|
||||
container: Mapped[Optional["JobContainerAssociation"]] = relationship(back_populates="job", uselist=False)
|
||||
@@ -1685,6 +1685,27 @@ class Job(Base, JobLike, UsesCreateAndUpdateTime, Dictifiable, Serializable):
|
||||
back_populates="job"
|
||||
)
|
||||
|
||||
@property
|
||||
def effective_workflow_invocation_step(self) -> Optional["WorkflowInvocationStep"]:
|
||||
"""The WorkflowInvocationStep backing this job, including mapped steps.
|
||||
|
||||
For non-mapped steps this is the direct ``workflow_invocation_step`` back-ref.
|
||||
For mapped steps ``WorkflowInvocationStep.job_id`` is NULL — the step points
|
||||
at an ``ImplicitCollectionJobs`` instead, and each job is linked to that ICJ
|
||||
via ``ImplicitCollectionJobsJobAssociation``. Resolve it by querying
|
||||
``WorkflowInvocationStep`` using the ICJ id.
|
||||
"""
|
||||
if self.workflow_invocation_step is not None:
|
||||
return self.workflow_invocation_step
|
||||
icj_assoc = self.implicit_collection_jobs_association
|
||||
if icj_assoc is None:
|
||||
return None
|
||||
icj_id = icj_assoc.implicit_collection_jobs_id
|
||||
session = required_object_session(self)
|
||||
return session.execute(
|
||||
select(WorkflowInvocationStep).where(WorkflowInvocationStep.implicit_collection_jobs_id == icj_id)
|
||||
).scalar_one_or_none()
|
||||
|
||||
dict_collection_visible_keys = [
|
||||
"id",
|
||||
"state",
|
||||
@@ -2665,11 +2686,15 @@ class JobToInputDatasetAssociation(Base, RepresentById):
|
||||
|
||||
id: Mapped[int] = mapped_column(primary_key=True)
|
||||
job_id: Mapped[int] = mapped_column(ForeignKey("job.id"), index=True, nullable=True)
|
||||
dataset_id: Mapped[int] = mapped_column(ForeignKey("history_dataset_association.id"), index=True, nullable=True)
|
||||
dataset_id: Mapped[Optional[int]] = mapped_column(
|
||||
ForeignKey("history_dataset_association.id"), index=True, nullable=True
|
||||
)
|
||||
dataset_version: Mapped[Optional[int]]
|
||||
name: Mapped[str] = mapped_column(String(255), nullable=True)
|
||||
adapter: Mapped[Optional[dict[str, Any]]] = mapped_column(JSONType, nullable=True)
|
||||
dataset: Mapped["HistoryDatasetAssociation"] = relationship(lazy="joined", back_populates="dependent_jobs")
|
||||
dataset: Mapped[Optional["HistoryDatasetAssociation"]] = relationship(
|
||||
lazy="joined", back_populates="dependent_jobs"
|
||||
)
|
||||
job: Mapped["Job"] = relationship(back_populates="input_datasets")
|
||||
|
||||
def __init__(self, name, dataset, adapter_json=None):
|
||||
|
||||
@@ -633,10 +633,16 @@ class FileParameter(MetadataParameter):
|
||||
if wrapped_value:
|
||||
return wrapped_value
|
||||
else:
|
||||
# If we've simultaneously copied the dataset and we've changed the datatype on the
|
||||
# copy we may not have committed the MetadataFile yet, so we need to commit the session.
|
||||
# TODO: It would be great if we can avoid the commit in the future.
|
||||
session.commit()
|
||||
# If we've simultaneously copied the dataset and we've changed the datatype on the
|
||||
# copy we may not have flushed the pending MetadataFile to the DB yet, so the
|
||||
# select above may have missed it. Flush (NOT commit) so the pending INSERT is
|
||||
# visible to the retry SELECT within this session without leaking every other
|
||||
# pending change in the session to concurrent readers. Committing here caused
|
||||
# https://github.com/galaxyproject/galaxy/issues/22194: a mid-loop wrap() inside
|
||||
# JobWrapper.finish() flushed an intermedia expression.json's dataset's
|
||||
# Dataset.state = OK to the DB before exec_after_process had replaced the file,
|
||||
# exposing a globall inconsistent state to the workflow scheduler.
|
||||
session.flush()
|
||||
return session.execute(select(galaxy.model.MetadataFile).filter_by(uuid=value)).scalar_one_or_none()
|
||||
|
||||
def make_copy(self, value, target_context: MetadataCollection, source_context):
|
||||
|
||||
@@ -12,6 +12,14 @@ warnings.filterwarnings("ignore", message=r"[\n.]DEPRECATION: Python 2", module=
|
||||
|
||||
from galaxy.util import requests
|
||||
|
||||
try:
|
||||
from cwl_utils.types import CWLObjectType
|
||||
except ImportError:
|
||||
try:
|
||||
from cwltool.utils import CWLObjectType # type: ignore[attr-defined, unused-ignore]
|
||||
except ImportError:
|
||||
CWLObjectType = object # type: ignore[assignment, misc]
|
||||
|
||||
try:
|
||||
from cwltool import (
|
||||
main,
|
||||
@@ -67,13 +75,11 @@ except ImportError:
|
||||
|
||||
try:
|
||||
from cwltool.utils import (
|
||||
CWLObjectType,
|
||||
JobsType,
|
||||
normalizeFilesDirs,
|
||||
visit_class,
|
||||
)
|
||||
except ImportError:
|
||||
CWLObjectType = object # type: ignore[assignment, misc]
|
||||
JobsType = object # type: ignore[misc, unused-ignore]
|
||||
visit_class = None # type: ignore[assignment]
|
||||
normalizeFilesDirs = None # type: ignore[assignment]
|
||||
|
||||
@@ -259,6 +259,122 @@ def _check_pattern(node):
|
||||
return True
|
||||
|
||||
|
||||
class OutputsStructuredLikeReference(Linter):
|
||||
@classmethod
|
||||
def lint(cls, tool_source: "ToolSource", lint_ctx: "LintContext"):
|
||||
tool_xml = getattr(tool_source, "xml_tree", None)
|
||||
if not tool_xml:
|
||||
return
|
||||
param_qualified_paths = _collect_param_qualified_paths(tool_xml)
|
||||
for output in tool_xml.findall("./outputs/collection[@structured_like]"):
|
||||
structured_like = output.attrib["structured_like"]
|
||||
_check_unqualified_reference(
|
||||
lint_ctx, cls.name(), output, structured_like, "structured_like", param_qualified_paths
|
||||
)
|
||||
|
||||
|
||||
class OutputsFormatSourceReference(Linter):
|
||||
@classmethod
|
||||
def lint(cls, tool_source: "ToolSource", lint_ctx: "LintContext"):
|
||||
tool_xml = getattr(tool_source, "xml_tree", None)
|
||||
if not tool_xml:
|
||||
return
|
||||
param_qualified_paths = _collect_param_qualified_paths(tool_xml)
|
||||
output_names = {
|
||||
o.attrib["name"]
|
||||
for o in tool_xml.findall("./outputs/data[@name]") + tool_xml.findall("./outputs/collection[@name]")
|
||||
}
|
||||
for output in tool_xml.findall("./outputs/data[@format_source]") + tool_xml.findall(
|
||||
"./outputs/collection[@format_source]"
|
||||
):
|
||||
format_source = output.attrib["format_source"]
|
||||
# format_source can reference other outputs, skip if it matches an output name
|
||||
if format_source in output_names:
|
||||
continue
|
||||
_check_unqualified_reference(
|
||||
lint_ctx, cls.name(), output, format_source, "format_source", param_qualified_paths
|
||||
)
|
||||
|
||||
|
||||
def _check_unqualified_reference(
|
||||
lint_ctx: "LintContext",
|
||||
linter_name: str,
|
||||
node: "Element",
|
||||
ref_value: str,
|
||||
attr_name: str,
|
||||
param_qualified_paths: dict,
|
||||
):
|
||||
if "|" in ref_value:
|
||||
return
|
||||
# Check if it matches a top-level param directly
|
||||
top_level_match = any(qp == ref_value for paths in param_qualified_paths.values() for qp in paths)
|
||||
if top_level_match:
|
||||
return
|
||||
matches = param_qualified_paths.get(ref_value, [])
|
||||
output_name = node.attrib.get("name", "unknown")
|
||||
if len(matches) == 1:
|
||||
lint_ctx.warn(
|
||||
f"Output '{output_name}' uses unqualified {attr_name}='{ref_value}'. "
|
||||
f"Use the qualified name '{matches[0]}'.",
|
||||
linter=linter_name,
|
||||
node=node,
|
||||
)
|
||||
elif len(matches) > 1:
|
||||
lint_ctx.warn(
|
||||
f"Output '{output_name}' uses ambiguous unqualified {attr_name}='{ref_value}' "
|
||||
f"matching multiple inputs: {', '.join(matches)}. Use a qualified name.",
|
||||
linter=linter_name,
|
||||
node=node,
|
||||
)
|
||||
else:
|
||||
lint_ctx.error(
|
||||
f"Output '{output_name}' references {attr_name}='{ref_value}' "
|
||||
f"which does not match any input parameter.",
|
||||
linter=linter_name,
|
||||
node=node,
|
||||
)
|
||||
|
||||
|
||||
def _collect_param_qualified_paths(tool_xml: "ElementTree") -> dict:
|
||||
"""Build a map of unqualified param name -> list of qualified paths."""
|
||||
param_paths: dict = {}
|
||||
parent_map = {child: parent for parent in tool_xml.iter() for child in parent}
|
||||
for param in tool_xml.findall("./inputs//param"):
|
||||
name = param.attrib.get("name")
|
||||
if not name:
|
||||
argument = param.attrib.get("argument")
|
||||
if argument:
|
||||
name = argument.lstrip("-").replace("-", "_")
|
||||
if not name:
|
||||
continue
|
||||
qualified = _get_qualified_name(param, parent_map)
|
||||
param_paths.setdefault(name, []).append(qualified)
|
||||
return param_paths
|
||||
|
||||
|
||||
def _get_qualified_name(param_elem: "Element", parent_map: dict) -> str:
|
||||
"""Walk up the XML tree to build the qualified path for a param element."""
|
||||
name = param_elem.attrib.get("name")
|
||||
if not name:
|
||||
argument = param_elem.attrib.get("argument")
|
||||
if argument:
|
||||
name = argument.lstrip("-").replace("-", "_")
|
||||
parts = [name] if name else []
|
||||
current = param_elem
|
||||
while True:
|
||||
parent = parent_map.get(current)
|
||||
if parent is None:
|
||||
break
|
||||
if parent.tag in ("conditional", "section"):
|
||||
parent_name = parent.attrib.get("name")
|
||||
if parent_name:
|
||||
parts.insert(0, parent_name)
|
||||
elif parent.tag in ("inputs", "tool"):
|
||||
break
|
||||
current = parent
|
||||
return "|".join(parts)
|
||||
|
||||
|
||||
def _has_tool_provided_metadata(tool_xml: "ElementTree") -> bool:
|
||||
outputs = tool_xml.find("./outputs")
|
||||
if outputs is not None:
|
||||
|
||||
@@ -964,6 +964,7 @@ class DefaultToolAction(ToolAction):
|
||||
job.user = trans.user
|
||||
if history:
|
||||
job.history_id = model.cached_id(history)
|
||||
job.history = history
|
||||
job.tool_id = tool.id
|
||||
try:
|
||||
# For backward compatibility, some tools may not have versions yet.
|
||||
|
||||
@@ -512,25 +512,26 @@ class ExecutionTracker:
|
||||
return output_collection_name
|
||||
|
||||
def sliced_input_collection_structure(self, input_name):
|
||||
unqualified_recurse = Version(str(self.tool.profile)) < Version("18.09") and "|" not in input_name
|
||||
unqualified_recurse = Version(str(self.tool.profile)) < Version("26.0") and "|" not in input_name
|
||||
|
||||
def find_collection(input_dict, input_name):
|
||||
def find_collection(input_dict, input_name, path_prefix=""):
|
||||
for key, value in input_dict.items():
|
||||
if key == input_name:
|
||||
return value
|
||||
return value, f"{path_prefix}{key}"
|
||||
if isinstance(value, dict):
|
||||
if "|" in input_name:
|
||||
prefix, rest_input_name = input_name.split("|", 1)
|
||||
if key == prefix:
|
||||
return find_collection(value, rest_input_name)
|
||||
return find_collection(value, rest_input_name, f"{path_prefix}{key}|")
|
||||
elif unqualified_recurse:
|
||||
# Looking for "input1" instead of "cond|input1" for instance.
|
||||
# See discussion on https://github.com/galaxyproject/galaxy/issues/6157.
|
||||
unqualified_match = find_collection(value, input_name)
|
||||
unqualified_match, qualified_path = find_collection(value, input_name, f"{path_prefix}{key}|")
|
||||
if unqualified_match:
|
||||
return unqualified_match
|
||||
return unqualified_match, qualified_path
|
||||
return None, None
|
||||
|
||||
input_collection = find_collection(self.example_params, input_name)
|
||||
input_collection, qualified_name = find_collection(self.example_params, input_name)
|
||||
if input_collection is None:
|
||||
raise Exception("Failed to find referenced collection in inputs.")
|
||||
|
||||
@@ -543,10 +544,10 @@ class ExecutionTracker:
|
||||
)
|
||||
)
|
||||
subcollection_mapping_type = None
|
||||
if self.is_implicit_input(input_name):
|
||||
if self.is_implicit_input(qualified_name):
|
||||
collection_info = self.collection_info
|
||||
assert collection_info
|
||||
subcollection_mapping_type = collection_info.subcollection_mapping_type(input_name)
|
||||
subcollection_mapping_type = collection_info.subcollection_mapping_type(qualified_name)
|
||||
|
||||
return get_structure(
|
||||
input_collection, collection_type_description, leaf_subcollection_type=subcollection_mapping_type
|
||||
|
||||
@@ -6,13 +6,19 @@ import logging
|
||||
import os
|
||||
import re
|
||||
import shutil
|
||||
from typing import Union
|
||||
|
||||
from galaxy import (
|
||||
exceptions,
|
||||
util,
|
||||
)
|
||||
from galaxy.managers.context import ProvidesAppContext
|
||||
from galaxy.model import Job
|
||||
from galaxy.model import (
|
||||
Job,
|
||||
JobToOutputDatasetAssociation,
|
||||
JobToOutputLibraryDatasetAssociation,
|
||||
)
|
||||
from galaxy.structured_app import MinimalManagerApp
|
||||
from galaxy.web import (
|
||||
expose_api_anonymous_and_sessionless,
|
||||
expose_api_raw_anonymous_and_sessionless,
|
||||
@@ -35,7 +41,7 @@ class JobFilesAPIController(BaseGalaxyAPIController):
|
||||
"""
|
||||
|
||||
@expose_api_raw_anonymous_and_sessionless
|
||||
def index(self, trans: ProvidesAppContext, job_id, **kwargs):
|
||||
def index(self, trans: ProvidesAppContext, job_id: str, **kwargs):
|
||||
"""
|
||||
GET /api/jobs/{job_id}/files
|
||||
|
||||
@@ -74,7 +80,7 @@ class JobFilesAPIController(BaseGalaxyAPIController):
|
||||
raise
|
||||
|
||||
@expose_api_anonymous_and_sessionless
|
||||
def create(self, trans, job_id, payload, **kwargs):
|
||||
def create(self, trans: ProvidesAppContext, job_id: str, payload, **kwargs):
|
||||
"""
|
||||
create( self, trans, job_id, payload, **kwargs )
|
||||
* POST /api/jobs/{job_id}/files
|
||||
@@ -180,7 +186,7 @@ class JobFilesAPIController(BaseGalaxyAPIController):
|
||||
return None
|
||||
|
||||
@expose_api_anonymous_and_sessionless
|
||||
def tus_hooks(self, trans, **kwds):
|
||||
def tus_hooks(self, trans: ProvidesAppContext, **kwds):
|
||||
"""No-op but if hook specified the way we do for user upload it would hit this action.
|
||||
|
||||
Exposed as PATCH /api/job_files/tus_hooks and documented in the docstring for
|
||||
@@ -188,7 +194,7 @@ class JobFilesAPIController(BaseGalaxyAPIController):
|
||||
"""
|
||||
pass
|
||||
|
||||
def __authorize_job_access(self, trans, encoded_job_id, **kwargs):
|
||||
def __authorize_job_access(self, trans: ProvidesAppContext, encoded_job_id: str, **kwargs):
|
||||
for key in ["path", "job_key"]:
|
||||
if key not in kwargs:
|
||||
error_message = f"Job files action requires a valid '{key}'."
|
||||
@@ -201,12 +207,13 @@ class JobFilesAPIController(BaseGalaxyAPIController):
|
||||
|
||||
# Verify job is active. Don't update the contents of complete jobs.
|
||||
job = trans.sa_session.get(Job, job_id)
|
||||
assert job
|
||||
if job.state not in Job.non_ready_states:
|
||||
error_message = "Attempting to read or modify the files of a job that has already completed."
|
||||
raise exceptions.ItemAccessibilityException(error_message)
|
||||
return job
|
||||
|
||||
def __check_job_can_write_to_path(self, trans, job, path):
|
||||
def __check_job_can_write_to_path(self, trans: ProvidesAppContext, job: Job, path: str):
|
||||
"""Verify an idealized job runner should actually be able to write to
|
||||
the specified path - it must be a dataset output, a dataset "extra
|
||||
file", or a some place in the working directory of this job.
|
||||
@@ -219,23 +226,25 @@ class JobFilesAPIController(BaseGalaxyAPIController):
|
||||
if not in_work_dir and not self.__is_output_dataset_path(job, path):
|
||||
raise exceptions.ItemAccessibilityException("Job is not authorized to write to supplied path.")
|
||||
|
||||
def __is_output_dataset_path(self, job, path):
|
||||
def __is_output_dataset_path(self, job: Job, path: str):
|
||||
"""Check if is an output path for this job or a file in the an
|
||||
output's extra files path.
|
||||
"""
|
||||
da_lists = [job.output_datasets, job.output_library_datasets]
|
||||
for da_list in da_lists:
|
||||
for job_dataset_association in da_list:
|
||||
dataset = job_dataset_association.dataset
|
||||
if not dataset:
|
||||
continue
|
||||
if os.path.abspath(dataset.get_file_name()) == os.path.abspath(path):
|
||||
return True
|
||||
elif util.in_directory(path, dataset.extra_files_path):
|
||||
return True
|
||||
all_output_assocs: list[Union[JobToOutputDatasetAssociation, JobToOutputLibraryDatasetAssociation]] = [
|
||||
*job.output_datasets,
|
||||
*job.output_library_datasets,
|
||||
]
|
||||
for assoc in all_output_assocs:
|
||||
dataset = assoc.dataset
|
||||
if not dataset:
|
||||
continue
|
||||
if os.path.abspath(dataset.get_file_name()) == os.path.abspath(path):
|
||||
return True
|
||||
elif util.in_directory(path, dataset.extra_files_path):
|
||||
return True
|
||||
return False
|
||||
|
||||
def __in_working_directory(self, job, path, app):
|
||||
def __in_working_directory(self, job: Job, path: str, app: MinimalManagerApp):
|
||||
working_directory = app.object_store.get_filename(
|
||||
job, base_dir="job_work", dir_only=True, extra_dir=str(job.id)
|
||||
)
|
||||
|
||||
@@ -339,7 +339,7 @@ class FastAPIJobs:
|
||||
for job_input_assoc in job.input_datasets:
|
||||
input_dataset_instance = job_input_assoc.dataset
|
||||
if input_dataset_instance is None:
|
||||
continue # type: ignore[unreachable] # TODO if job_input_assoc.dataset is indeed never None, remove the above check
|
||||
continue
|
||||
if input_dataset_instance.get_total_size() == 0:
|
||||
has_empty_inputs = True
|
||||
input_instance_id = input_dataset_instance.id
|
||||
|
||||
@@ -74,7 +74,7 @@ from galaxy_test.base.testcase import FunctionalTestCase
|
||||
try:
|
||||
from galaxy_test.driver.driver_util import GalaxyTestDriver
|
||||
except ImportError:
|
||||
GalaxyTestDriver = None # type: ignore[misc,assignment]
|
||||
GalaxyTestDriver = None # type: ignore[assignment, misc, unused-ignore]
|
||||
|
||||
|
||||
def _load_config_file() -> None:
|
||||
|
||||
@@ -19,6 +19,7 @@ from sqlalchemy import (
|
||||
|
||||
import tool_shed.util.shed_util_common as suc
|
||||
from galaxy.exceptions import (
|
||||
AuthenticationRequired,
|
||||
InsufficientPermissionsException,
|
||||
ObjectNotFound,
|
||||
RequestParameterInvalidException,
|
||||
@@ -134,7 +135,7 @@ class FastAPIUsers:
|
||||
def current(self, trans: SessionRequestContext = DependsOnTrans) -> User:
|
||||
user = trans.user
|
||||
if not user:
|
||||
raise ObjectNotFound()
|
||||
raise AuthenticationRequired()
|
||||
|
||||
return get_api_user(trans.app, user)
|
||||
|
||||
|
||||
@@ -41,7 +41,7 @@ install_requires =
|
||||
galaxy-tool-util[cwl,edam]
|
||||
galaxy-tool-shed-schema
|
||||
galaxy-tours
|
||||
galaxy-util[image-util]
|
||||
galaxy-util[image-util,jstree]
|
||||
galaxy-web-framework
|
||||
galaxy-web-stack
|
||||
Beaker
|
||||
@@ -91,6 +91,8 @@ test =
|
||||
pykwalify
|
||||
pytest
|
||||
pytest-asyncio
|
||||
pytest-mock
|
||||
responses
|
||||
testfixtures
|
||||
|
||||
[options.entry_points]
|
||||
|
||||
@@ -69,6 +69,7 @@ python_requires = >=3.10
|
||||
|
||||
[options.extras_require]
|
||||
test =
|
||||
galaxy-config
|
||||
pytest
|
||||
roc-validator
|
||||
|
||||
|
||||
@@ -45,6 +45,7 @@ test =
|
||||
galaxy-test-base
|
||||
pytest
|
||||
gcsfs
|
||||
responses
|
||||
s3fs>=2023.1.0
|
||||
|
||||
[options.entry_points]
|
||||
|
||||
@@ -20,6 +20,7 @@ tours
|
||||
auth
|
||||
job_execution
|
||||
test_api
|
||||
test_selenium
|
||||
web_framework
|
||||
web_stack
|
||||
|
||||
|
||||
@@ -31,6 +31,7 @@ version = 26.1.dev0
|
||||
[options]
|
||||
include_package_data = True
|
||||
install_requires =
|
||||
galaxy-tool-util-models
|
||||
galaxy-util
|
||||
pydantic[email]>=2.7.4
|
||||
packages = find:
|
||||
|
||||
+14
-15
@@ -36,32 +36,20 @@ fi
|
||||
# Change to packages directory.
|
||||
cd "$(dirname "$0")"
|
||||
|
||||
# Use a throw-away virtualenv
|
||||
TEST_PYTHON=${TEST_PYTHON:-"python3"}
|
||||
TEST_ENV_DIR=${TEST_ENV_DIR:-$(mktemp -d -t gxpkgtestenvXXXXXX)}
|
||||
|
||||
if command -v uv >/dev/null; then
|
||||
uv venv --python "$TEST_PYTHON" "$TEST_ENV_DIR"
|
||||
VENV_CMD="uv venv --python $TEST_PYTHON"
|
||||
PIP_CMD="$(command -v uv) pip"
|
||||
BUILD_WHEEL_CMD="$(command -v uv) build"
|
||||
TWINE_CMD="$(command -v uvx) twine"
|
||||
else
|
||||
"$TEST_PYTHON" -m venv "$TEST_ENV_DIR"
|
||||
VENV_CMD="$TEST_PYTHON -m venv"
|
||||
PIP_CMD='python -m pip'
|
||||
BUILD_WHEEL_CMD='python -m build'
|
||||
TWINE_CMD=twine
|
||||
fi
|
||||
|
||||
# shellcheck disable=SC1091
|
||||
. "${TEST_ENV_DIR}/bin/activate"
|
||||
if [ "${PIP_CMD}" = 'python -m pip' ]; then
|
||||
${PIP_CMD} install --upgrade build pip setuptools twine wheel
|
||||
fi
|
||||
if [ $FOR_PULSAR -eq 0 ]; then
|
||||
# shellcheck disable=SC2086 - word splitting is intentional for PIP_EXTRA_ARGS
|
||||
${PIP_CMD} install ${PIP_EXTRA_ARGS} -r ../lib/galaxy/dependencies/pinned-typecheck-requirements.txt
|
||||
fi
|
||||
|
||||
# Ensure ordered by dependency DAG
|
||||
while read -r package_dir || [ -n "$package_dir" ]; do # https://stackoverflow.com/questions/12916352/shell-script-read-missing-last-line
|
||||
# Ignore empty lines
|
||||
@@ -81,8 +69,17 @@ while read -r package_dir || [ -n "$package_dir" ]; do # https://stackoverflow.
|
||||
|
||||
cd "$package_dir"
|
||||
|
||||
# Use a throw-away virtualenv
|
||||
TEST_ENV_DIR=$(mktemp -d -t gxpkgtestenvXXXXXX)
|
||||
${VENV_CMD} "$TEST_ENV_DIR"
|
||||
# shellcheck disable=SC1091
|
||||
. "${TEST_ENV_DIR}/bin/activate"
|
||||
if [ "${PIP_CMD}" = 'python -m pip' ]; then
|
||||
${PIP_CMD} install --upgrade build pip setuptools twine wheel
|
||||
fi
|
||||
|
||||
# Install extras (if needed)
|
||||
# shellcheck disable=SC2086 - word splitting is intentional for PIP_EXTRA_ARGS
|
||||
# shellcheck disable=SC2086 # word splitting is intentional for PIP_EXTRA_ARGS
|
||||
if [ "$package_dir" = "util" ]; then
|
||||
${PIP_CMD} install ${PIP_EXTRA_ARGS} '.[image-util,template,jstree,config-template,test]'
|
||||
elif [ "$package_dir" = "tool_util" ]; then
|
||||
@@ -101,6 +98,8 @@ while read -r package_dir || [ -n "$package_dir" ]; do # https://stackoverflow.
|
||||
# Ignore exit code 5 (no tests ran)
|
||||
pytest "${marker_args[@]}" . || test $? -eq 5
|
||||
if [ $FOR_PULSAR -eq 0 ]; then
|
||||
# shellcheck disable=SC2086 # word splitting is intentional for PIP_EXTRA_ARGS
|
||||
${PIP_CMD} install ${PIP_EXTRA_ARGS} -r ../../lib/galaxy/dependencies/pinned-typecheck-requirements.txt
|
||||
# make mypy uses uv now and so this legacy code should just run mypy
|
||||
# directly to use the venv we have already activated
|
||||
mypy .
|
||||
|
||||
@@ -64,6 +64,7 @@ python_requires = >=3.10
|
||||
|
||||
[options.extras_require]
|
||||
test =
|
||||
galaxy-test-base
|
||||
pytest
|
||||
|
||||
[options.packages.find]
|
||||
|
||||
@@ -34,7 +34,7 @@ version = 26.1.dev0
|
||||
include_package_data = True
|
||||
install_requires =
|
||||
galaxy-tool-util-models
|
||||
galaxy-util[image-util]>=22.1
|
||||
galaxy-util[template,image-util]>=22.1
|
||||
conda-package-streaming
|
||||
jsonschema
|
||||
lxml!=4.2.2
|
||||
@@ -66,6 +66,7 @@ console_scripts =
|
||||
|
||||
[options.extras_require]
|
||||
cwl =
|
||||
cwl-utils
|
||||
cwltool>=3.1.20230624081518
|
||||
mulled =
|
||||
jinja2
|
||||
|
||||
@@ -32,6 +32,8 @@ version = 26.1.dev0
|
||||
include_package_data = True
|
||||
install_requires =
|
||||
galaxy-navigation
|
||||
galaxy-schema
|
||||
galaxy-util
|
||||
pydantic>=2.7.4
|
||||
PyYAML
|
||||
packages = find:
|
||||
@@ -41,6 +43,10 @@ python_requires = >=3.10
|
||||
console_scripts =
|
||||
gx-validate-tours = galaxy.tours.validate:main
|
||||
|
||||
[options.extras_require]
|
||||
test =
|
||||
pytest
|
||||
|
||||
[options.packages.find]
|
||||
exclude =
|
||||
tests*
|
||||
|
||||
@@ -56,6 +56,7 @@ template =
|
||||
fissix;python_version>='3.13'
|
||||
future>=1.0.0
|
||||
config-template =
|
||||
galaxy-tool-util-models
|
||||
Jinja2
|
||||
pydantic>=2.7.4
|
||||
test =
|
||||
|
||||
@@ -141,6 +141,32 @@ text_input1: |
|
||||
samp2\t20.0
|
||||
"""
|
||||
|
||||
WORKFLOW_WITH_MAPPED_COLLECTION_OUTPUT = """
|
||||
class: GalaxyWorkflow
|
||||
inputs:
|
||||
input_collection:
|
||||
type: collection
|
||||
collection_type: list
|
||||
outputs:
|
||||
wf_output_1:
|
||||
outputSource: cat_mapped/out_file1
|
||||
steps:
|
||||
cat_mapped:
|
||||
tool_id: cat
|
||||
in:
|
||||
input1: input_collection
|
||||
"""
|
||||
|
||||
WORKFLOW_MAPPED_COLLECTION_TEST_DATA = """
|
||||
input_collection:
|
||||
collection_type: list
|
||||
elements:
|
||||
- identifier: el1
|
||||
content: "data 1"
|
||||
- identifier: el2
|
||||
content: "data 2"
|
||||
"""
|
||||
|
||||
|
||||
def assert_storage_name_is(storage_dict: dict[str, Any], name: str):
|
||||
storage_name = storage_dict["name"]
|
||||
@@ -366,6 +392,67 @@ class TestObjectStoreSelectionWithPreferredObjectStoresIntegration(BaseObjectSto
|
||||
assert_storage_name_is(output_info, "Static Storage")
|
||||
assert_storage_name_is(intermediate_dict, "Dynamic EBS")
|
||||
|
||||
def test_workflow_mapped_collection_objectstore_selection(self):
|
||||
# Regression for https://github.com/galaxyproject/galaxy/issues/21846
|
||||
with self.dataset_populator.test_history() as history_id:
|
||||
element_storages = self._run_workflow_with_mapped_collection(
|
||||
history_id,
|
||||
WORKFLOW_WITH_MAPPED_COLLECTION_OUTPUT,
|
||||
WORKFLOW_MAPPED_COLLECTION_TEST_DATA,
|
||||
)
|
||||
for storage in element_storages:
|
||||
assert_storage_name_is(storage, "Default Store")
|
||||
|
||||
with self.dataset_populator.test_history() as history_id:
|
||||
element_storages = self._run_workflow_with_mapped_collection(
|
||||
history_id,
|
||||
WORKFLOW_WITH_MAPPED_COLLECTION_OUTPUT,
|
||||
WORKFLOW_MAPPED_COLLECTION_TEST_DATA,
|
||||
extra_invocation_kwds={"preferred_object_store_id": "static"},
|
||||
)
|
||||
for storage in element_storages:
|
||||
assert_storage_name_is(storage, "Static Storage")
|
||||
|
||||
def test_workflow_mapped_collection_objectstore_selection_split(self):
|
||||
# Regression for https://github.com/galaxyproject/galaxy/issues/21846
|
||||
with self.dataset_populator.test_history() as history_id:
|
||||
element_storages = self._run_workflow_with_mapped_collection(
|
||||
history_id,
|
||||
WORKFLOW_WITH_MAPPED_COLLECTION_OUTPUT,
|
||||
WORKFLOW_MAPPED_COLLECTION_TEST_DATA,
|
||||
extra_invocation_kwds={
|
||||
"preferred_outputs_object_store_id": "static",
|
||||
"preferred_intermediate_object_store_id": "dynamic_ebs",
|
||||
},
|
||||
)
|
||||
for storage in element_storages:
|
||||
assert_storage_name_is(storage, "Static Storage")
|
||||
|
||||
def _run_workflow_with_mapped_collection(
|
||||
self,
|
||||
history_id: str,
|
||||
workflow: str,
|
||||
test_data: str,
|
||||
extra_invocation_kwds: Optional[dict[str, Any]] = None,
|
||||
):
|
||||
self.workflow_populator.run_workflow(
|
||||
workflow,
|
||||
test_data=test_data,
|
||||
history_id=history_id,
|
||||
extra_invocation_kwds=extra_invocation_kwds,
|
||||
)
|
||||
# Find the implicit collection in the history and inspect element storage
|
||||
history_contents = self.dataset_populator.get_history_contents(history_id, data={"v": "dev"})
|
||||
hdca = None
|
||||
for entry in history_contents:
|
||||
if entry.get("history_content_type") == "dataset_collection" and entry.get("visible", True):
|
||||
hdca = entry
|
||||
assert hdca is not None, "No collection found in history"
|
||||
hdca_details = self.dataset_populator.get_history_collection_details(history_id, content_id=hdca["id"])
|
||||
elements = hdca_details["elements"]
|
||||
assert len(elements) > 0, "Collection has no elements"
|
||||
return [self._storage_info(element["object"]) for element in elements]
|
||||
|
||||
def _run_workflow_with_collections_1(self, history_id: str, extra_invocation_kwds: Optional[dict[str, Any]] = None):
|
||||
wf_run = self.workflow_populator.run_workflow(
|
||||
WORKFLOW_WITH_COLLECTIONS_1,
|
||||
|
||||
@@ -80,7 +80,7 @@ class TestIsEmailBanned:
|
||||
assert is_email_banned("ab+bar@gmail.com", "_", rules)
|
||||
assert not is_email_banned("ab-bar@gmail.com", "_", rules) # different sub-addressing delimiter
|
||||
|
||||
def test_no_canonical_rules(self, monkeypatch, appconfig):
|
||||
def test_no_canonical_rules(self, monkeypatch):
|
||||
"""No rules loaded."""
|
||||
monkeypatch.setattr(validate_user_input, "_read_email_ban_list", lambda a: self.mock_ban_list)
|
||||
|
||||
@@ -97,7 +97,7 @@ class TestIsEmailBanned:
|
||||
assert not is_email_banned("a.b@gmail.com", "_", rules)
|
||||
assert not is_email_banned("ab+bar@gmail.com", "_", rules)
|
||||
|
||||
def test_custom_canonical_rules(self, monkeypatch, appconfig):
|
||||
def test_custom_canonical_rules(self, monkeypatch):
|
||||
"""No rules loaded."""
|
||||
monkeypatch.setattr(validate_user_input, "_read_email_ban_list", lambda a: self.mock_ban_list)
|
||||
|
||||
|
||||
@@ -711,6 +711,89 @@ OUTPUTS_FILTER_EXPRESSION = """
|
||||
</tool>
|
||||
"""
|
||||
|
||||
OUTPUTS_STRUCTURED_LIKE_UNQUALIFIED = """
|
||||
<tool id="id" name="name">
|
||||
<inputs>
|
||||
<conditional name="cond">
|
||||
<param name="cond_param" type="select">
|
||||
<option value="paired">Paired</option>
|
||||
</param>
|
||||
<when value="paired">
|
||||
<param name="input1" type="data_collection" collection_type="paired" format="data" />
|
||||
</when>
|
||||
</conditional>
|
||||
</inputs>
|
||||
<outputs>
|
||||
<collection name="list_output" structured_like="input1" type="paired" inherit_format="true" />
|
||||
</outputs>
|
||||
</tool>
|
||||
"""
|
||||
|
||||
OUTPUTS_STRUCTURED_LIKE_QUALIFIED = """
|
||||
<tool id="id" name="name">
|
||||
<inputs>
|
||||
<conditional name="cond">
|
||||
<param name="cond_param" type="select">
|
||||
<option value="paired">Paired</option>
|
||||
</param>
|
||||
<when value="paired">
|
||||
<param name="input1" type="data_collection" collection_type="paired" format="data" />
|
||||
</when>
|
||||
</conditional>
|
||||
</inputs>
|
||||
<outputs>
|
||||
<collection name="list_output" structured_like="cond|input1" type="paired" inherit_format="true" />
|
||||
</outputs>
|
||||
</tool>
|
||||
"""
|
||||
|
||||
OUTPUTS_STRUCTURED_LIKE_MISSING = """
|
||||
<tool id="id" name="name">
|
||||
<inputs>
|
||||
<param name="input2" type="data" format="data" />
|
||||
</inputs>
|
||||
<outputs>
|
||||
<collection name="list_output" structured_like="nonexistent" type="paired" inherit_format="true" />
|
||||
</outputs>
|
||||
</tool>
|
||||
"""
|
||||
|
||||
OUTPUTS_FORMAT_SOURCE_UNQUALIFIED = """
|
||||
<tool id="id" name="name">
|
||||
<inputs>
|
||||
<conditional name="cond">
|
||||
<param name="cond_param" type="select">
|
||||
<option value="yes">Yes</option>
|
||||
</param>
|
||||
<when value="yes">
|
||||
<param name="input1" type="data" format="data" />
|
||||
</when>
|
||||
</conditional>
|
||||
</inputs>
|
||||
<outputs>
|
||||
<data name="output1" format_source="input1" />
|
||||
</outputs>
|
||||
</tool>
|
||||
"""
|
||||
|
||||
OUTPUTS_FORMAT_SOURCE_QUALIFIED = """
|
||||
<tool id="id" name="name">
|
||||
<inputs>
|
||||
<conditional name="cond">
|
||||
<param name="cond_param" type="select">
|
||||
<option value="yes">Yes</option>
|
||||
</param>
|
||||
<when value="yes">
|
||||
<param name="input1" type="data" format="data" />
|
||||
</when>
|
||||
</conditional>
|
||||
</inputs>
|
||||
<outputs>
|
||||
<data name="output1" format_source="cond|input1" />
|
||||
</outputs>
|
||||
</tool>
|
||||
"""
|
||||
|
||||
# tool xml for repeats linter
|
||||
REPEATS = """
|
||||
<tool id="id" name="name">
|
||||
@@ -1881,6 +1964,41 @@ def test_outputs_filter_expression(lint_ctx):
|
||||
assert not lint_ctx.error_messages
|
||||
|
||||
|
||||
def test_outputs_structured_like_unqualified(lint_ctx):
|
||||
tool_source = get_xml_tool_source(OUTPUTS_STRUCTURED_LIKE_UNQUALIFIED)
|
||||
run_lint_module(lint_ctx, output, tool_source)
|
||||
assert "unqualified structured_like='input1'" in lint_ctx.warn_messages
|
||||
assert "cond|input1" in lint_ctx.warn_messages
|
||||
assert "structured_like" not in lint_ctx.error_messages
|
||||
|
||||
|
||||
def test_outputs_structured_like_qualified(lint_ctx):
|
||||
tool_source = get_xml_tool_source(OUTPUTS_STRUCTURED_LIKE_QUALIFIED)
|
||||
run_lint_module(lint_ctx, output, tool_source)
|
||||
assert "structured_like" not in lint_ctx.warn_messages
|
||||
assert "structured_like" not in lint_ctx.error_messages
|
||||
|
||||
|
||||
def test_outputs_structured_like_missing(lint_ctx):
|
||||
tool_source = get_xml_tool_source(OUTPUTS_STRUCTURED_LIKE_MISSING)
|
||||
run_lint_module(lint_ctx, output, tool_source)
|
||||
assert "does not match any input" in lint_ctx.error_messages
|
||||
|
||||
|
||||
def test_outputs_format_source_unqualified(lint_ctx):
|
||||
tool_source = get_xml_tool_source(OUTPUTS_FORMAT_SOURCE_UNQUALIFIED)
|
||||
run_lint_module(lint_ctx, output, tool_source)
|
||||
assert "unqualified format_source='input1'" in lint_ctx.warn_messages
|
||||
assert "cond|input1" in lint_ctx.warn_messages
|
||||
|
||||
|
||||
def test_outputs_format_source_qualified(lint_ctx):
|
||||
tool_source = get_xml_tool_source(OUTPUTS_FORMAT_SOURCE_QUALIFIED)
|
||||
run_lint_module(lint_ctx, output, tool_source)
|
||||
assert "format_source" not in lint_ctx.warn_messages
|
||||
assert "format_source" not in lint_ctx.error_messages
|
||||
|
||||
|
||||
def test_stdio_default_for_default_profile(lint_ctx):
|
||||
tool_source = get_xml_tool_source(STDIO_DEFAULT_FOR_DEFAULT_PROFILE)
|
||||
run_lint_module(lint_ctx, stdio, tool_source)
|
||||
@@ -2425,7 +2543,7 @@ def test_skip_by_module(lint_ctx):
|
||||
def test_list_linters():
|
||||
linter_names = Linter.list_listers()
|
||||
# make sure to add/remove a test for new/removed linters if this number changes
|
||||
assert len(linter_names) == 143
|
||||
assert len(linter_names) == 145
|
||||
assert "Linter" not in linter_names
|
||||
# make sure that linters from all modules are available
|
||||
for prefix in [
|
||||
|
||||
Reference in New Issue
Block a user