Merge pull request #9610 from mvdbeek/store_tool_source

Startup speedups (XML doc cache, delay non-essential tool param parsing & indexing)
This commit is contained in:
John Chilton
2020-04-27 14:54:32 -04:00
committed by GitHub
29 changed files with 378 additions and 129 deletions
+25
View File
@@ -1002,6 +1002,31 @@
:Type: str
~~~~~~~~~~~~~~~~~~~~~~~
``tool_cache_data_dir``
~~~~~~~~~~~~~~~~~~~~~~~
:Description:
Tool related caching. Fully expanded tools and metadata will be
stored at this path. Per tool_conf cache locations can be
configured in (shed_)tool_conf.xml files using the
tool_cache_data_dir attribute.
:Default: ``tool_cache``
:Type: str
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
``delay_tool_initialization``
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
:Description:
Set this to true to delay parsing of tool inputs and outputs until
they are needed. This results in faster startup times but uses
more memory when using forked Galaxy processes.
:Default: ``false``
:Type: bool
~~~~~~~~~~~~~~~~~~~~~~~
``citation_cache_type``
~~~~~~~~~~~~~~~~~~~~~~~
+9 -1
View File
@@ -24,7 +24,10 @@ from galaxy.managers.users import UserManager
from galaxy.managers.workflows import WorkflowsManager
from galaxy.model.database_heartbeat import DatabaseHeartbeat
from galaxy.model.tags import GalaxyTagHandler
from galaxy.queue_worker import GalaxyQueueWorker
from galaxy.queue_worker import (
GalaxyQueueWorker,
send_local_control_task,
)
from galaxy.tool_shed.galaxy_install.installed_repository_manager import InstalledRepositoryManager
from galaxy.tool_shed.galaxy_install.update_repository_manager import UpdateRepositoryManager
from galaxy.tool_util.deps.views import DependencyResolversView
@@ -65,6 +68,8 @@ class UniverseApplication(config.ConfiguresGalaxyMixin):
logging.basicConfig(level=logging.DEBUG)
log.debug("python path is: %s", ", ".join(sys.path))
self.name = 'galaxy'
# is_webapp will be set to true when building WSGI app
self.is_webapp = False
self.startup_timer = ExecutionTimer()
self.new_installation = False
# Read config file and check for errors
@@ -244,6 +249,9 @@ class UniverseApplication(config.ConfiguresGalaxyMixin):
# Start web stack message handling
self.application_stack.register_postfork_function(self.application_stack.start)
self.application_stack.register_postfork_function(self.queue_worker.bind_and_start)
# Delay toolbox index until after startup
self.application_stack.register_postfork_function(lambda: send_local_control_task(self, 'rebuild_toolbox_search_index'))
self.model.engine.dispose()
+1 -2
View File
@@ -1031,8 +1031,7 @@ class ConfiguresGalaxyMixin(object):
self.container_finder = containers.ContainerFinder(app_info, mulled_resolution_cache=mulled_resolution_cache)
self._set_enabled_container_types()
index_help = getattr(self.config, "index_tool_help", True)
self.toolbox_search = galaxy.tools.search.ToolBoxSearch(self.toolbox, index_help)
self.reindex_tool_search()
self.toolbox_search = galaxy.tools.search.ToolBoxSearch(self.toolbox, index_dir=self.config.tool_search_index_dir, index_help=index_help)
def reindex_tool_search(self):
# Call this when tools are added or removed.
@@ -588,6 +588,17 @@ galaxy:
# generated commands run in sh.
#default_job_shell: /bin/bash
# Tool related caching. Fully expanded tools and metadata will be
# stored at this path. Per tool_conf cache locations can be configured
# in (shed_)tool_conf.xml files using the tool_cache_data_dir
# attribute.
#tool_cache_data_dir: tool_cache
# Set this to true to delay parsing of tool inputs and outputs until
# they are needed. This results in faster startup times but uses more
# memory when using forked Galaxy processes.
#delay_tool_initialization: false
# Citation related caching. Tool citations information maybe fetched
# from external sources such as https://doi.org/ by Galaxy - the
# following parameters can be used to control the caching used to
@@ -62,6 +62,7 @@ dictobj = "*"
nose = "*"
Parsley = "*"
six = "*"
sortedcontainers = "*"
Whoosh = "*"
galaxy_sequence_utils = "*"
"h5py" = "!=2.7.0, !=2.7.1"
@@ -171,6 +171,7 @@ shellescape==3.4.1
simplejson==3.17.0
six==1.11.0
social-auth-core[openidconnect]==3.3.0
sortedcontainers==2.1.0
sqlalchemy-migrate==0.13.0
sqlalchemy-utils==0.36.3
sqlalchemy==1.3.16
+5 -2
View File
@@ -247,8 +247,11 @@ def reload_tool_data_tables(app, **kwargs):
def rebuild_toolbox_search_index(app, **kwargs):
if app.toolbox_search.index_count < app.toolbox._reload_count:
app.reindex_tool_search()
if app.is_webapp:
if app.toolbox_search.index_count < app.toolbox._reload_count:
app.reindex_tool_search()
else:
log.debug("App is not a webapp, not building a search index")
def reload_job_rules(app, **kwargs):
@@ -78,6 +78,7 @@ class ToolValidator(object):
)
try:
tool = create_tool_from_source(config_file=full_path, app=self.app, tool_source=tool_source, repository_id=repository_id, allow_code_files=False)
tool.assert_finalized(raise_if_invalid=True)
valid = True
error_message = None
except KeyError as e:
+2 -2
View File
@@ -14,14 +14,14 @@ from ..fetcher import ToolLocationFetcher
log = logging.getLogger(__name__)
def get_tool_source(config_file=None, xml_tree=None, enable_beta_formats=True, tool_location_fetcher=None):
def get_tool_source(config_file=None, xml_tree=None, enable_beta_formats=True, tool_location_fetcher=None, macro_paths=None):
"""Return a ToolSource object corresponding to supplied source.
The supplied source may be specified as a file path (using the config_file
parameter) or as an XML object loaded with load_tool_with_refereces.
"""
if xml_tree is not None:
return XmlToolSource(xml_tree, source_path=config_file)
return XmlToolSource(xml_tree, source_path=config_file, macro_paths=macro_paths)
elif config_file is None:
raise ValueError("get_tool_source called with invalid config_file None.")
+3 -4
View File
@@ -50,6 +50,9 @@ class XmlToolSource(ToolSource):
self._macro_paths = macro_paths or []
self.legacy_defaults = self.parse_profile() == "16.01"
def to_string(self):
return xml_to_string(self.root)
def parse_version(self):
return self.root.get("version", None)
@@ -711,21 +714,18 @@ def __expand_input_elems(root_elem, prefix=""):
new_prefix = __prefix_join(prefix, name, index=index)
__expand_input_elems(repeat_elem, new_prefix)
__pull_up_params(root_elem, repeat_elem)
root_elem.remove(repeat_elem)
cond_elems = root_elem.findall('conditional')
for cond_elem in cond_elems:
new_prefix = __prefix_join(prefix, cond_elem.get("name"))
__expand_input_elems(cond_elem, new_prefix)
__pull_up_params(root_elem, cond_elem)
root_elem.remove(cond_elem)
section_elems = root_elem.findall('section')
for section_elem in section_elems:
new_prefix = __prefix_join(prefix, section_elem.get("name"))
__expand_input_elems(section_elem, new_prefix)
__pull_up_params(root_elem, section_elem)
root_elem.remove(section_elem)
def __append_prefix_to_params(elem, prefix):
@@ -736,7 +736,6 @@ def __append_prefix_to_params(elem, prefix):
def __pull_up_params(parent_elem, child_elem):
for param_elem in child_elem.findall('param'):
parent_elem.append(param_elem)
child_elem.remove(param_elem)
def __prefix_join(prefix, name, index=None):
+80 -14
View File
@@ -57,6 +57,7 @@ from galaxy.tools.actions import DefaultToolAction
from galaxy.tools.actions.data_manager import DataManagerToolAction
from galaxy.tools.actions.data_source import DataSourceToolAction
from galaxy.tools.actions.model_operations import ModelOperationToolAction
from galaxy.tools.cache import create_cache_region
from galaxy.tools.parameters import (
check_param,
params_from_strings,
@@ -114,6 +115,7 @@ from .execute import (
MappingParameters,
)
log = logging.getLogger(__name__)
REQUIRES_JS_RUNTIME_MESSAGE = ("The tool [%s] requires a nodejs runtime to execute "
@@ -244,6 +246,9 @@ class ToolBox(BaseGalaxyToolBox):
def __init__(self, config_filenames, tool_root_dir, app, save_integrated_tool_panel=True):
self._reload_count = 0
self.tool_location_fetcher = ToolLocationFetcher()
self.cache_regions = {}
if not os.path.exists(app.config.tool_cache_data_dir):
os.makedirs(app.config.tool_cache_data_dir)
# This is here to deal with the old default value, which doesn't make
# sense in an "installed Galaxy" world.
# FIXME: ./
@@ -282,9 +287,22 @@ class ToolBox(BaseGalaxyToolBox):
# Deprecated method, TODO - eliminate calls to this in test/.
return self._tools_by_id
def create_tool(self, config_file, **kwds):
def get_cache_region(self, tool_cache_data_dir):
if tool_cache_data_dir not in self.cache_regions:
self.cache_regions[tool_cache_data_dir] = create_cache_region(tool_cache_data_dir)
return self.cache_regions[tool_cache_data_dir]
def create_tool(self, config_file, tool_cache_data_dir=None, **kwds):
cache = self.get_cache_region(tool_cache_data_dir or self.app.config.tool_cache_data_dir)
tool_source = cache.get_or_create(config_file, creator=self.get_expanded_tool_source, expiration_time=-1, creator_args=((config_file,), {}))
tool = self._create_tool_from_source(tool_source, config_file=config_file, **kwds)
if not self.app.config.delay_tool_initialization:
tool.assert_finalized(raise_if_invalid=True)
return tool
def get_expanded_tool_source(self, config_file):
try:
tool_source = get_tool_source(
return get_tool_source(
config_file,
enable_beta_formats=getattr(self.app.config, "enable_beta_tool_formats", False),
tool_location_fetcher=self.tool_location_fetcher,
@@ -293,7 +311,6 @@ class ToolBox(BaseGalaxyToolBox):
# capture and log parsing errors
global_tool_errors.add_error(config_file, "Tool XML parsing", e)
raise e
return self._create_tool_from_source(tool_source, config_file=config_file, **kwds)
def _create_tool_from_source(self, tool_source, **kwds):
return create_tool_from_source(self.app, tool_source, **kwds)
@@ -447,7 +464,6 @@ class Tool(Dictifiable):
self.repository_id = repository_id
self._allow_code_files = allow_code_files
# setup initial attribute values
self.inputs = OrderedDict()
self.stdio_exit_codes = list()
self.stdio_regexes = list()
self.inputs_by_page = list()
@@ -492,6 +508,9 @@ class Tool(Dictifiable):
self.populate_resource_parameters(tool_source)
self.tool_errors = None
# Parse XML element containing configuration
self.tool_source = tool_source
self._is_workflow_compatible = None
self.finalized = False
try:
self.parse(tool_source, guid=guid, dynamic=dynamic)
except Exception as e:
@@ -502,6 +521,50 @@ class Tool(Dictifiable):
if self.app.name == 'galaxy':
self.job_search = JobSearch(app=self.app)
def __getattr__(self, name):
lazy_attributes = {
'action',
'check_values',
'display_by_page',
'enctype',
'has_multiple_pages',
'inputs',
'inputs_by_page',
'last_page',
'method',
'npages',
'nginx_upload',
'target',
'template_macro_params',
'outputs',
'output_collections'
}
if name in lazy_attributes:
self.assert_finalized()
return getattr(self, name)
raise AttributeError(name)
def assert_finalized(self, raise_if_invalid=False):
if self.finalized is False:
try:
self.parse_inputs(self.tool_source)
self.parse_outputs(self.tool_source)
self.finalized = True
except Exception:
toolbox = getattr(self.app, 'toolbox', None)
if toolbox:
toolbox.remove_tool_by_id(self.id)
if raise_if_invalid:
raise
else:
log.warning("An error occured while parsing the tool wrapper xml, the tool is not functional", exc_info=True)
def remove_from_cache(self):
source_path = self.tool_source._source_path
if source_path:
for region in self.app.toolbox.cache_regions.values():
region.delete(source_path)
@property
def history_manager(self):
return self.app.history_manager
@@ -528,7 +591,7 @@ class Tool(Dictifiable):
def tool_versions(self):
# If we have versions, return them.
if self.lineage:
return self.lineage.tool_versions
return list(self.lineage.tool_versions)
else:
return []
@@ -803,15 +866,9 @@ class Tool(Dictifiable):
self.provided_metadata_file = tool_source.parse_provided_metadata_file()
self.provided_metadata_style = tool_source.parse_provided_metadata_style()
# Parse tool inputs (if there are any required)
self.parse_inputs(tool_source)
# Parse tool help
self.parse_help(tool_source)
# Description of outputs produced by an invocation of the tool
self.parse_outputs(tool_source)
# Parse result handling for tool exit codes and stdout/stderr messages:
self.parse_stdio(tool_source)
@@ -844,8 +901,6 @@ class Tool(Dictifiable):
self.citations = self._parse_citations(tool_source)
self.xrefs = tool_source.parse_xrefs()
# Determine if this tool can be used in workflows
self.is_workflow_compatible = self.check_workflow_compatible(tool_source)
self.__parse_trackster_conf(tool_source)
# Record macro paths so we can reload a tool if any of its macro has changes
self._macro_paths = tool_source.macro_paths()
@@ -930,6 +985,7 @@ class Tool(Dictifiable):
@property
def tests(self):
self.assert_finalized()
if not self.__tests_populated:
tests_source = self.__tests_source
if tests_source:
@@ -1015,6 +1071,7 @@ class Tool(Dictifiable):
This implementation supports multiple pages and grouping constructs.
"""
# Load parameters (optional)
self.inputs = OrderedDict()
pages = tool_source.parse_input_pages()
enctypes = set()
if pages.inputs_defined:
@@ -1352,6 +1409,15 @@ class Tool(Dictifiable):
else:
return self.outputs.get(name, None)
@property
def is_workflow_compatible(self):
is_workflow_compatible = self._is_workflow_compatible
if is_workflow_compatible is None:
is_workflow_compatible = self.check_workflow_compatible(self.tool_source)
if self.finalized:
self._is_workflow_compatible = is_workflow_compatible
return is_workflow_compatible
def check_workflow_compatible(self, tool_source):
"""
Determine if a tool can be used in workflows. External tools and the
@@ -1359,7 +1425,7 @@ class Tool(Dictifiable):
"""
# Multiple page tools are not supported -- we're eliminating most
# of these anyway
if self.has_multiple_pages:
if self.finalized and self.has_multiple_pages:
return False
# This is probably the best bet for detecting external web tools
# right now
+121 -13
View File
@@ -1,19 +1,104 @@
import json
import logging
import os
from collections import defaultdict
from threading import Lock
from dogpile.cache import make_region
from dogpile.cache.api import (
CachedValue,
NO_VALUE,
)
from dogpile.cache.backends.file import AbstractFileLock
from dogpile.cache.proxy import ProxyBackend
from dogpile.util import ReadWriteMutex
from lxml import etree
from sqlalchemy.orm import (
defer,
joinedload,
)
from galaxy.tool_util.parser import get_tool_source
from galaxy.util import unicodify
from galaxy.util.hash_util import md5_hash_file
log = logging.getLogger(__name__)
class JSONBackend(ProxyBackend):
def set(self, key, value):
with self.proxied._dbm_file(True) as dbm:
dbm[key] = json.dumps({'metadata': value.metadata, 'payload': self.value_encode(value), 'macro_paths': value.payload.macro_paths()})
def get(self, key):
with self.proxied._dbm_file(False) as dbm:
if hasattr(dbm, "get"):
value = dbm.get(key, NO_VALUE)
else:
# gdbm objects lack a .get method
try:
value = dbm[key]
except KeyError:
value = NO_VALUE
if value is not NO_VALUE:
value = self.value_decode(key, value)
return value
def value_decode(self, k, v):
if not v or v is NO_VALUE:
return NO_VALUE
# v is returned as bytestring, so we need to `unicodify` on python < 3.6 before we can use json.loads
v = json.loads(unicodify(v))
payload = get_tool_source(
config_file=k,
xml_tree=etree.ElementTree(etree.fromstring(v['payload'].encode('utf-8'))),
macro_paths=v['macro_paths']
)
return CachedValue(metadata=v['metadata'], payload=payload)
def value_encode(self, v):
return unicodify(v.payload.to_string())
class MutexLock(AbstractFileLock):
def __init__(self, filename):
self.mutex = ReadWriteMutex()
def acquire_read_lock(self, wait):
# No need for read lock. It is supposed to prevent the "dogpile" effect
# where multiple functions each create the cached resource, but I don't
# think we care.
return True
def acquire_write_lock(self, wait):
ret = self.mutex.acquire_write_lock(wait)
return wait or ret
def release_read_lock(self):
return True
def release_write_lock(self):
return self.mutex.release_write_lock()
def create_cache_region(tool_cache_data_dir):
if not os.path.exists(tool_cache_data_dir):
os.makedirs(tool_cache_data_dir)
region = make_region()
region.configure(
'dogpile.cache.dbm',
arguments={
"filename": os.path.join(tool_cache_data_dir, "cache.dbm"),
"lock_factory": MutexLock,
},
expiration_time=-1,
wrap=[JSONBackend],
replace_existing_backend=True,
)
return region
class ToolCache(object):
"""
Cache tool definitions to allow quickly reloading the whole
@@ -26,10 +111,16 @@ class ToolCache(object):
self._tools_by_path = {}
self._tool_paths_by_id = {}
self._macro_paths_by_id = {}
self._mod_time_by_path = {}
self._new_tool_ids = set()
self._removed_tool_ids = set()
self._removed_tools_by_path = {}
self._hashes_initialized = False
def assert_hashes_initialized(self):
if not self._hashes_initialized:
for tool_hash in self._hash_by_tool_paths.values():
tool_hash.hash
self._hashes_initialized = True
def cleanup(self):
"""
@@ -40,14 +131,16 @@ class ToolCache(object):
removed_tool_ids = []
try:
with self._lock:
paths_to_cleanup = {path: tool.all_ids for path, tool in self._tools_by_path.items() if self._should_cleanup(path)}
for config_filename, tool_ids in paths_to_cleanup.items():
paths_to_cleanup = {(path, tool) for path, tool in self._tools_by_path.items() if self._should_cleanup(path)}
for config_filename, tool in paths_to_cleanup:
tool.remove_from_cache()
del self._hash_by_tool_paths[config_filename]
if os.path.exists(config_filename):
# This tool has probably been broken while editing on disk
# We record it here, so that we can recover it
self._removed_tools_by_path[config_filename] = self._tools_by_path[config_filename]
del self._tools_by_path[config_filename]
tool_ids = tool.all_ids
for tool_id in tool_ids:
if tool_id in self._tool_paths_by_id:
del self._tool_paths_by_id[tool_id]
@@ -70,13 +163,14 @@ class ToolCache(object):
if not os.path.exists(config_filename):
return True
new_mtime = os.path.getmtime(config_filename)
if self._mod_time_by_path.get(config_filename) < new_mtime:
if md5_hash_file(config_filename) != self._hash_by_tool_paths.get(config_filename):
tool_hash = self._hash_by_tool_paths.get(config_filename)
if tool_hash.modtime < new_mtime:
if md5_hash_file(config_filename) != tool_hash.hash:
return True
tool = self._tools_by_path[config_filename]
for macro_path in tool._macro_paths:
new_mtime = os.path.getmtime(macro_path)
if self._mod_time_by_path.get(macro_path) < new_mtime:
if self._hash_by_tool_paths.get(macro_path).modtime < new_mtime:
return True
return False
@@ -98,23 +192,21 @@ class ToolCache(object):
del self._hash_by_tool_paths[config_filename]
del self._tool_paths_by_id[tool_id]
del self._tools_by_path[config_filename]
del self._mod_time_by_path[config_filename]
if tool_id in self._new_tool_ids:
self._new_tool_ids.remove(tool_id)
def cache_tool(self, config_filename, tool):
tool_hash = md5_hash_file(config_filename)
if tool_hash is None:
return
tool_id = str(tool.id)
# We defer hashing of the config file if we haven't called assert_hashes_initialized.
# This allows startup to occur without having to read in and hash all tool and macro files
lazy_hash = not self._hashes_initialized
with self._lock:
self._hash_by_tool_paths[config_filename] = tool_hash
self._mod_time_by_path[config_filename] = os.path.getmtime(config_filename)
self._hash_by_tool_paths[config_filename] = ToolHash(config_filename, lazy_hash=lazy_hash)
self._tool_paths_by_id[tool_id] = config_filename
self._tools_by_path[config_filename] = tool
self._new_tool_ids.add(tool_id)
for macro_path in tool._macro_paths:
self._mod_time_by_path[macro_path] = os.path.getmtime(macro_path)
self._hash_by_tool_paths[macro_path] = ToolHash(macro_path, lazy_hash=lazy_hash)
if tool_id not in self._macro_paths_by_id:
self._macro_paths_by_id[tool_id] = {macro_path}
else:
@@ -130,6 +222,22 @@ class ToolCache(object):
self._removed_tools_by_path = {}
class ToolHash(object):
def __init__(self, path, modtime=None, lazy_hash=False):
self.path = path
self.modtime = modtime or os.path.getmtime(path)
self._tool_hash = None
if not lazy_hash:
self.hash
@property
def hash(self):
if self._tool_hash is None:
self._tool_hash = md5_hash_file(self.path)
return self._tool_hash
class ToolShedRepositoryCache(object):
"""
Cache installed ToolShedRepository objects.
+37 -30
View File
@@ -5,24 +5,24 @@ or searching related parts it is deeply recommended to read
through the library docs at https://whoosh.readthedocs.io.
"""
import logging
import os
import re
import tempfile
from whoosh import analysis
from whoosh import (
analysis,
index,
)
from whoosh.analysis import StandardAnalyzer
from whoosh.fields import (
ID,
KEYWORD,
Schema,
STORED,
TEXT
)
from whoosh.filedb.filestore import (
FileStorage,
RamStorage
)
from whoosh.qparser import MultifieldParser
from whoosh.qparser import OrGroup
from whoosh.scoring import BM25F
from whoosh.writing import AsyncWriter
from galaxy.util import ExecutionTimer
from galaxy.web.framework.helpers import to_unicode
@@ -30,14 +30,27 @@ from galaxy.web.framework.helpers import to_unicode
log = logging.getLogger(__name__)
def get_or_create_index(index_dir, schema):
if not os.path.exists(index_dir):
os.makedirs(index_dir)
if index.exists_in(index_dir):
idx = index.open_dir(index_dir)
try:
assert idx.schema == schema
return idx
except AssertionError:
log.warning("Index at '%s' uses outdated schema, creating new index", index_dir)
return index.create_in(index_dir, schema=schema)
class ToolBoxSearch(object):
"""
Support searching tools in a toolbox. This implementation uses
the Whoosh search library.
"""
def __init__(self, toolbox, index_help=True):
self.schema = Schema(id=STORED,
def __init__(self, toolbox, index_dir=None, index_help=True):
self.schema = Schema(id=ID(stored=True),
stub=KEYWORD,
name=TEXT(analyzer=analysis.SimpleAnalyzer()),
description=TEXT,
@@ -45,8 +58,9 @@ class ToolBoxSearch(object):
help=TEXT,
labels=KEYWORD)
self.rex = analysis.RegexTokenizer()
self.index_dir = index_dir
self.toolbox = toolbox
self.storage, self.index = self._index_setup()
self.index = self._index_setup()
# We keep track of how many times the tool index has been rebuilt.
# We start at -1, so that after the first index the count is at 0,
# which is the same as the toolbox reload count. This way we can skip
@@ -54,11 +68,7 @@ class ToolBoxSearch(object):
self.index_count = -1
def _index_setup(self):
RamStorage.temp_storage = _temp_storage
# Works around https://bitbucket.org/mchaput/whoosh/issues/391/race-conditions-with-temp-storage
storage = RamStorage()
index = storage.create_index(self.schema)
return storage, index
return get_or_create_index(index_dir=self.index_dir, schema=self.schema)
def build_index(self, tool_cache, index_help=True):
"""
@@ -68,15 +78,18 @@ class ToolBoxSearch(object):
log.debug('Starting to build toolbox index.')
self.index_count += 1
execution_timer = ExecutionTimer()
writer = self.index.writer()
for tool_id in tool_cache._removed_tool_ids:
writer.delete_by_term('id', tool_id)
for tool_id in tool_cache._new_tool_ids:
tool = tool_cache.get_tool_by_id(tool_id)
if tool and tool.is_latest_version:
add_doc_kwds = self._create_doc(tool_id=tool_id, tool=tool, index_help=index_help)
writer.add_document(**add_doc_kwds)
writer.commit()
with self.index.reader() as reader:
# Index ocasionally contains empty stored fields
indexed_tool_ids = {f['id'] for f in reader.all_stored_fields() if f}
tool_ids_to_remove = (indexed_tool_ids - set(tool_cache._tool_paths_by_id.keys())).union(tool_cache._removed_tool_ids)
with AsyncWriter(self.index) as writer:
for tool_id in tool_ids_to_remove:
writer.delete_by_term('id', tool_id)
for tool_id in tool_cache._new_tool_ids - indexed_tool_ids:
tool = tool_cache.get_tool_by_id(tool_id)
if tool and tool.is_latest_version:
add_doc_kwds = self._create_doc(tool_id=tool_id, tool=tool, index_help=index_help)
writer.add_document(**add_doc_kwds)
log.debug("Toolbox index finished %s", execution_timer)
def _create_doc(self, tool_id, tool, index_help=True):
@@ -181,9 +194,3 @@ class ToolBoxSearch(object):
hits_with_score = sorted(hits_with_score.items(), key=lambda x: x[1], reverse=True)
# Return the tool ids
return [item[0] for item in hits_with_score[0:int(tool_search_limit)]]
def _temp_storage(self, name=None):
path = tempfile.mkdtemp()
tempstore = FileStorage(path)
return tempstore.create()
+20 -30
View File
@@ -11,7 +11,6 @@ from errno import ENOENT
from xml.etree.ElementTree import ParseError
from markupsafe import escape
from six import iteritems
from six.moves.urllib.parse import urlparse
from galaxy.exceptions import (
@@ -106,13 +105,8 @@ class AbstractToolBox(Dictifiable, ManagesIntegratedToolPanelMixin):
# (e.g., shed_tool_conf.xml) files include the tool_path attribute within the <toolbox> tag.
self._tool_root_dir = tool_root_dir
self.app = app
if hasattr(self.app, 'watchers'):
self._tool_watcher = self.app.watchers.tool_watcher
self._tool_config_watcher = self.app.watchers.tool_config_watcher
else:
# Toolbox is loaded but not used during toolshed tests
self._tool_watcher = None
self._tool_config_watcher = None
self._tool_watcher = self.app.watchers.tool_watcher
self._tool_config_watcher = self.app.watchers.tool_config_watcher
self._filter_factory = FilterFactory(self)
self._tool_tag_manager = tool_tag_manager(app)
self._init_tools_from_configs(config_filenames)
@@ -197,6 +191,7 @@ class AbstractToolBox(Dictifiable, ManagesIntegratedToolPanelMixin):
return
raise
tool_path = tool_conf_source.parse_tool_path()
tool_cache_data_dir = tool_conf_source.parse_tool_cache_data_dir()
parsing_shed_tool_conf = tool_conf_source.is_shed_tool_conf()
if parsing_shed_tool_conf:
# Keep an in-memory list of xml elements to enable persistence of the changing tool config.
@@ -214,10 +209,10 @@ class AbstractToolBox(Dictifiable, ManagesIntegratedToolPanelMixin):
self.load_item(
item,
tool_path=tool_path,
tool_cache_data_dir=tool_cache_data_dir,
load_panel_dict=load_panel_dict,
guid=item.get('guid'),
index=index,
internal=True
)
if parsing_shed_tool_conf:
@@ -245,29 +240,24 @@ class AbstractToolBox(Dictifiable, ManagesIntegratedToolPanelMixin):
self._tools_by_uuid[dynamic_tool.uuid] = tool
return tool
def load_item(self, item, tool_path, panel_dict=None, integrated_panel_dict=None, load_panel_dict=True, guid=None, index=None, internal=False):
def load_item(self, item, tool_path, panel_dict=None, integrated_panel_dict=None, load_panel_dict=True, guid=None, index=None, tool_cache_data_dir=None):
with self.app._toolbox_lock:
item = ensure_tool_conf_item(item)
item_type = item.type
if item_type not in ['tool', 'section'] and not internal:
# External calls from tool shed code cannot load labels or tool
# directories.
return
if panel_dict is None:
panel_dict = self._tool_panel
if integrated_panel_dict is None:
integrated_panel_dict = self._integrated_tool_panel
if item_type == 'tool':
self._load_tool_tag_set(item, panel_dict=panel_dict, integrated_panel_dict=integrated_panel_dict, tool_path=tool_path, load_panel_dict=load_panel_dict, guid=guid, index=index, internal=internal)
self._load_tool_tag_set(item, panel_dict=panel_dict, integrated_panel_dict=integrated_panel_dict, tool_path=tool_path, load_panel_dict=load_panel_dict, guid=guid, index=index, tool_cache_data_dir=tool_cache_data_dir)
elif item_type == 'workflow':
self._load_workflow_tag_set(item, panel_dict=panel_dict, integrated_panel_dict=integrated_panel_dict, load_panel_dict=load_panel_dict, index=index)
elif item_type == 'section':
self._load_section_tag_set(item, tool_path=tool_path, load_panel_dict=load_panel_dict, index=index, internal=internal)
self._load_section_tag_set(item, tool_path=tool_path, load_panel_dict=load_panel_dict, index=index, tool_cache_data_dir=tool_cache_data_dir)
elif item_type == 'label':
self._load_label_tag_set(item, panel_dict=panel_dict, integrated_panel_dict=integrated_panel_dict, load_panel_dict=load_panel_dict, index=index)
elif item_type == 'tool_dir':
self._load_tooldir_tag_set(item, panel_dict, tool_path, integrated_panel_dict, load_panel_dict=load_panel_dict)
self._load_tooldir_tag_set(item, panel_dict, tool_path, integrated_panel_dict, load_panel_dict=load_panel_dict, tool_cache_data_dir=tool_cache_data_dir)
def get_shed_config_dict_by_filename(self, filename):
filename = os.path.abspath(filename)
@@ -602,7 +592,7 @@ class AbstractToolBox(Dictifiable, ManagesIntegratedToolPanelMixin):
return []
def tools(self):
return iteritems(self._tools_by_id)
return self._tools_by_id.copy().items()
def dynamic_confs(self, include_migrated_tool_conf=False):
confs = []
@@ -638,7 +628,7 @@ class AbstractToolBox(Dictifiable, ManagesIntegratedToolPanelMixin):
def _path_template_kwds(self):
return {}
def _load_tool_tag_set(self, item, panel_dict, integrated_panel_dict, tool_path, load_panel_dict, guid=None, index=None, internal=False):
def _load_tool_tag_set(self, item, panel_dict, integrated_panel_dict, tool_path, load_panel_dict, guid=None, index=None, tool_cache_data_dir=None):
try:
path_template = item.get("file")
template_kwds = self._path_template_kwds()
@@ -664,9 +654,9 @@ class AbstractToolBox(Dictifiable, ManagesIntegratedToolPanelMixin):
# The shed tool is in the install database
# Only load tools if the repository is not deactivated or uninstalled.
can_load_into_panel_dict = not tool_shed_repository.deleted
tool = self.load_tool(concrete_path, guid=guid, tool_shed_repository=tool_shed_repository, use_cached=False)
tool = self.load_tool(concrete_path, guid=guid, tool_shed_repository=tool_shed_repository, use_cached=False, tool_cache_data_dir=tool_cache_data_dir)
if not tool: # tool was not in cache and is not a tool shed tool.
tool = self.load_tool(concrete_path, use_cached=False)
tool = self.load_tool(concrete_path, use_cached=False, tool_cache_data_dir=tool_cache_data_dir)
if string_as_bool(item.get('hidden', False)):
tool.hidden = True
key = 'tool_%s' % str(tool.id)
@@ -772,7 +762,7 @@ class AbstractToolBox(Dictifiable, ManagesIntegratedToolPanelMixin):
panel_dict[key] = label
integrated_panel_dict.update_or_append(index, key, label)
def _load_section_tag_set(self, item, tool_path, load_panel_dict, index=None, internal=False):
def _load_section_tag_set(self, item, tool_path, load_panel_dict, index=None, tool_cache_data_dir=None):
key = item.get("id")
if key in self._tool_panel:
section = self._tool_panel[key]
@@ -795,7 +785,7 @@ class AbstractToolBox(Dictifiable, ManagesIntegratedToolPanelMixin):
load_panel_dict=load_panel_dict,
guid=sub_item.get('guid'),
index=sub_index,
internal=internal,
tool_cache_data_dir=tool_cache_data_dir,
)
# Ensure each tool's section is stored
@@ -810,16 +800,16 @@ class AbstractToolBox(Dictifiable, ManagesIntegratedToolPanelMixin):
# Always load sections into the integrated_tool_panel.
self._integrated_tool_panel.update_or_append(index, key, integrated_section)
def _load_tooldir_tag_set(self, item, elems, tool_path, integrated_elems, load_panel_dict):
def _load_tooldir_tag_set(self, item, elems, tool_path, integrated_elems, load_panel_dict, tool_cache_data_dir=None):
directory = os.path.join(tool_path, item.get("dir"))
recursive = string_as_bool(item.get("recursive", True))
self.__watch_directory(directory, elems, integrated_elems, load_panel_dict, recursive, force_watch=True)
self.__watch_directory(directory, elems, integrated_elems, load_panel_dict, recursive, force_watch=True, tool_cache_data_dir=tool_cache_data_dir)
def __watch_directory(self, directory, elems, integrated_elems, load_panel_dict, recursive, force_watch=False):
def __watch_directory(self, directory, elems, integrated_elems, load_panel_dict, recursive, force_watch=False, tool_cache_data_dir=None):
def quick_load(tool_file, async_load=True):
try:
tool = self.load_tool(tool_file)
tool = self.load_tool(tool_file, tool_cache_data_dir)
self.__add_tool(tool, load_panel_dict, elems)
# Always load the tool into the integrated_panel_dict, or it will not be included in the integrated_tool_panel.xml file.
key = 'tool_%s' % str(tool.id)
@@ -850,7 +840,7 @@ class AbstractToolBox(Dictifiable, ManagesIntegratedToolPanelMixin):
if (tool_loaded or force_watch) and self._tool_watcher:
self._tool_watcher.watch_directory(directory, quick_load)
def load_tool(self, config_file, guid=None, tool_shed_repository=None, use_cached=False, **kwds):
def load_tool(self, config_file, guid=None, tool_shed_repository=None, use_cached=False, tool_cache_data_dir=None, **kwds):
"""Load a single tool from the file named by `config_file` and return an instance of `Tool`."""
# Parse XML configuration file and get the root element
tool = None
@@ -858,7 +848,7 @@ class AbstractToolBox(Dictifiable, ManagesIntegratedToolPanelMixin):
tool = self.load_tool_from_cache(config_file)
if not tool or guid and guid != tool.guid:
try:
tool = self.create_tool(config_file=config_file, tool_shed_repository=tool_shed_repository, guid=guid, **kwds)
tool = self.create_tool(config_file=config_file, tool_shed_repository=tool_shed_repository, guid=guid, tool_cache_data_dir=tool_cache_data_dir, **kwds)
except Exception:
# If the tool is broken but still exists we can load it from the cache
tool = self.load_tool_from_cache(config_file, recover_tool=True)
@@ -1,6 +1,7 @@
import threading
import packaging.version
from sortedcontainers import SortedSet
from galaxy.util.tool_version import remove_version_from_guid
@@ -39,11 +40,7 @@ class ToolLineage(object):
def __init__(self, tool_id, **kwds):
self.tool_id = tool_id
self._tool_versions = set()
@property
def tool_versions(self):
return sorted(self._tool_versions, key=packaging.version.parse)
self.tool_versions = SortedSet(key=packaging.version.parse)
@property
def tool_ids(self):
@@ -64,7 +61,7 @@ class ToolLineage(object):
def register_version(self, tool_version):
assert tool_version is not None
self._tool_versions.add(str(tool_version))
self.tool_versions.add(str(tool_version))
def get_versions(self):
"""
+6
View File
@@ -43,6 +43,9 @@ class XmlToolConfSource(ToolConfSource):
def parse_tool_path(self):
return self.root.get('tool_path')
def parse_tool_cache_data_dir(self):
return self.root.get('tool_cache_data_dir')
def parse_items(self):
return [ensure_tool_conf_item(_) for _ in self.root]
@@ -65,6 +68,9 @@ class YamlToolConfSource(ToolConfSource):
def parse_tool_path(self):
return self.as_dict.get('tool_path')
def parse_tool_cache_data_dir(self):
return self.as_dict.get('tool_cache_data_dir')
def parse_items(self):
return [ToolConfItem.from_dict(_) for _ in self.as_dict.get('items')]
+2
View File
@@ -99,6 +99,8 @@ class ToolConfWatcher(object):
def check(self):
"""Check for changes in self.paths or self.cache and call the event handler."""
hashes = {}
if self.cache:
self.cache.assert_hashes_initialized()
while self._active and not self.exit.isSet():
do_reload = False
drop_on_next_loop = set()
+5 -1
View File
@@ -5,13 +5,17 @@ import os
from collections import OrderedDict
import yaml
try:
from yaml import CSafeLoader as SafeLoader
except ImportError:
from yaml import SafeLoader
from yaml.constructor import ConstructorError
log = logging.getLogger(__name__)
class OrderedLoader(yaml.SafeLoader):
class OrderedLoader(SafeLoader):
# This class was pulled out of ordered_load() for the sake of
# mocking __init__ in a unit test.
def __init__(self, stream):
+1
View File
@@ -79,6 +79,7 @@ class WebApplication(base.WebApplication):
def __init__(self, galaxy_app, session_cookie='galaxysession', name=None):
self.name = name
base.WebApplication.__init__(self)
galaxy_app.is_webapp = True
self.set_transaction_factory(lambda e: self.transaction_chooser(e, galaxy_app, session_cookie))
# Mako support
self.mako_template_lookup = self.create_mako_template_lookup(galaxy_app, name)
-1
View File
@@ -204,7 +204,6 @@ uwsgi_app_factory = uwsgi_app
def postfork_setup():
from galaxy.app import app
app.queue_worker.bind_and_start()
app.application_stack.log_startup()
@@ -749,6 +749,33 @@ mapping:
should be disabled. Containerized jobs always use /bin/sh - so more maximum
portability tool authors should assume generated commands run in sh.
tool_cache_data_dir:
type: str
default: tool_cache
path_resolves_to: data_dir
required: false
desc: |
Tool related caching. Fully expanded tools and metadata will be stored at this path.
Per tool_conf cache locations can be configured in (shed_)tool_conf.xml files using
the tool_cache_data_dir attribute.
tool_search_index_dir:
type: str
default: tool_search_index
path_resolves_to: data_dir
required: false
desc:
Directory in which the toolbox search index is stored.
delay_tool_initialization:
type: bool
default: false
required: false
desc: |
Set this to true to delay parsing of tool inputs and outputs until they are needed.
This results in faster startup times but uses more memory when using forked Galaxy
processes.
citation_cache_type:
type: str
default: file
+4 -3
View File
@@ -1251,9 +1251,10 @@ class ToolModule(WorkflowModule):
self.tool_version = tool_version
self.tool_uuid = tool_uuid
self.tool = trans.app.toolbox.get_tool(tool_id, tool_version=tool_version, exact=exact_tools, tool_uuid=tool_uuid)
if self.tool and tool_version and exact_tools and str(self.tool.version) != str(tool_version):
log.info("Exact tool specified during workflow module creation for [%s] but couldn't find correct version [%s]." % (tool_id, tool_version))
self.tool = None
if self.tool:
if tool_version and exact_tools and str(self.tool.version) != str(tool_version):
log.info("Exact tool specified during workflow module creation for [%s] but couldn't find correct version [%s]." % (tool_id, tool_version))
self.tool = None
self.post_job_actions = {}
self.runtime_post_job_actions = {}
self.workflow_outputs = []
+4 -15
View File
@@ -2,11 +2,11 @@ import logging
import os
from mercurial import hg, ui
from whoosh import index
from whoosh.writing import AsyncWriter
import tool_shed.webapp.model.mapping as ts_mapping
from galaxy.tool_util.loader_directory import load_tool_elements_from_path
from galaxy.tools.search import get_or_create_index
from galaxy.util import (
directory_hash_id,
ExecutionTimer,
@@ -21,24 +21,13 @@ from tool_shed.webapp.search.tool_search import schema as tool_schema
log = logging.getLogger(__name__)
def get_or_create_index(whoosh_index_dir):
def _get_or_create_index(whoosh_index_dir):
tool_index_dir = os.path.join(whoosh_index_dir, 'tools')
if not os.path.exists(whoosh_index_dir):
os.makedirs(whoosh_index_dir)
if not os.path.exists(tool_index_dir):
os.makedirs(tool_index_dir)
return _get_or_create_index(whoosh_index_dir, repo_schema), _get_or_create_index(tool_index_dir, tool_schema)
def _get_or_create_index(index_dir, schema):
if index.exists_in(index_dir):
idx = index.open_dir(index_dir)
try:
assert idx.schema == schema
return idx
except AssertionError:
log.warning("Index at '%s' uses outdated schema, creating new index", index_dir)
return index.create_in(index_dir, schema=schema)
return get_or_create_index(whoosh_index_dir, repo_schema), get_or_create_index(tool_index_dir, tool_schema)
def build_index(whoosh_index_dir, file_path, hgweb_config_dir, dburi, **kwargs):
@@ -50,7 +39,7 @@ def build_index(whoosh_index_dir, file_path, hgweb_config_dir, dburi, **kwargs):
"""
model = ts_mapping.init(file_path, dburi, engine_options={}, create_tables=False)
sa_session = model.context.current
repo_index, tool_index = get_or_create_index(whoosh_index_dir)
repo_index, tool_index = _get_or_create_index(whoosh_index_dir)
repo_index_writer = AsyncWriter(repo_index)
tool_index_writer = AsyncWriter(tool_index)
+2 -3
View File
@@ -8,7 +8,6 @@ import galaxy.tools.data
import tool_shed.repository_registry
import tool_shed.repository_types.registry
import tool_shed.webapp.model
from galaxy import tools
from galaxy.config import configure_logging
from galaxy.model.tags import CommunityTagHandler
from galaxy.security import idencoding
@@ -27,6 +26,8 @@ class UniverseApplication(object):
def __init__(self, **kwd):
log.debug("python path is: %s", ", ".join(sys.path))
self.name = "tool_shed"
# will be overwritten when building WSGI app
self.is_webapp = False
# Read the tool_shed.ini configuration file and check for errors.
self.config = config.Configuration(**kwd)
self.config.check()
@@ -64,9 +65,7 @@ class UniverseApplication(object):
# Citation manager needed to load tools.
from galaxy.managers.citations import CitationsManager
self.citations_manager = CitationsManager(self)
# The Tool Shed makes no use of a Galaxy toolbox, but this attribute is still required.
self.use_tool_dependency_resolution = False
self.toolbox = tools.ToolBox([], self.config.tool_path, self)
# Initialize the Tool Shed security agent.
self.security_agent = self.model.security_agent
# The Tool Shed makes no use of a quota, but this attribute is still required.
-1
View File
@@ -108,7 +108,6 @@ def load_galaxy_app(
**kwds
)
app.database_heartbeat.start()
app.queue_worker.bind_and_start()
app.application_stack.log_startup()
return app
+2
View File
@@ -93,10 +93,12 @@ class ExpectedValues:
'shed_tool_data_path': self._in_root_dir('tool-data'),
'shed_tool_data_table_config': self._in_managed_config_dir('shed_tool_data_table_conf.xml'),
'template_cache_path': self._in_data_dir('compiled_templates'),
'tool_cache_data_dir': self._in_data_dir('tool_cache'),
'tool_config_file': self._in_sample_dir('tool_conf.xml.sample'),
'tool_data_path': self._in_root_dir('tool-data'),
'tool_data_table_config_path': self._in_sample_dir('tool_data_table_conf.xml.sample'),
'tool_path': self._in_root_dir('tools'),
'tool_search_index_dir': self._in_data_dir('tool_search_index'),
'tool_sheds_config_file': self._in_config_dir('tool_sheds_conf.xml'),
'tool_test_data_directories': self._in_root_dir('test-data'),
'user_preferences_extra_conf_path': self._in_config_dir('user_preferences_extra_conf.yml'),
+1
View File
@@ -99,6 +99,7 @@ class UsesTools(object):
tool_source = get_tool_source(self.tool_file)
try:
self.tool = create_tool_from_source(self.app, tool_source, config_file=self.tool_file)
self.tool.assert_finalized()
except Exception:
self.tool = None
if getattr(self, "tool_action", None and self.tool):
+2
View File
@@ -170,6 +170,8 @@ class MockAppConfig(Bunch):
# set by MockDir
self.root = root
self.tool_cache_data_dir = os.path.join(root, 'tool_cache')
self.delay_tool_initialization = True
self.config_file = None
+2 -1
View File
@@ -456,7 +456,8 @@ def __mock_tool(
output_type='data')},
params_from_strings=mock.Mock(),
check_and_update_param_values=mock.Mock(),
to_json=_to_json
to_json=_to_json,
assert_finalized=lambda: None,
)
return tool