Progress toward fixing extended metadata with Pulsar

This commit is contained in:
John Chilton
2022-11-02 08:59:12 -04:00
parent d964e4f1ca
commit 89394bd8b9
8 changed files with 44 additions and 14 deletions
+12 -4
View File
@@ -82,7 +82,10 @@ def build_command(
__handle_task_splitting(commands_builder, job_wrapper)
externalized = False
for_pulsar = "pulsar_version" in remote_command_params
if (container and modify_command_for_container) or job_wrapper.commands_in_new_shell:
externalized = True
if container and modify_command_for_container:
# Many Docker containers do not have /bin/bash.
external_command_shell = container.shell
@@ -103,10 +106,13 @@ def build_command(
else:
commands_builder = CommandsBuilder(externalized_commands)
for_pulsar = "script_directory" in remote_command_params
if not for_pulsar:
wrap_stdio = externalized or not for_pulsar
if wrap_stdio:
# Galaxy writes I/O files to outputs, Pulsar uses metadata. metadata seems like
# it should be preferred - at least if the directory exists.
io_directory = "../metadata" if for_pulsar else "../outputs"
commands_builder.capture_stdout_stderr(
"../outputs/tool_stdout", "../outputs/tool_stderr", stream_stdout_stderr=stream_stdout_stderr
f"{io_directory}/tool_stdout", f"{io_directory}/tool_stderr", stream_stdout_stderr=stream_stdout_stderr
)
# Don't need to create a separate tool working directory for Pulsar
@@ -201,7 +207,7 @@ def __externalize_commands(
# doesn't need to mount the job directory (rw) and then eliminate this hack
# (or restrict to older Pulsar versions).
# https://github.com/galaxyproject/galaxy/pull/8449
for_pulsar = "script_directory" in remote_command_params
for_pulsar = "pulsar_version" in remote_command_params
if for_pulsar:
commands = f"{shell} {join(remote_command_params['script_directory'], script_name)}"
log.info(f"Built script [{local_container_script}] for tool command [{tool_commands}]")
@@ -260,6 +266,7 @@ def __handle_metadata(
config_file = metadata_kwds.get("config_file", None)
datatypes_config = metadata_kwds.get("datatypes_config", None)
compute_tmp_dir = metadata_kwds.get("compute_tmp_dir", None)
version_path = job_wrapper.job_io.version_path
resolve_metadata_dependencies = job_wrapper.commands_in_new_shell
metadata_command = (
job_wrapper.setup_external_metadata(
@@ -272,6 +279,7 @@ def __handle_metadata(
config_file=config_file,
datatypes_config=datatypes_config,
compute_tmp_dir=compute_tmp_dir,
compute_version_path=version_path,
resolve_metadata_dependencies=resolve_metadata_dependencies,
use_bin=job_wrapper.use_metadata_binary,
kwds={"overwrite": False},
+11 -4
View File
@@ -485,15 +485,17 @@ class PulsarJobRunner(AsynchronousJobRunner):
remote_working_directory = remote_job_config["working_directory"]
remote_job_directory = os.path.abspath(os.path.join(remote_working_directory, os.path.pardir))
remote_tool_directory = os.path.abspath(os.path.join(remote_job_directory, "tool_files"))
pulsar_version = PulsarJobRunner.pulsar_version(remote_job_config)
remote_command_params = dict(
working_directory=remote_job_config["metadata_directory"],
script_directory=remote_job_directory,
metadata_kwds=metadata_kwds,
dependency_resolution=dependency_resolution,
pulsar_version=pulsar_version,
)
# TODO: Following directories work for Pulsar, always worked for Pulsar - but should be
# calculated at some other level.
job_wrapper.disable_commands_in_new_shell()
rewrite_paths = not PulsarJobRunner.__rewrite_parameters(client)
if pulsar_version < packaging.version.parse("0.14.999") and rewrite_paths:
job_wrapper.disable_commands_in_new_shell()
container = None
if remote_container is None:
container = self._find_container(
@@ -794,11 +796,16 @@ class PulsarJobRunner(AsynchronousJobRunner):
)
return client_outputs
@staticmethod
def pulsar_version(remote_job_config):
pulsar_version = packaging.version.parse(remote_job_config.get("pulsar_version", "0.6.0"))
return pulsar_version
@staticmethod
def check_job_config(remote_job_config, check_features=None):
check_features = check_features or {}
# 0.6.0 was newest Pulsar version that did not report it's version.
pulsar_version = packaging.version.parse(remote_job_config.get("pulsar_version", "0.6.0"))
pulsar_version = PulsarJobRunner.pulsar_version(remote_job_config)
needed_version = packaging.version.parse("0.0.0")
log.info(f"pulsar_version is {pulsar_version}")
for feature, needed in list(check_features.items()) + [("_default_", True)]:
+3
View File
@@ -137,6 +137,7 @@ class PortableDirectoryMetadataGenerator(MetadataCollectionStrategy):
job_metadata=None,
provided_metadata_style=None,
compute_tmp_dir=None,
compute_version_path=None,
include_command=True,
max_metadata_value_size=0,
max_discovered_files=None,
@@ -240,6 +241,8 @@ class PortableDirectoryMetadataGenerator(MetadataCollectionStrategy):
)
metadata_params["job_params"] = job.raw_param_dict()
metadata_params["output_collections"] = output_collections
if compute_version_path:
metadata_params["compute_version_path"] = compute_version_path
with open(metadata_params_path, "w") as f:
json.dump(metadata_params, f)
+5 -4
View File
@@ -208,15 +208,16 @@ def set_metadata_portable(
outputs_directory = os.path.join(tool_job_working_directory, "outputs")
if not os.path.exists(outputs_directory):
outputs_directory = tool_job_working_directory
metadata_directory = os.path.join(tool_job_working_directory, "metadata")
# TODO: constants...
locations = [
(metadata_directory, "tool_"),
(outputs_directory, "tool_"),
(tool_job_working_directory, ""),
(outputs_directory, ""), # # Pulsar style output directory? Was this ever used - did this ever work?
]
for directory, prefix in locations:
if os.path.exists(os.path.join(directory, f"{prefix}stdout")):
if directory and os.path.exists(os.path.join(directory, f"{prefix}stdout")):
with open(os.path.join(directory, f"{prefix}stdout"), "rb") as f:
tool_stdout = f.read(MAX_STDIO_READ_BYTES)
with open(os.path.join(directory, f"{prefix}stderr"), "rb") as f:
@@ -261,9 +262,9 @@ def set_metadata_portable(
else:
final_job_state = Job.states.ERROR
version_string_path = os.path.join("outputs", COMMAND_VERSION_FILENAME)
default_version_string_path = os.path.join("outputs", COMMAND_VERSION_FILENAME)
version_string_path = metadata_params.get("compute_version_path", default_version_string_path)
version_string = collect_shrinked_content_from_path(version_string_path)
expression_context = ExpressionContext(dict(stdout=tool_stdout[:255], stderr=tool_stderr[:255]))
# Load outputs.
@@ -18,6 +18,9 @@ cp '$input' '$output'
<assert_command>
<not_has_text text="VERSION"/>
</assert_command>
<assert_command_version>
<has_text text="4.0.0" />
</assert_command_version>
</test>
</tests>
</tool>
@@ -25,6 +25,8 @@ instance = integration_util.integration_module_instance(EmbeddedAndExtendedMetad
test_tools = integration_util.integration_tool_runner(
[
"version_command_plain",
"version_command_tool_dir",
"simple_constructs",
"metadata_bam",
# "job_properties", # https://github.com/galaxyproject/galaxy/issues/11813
+5 -1
View File
@@ -253,7 +253,11 @@ class MockJobWrapper:
@property
def job_io(self):
return Bunch(get_output_fnames=lambda: ["output1"], check_job_script_integrity=False)
return Bunch(
get_output_fnames=lambda: ["output1"],
check_job_script_integrity=False,
version_path=None,
)
@property
def is_cwl_job(self):
+3 -1
View File
@@ -191,7 +191,9 @@ class MockJobWrapper:
@property
def job_io(self):
return bunch.Bunch(get_output_fnames=lambda: [], check_job_script_integrity=False)
return bunch.Bunch(
get_output_fnames=lambda: [], check_job_script_integrity=False, version_path="/tmp/version_path"
)
def get_job(self):
return self.job