mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge pull request #8114 from jmchilton/galaxy-job-execution
Add galaxy-job-execution package with galaxy-set-metadata script.
This commit is contained in:
+61
@@ -4,6 +4,7 @@ import logging
|
||||
import operator
|
||||
import os
|
||||
import re
|
||||
from tempfile import NamedTemporaryFile
|
||||
|
||||
import galaxy.model
|
||||
from galaxy.model.dataset_collections import builder
|
||||
@@ -463,3 +464,63 @@ def _compose(f, g):
|
||||
|
||||
DEFAULT_DATASET_COLLECTOR = DatasetCollector(DEFAULT_DATASET_COLLECTOR_DESCRIPTION)
|
||||
DEFAULT_TOOL_PROVIDED_DATASET_COLLECTOR = ToolMetadataDatasetCollector(ToolProvidedMetadataDatasetCollection())
|
||||
|
||||
|
||||
def read_exit_code_from(exit_code_file, id_tag):
|
||||
"""Read exit code reported for a Galaxy job."""
|
||||
try:
|
||||
# This should be an 8-bit exit code, but read ahead anyway:
|
||||
exit_code_str = open(exit_code_file, "r").read(32)
|
||||
except Exception:
|
||||
# By default, the exit code is 0, which typically indicates success.
|
||||
exit_code_str = "0"
|
||||
|
||||
try:
|
||||
# Decode the exit code. If it's bogus, then just use 0.
|
||||
exit_code = int(exit_code_str)
|
||||
except ValueError:
|
||||
galaxy_id_tag = id_tag
|
||||
log.warning("(%s) Exit code '%s' invalid. Using 0." % (galaxy_id_tag, exit_code_str))
|
||||
exit_code = 0
|
||||
|
||||
return exit_code
|
||||
|
||||
|
||||
def default_exit_code_file(files_dir, id_tag):
|
||||
return os.path.join(files_dir, 'galaxy_%s.ec' % id_tag)
|
||||
|
||||
|
||||
def collect_extra_files(object_store, dataset, job_working_directory):
|
||||
store_by = getattr(object_store, "store_by", "id")
|
||||
file_name = "dataset_%s_files" % getattr(dataset.dataset, store_by)
|
||||
temp_file_path = os.path.join(job_working_directory, file_name)
|
||||
extra_dir = None
|
||||
try:
|
||||
# This skips creation of directories - object store
|
||||
# automatically creates them. However, empty directories will
|
||||
# not be created in the object store at all, which might be a
|
||||
# problem.
|
||||
for root, dirs, files in os.walk(temp_file_path):
|
||||
extra_dir = root.replace(job_working_directory, '', 1).lstrip(os.path.sep)
|
||||
for f in files:
|
||||
object_store.update_from_file(
|
||||
dataset.dataset,
|
||||
extra_dir=extra_dir,
|
||||
alt_name=f,
|
||||
file_name=os.path.join(root, f),
|
||||
create=True,
|
||||
preserve_symlinks=True
|
||||
)
|
||||
except Exception as e:
|
||||
log.debug("Error in collect_associated_files: %s" % (e))
|
||||
|
||||
# Handle composite datatypes of auto_primary_file type
|
||||
if dataset.datatype.composite_type == 'auto_primary_file' and not dataset.has_data():
|
||||
try:
|
||||
with NamedTemporaryFile(mode='w') as temp_fh:
|
||||
temp_fh.write(dataset.datatype.generate_primary_file(dataset))
|
||||
temp_fh.flush()
|
||||
object_store.update_from_file(dataset.dataset, file_name=temp_fh.name, create=True)
|
||||
dataset.set_size()
|
||||
except Exception as e:
|
||||
log.warning('Unable to generate primary composite file automatically for %s: %s', dataset.dataset.id, e)
|
||||
+97
-119
@@ -16,7 +16,6 @@ import time
|
||||
import traceback
|
||||
from abc import ABCMeta, abstractmethod
|
||||
from json import loads
|
||||
from tempfile import NamedTemporaryFile
|
||||
from xml.etree import ElementTree
|
||||
|
||||
import six
|
||||
@@ -27,20 +26,25 @@ import galaxy
|
||||
from galaxy import model, util
|
||||
from galaxy.datatypes import sniff
|
||||
from galaxy.exceptions import ObjectInvalid, ObjectNotFound
|
||||
from galaxy.job_execution.datasets import (
|
||||
DatasetPath,
|
||||
NullDatasetPathRewriter,
|
||||
OutputsToWorkingDirectoryPathRewriter,
|
||||
TaskPathRewriter
|
||||
)
|
||||
from galaxy.job_execution.output_collect import collect_extra_files
|
||||
from galaxy.jobs.actions.post import ActionBox
|
||||
from galaxy.jobs.mapper import JobMappingException, JobRunnerMapper
|
||||
from galaxy.jobs.runners import BaseJobRunner, JobState
|
||||
from galaxy.metadata import get_metadata_compute_strategy
|
||||
from galaxy.objectstore import ObjectStorePopulator
|
||||
from galaxy.tool_util.deps import requirements
|
||||
from galaxy.tool_util.output_checker import check_output, DETECTED_JOB_STATE
|
||||
from galaxy.util import safe_makedirs, unicodify
|
||||
from galaxy.util.bunch import Bunch
|
||||
from galaxy.util.expressions import ExpressionContext
|
||||
from galaxy.util.xml_macros import load
|
||||
from galaxy.web.stack.handlers import ConfiguresHandlers
|
||||
from .datasets import (DatasetPath, NullDatasetPathRewriter,
|
||||
OutputsToWorkingDirectoryPathRewriter, TaskPathRewriter)
|
||||
from .output_checker import check_output, DETECTED_JOB_STATE
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
@@ -871,7 +875,7 @@ class JobWrapper(HasResourceParameters):
|
||||
self.output_hdas_and_paths = None
|
||||
self.tool_provided_job_metadata = None
|
||||
# Wrapper holding the info required to restore and clean up from files used for setting metadata externally
|
||||
self.external_output_metadata = get_metadata_compute_strategy(self.app, job.id)
|
||||
self.external_output_metadata = get_metadata_compute_strategy(self.app.config, job.id)
|
||||
self.job_runner_mapper = JobRunnerMapper(self, queue.dispatcher.url_to_destination, self.app.job_config)
|
||||
self.params = None
|
||||
if job.params:
|
||||
@@ -917,6 +921,10 @@ class JobWrapper(HasResourceParameters):
|
||||
def requires_containerization(self):
|
||||
return util.asbool(self.get_destination_configuration("require_container", "False"))
|
||||
|
||||
@property
|
||||
def use_metadata_binary(self):
|
||||
return util.asbool(self.get_destination_configuration('use_metadata_binary', "False"))
|
||||
|
||||
def can_split(self):
|
||||
# Should the job handler split this job up?
|
||||
return self.app.config.use_tasked_jobs and self.tool.parallelism
|
||||
@@ -1390,6 +1398,87 @@ class JobWrapper(HasResourceParameters):
|
||||
job.object_store_id = object_store_populator.object_store_id
|
||||
self._setup_working_directory(job=job)
|
||||
|
||||
def _finish_dataset(self, output_name, dataset, job, context, final_job_state, remote_metadata_directory):
|
||||
implicit_collection_jobs = job.implicit_collection_jobs_association
|
||||
purged = dataset.dataset.purged
|
||||
if not purged and dataset.dataset.external_filename is None:
|
||||
trynum = 0
|
||||
while trynum < self.app.config.retry_job_output_collection:
|
||||
try:
|
||||
# Attempt to short circuit NFS attribute caching
|
||||
os.stat(dataset.dataset.file_name)
|
||||
os.chown(dataset.dataset.file_name, os.getuid(), -1)
|
||||
trynum = self.app.config.retry_job_output_collection
|
||||
except (OSError, ObjectNotFound) as e:
|
||||
trynum += 1
|
||||
log.warning('Error accessing dataset with ID %i, will retry: %s', dataset.dataset.id, e)
|
||||
time.sleep(2)
|
||||
if getattr(dataset, "hidden_beneath_collection_instance", None):
|
||||
dataset.visible = False
|
||||
dataset.blurb = 'done'
|
||||
dataset.peek = 'no peek'
|
||||
dataset.info = (dataset.info or '')
|
||||
if context['stdout'].strip():
|
||||
# Ensure white space between entries
|
||||
dataset.info = dataset.info.rstrip() + "\n" + context['stdout'].strip()
|
||||
if context['stderr'].strip():
|
||||
# Ensure white space between entries
|
||||
dataset.info = dataset.info.rstrip() + "\n" + context['stderr'].strip()
|
||||
dataset.tool_version = self.version_string
|
||||
dataset.set_size()
|
||||
if 'uuid' in context:
|
||||
dataset.dataset.uuid = context['uuid']
|
||||
self.__update_output(job, dataset)
|
||||
if not purged:
|
||||
collect_extra_files(self.object_store, dataset, self.working_directory)
|
||||
if job.states.ERROR == final_job_state:
|
||||
dataset.blurb = "error"
|
||||
if not implicit_collection_jobs:
|
||||
# Only unhide dataset outputs that are not part of a implicit collection
|
||||
dataset.mark_unhidden()
|
||||
elif not purged:
|
||||
# If the tool was expected to set the extension, attempt to retrieve it
|
||||
if dataset.ext == 'auto':
|
||||
dataset.extension = context.get('ext', 'data')
|
||||
dataset.init_meta(copy_from=dataset)
|
||||
# if a dataset was copied, it won't appear in our dictionary:
|
||||
# either use the metadata from originating output dataset, or call set_meta on the copies
|
||||
# it would be quicker to just copy the metadata from the originating output dataset,
|
||||
# but somewhat trickier (need to recurse up the copied_from tree), for now we'll call set_meta()
|
||||
retry_internally = util.asbool(self.get_destination_configuration("retry_metadata_internally", True))
|
||||
metadata_set_successfully = self.external_output_metadata.external_metadata_set_successfully(dataset, output_name, self.sa_session, working_directory=self.working_directory)
|
||||
if retry_internally and not metadata_set_successfully:
|
||||
# If Galaxy was expected to sniff type and didn't - do so.
|
||||
if dataset.ext == "_sniff_":
|
||||
extension = sniff.handle_uploaded_dataset_file(dataset.dataset.file_name, self.app.datatypes_registry)
|
||||
dataset.extension = extension
|
||||
|
||||
# call datatype.set_meta directly for the initial set_meta call during dataset creation
|
||||
dataset.datatype.set_meta(dataset, overwrite=False)
|
||||
elif (job.states.ERROR != final_job_state and not metadata_set_successfully):
|
||||
dataset._state = model.Dataset.states.FAILED_METADATA
|
||||
else:
|
||||
self.external_output_metadata.load_metadata(dataset, output_name, self.sa_session, working_directory=self.working_directory, remote_metadata_directory=remote_metadata_directory)
|
||||
line_count = context.get('line_count', None)
|
||||
try:
|
||||
# Certain datatype's set_peek methods contain a line_count argument
|
||||
dataset.set_peek(line_count=line_count)
|
||||
except TypeError:
|
||||
# ... and others don't
|
||||
dataset.set_peek()
|
||||
else:
|
||||
# Handle purged datasets.
|
||||
dataset.blurb = "empty"
|
||||
if dataset.ext == 'auto':
|
||||
dataset.extension = context.get('ext', 'txt')
|
||||
|
||||
for context_key in TOOL_PROVIDED_JOB_METADATA_KEYS:
|
||||
if context_key in context:
|
||||
context_value = context[context_key]
|
||||
setattr(dataset, context_key, context_value)
|
||||
|
||||
self.sa_session.add(dataset)
|
||||
|
||||
def finish(
|
||||
self,
|
||||
tool_stdout,
|
||||
@@ -1477,101 +1566,15 @@ class JobWrapper(HasResourceParameters):
|
||||
return self.fail("Job %s's output dataset(s) could not be read" % job.id)
|
||||
|
||||
job_context = ExpressionContext(dict(stdout=job.stdout, stderr=job.stderr))
|
||||
implicit_collection_jobs = job.implicit_collection_jobs_association
|
||||
for dataset_assoc in job.output_datasets + job.output_library_datasets:
|
||||
context = self.get_dataset_finish_context(job_context, dataset_assoc)
|
||||
# should this also be checking library associations? - can a library item be added from a history before the job has ended? -
|
||||
# lets not allow this to occur
|
||||
# need to update all associated output hdas, i.e. history was shared with job running
|
||||
for dataset in dataset_assoc.dataset.dataset.history_associations + dataset_assoc.dataset.dataset.library_associations:
|
||||
purged = dataset.dataset.purged
|
||||
if not purged and dataset.dataset.external_filename is None:
|
||||
trynum = 0
|
||||
while trynum < self.app.config.retry_job_output_collection:
|
||||
try:
|
||||
# Attempt to short circuit NFS attribute caching
|
||||
os.stat(dataset.dataset.file_name)
|
||||
os.chown(dataset.dataset.file_name, os.getuid(), -1)
|
||||
trynum = self.app.config.retry_job_output_collection
|
||||
except (OSError, ObjectNotFound) as e:
|
||||
trynum += 1
|
||||
log.warning('Error accessing dataset with ID %i, will retry: %s', dataset.dataset.id, e)
|
||||
time.sleep(2)
|
||||
if getattr(dataset, "hidden_beneath_collection_instance", None):
|
||||
dataset.visible = False
|
||||
dataset.blurb = 'done'
|
||||
dataset.peek = 'no peek'
|
||||
dataset.info = (dataset.info or '')
|
||||
if context['stdout'].strip():
|
||||
# Ensure white space between entries
|
||||
dataset.info = dataset.info.rstrip() + "\n" + context['stdout'].strip()
|
||||
if context['stderr'].strip():
|
||||
# Ensure white space between entries
|
||||
dataset.info = dataset.info.rstrip() + "\n" + context['stderr'].strip()
|
||||
dataset.tool_version = self.version_string
|
||||
dataset.set_size()
|
||||
if 'uuid' in context:
|
||||
dataset.dataset.uuid = context['uuid']
|
||||
self.__update_output(job, dataset)
|
||||
if not purged:
|
||||
self._collect_extra_files(dataset.dataset, self.working_directory)
|
||||
# Handle composite datatypes of auto_primary_file type
|
||||
if dataset.datatype.composite_type == 'auto_primary_file' and not dataset.has_data():
|
||||
try:
|
||||
with NamedTemporaryFile(mode='w') as temp_fh:
|
||||
temp_fh.write(dataset.datatype.generate_primary_file(dataset))
|
||||
temp_fh.flush()
|
||||
self.object_store.update_from_file(dataset.dataset, file_name=temp_fh.name, create=True)
|
||||
dataset.set_size()
|
||||
except Exception as e:
|
||||
log.warning('Unable to generate primary composite file automatically for %s: %s', dataset.dataset.id, e)
|
||||
if job.states.ERROR == final_job_state:
|
||||
dataset.blurb = "error"
|
||||
if not implicit_collection_jobs:
|
||||
# Only unhide dataset outputs that are not part of a implicit collection
|
||||
dataset.mark_unhidden()
|
||||
elif not purged:
|
||||
# If the tool was expected to set the extension, attempt to retrieve it
|
||||
if dataset.ext == 'auto':
|
||||
dataset.extension = context.get('ext', 'data')
|
||||
dataset.init_meta(copy_from=dataset)
|
||||
# if a dataset was copied, it won't appear in our dictionary:
|
||||
# either use the metadata from originating output dataset, or call set_meta on the copies
|
||||
# it would be quicker to just copy the metadata from the originating output dataset,
|
||||
# but somewhat trickier (need to recurse up the copied_from tree), for now we'll call set_meta()
|
||||
retry_internally = util.asbool(self.get_destination_configuration("retry_metadata_internally", True))
|
||||
metadata_set_successfully = self.external_output_metadata.external_metadata_set_successfully(dataset, dataset_assoc.name, self.sa_session, working_directory=self.working_directory)
|
||||
if retry_internally and not metadata_set_successfully:
|
||||
# If Galaxy was expected to sniff type and didn't - do so.
|
||||
if dataset.ext == "_sniff_":
|
||||
extension = sniff.handle_uploaded_dataset_file(dataset.dataset.file_name, self.app.datatypes_registry)
|
||||
dataset.extension = extension
|
||||
|
||||
# call datatype.set_meta directly for the initial set_meta call during dataset creation
|
||||
dataset.datatype.set_meta(dataset, overwrite=False)
|
||||
elif (job.states.ERROR != final_job_state and not metadata_set_successfully):
|
||||
dataset._state = model.Dataset.states.FAILED_METADATA
|
||||
else:
|
||||
self.external_output_metadata.load_metadata(dataset, dataset_assoc.name, self.sa_session, working_directory=self.working_directory, remote_metadata_directory=remote_metadata_directory)
|
||||
line_count = context.get('line_count', None)
|
||||
try:
|
||||
# Certain datatype's set_peek methods contain a line_count argument
|
||||
dataset.set_peek(line_count=line_count)
|
||||
except TypeError:
|
||||
# ... and others don't
|
||||
dataset.set_peek()
|
||||
else:
|
||||
# Handle purged datasets.
|
||||
dataset.blurb = "empty"
|
||||
if dataset.ext == 'auto':
|
||||
dataset.extension = context.get('ext', 'txt')
|
||||
|
||||
for context_key in TOOL_PROVIDED_JOB_METADATA_KEYS:
|
||||
if context_key in context:
|
||||
context_value = context[context_key]
|
||||
setattr(dataset, context_key, context_value)
|
||||
|
||||
self.sa_session.add(dataset)
|
||||
self._finish_dataset(
|
||||
dataset_assoc.name, dataset, job, context, final_job_state, remote_metadata_directory
|
||||
)
|
||||
if job.states.ERROR == final_job_state:
|
||||
log.debug("(%s) setting dataset %s state to ERROR", job.id, dataset_assoc.dataset.dataset.id)
|
||||
# TODO: This is where the state is being set to error. Change it!
|
||||
@@ -1702,31 +1705,6 @@ class JobWrapper(HasResourceParameters):
|
||||
except Exception:
|
||||
log.exception("Unable to cleanup job %d", self.job_id)
|
||||
|
||||
def _collect_extra_files(self, dataset, job_working_directory):
|
||||
object_store = self.app.object_store
|
||||
store_by = getattr(object_store, "store_by", "id")
|
||||
file_name = "dataset_%s_files" % getattr(dataset, store_by)
|
||||
temp_file_path = os.path.join(job_working_directory, file_name)
|
||||
extra_dir = None
|
||||
try:
|
||||
# This skips creation of directories - object store
|
||||
# automatically creates them. However, empty directories will
|
||||
# not be created in the object store at all, which might be a
|
||||
# problem.
|
||||
for root, dirs, files in os.walk(temp_file_path):
|
||||
extra_dir = root.replace(job_working_directory, '', 1).lstrip(os.path.sep)
|
||||
for f in files:
|
||||
self.object_store.update_from_file(
|
||||
dataset,
|
||||
extra_dir=extra_dir,
|
||||
alt_name=f,
|
||||
file_name=os.path.join(root, f),
|
||||
create=True,
|
||||
preserve_symlinks=True
|
||||
)
|
||||
except Exception as e:
|
||||
log.debug("Error in collect_associated_files: %s" % (e))
|
||||
|
||||
def _collect_metrics(self, has_metrics, job_metrics_directory=None):
|
||||
job = has_metrics.get_job()
|
||||
job_metrics_directory = job_metrics_directory or self.working_directory
|
||||
|
||||
@@ -204,6 +204,7 @@ def __handle_metadata(commands_builder, job_wrapper, runner, remote_command_para
|
||||
datatypes_config=datatypes_config,
|
||||
compute_tmp_dir=compute_tmp_dir,
|
||||
resolve_metadata_dependencies=resolve_metadata_dependencies,
|
||||
use_bin=job_wrapper.use_metadata_binary,
|
||||
kwds={'overwrite': False}
|
||||
) or ''
|
||||
metadata_command = metadata_command.strip()
|
||||
|
||||
@@ -18,8 +18,8 @@ from six.moves.queue import (
|
||||
|
||||
import galaxy.jobs
|
||||
from galaxy import model
|
||||
from galaxy.job_execution.output_collect import default_exit_code_file, read_exit_code_from
|
||||
from galaxy.jobs.command_factory import build_command
|
||||
from galaxy.jobs.output_checker import DETECTED_JOB_STATE
|
||||
from galaxy.jobs.runners.util.env import env_to_statement
|
||||
from galaxy.jobs.runners.util.job_script import (
|
||||
job_script,
|
||||
@@ -29,6 +29,7 @@ from galaxy.tool_util.deps.dependencies import (
|
||||
JobInfo,
|
||||
ToolInfo
|
||||
)
|
||||
from galaxy.tool_util.output_checker import DETECTED_JOB_STATE
|
||||
from galaxy.util import (
|
||||
DATABASE_MAX_STRING_SIZE,
|
||||
ExecutionTimer,
|
||||
@@ -544,7 +545,7 @@ class JobState(object):
|
||||
self.job_file = JobState.default_job_file(files_dir, id_tag)
|
||||
self.output_file = os.path.join(files_dir, 'galaxy_%s.o' % id_tag)
|
||||
self.error_file = os.path.join(files_dir, 'galaxy_%s.e' % id_tag)
|
||||
self.exit_code_file = os.path.join(files_dir, 'galaxy_%s.ec' % id_tag)
|
||||
self.exit_code_file = default_exit_code_file(files_dir, id_tag)
|
||||
job_name = 'g%s' % id_tag
|
||||
if self.job_wrapper.tool.old_id:
|
||||
job_name += '_%s' % self.job_wrapper.tool.old_id
|
||||
@@ -556,27 +557,8 @@ class JobState(object):
|
||||
def default_job_file(files_dir, id_tag):
|
||||
return os.path.join(files_dir, 'galaxy_%s.sh' % id_tag)
|
||||
|
||||
@staticmethod
|
||||
def default_exit_code_file(files_dir, id_tag):
|
||||
return os.path.join(files_dir, 'galaxy_%s.ec' % id_tag)
|
||||
|
||||
def read_exit_code(self):
|
||||
try:
|
||||
# This should be an 8-bit exit code, but read ahead anyway:
|
||||
exit_code_str = open(self.exit_code_file, "r").read(32)
|
||||
except Exception:
|
||||
# By default, the exit code is 0, which typically indicates success.
|
||||
exit_code_str = "0"
|
||||
|
||||
try:
|
||||
# Decode the exit code. If it's bogus, then just use 0.
|
||||
exit_code = int(exit_code_str)
|
||||
except ValueError:
|
||||
galaxy_id_tag = self.job_wrapper.get_id_tag()
|
||||
log.warning("(%s) Exit code '%s' invalid. Using 0." % (galaxy_id_tag, exit_code_str))
|
||||
exit_code = 0
|
||||
|
||||
return exit_code
|
||||
return read_exit_code_from(self.exit_code_file, self.job_wrapper.get_id_tag())
|
||||
|
||||
def cleanup(self):
|
||||
for file in [getattr(self, a) for a in self.cleanup_file_attributes if hasattr(self, a)]:
|
||||
|
||||
@@ -11,6 +11,7 @@ import threading
|
||||
from time import sleep
|
||||
|
||||
from galaxy import model
|
||||
from galaxy.job_execution.output_collect import default_exit_code_file
|
||||
from galaxy.util import (
|
||||
asbool,
|
||||
)
|
||||
@@ -65,7 +66,7 @@ class LocalJobRunner(BaseJobRunner):
|
||||
|
||||
job_id = job_wrapper.get_id_tag()
|
||||
job_file = JobState.default_job_file(job_wrapper.working_directory, job_id)
|
||||
exit_code_path = JobState.default_exit_code_file(job_wrapper.working_directory, job_id)
|
||||
exit_code_path = default_exit_code_file(job_wrapper.working_directory, job_id)
|
||||
job_script_props = {
|
||||
'slots_statement': slots_statement,
|
||||
'command': command_line,
|
||||
@@ -137,7 +138,7 @@ class LocalJobRunner(BaseJobRunner):
|
||||
|
||||
job_destination = job_wrapper.job_destination
|
||||
job_state = JobState(job_wrapper, job_destination)
|
||||
job_state.exit_code_file = JobState.default_exit_code_file(job_wrapper.working_directory, job_id)
|
||||
job_state.exit_code_file = default_exit_code_file(job_wrapper.working_directory, job_id)
|
||||
job_state.stop_job = False
|
||||
self._finish_or_resubmit_job(job_state, stdout, stderr, job_id=job_id)
|
||||
|
||||
|
||||
@@ -699,7 +699,7 @@ class PulsarJobRunner(AsynchronousJobRunner):
|
||||
|
||||
def __build_metadata_configuration(self, client, job_wrapper, remote_metadata, remote_job_config):
|
||||
metadata_kwds = {}
|
||||
if remote_metadata:
|
||||
if remote_metadata and not job_wrapper.use_metadata_binary:
|
||||
remote_system_properties = remote_job_config.get("system_properties", {})
|
||||
remote_galaxy_home = remote_system_properties.get("galaxy_home", None)
|
||||
if not remote_galaxy_home:
|
||||
|
||||
@@ -20,8 +20,8 @@ log = getLogger(__name__)
|
||||
SET_METADATA_SCRIPT = 'from galaxy_ext.metadata.set_metadata import set_metadata; set_metadata()'
|
||||
|
||||
|
||||
def get_metadata_compute_strategy(app, job_id):
|
||||
metadata_strategy = app.config.metadata_strategy
|
||||
def get_metadata_compute_strategy(config, job_id):
|
||||
metadata_strategy = config.metadata_strategy
|
||||
if metadata_strategy == "legacy":
|
||||
return JobExternalOutputMetadataWrapper(job_id)
|
||||
else:
|
||||
@@ -46,7 +46,7 @@ class MetadataCollectionStrategy(object):
|
||||
@abc.abstractmethod
|
||||
def setup_external_metadata(self, datasets_dict, sa_session, exec_dir=None,
|
||||
tmp_dir=None, dataset_files_path=None,
|
||||
output_fnames=None, config_root=None,
|
||||
output_fnames=None, config_root=None, use_bin=False,
|
||||
config_file=None, datatypes_config=None,
|
||||
job_metadata=None, compute_tmp_dir=None,
|
||||
include_command=True, max_metadata_value_size=0,
|
||||
@@ -99,7 +99,7 @@ class PortableDirectoryMetadataGenerator(MetadataCollectionStrategy):
|
||||
|
||||
def setup_external_metadata(self, datasets_dict, sa_session, exec_dir=None,
|
||||
tmp_dir=None, dataset_files_path=None,
|
||||
output_fnames=None, config_root=None,
|
||||
output_fnames=None, config_root=None, use_bin=False,
|
||||
config_file=None, datatypes_config=None,
|
||||
job_metadata=None, compute_tmp_dir=None,
|
||||
include_command=True, max_metadata_value_size=0,
|
||||
@@ -147,9 +147,12 @@ class PortableDirectoryMetadataGenerator(MetadataCollectionStrategy):
|
||||
if include_command:
|
||||
# return command required to build
|
||||
script_path = os.path.join(metadata_dir, "set.py")
|
||||
with open(script_path, "w") as f:
|
||||
f.write(SET_METADATA_SCRIPT)
|
||||
return 'python "metadata/set.py"'
|
||||
if use_bin:
|
||||
return "galaxy-set-metadata"
|
||||
else:
|
||||
with open(script_path, "w") as f:
|
||||
f.write(SET_METADATA_SCRIPT)
|
||||
return 'python "metadata/set.py"'
|
||||
else:
|
||||
# return args to galaxy_ext.metadata.set_metadata required to build
|
||||
return ''
|
||||
@@ -214,7 +217,7 @@ class JobExternalOutputMetadataWrapper(MetadataCollectionStrategy):
|
||||
|
||||
def setup_external_metadata(self, datasets_dict, sa_session, exec_dir=None,
|
||||
tmp_dir=None, dataset_files_path=None,
|
||||
output_fnames=None, config_root=None,
|
||||
output_fnames=None, config_root=None, use_bin=False,
|
||||
config_file=None, datatypes_config=None,
|
||||
job_metadata=None, compute_tmp_dir=None,
|
||||
include_command=True, max_metadata_value_size=0,
|
||||
@@ -296,6 +299,7 @@ class JobExternalOutputMetadataWrapper(MetadataCollectionStrategy):
|
||||
job_metadata,
|
||||
" ".join(map(__metadata_files_list_to_cmd_line, metadata_files_list)),
|
||||
max_metadata_value_size)
|
||||
assert not use_bin
|
||||
if include_command:
|
||||
# return command required to build
|
||||
fd, fp = tempfile.mkstemp(suffix='.py', dir=tmp_dir, prefix="set_metadata_")
|
||||
|
||||
@@ -0,0 +1,219 @@
|
||||
"""
|
||||
Execute an external process to set_meta() on a provided list of pickled datasets.
|
||||
|
||||
This was formerly scripts/set_metadata.py and expects these arguments:
|
||||
|
||||
%prog datatypes_conf.xml job_metadata_file metadata_in,metadata_kwds,metadata_out,metadata_results_code,output_filename_override,metadata_override... max_metadata_value_size
|
||||
|
||||
Galaxy should be importable on sys.path and output_filename_override should be
|
||||
set to the path of the dataset on which metadata is being set
|
||||
(output_filename_override could previously be left empty and the path would be
|
||||
constructed automatically).
|
||||
"""
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
|
||||
from six.moves import cPickle
|
||||
from sqlalchemy.orm import clear_mappers
|
||||
|
||||
import galaxy.model.mapping # need to load this before we unpickle, in order to setup properties assigned by the mappers
|
||||
from galaxy.model.custom_types import total_size
|
||||
from galaxy.tool_util.provided_metadata import parse_tool_provided_metadata
|
||||
from galaxy.util import stringify_dictionary_keys
|
||||
|
||||
logging.basicConfig()
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
galaxy.model.Job() # this looks REAL stupid, but it is REQUIRED in order for SA to insert parameters into the classes defined by the mappers --> it appears that instantiating ANY mapper'ed class would suffice here
|
||||
|
||||
|
||||
def set_meta_with_tool_provided(dataset_instance, file_dict, set_meta_kwds, datatypes_registry, max_metadata_value_size):
|
||||
# This method is somewhat odd, in that we set the metadata attributes from tool,
|
||||
# then call set_meta, then set metadata attributes from tool again.
|
||||
# This is intentional due to interplay of overwrite kwd, the fact that some metadata
|
||||
# parameters may rely on the values of others, and that we are accepting the
|
||||
# values provided by the tool as Truth.
|
||||
extension = dataset_instance.extension
|
||||
if extension == "_sniff_":
|
||||
try:
|
||||
from galaxy.datatypes import sniff
|
||||
extension = sniff.handle_uploaded_dataset_file(dataset_instance.dataset.external_filename, datatypes_registry)
|
||||
# We need to both set the extension so it is available to set_meta
|
||||
# and record it in the metadata so it can be reloaded on the server
|
||||
# side and the model updated (see MetadataCollection.{from,to}_JSON_dict)
|
||||
dataset_instance.extension = extension
|
||||
# Set special metadata property that will reload this on server side.
|
||||
setattr(dataset_instance.metadata, "__extension__", extension)
|
||||
except Exception:
|
||||
log.exception("Problem sniffing datatype.")
|
||||
|
||||
for metadata_name, metadata_value in file_dict.get('metadata', {}).items():
|
||||
setattr(dataset_instance.metadata, metadata_name, metadata_value)
|
||||
dataset_instance.datatype.set_meta(dataset_instance, **set_meta_kwds)
|
||||
for metadata_name, metadata_value in file_dict.get('metadata', {}).items():
|
||||
setattr(dataset_instance.metadata, metadata_name, metadata_value)
|
||||
|
||||
if max_metadata_value_size:
|
||||
for k, v in list(dataset_instance.metadata.items()):
|
||||
if total_size(v) > max_metadata_value_size:
|
||||
log.info("Key %s too large for metadata, discarding" % k)
|
||||
dataset_instance.metadata.remove_key(k)
|
||||
|
||||
|
||||
def set_metadata():
|
||||
if len(sys.argv) == 1:
|
||||
set_metadata_portable()
|
||||
else:
|
||||
set_metadata_legacy()
|
||||
|
||||
|
||||
def set_metadata_portable():
|
||||
import galaxy.model
|
||||
galaxy.model.metadata.MetadataTempFile.tmp_dir = tool_job_working_directory = os.path.abspath(os.getcwd())
|
||||
|
||||
metadata_params_path = os.path.join("metadata", "params.json")
|
||||
try:
|
||||
with open(metadata_params_path, "r") as f:
|
||||
metadata_params = json.load(f)
|
||||
except IOError:
|
||||
raise Exception("Failed to find metadata/params.json from cwd [%s]" % tool_job_working_directory)
|
||||
datatypes_config = metadata_params["datatypes_config"]
|
||||
job_metadata = metadata_params["job_metadata"]
|
||||
max_metadata_value_size = metadata_params.get("max_metadata_value_size") or 0
|
||||
outputs = metadata_params["outputs"]
|
||||
|
||||
datatypes_registry = validate_and_load_datatypes_config(datatypes_config)
|
||||
tool_provided_metadata = load_job_metadata(job_metadata)
|
||||
|
||||
def set_meta(new_dataset_instance, file_dict):
|
||||
set_meta_with_tool_provided(new_dataset_instance, file_dict, set_meta_kwds, datatypes_registry, max_metadata_value_size)
|
||||
|
||||
for output_name, output_dict in outputs.items():
|
||||
filename_in = os.path.join("metadata/metadata_in_%s" % output_name)
|
||||
filename_kwds = os.path.join("metadata/metadata_kwds_%s" % output_name)
|
||||
filename_out = os.path.join("metadata/metadata_out_%s" % output_name)
|
||||
filename_results_code = os.path.join("metadata/metadata_results_%s" % output_name)
|
||||
override_metadata = os.path.join("metadata/metadata_override_%s" % output_name)
|
||||
dataset_filename_override = output_dict["filename_override"]
|
||||
|
||||
# Same block as below...
|
||||
set_meta_kwds = stringify_dictionary_keys(json.load(open(filename_kwds))) # load kwds; need to ensure our keywords are not unicode
|
||||
try:
|
||||
dataset = cPickle.load(open(filename_in, 'rb')) # load DatasetInstance
|
||||
dataset.dataset.external_filename = dataset_filename_override
|
||||
store_by = metadata_params.get("object_store_store_by", "id")
|
||||
extra_files_dir_name = "dataset_%s_files" % getattr(dataset.dataset, store_by)
|
||||
files_path = os.path.abspath(os.path.join(tool_job_working_directory, extra_files_dir_name))
|
||||
dataset.dataset.external_extra_files_path = files_path
|
||||
file_dict = tool_provided_metadata.get_dataset_meta(output_name, dataset.dataset.id)
|
||||
if 'ext' in file_dict:
|
||||
dataset.extension = file_dict['ext']
|
||||
# Metadata FileParameter types may not be writable on a cluster node, and are therefore temporarily substituted with MetadataTempFiles
|
||||
override_metadata = json.load(open(override_metadata))
|
||||
for metadata_name, metadata_file_override in override_metadata:
|
||||
if galaxy.datatypes.metadata.MetadataTempFile.is_JSONified_value(metadata_file_override):
|
||||
metadata_file_override = galaxy.datatypes.metadata.MetadataTempFile.from_JSON(metadata_file_override)
|
||||
setattr(dataset.metadata, metadata_name, metadata_file_override)
|
||||
set_meta(dataset, file_dict)
|
||||
dataset.metadata.to_JSON_dict(filename_out) # write out results of set_meta
|
||||
json.dump((True, 'Metadata has been set successfully'), open(filename_results_code, 'wt+')) # setting metadata has succeeded
|
||||
except Exception as e:
|
||||
json.dump((False, str(e)), open(filename_results_code, 'wt+')) # setting metadata has failed somehow
|
||||
|
||||
write_job_metadata(tool_job_working_directory, job_metadata, set_meta, tool_provided_metadata)
|
||||
|
||||
|
||||
def set_metadata_legacy():
|
||||
import galaxy.model
|
||||
galaxy.model.metadata.MetadataTempFile.tmp_dir = tool_job_working_directory = os.path.abspath(os.getcwd())
|
||||
|
||||
# This is ugly, but to transition from existing jobs without this parameter
|
||||
# to ones with, smoothly, it has to be the last optional parameter and we
|
||||
# have to sniff it.
|
||||
try:
|
||||
max_metadata_value_size = int(sys.argv[-1])
|
||||
sys.argv = sys.argv[:-1]
|
||||
except ValueError:
|
||||
max_metadata_value_size = 0
|
||||
# max_metadata_value_size is unspecified and should be 0
|
||||
|
||||
# Set up datatypes registry
|
||||
datatypes_config = sys.argv.pop(1)
|
||||
datatypes_registry = validate_and_load_datatypes_config(datatypes_config)
|
||||
|
||||
job_metadata = sys.argv.pop(1)
|
||||
tool_provided_metadata = load_job_metadata(job_metadata)
|
||||
|
||||
def set_meta(new_dataset_instance, file_dict):
|
||||
set_meta_with_tool_provided(new_dataset_instance, file_dict, set_meta_kwds, datatypes_registry, max_metadata_value_size)
|
||||
|
||||
for filenames in sys.argv[1:]:
|
||||
fields = filenames.split(',')
|
||||
filename_in = fields.pop(0)
|
||||
filename_kwds = fields.pop(0)
|
||||
filename_out = fields.pop(0)
|
||||
filename_results_code = fields.pop(0)
|
||||
dataset_filename_override = fields.pop(0)
|
||||
override_metadata = fields.pop(0)
|
||||
set_meta_kwds = stringify_dictionary_keys(json.load(open(filename_kwds))) # load kwds; need to ensure our keywords are not unicode
|
||||
try:
|
||||
dataset = cPickle.load(open(filename_in, 'rb')) # load DatasetInstance
|
||||
dataset.dataset.external_filename = dataset_filename_override
|
||||
files_path = os.path.abspath(os.path.join(tool_job_working_directory, "dataset_%s_files" % (dataset.dataset.id)))
|
||||
dataset.dataset.external_extra_files_path = files_path
|
||||
file_dict = tool_provided_metadata.get_dataset_meta(None, dataset.dataset.id)
|
||||
if 'ext' in file_dict:
|
||||
dataset.extension = file_dict['ext']
|
||||
# Metadata FileParameter types may not be writable on a cluster node, and are therefore temporarily substituted with MetadataTempFiles
|
||||
override_metadata = json.load(open(override_metadata))
|
||||
for metadata_name, metadata_file_override in override_metadata:
|
||||
if galaxy.datatypes.metadata.MetadataTempFile.is_JSONified_value(metadata_file_override):
|
||||
metadata_file_override = galaxy.datatypes.metadata.MetadataTempFile.from_JSON(metadata_file_override)
|
||||
setattr(dataset.metadata, metadata_name, metadata_file_override)
|
||||
set_meta(dataset, file_dict)
|
||||
dataset.metadata.to_JSON_dict(filename_out) # write out results of set_meta
|
||||
json.dump((True, 'Metadata has been set successfully'), open(filename_results_code, 'wt+')) # setting metadata has succeeded
|
||||
except Exception as e:
|
||||
json.dump((False, str(e)), open(filename_results_code, 'wt+')) # setting metadata has failed somehow
|
||||
|
||||
write_job_metadata(tool_job_working_directory, job_metadata, set_meta, tool_provided_metadata)
|
||||
|
||||
|
||||
def validate_and_load_datatypes_config(datatypes_config):
|
||||
galaxy_root = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir, os.pardir))
|
||||
|
||||
if not os.path.exists(datatypes_config):
|
||||
# Hack for Pulsar on usegalaxy.org, drop ASAP.
|
||||
datatypes_config = "configs/registry.xml"
|
||||
|
||||
if not os.path.exists(datatypes_config):
|
||||
print("Metadata setting failed because registry.xml [%s] could not be found. You may retry setting metadata." % datatypes_config)
|
||||
sys.exit(1)
|
||||
import galaxy.datatypes.registry
|
||||
datatypes_registry = galaxy.datatypes.registry.Registry()
|
||||
datatypes_registry.load_datatypes(root_dir=galaxy_root, config=datatypes_config)
|
||||
galaxy.model.set_datatypes_registry(datatypes_registry)
|
||||
return datatypes_registry
|
||||
|
||||
|
||||
def load_job_metadata(job_metadata):
|
||||
return parse_tool_provided_metadata(job_metadata)
|
||||
|
||||
|
||||
def write_job_metadata(tool_job_working_directory, job_metadata, set_meta, tool_provided_metadata):
|
||||
for i, file_dict in enumerate(tool_provided_metadata.get_new_datasets_for_metadata_collection(), start=1):
|
||||
filename = file_dict["filename"]
|
||||
new_dataset_filename = os.path.join(tool_job_working_directory, "working", filename)
|
||||
new_dataset = galaxy.model.Dataset(id=-i, external_filename=new_dataset_filename)
|
||||
extra_files = file_dict.get('extra_files', None)
|
||||
if extra_files is not None:
|
||||
new_dataset._extra_files_path = os.path.join(tool_job_working_directory, "working", extra_files)
|
||||
new_dataset.state = new_dataset.states.OK
|
||||
new_dataset_instance = galaxy.model.HistoryDatasetAssociation(id=-i, dataset=new_dataset, extension=file_dict.get('ext', 'data'))
|
||||
set_meta(new_dataset_instance, file_dict)
|
||||
file_dict['metadata'] = json.loads(new_dataset_instance.metadata.to_JSON_dict()) # storing metadata in external form, need to turn back into dict, then later jsonify
|
||||
|
||||
tool_provided_metadata.rewrite()
|
||||
clear_mappers()
|
||||
@@ -25,6 +25,7 @@ from galaxy import (
|
||||
exceptions,
|
||||
model
|
||||
)
|
||||
from galaxy.job_execution import output_collect
|
||||
from galaxy.managers.jobs import JobSearch
|
||||
from galaxy.metadata import get_metadata_compute_strategy
|
||||
from galaxy.model.tags import GalaxyTagHandler
|
||||
@@ -44,6 +45,7 @@ from galaxy.tool_util.parser import (
|
||||
ToolOutputCollectionPart
|
||||
)
|
||||
from galaxy.tool_util.parser.xml import XmlPageSource
|
||||
from galaxy.tool_util.provided_metadata import parse_tool_provided_metadata
|
||||
from galaxy.tools import expressions
|
||||
from galaxy.tools.actions import DefaultToolAction
|
||||
from galaxy.tools.actions.data_manager import DataManagerToolAction
|
||||
@@ -57,7 +59,6 @@ from galaxy.tools.parameters import (
|
||||
populate_state,
|
||||
visit_input_values
|
||||
)
|
||||
from galaxy.tools.parameters import output_collect
|
||||
from galaxy.tools.parameters.basic import (
|
||||
BaseURLToolParameter,
|
||||
DataCollectionToolParameter,
|
||||
@@ -101,7 +102,6 @@ from .execute import (
|
||||
execute as execute_job,
|
||||
MappingParameters,
|
||||
)
|
||||
from .provided_metadata import parse_tool_provided_metadata
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
@@ -1737,6 +1737,8 @@ class Tool(Dictifiable):
|
||||
Find any additional datasets generated by a tool and attach (for
|
||||
cases where number of outputs is not known in advance).
|
||||
"""
|
||||
# given the job_execution import is the only one, probably makes sense to refactor this out
|
||||
# into job_wrapper.
|
||||
tool = self
|
||||
permission_provider = output_collect.PermissionProvider(inp_data, tool.app.security_agent, job)
|
||||
metadata_source_provider = output_collect.MetadataSourceProvider(inp_data)
|
||||
@@ -2392,7 +2394,7 @@ class SetMetadataTool(Tool):
|
||||
job, base_dir='job_work', dir_only=True, obj_dir=True
|
||||
)
|
||||
for name, dataset in inp_data.items():
|
||||
external_metadata = get_metadata_compute_strategy(app, job.id)
|
||||
external_metadata = get_metadata_compute_strategy(app.config, job.id)
|
||||
sa_session = app.model.context
|
||||
if external_metadata.external_metadata_set_successfully(dataset, name, sa_session, working_directory=working_directory):
|
||||
external_metadata.load_metadata(dataset, name, sa_session, working_directory=working_directory)
|
||||
|
||||
@@ -2,7 +2,7 @@ import logging
|
||||
import os
|
||||
from json import dumps
|
||||
|
||||
from galaxy.jobs.datasets import DatasetPath
|
||||
from galaxy.job_execution.datasets import DatasetPath
|
||||
from galaxy.metadata import get_metadata_compute_strategy
|
||||
from galaxy.util.odict import odict
|
||||
from . import ToolAction
|
||||
@@ -78,7 +78,7 @@ class SetMetadataToolAction(ToolAction):
|
||||
job_working_dir = app.object_store.get_filename(job, base_dir='job_work', dir_only=True, extra_dir=str(job.id))
|
||||
datatypes_config = os.path.join(job_working_dir, 'registry.xml')
|
||||
app.datatypes_registry.to_xml_file(path=datatypes_config)
|
||||
external_metadata_wrapper = get_metadata_compute_strategy(app, job.id)
|
||||
external_metadata_wrapper = get_metadata_compute_strategy(app.config, job.id)
|
||||
output_datatasets_dict = {
|
||||
dataset_name: dataset,
|
||||
}
|
||||
|
||||
@@ -7,7 +7,7 @@ import tempfile
|
||||
from six import string_types
|
||||
|
||||
from galaxy import model
|
||||
from galaxy.jobs.datasets import dataset_path_rewrites
|
||||
from galaxy.job_execution.datasets import dataset_path_rewrites
|
||||
from galaxy.model.none_like import NoneDataset
|
||||
from galaxy.tools import global_tool_errors
|
||||
from galaxy.tools.parameters import (
|
||||
|
||||
@@ -1 +0,0 @@
|
||||
../../galaxy/tools/provided_metadata.py
|
||||
@@ -10,217 +10,13 @@ set to the path of the dataset on which metadata is being set
|
||||
(output_filename_override could previously be left empty and the path would be
|
||||
constructed automatically).
|
||||
"""
|
||||
import json
|
||||
import logging
|
||||
|
||||
import os
|
||||
import sys
|
||||
|
||||
# insert *this* galaxy before all others on sys.path
|
||||
sys.path.insert(1, os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir)))
|
||||
|
||||
from six.moves import cPickle
|
||||
from sqlalchemy.orm import clear_mappers
|
||||
from galaxy.metadata.set_metadata import set_metadata
|
||||
|
||||
import galaxy.model.mapping # need to load this before we unpickle, in order to setup properties assigned by the mappers
|
||||
from galaxy.model.custom_types import total_size
|
||||
from galaxy.util import stringify_dictionary_keys
|
||||
from ._provided_metadata import parse_tool_provided_metadata
|
||||
|
||||
# ensure supported version
|
||||
assert sys.version_info[:2] >= (2, 7), 'Python version must be at least 2.7, this is: %s' % sys.version
|
||||
|
||||
logging.basicConfig()
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
galaxy.model.Job() # this looks REAL stupid, but it is REQUIRED in order for SA to insert parameters into the classes defined by the mappers --> it appears that instantiating ANY mapper'ed class would suffice here
|
||||
|
||||
|
||||
def set_meta_with_tool_provided(dataset_instance, file_dict, set_meta_kwds, datatypes_registry, max_metadata_value_size):
|
||||
# This method is somewhat odd, in that we set the metadata attributes from tool,
|
||||
# then call set_meta, then set metadata attributes from tool again.
|
||||
# This is intentional due to interplay of overwrite kwd, the fact that some metadata
|
||||
# parameters may rely on the values of others, and that we are accepting the
|
||||
# values provided by the tool as Truth.
|
||||
extension = dataset_instance.extension
|
||||
if extension == "_sniff_":
|
||||
try:
|
||||
from galaxy.datatypes import sniff
|
||||
extension = sniff.handle_uploaded_dataset_file(dataset_instance.dataset.external_filename, datatypes_registry)
|
||||
# We need to both set the extension so it is available to set_meta
|
||||
# and record it in the metadata so it can be reloaded on the server
|
||||
# side and the model updated (see MetadataCollection.{from,to}_JSON_dict)
|
||||
dataset_instance.extension = extension
|
||||
# Set special metadata property that will reload this on server side.
|
||||
setattr(dataset_instance.metadata, "__extension__", extension)
|
||||
except Exception:
|
||||
log.exception("Problem sniffing datatype.")
|
||||
|
||||
for metadata_name, metadata_value in file_dict.get('metadata', {}).items():
|
||||
setattr(dataset_instance.metadata, metadata_name, metadata_value)
|
||||
dataset_instance.datatype.set_meta(dataset_instance, **set_meta_kwds)
|
||||
for metadata_name, metadata_value in file_dict.get('metadata', {}).items():
|
||||
setattr(dataset_instance.metadata, metadata_name, metadata_value)
|
||||
|
||||
if max_metadata_value_size:
|
||||
for k, v in list(dataset_instance.metadata.items()):
|
||||
if total_size(v) > max_metadata_value_size:
|
||||
log.info("Key %s too large for metadata, discarding" % k)
|
||||
dataset_instance.metadata.remove_key(k)
|
||||
|
||||
|
||||
def set_metadata():
|
||||
if len(sys.argv) == 1:
|
||||
set_metadata_portable()
|
||||
else:
|
||||
set_metadata_legacy()
|
||||
|
||||
|
||||
def set_metadata_portable():
|
||||
import galaxy.model
|
||||
tool_job_working_directory = os.path.abspath(os.getcwd())
|
||||
galaxy.model.metadata.MetadataTempFile.tmp_dir = os.path.join(tool_job_working_directory, "metadata")
|
||||
|
||||
metadata_params_path = os.path.join("metadata", "params.json")
|
||||
try:
|
||||
with open(metadata_params_path, "r") as f:
|
||||
metadata_params = json.load(f)
|
||||
except IOError:
|
||||
raise Exception("Failed to find metadata/params.json from cwd [%s]" % tool_job_working_directory)
|
||||
datatypes_config = metadata_params["datatypes_config"]
|
||||
job_metadata = metadata_params["job_metadata"]
|
||||
max_metadata_value_size = metadata_params.get("max_metadata_value_size") or 0
|
||||
outputs = metadata_params["outputs"]
|
||||
|
||||
datatypes_registry = validate_and_load_datatypes_config(datatypes_config)
|
||||
tool_provided_metadata = load_job_metadata(job_metadata)
|
||||
|
||||
def set_meta(new_dataset_instance, file_dict):
|
||||
set_meta_with_tool_provided(new_dataset_instance, file_dict, set_meta_kwds, datatypes_registry, max_metadata_value_size)
|
||||
|
||||
for output_name, output_dict in outputs.items():
|
||||
filename_in = os.path.join("metadata/metadata_in_%s" % output_name)
|
||||
filename_kwds = os.path.join("metadata/metadata_kwds_%s" % output_name)
|
||||
filename_out = os.path.join("metadata/metadata_out_%s" % output_name)
|
||||
filename_results_code = os.path.join("metadata/metadata_results_%s" % output_name)
|
||||
override_metadata = os.path.join("metadata/metadata_override_%s" % output_name)
|
||||
dataset_filename_override = output_dict["filename_override"]
|
||||
|
||||
# Same block as below...
|
||||
set_meta_kwds = stringify_dictionary_keys(json.load(open(filename_kwds))) # load kwds; need to ensure our keywords are not unicode
|
||||
try:
|
||||
dataset = cPickle.load(open(filename_in, 'rb')) # load DatasetInstance
|
||||
dataset.dataset.external_filename = dataset_filename_override
|
||||
store_by = metadata_params.get("object_store_store_by", "id")
|
||||
extra_files_dir_name = "dataset_%s_files" % getattr(dataset.dataset, store_by)
|
||||
files_path = os.path.abspath(os.path.join(tool_job_working_directory, extra_files_dir_name))
|
||||
dataset.dataset.external_extra_files_path = files_path
|
||||
file_dict = tool_provided_metadata.get_dataset_meta(output_name, dataset.dataset.id)
|
||||
if 'ext' in file_dict:
|
||||
dataset.extension = file_dict['ext']
|
||||
# Metadata FileParameter types may not be writable on a cluster node, and are therefore temporarily substituted with MetadataTempFiles
|
||||
override_metadata = json.load(open(override_metadata))
|
||||
for metadata_name, metadata_file_override in override_metadata:
|
||||
if galaxy.datatypes.metadata.MetadataTempFile.is_JSONified_value(metadata_file_override):
|
||||
metadata_file_override = galaxy.datatypes.metadata.MetadataTempFile.from_JSON(metadata_file_override)
|
||||
setattr(dataset.metadata, metadata_name, metadata_file_override)
|
||||
set_meta(dataset, file_dict)
|
||||
dataset.metadata.to_JSON_dict(filename_out) # write out results of set_meta
|
||||
json.dump((True, 'Metadata has been set successfully'), open(filename_results_code, 'wt+')) # setting metadata has succeeded
|
||||
except Exception as e:
|
||||
json.dump((False, str(e)), open(filename_results_code, 'wt+')) # setting metadata has failed somehow
|
||||
|
||||
write_job_metadata(tool_job_working_directory, job_metadata, set_meta, tool_provided_metadata)
|
||||
|
||||
|
||||
def set_metadata_legacy():
|
||||
import galaxy.model
|
||||
galaxy.model.metadata.MetadataTempFile.tmp_dir = tool_job_working_directory = os.path.abspath(os.getcwd())
|
||||
|
||||
# This is ugly, but to transition from existing jobs without this parameter
|
||||
# to ones with, smoothly, it has to be the last optional parameter and we
|
||||
# have to sniff it.
|
||||
try:
|
||||
max_metadata_value_size = int(sys.argv[-1])
|
||||
sys.argv = sys.argv[:-1]
|
||||
except ValueError:
|
||||
max_metadata_value_size = 0
|
||||
# max_metadata_value_size is unspecified and should be 0
|
||||
|
||||
# Set up datatypes registry
|
||||
datatypes_config = sys.argv.pop(1)
|
||||
datatypes_registry = validate_and_load_datatypes_config(datatypes_config)
|
||||
|
||||
job_metadata = sys.argv.pop(1)
|
||||
tool_provided_metadata = load_job_metadata(job_metadata)
|
||||
|
||||
def set_meta(new_dataset_instance, file_dict):
|
||||
set_meta_with_tool_provided(new_dataset_instance, file_dict, set_meta_kwds, datatypes_registry, max_metadata_value_size)
|
||||
|
||||
for filenames in sys.argv[1:]:
|
||||
fields = filenames.split(',')
|
||||
filename_in = fields.pop(0)
|
||||
filename_kwds = fields.pop(0)
|
||||
filename_out = fields.pop(0)
|
||||
filename_results_code = fields.pop(0)
|
||||
dataset_filename_override = fields.pop(0)
|
||||
override_metadata = fields.pop(0)
|
||||
set_meta_kwds = stringify_dictionary_keys(json.load(open(filename_kwds))) # load kwds; need to ensure our keywords are not unicode
|
||||
try:
|
||||
dataset = cPickle.load(open(filename_in, 'rb')) # load DatasetInstance
|
||||
dataset.dataset.external_filename = dataset_filename_override
|
||||
files_path = os.path.abspath(os.path.join(tool_job_working_directory, "dataset_%s_files" % (dataset.dataset.id)))
|
||||
dataset.dataset.external_extra_files_path = files_path
|
||||
file_dict = tool_provided_metadata.get_dataset_meta(None, dataset.dataset.id)
|
||||
if 'ext' in file_dict:
|
||||
dataset.extension = file_dict['ext']
|
||||
# Metadata FileParameter types may not be writable on a cluster node, and are therefore temporarily substituted with MetadataTempFiles
|
||||
override_metadata = json.load(open(override_metadata))
|
||||
for metadata_name, metadata_file_override in override_metadata:
|
||||
if galaxy.datatypes.metadata.MetadataTempFile.is_JSONified_value(metadata_file_override):
|
||||
metadata_file_override = galaxy.datatypes.metadata.MetadataTempFile.from_JSON(metadata_file_override)
|
||||
setattr(dataset.metadata, metadata_name, metadata_file_override)
|
||||
set_meta(dataset, file_dict)
|
||||
dataset.metadata.to_JSON_dict(filename_out) # write out results of set_meta
|
||||
json.dump((True, 'Metadata has been set successfully'), open(filename_results_code, 'wt+')) # setting metadata has succeeded
|
||||
except Exception as e:
|
||||
json.dump((False, str(e)), open(filename_results_code, 'wt+')) # setting metadata has failed somehow
|
||||
|
||||
write_job_metadata(tool_job_working_directory, job_metadata, set_meta, tool_provided_metadata)
|
||||
|
||||
|
||||
def validate_and_load_datatypes_config(datatypes_config):
|
||||
galaxy_root = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir, os.pardir))
|
||||
|
||||
if not os.path.exists(datatypes_config):
|
||||
# Hack for Pulsar on usegalaxy.org, drop ASAP.
|
||||
datatypes_config = "configs/registry.xml"
|
||||
|
||||
if not os.path.exists(datatypes_config):
|
||||
print("Metadata setting failed because registry.xml [%s] could not be found. You may retry setting metadata." % datatypes_config)
|
||||
sys.exit(1)
|
||||
import galaxy.datatypes.registry
|
||||
datatypes_registry = galaxy.datatypes.registry.Registry()
|
||||
datatypes_registry.load_datatypes(root_dir=galaxy_root, config=datatypes_config)
|
||||
galaxy.model.set_datatypes_registry(datatypes_registry)
|
||||
return datatypes_registry
|
||||
|
||||
|
||||
def load_job_metadata(job_metadata):
|
||||
return parse_tool_provided_metadata(job_metadata)
|
||||
|
||||
|
||||
def write_job_metadata(tool_job_working_directory, job_metadata, set_meta, tool_provided_metadata):
|
||||
for i, file_dict in enumerate(tool_provided_metadata.get_new_datasets_for_metadata_collection(), start=1):
|
||||
filename = file_dict["filename"]
|
||||
new_dataset_filename = os.path.join(tool_job_working_directory, "working", filename)
|
||||
new_dataset = galaxy.model.Dataset(id=-i, external_filename=new_dataset_filename)
|
||||
extra_files = file_dict.get('extra_files', None)
|
||||
if extra_files is not None:
|
||||
new_dataset._extra_files_path = os.path.join(tool_job_working_directory, "working", extra_files)
|
||||
new_dataset.state = new_dataset.states.OK
|
||||
new_dataset_instance = galaxy.model.HistoryDatasetAssociation(id=-i, dataset=new_dataset, extension=file_dict.get('ext', 'data'))
|
||||
set_meta(new_dataset_instance, file_dict)
|
||||
file_dict['metadata'] = json.loads(new_dataset_instance.metadata.to_JSON_dict()) # storing metadata in external form, need to turn back into dict, then later jsonify
|
||||
|
||||
tool_provided_metadata.rewrite()
|
||||
clear_mappers()
|
||||
__all__ = ('set_metadata', )
|
||||
|
||||
@@ -0,0 +1,12 @@
|
||||
.. :changelog:
|
||||
|
||||
History
|
||||
-------
|
||||
|
||||
.. to_doc
|
||||
|
||||
---------------------
|
||||
19.9.0.dev0
|
||||
---------------------
|
||||
|
||||
* Initial import from dev branch of Galaxy during 19.09 development cycle.
|
||||
Symlink
+1
@@ -0,0 +1 @@
|
||||
../../LICENSE.txt
|
||||
@@ -0,0 +1,2 @@
|
||||
include *.rst LICENSE
|
||||
|
||||
Symlink
+1
@@ -0,0 +1 @@
|
||||
../package.Makefile
|
||||
@@ -0,0 +1,14 @@
|
||||
|
||||
.. image:: https://badge.fury.io/py/galaxy-job-execution.svg
|
||||
:target: https://pypi.python.org/pypi/galaxy-job-execution/
|
||||
|
||||
|
||||
Overview
|
||||
--------
|
||||
|
||||
The Galaxy_ job execution runtime module.
|
||||
|
||||
* Free software: Academic Free License version 3.0
|
||||
* Code: https://github.com/galaxyproject/galaxy
|
||||
|
||||
.. _Galaxy: http://galaxyproject.org/
|
||||
@@ -0,0 +1 @@
|
||||
../package-dev-requirements.txt
|
||||
@@ -0,0 +1 @@
|
||||
__path__ = __import__('pkgutil').extend_path(__path__, __name__)
|
||||
@@ -0,0 +1 @@
|
||||
../../../lib/galaxy/job_execution
|
||||
+1
@@ -0,0 +1 @@
|
||||
../../../lib/galaxy/metadata
|
||||
@@ -0,0 +1,13 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
|
||||
__version__ = '19.9.0.dev0'
|
||||
|
||||
PROJECT_NAME = "galaxy-job-execution"
|
||||
PROJECT_OWNER = PROJECT_USERAME = "galaxyproject"
|
||||
PROJECT_URL = "https://github.com/galaxyproject/galaxy"
|
||||
PROJECT_AUTHOR = 'Galaxy Project and Community'
|
||||
PROJECT_DESCRIPTION = 'Galaxy Job Execution Runtime Utilities'
|
||||
PROJECT_EMAIL = 'jmchilton@gmail.com'
|
||||
RAW_CONTENT_URL = "https://raw.github.com/%s/%s/master/" % (
|
||||
PROJECT_USERAME, PROJECT_NAME
|
||||
)
|
||||
@@ -0,0 +1 @@
|
||||
galaxy-data
|
||||
Symlink
+1
@@ -0,0 +1 @@
|
||||
../build_scripts
|
||||
Symlink
+1
@@ -0,0 +1 @@
|
||||
../../setup.cfg
|
||||
@@ -0,0 +1,102 @@
|
||||
#!/usr/bin/env python
|
||||
# -*- coding: utf-8 -*-
|
||||
|
||||
import ast
|
||||
import os
|
||||
import re
|
||||
try:
|
||||
from setuptools import setup
|
||||
except ImportError:
|
||||
from distutils.core import setup
|
||||
|
||||
SOURCE_DIR = "galaxy"
|
||||
|
||||
_version_re = re.compile(r'__version__\s+=\s+(.*)')
|
||||
|
||||
with open('%s/project_galaxy_job_execution.py' % SOURCE_DIR, 'rb') as f:
|
||||
init_contents = f.read().decode('utf-8')
|
||||
|
||||
def get_var(var_name):
|
||||
pattern = re.compile(r'%s\s+=\s+(.*)' % var_name)
|
||||
match = pattern.search(init_contents).group(1)
|
||||
return str(ast.literal_eval(match))
|
||||
|
||||
version = get_var("__version__")
|
||||
PROJECT_NAME = get_var("PROJECT_NAME")
|
||||
PROJECT_URL = get_var("PROJECT_URL")
|
||||
PROJECT_AUTHOR = get_var("PROJECT_AUTHOR")
|
||||
PROJECT_EMAIL = get_var("PROJECT_EMAIL")
|
||||
PROJECT_DESCRIPTION = get_var("PROJECT_DESCRIPTION")
|
||||
|
||||
TEST_DIR = 'tests'
|
||||
PACKAGES = [
|
||||
'galaxy',
|
||||
'galaxy.job_execution',
|
||||
'galaxy.metadata',
|
||||
]
|
||||
ENTRY_POINTS = '''
|
||||
[console_scripts]
|
||||
galaxy-set-metadata=galaxy.metadata.set_metadata:set_metadata
|
||||
'''
|
||||
PACKAGE_DATA = {
|
||||
# Be sure to update MANIFEST.in for source dist.
|
||||
'galaxy': [
|
||||
],
|
||||
}
|
||||
PACKAGE_DIR = {
|
||||
SOURCE_DIR: SOURCE_DIR,
|
||||
}
|
||||
|
||||
readme = open('README.rst').read()
|
||||
history = open('HISTORY.rst').read().replace('.. :changelog:', '')
|
||||
|
||||
if os.path.exists("requirements.txt"):
|
||||
requirements = open("requirements.txt").read().split("\n")
|
||||
else:
|
||||
# In tox, it will cover them anyway.
|
||||
requirements = []
|
||||
|
||||
|
||||
test_requirements = [
|
||||
# TODO: put package test requirements here
|
||||
]
|
||||
|
||||
|
||||
setup(
|
||||
name=PROJECT_NAME,
|
||||
version=version,
|
||||
description=PROJECT_DESCRIPTION,
|
||||
long_description=readme + '\n\n' + history,
|
||||
long_description_content_type='text/x-rst',
|
||||
author=PROJECT_AUTHOR,
|
||||
author_email=PROJECT_EMAIL,
|
||||
url=PROJECT_URL,
|
||||
packages=PACKAGES,
|
||||
entry_points=ENTRY_POINTS,
|
||||
package_data=PACKAGE_DATA,
|
||||
package_dir=PACKAGE_DIR,
|
||||
include_package_data=True,
|
||||
install_requires=requirements,
|
||||
extras_require={},
|
||||
license="AFL",
|
||||
zip_safe=False,
|
||||
keywords='galaxy',
|
||||
classifiers=[
|
||||
'Development Status :: 5 - Production/Stable',
|
||||
'Intended Audience :: Developers',
|
||||
'Environment :: Console',
|
||||
'License :: OSI Approved :: Academic Free License (AFL)',
|
||||
'Operating System :: POSIX',
|
||||
'Topic :: Software Development',
|
||||
'Topic :: Software Development :: Code Generators',
|
||||
'Topic :: Software Development :: Testing',
|
||||
'Natural Language :: English',
|
||||
"Programming Language :: Python :: 2",
|
||||
'Programming Language :: Python :: 2.7',
|
||||
'Programming Language :: Python :: 3.5',
|
||||
'Programming Language :: Python :: 3.6',
|
||||
'Programming Language :: Python :: 3.7',
|
||||
],
|
||||
test_suite=TEST_DIR,
|
||||
tests_require=test_requirements
|
||||
)
|
||||
@@ -0,0 +1 @@
|
||||
../../../test/unit/jobs/test_datasets.py
|
||||
@@ -0,0 +1 @@
|
||||
../../../test/unit/tools/test_metadata.py
|
||||
+2
-1
@@ -21,10 +21,11 @@ PACKAGE_DIRS=(
|
||||
containers
|
||||
tool_util
|
||||
data
|
||||
job_execution
|
||||
)
|
||||
# containers has no tests, tool_util not yet working 100%,
|
||||
# data has many problems quota, tool shed install database, etc..
|
||||
RUN_TESTS=(1 1 1 0 0 0 0)
|
||||
RUN_TESTS=(1 1 1 0 0 0 0 1)
|
||||
|
||||
for ((i=0; i<${#PACKAGE_DIRS[@]}; i++)); do
|
||||
package_dir=${PACKAGE_DIRS[$i]}
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
../../../test/unit/tools/test_output_checker.py
|
||||
@@ -178,6 +178,7 @@ class MockJobWrapper(object):
|
||||
)
|
||||
)
|
||||
self.shell = "/bin/sh"
|
||||
self.use_metadata_binary = False
|
||||
|
||||
def get_command_line(self):
|
||||
return self.command_line
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
from galaxy.jobs.datasets import DatasetPath
|
||||
from galaxy.job_execution.datasets import DatasetPath
|
||||
|
||||
|
||||
def test_dataset_path():
|
||||
|
||||
@@ -144,6 +144,7 @@ class MockJobWrapper(object):
|
||||
self.shell = "/bin/bash"
|
||||
self.cleanup_job = "never"
|
||||
self.tmp_dir_creation_statement = ""
|
||||
self.use_metadata_binary = False
|
||||
|
||||
# Cruft for setting metadata externally, axe at some point.
|
||||
self.external_output_metadata = bunch.Bunch(
|
||||
|
||||
@@ -7,7 +7,7 @@ from galaxy import (
|
||||
util
|
||||
)
|
||||
from galaxy.tool_util.parser import output_collection_def
|
||||
from galaxy.tools.provided_metadata import LegacyToolProvidedMetadata, NullToolProvidedMetadata
|
||||
from galaxy.tool_util.provided_metadata import LegacyToolProvidedMetadata, NullToolProvidedMetadata
|
||||
from .. import tools_support
|
||||
|
||||
DEFAULT_TOOL_OUTPUT = "out1"
|
||||
|
||||
@@ -2,8 +2,8 @@ import os
|
||||
from unittest import TestCase
|
||||
from xml.etree.ElementTree import XML
|
||||
|
||||
from galaxy.job_execution.datasets import DatasetPath
|
||||
from galaxy.jobs import SimpleComputeEnvironment
|
||||
from galaxy.jobs.datasets import DatasetPath
|
||||
from galaxy.model import (
|
||||
Dataset,
|
||||
History,
|
||||
|
||||
@@ -3,7 +3,7 @@ import subprocess
|
||||
import unittest
|
||||
|
||||
from galaxy import model
|
||||
from galaxy.jobs.datasets import DatasetPath
|
||||
from galaxy.job_execution.datasets import DatasetPath
|
||||
from galaxy.metadata import get_metadata_compute_strategy
|
||||
from galaxy.objectstore import ObjectStorePopulator
|
||||
from .. import tools_support
|
||||
@@ -143,7 +143,7 @@ class MetadataTestCase(unittest.TestCase, tools_support.UsesApp, tools_support.U
|
||||
f.write(contents)
|
||||
|
||||
def metadata_command(self, output_datasets):
|
||||
metadata_compute_strategy = get_metadata_compute_strategy(self.app, self.job.id)
|
||||
metadata_compute_strategy = get_metadata_compute_strategy(self.app.config, self.job.id)
|
||||
self.metadata_compute_strategy = metadata_compute_strategy
|
||||
|
||||
exec_dir = None
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
from unittest import TestCase
|
||||
|
||||
from galaxy.jobs.output_checker import check_output, DETECTED_JOB_STATE
|
||||
from galaxy.tool_util.output_checker import check_output, DETECTED_JOB_STATE
|
||||
from galaxy.tool_util.parser.error_level import StdioErrorLevel
|
||||
from galaxy.tool_util.parser.interface import ToolStdioRegex
|
||||
from galaxy.util.bunch import Bunch
|
||||
@@ -3,7 +3,7 @@ import tempfile
|
||||
from xml.etree.ElementTree import XML
|
||||
|
||||
from galaxy.datatypes.metadata import MetadataSpecCollection
|
||||
from galaxy.jobs.datasets import DatasetPath
|
||||
from galaxy.job_execution.datasets import DatasetPath
|
||||
from galaxy.tools.parameters.basic import (
|
||||
DrillDownSelectToolParameter,
|
||||
IntegerToolParameter,
|
||||
|
||||
Reference in New Issue
Block a user