From 5a3069d283e25489d203ed298619a9eb64d3a274 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 6 Oct 2022 13:25:17 -0400 Subject: [PATCH] Pulsar infrastructure improvements for TES execution and MQ-less Kubernetes. Updating the Pulsar job runner/client for https://github.com/galaxyproject/pulsar/pull/302 and as described in the new docs at https://pulsar.readthedocs.io/en/latest/containers.html. --- .../dependencies/pinned-requirements.txt | 3 +- lib/galaxy/jobs/runners/pulsar.py | 48 ++++-- lib/galaxy_test/driver/driver_util.py | 8 +- lib/galaxy_test/driver/integration_util.py | 7 + mypy.ini | 2 +- packages/app/setup.cfg | 2 +- pyproject.toml | 2 +- ...ernetes_staging.py => test_coexecution.py} | 149 ++++++++++++++++-- test/integration/test_interactivetools_api.py | 10 +- 9 files changed, 193 insertions(+), 38 deletions(-) rename test/integration/{test_kubernetes_staging.py => test_coexecution.py} (61%) diff --git a/lib/galaxy/dependencies/pinned-requirements.txt b/lib/galaxy/dependencies/pinned-requirements.txt index 1c642fa6d3f..22d9f307f33 100644 --- a/lib/galaxy/dependencies/pinned-requirements.txt +++ b/lib/galaxy/dependencies/pinned-requirements.txt @@ -110,13 +110,14 @@ pkgutil-resolve-name==1.3.10 ; python_version >= "3.7" and python_version < "3.9 prompt-toolkit==3.0.31 ; python_version >= "3.7" and python_version < "3.11" prov==1.5.1 ; python_version >= "3.7" and python_version < "3.11" psutil==5.9.2 ; python_version >= "3.7" and python_version < "3.11" -pulsar-galaxy-lib==0.14.16 ; python_version >= "3.7" and python_version < "3.11" +pulsar-galaxy-lib==0.15.0.dev0 ; python_version >= "3.7" and python_version < "3.11" py==1.11.0 ; python_version >= "3.7" and python_version < "3.11" and implementation_name == "pypy" pyasn1==0.4.8 ; python_version >= "3.7" and python_version < "3.11" pycparser==2.21 ; python_version >= "3.7" and python_version < "3.11" pycryptodome==3.15.0 ; python_version >= "3.7" and python_version < "3.11" pydantic==1.10.2 ; python_version >= "3.7" and python_version < "3.11" pydantic[email]==1.10.2 ; python_version >= "3.7" and python_version < "3.11" +pydantic-tes==0.1.5 ; python_version >= "3.7" and python_version < "3.11" pydot==1.4.2 ; python_version >= "3.7" and python_version < "3.11" pyeventsystem==0.1.0 ; python_version >= "3.7" and python_version < "3.11" pyfaidx==0.7.1 ; python_version >= "3.7" and python_version < "3.11" diff --git a/lib/galaxy/jobs/runners/pulsar.py b/lib/galaxy/jobs/runners/pulsar.py index 04c6d945902..0dcc371bb8c 100644 --- a/lib/galaxy/jobs/runners/pulsar.py +++ b/lib/galaxy/jobs/runners/pulsar.py @@ -9,6 +9,10 @@ import os import re import subprocess from time import sleep +from typing import ( + Any, + Dict, +) import packaging.version import pulsar.core @@ -420,20 +424,22 @@ class PulsarJobRunner(AsynchronousJobRunner): guest_ports=job_wrapper.guest_ports, tool_directory_required_files=tool_directory_required_files, ) - job_id = pulsar_submit_job(client, client_job_description, remote_job_config) - log.info(f"Pulsar job submitted with job_id {job_id}") + external_job_id = pulsar_submit_job(client, client_job_description, remote_job_config) + log.info(f"Pulsar job submitted with job_id {external_job_id}") job = job_wrapper.get_job() # Set the job destination here (unlike other runners) because there are likely additional job destination # params from the Pulsar client. # Flush with change_state. - job_wrapper.set_job_destination(job_destination, external_id=job_id, flush=False, job=job) + job_wrapper.set_job_destination(job_destination, external_id=external_job_id, flush=False, job=job) job_wrapper.change_state(model.Job.states.QUEUED, job=job) except Exception: job_wrapper.fail("failure running job", exception=True) log.exception("failure running job %d", job_wrapper.job_id) return - pulsar_job_state = AsynchronousJobState(job_wrapper=job_wrapper, job_id=job_id, job_destination=job_destination) + pulsar_job_state = AsynchronousJobState( + job_wrapper=job_wrapper, job_id=external_job_id, job_destination=job_destination + ) pulsar_job_state.old_state = True pulsar_job_state.running = False self.monitor_job(pulsar_job_state) @@ -584,7 +590,7 @@ class PulsarJobRunner(AsynchronousJobRunner): def get_client_from_state(self, job_state): job_destination_params = job_state.job_destination.params - job_id = job_state.job_id + job_id = job_state.job_wrapper.job_id # we want the Galaxy ID here, job_state.job_id is the external one. return self.get_client(job_destination_params, job_id) def get_client(self, job_destination_params, job_id, env=None): @@ -910,6 +916,7 @@ class PulsarJobRunner(AsynchronousJobRunner): def __async_update(self, full_status): galaxy_job_id = None + remote_job_id = None try: remote_job_id = full_status["job_id"] if len(remote_job_id) == 32: @@ -924,7 +931,7 @@ class PulsarJobRunner(AsynchronousJobRunner): job_state = self._job_state(job, job_wrapper) self._update_job_state_for_status(job_state, full_status["status"], full_status=full_status) except Exception: - log.exception("Failed to update Pulsar job status for job_id %s", galaxy_job_id) + log.exception("Failed to update Pulsar job status for job_id (%s/%s)" % (galaxy_job_id, remote_job_id)) raise # Nothing else to do? - Attempt to fail the job? @@ -954,21 +961,20 @@ class PulsarMQJobRunner(PulsarJobRunner): ) -KUBERNETES_DESTINATION_DEFAULTS = { +DEFAULT_PULSAR_CONTAINER = "galaxy/pulsar-pod-staging:0.15.0.1" +COEXECUTION_DESTENTATION_DEFAULTS = { "default_file_action": "remote_transfer", "rewrite_parameters": "true", "jobs_directory": "/pulsar_staging", - "pulsar_container_image": "galaxy/pulsar-pod-staging:0.14.0", + "pulsar_container_image": DEFAULT_PULSAR_CONTAINER, "remote_container_handling": True, - "k8s_enabled": True, "url": PARAMETER_SPECIFICATION_IGNORED, "private_token": PARAMETER_SPECIFICATION_IGNORED, } -class PulsarKubernetesJobRunner(PulsarMQJobRunner): - destination_defaults = KUBERNETES_DESTINATION_DEFAULTS - poll = True # Poll so we can check API for pod IP for ITs. +class PulsarCoexecutionJobRunner(PulsarMQJobRunner): + destination_defaults = COEXECUTION_DESTENTATION_DEFAULTS def _populate_parameter_defaults(self, job_destination): super()._populate_parameter_defaults(job_destination) @@ -982,6 +988,24 @@ class PulsarKubernetesJobRunner(PulsarMQJobRunner): pulsar_app_config["staging_directory"] = params.get("jobs_directory") +KUBERNETES_DESTINATION_DEFAULTS: Dict[str, Any] = {"k8s_enabled": True, **COEXECUTION_DESTENTATION_DEFAULTS} + + +class PulsarKubernetesJobRunner(PulsarCoexecutionJobRunner): + destination_defaults = KUBERNETES_DESTINATION_DEFAULTS + poll = True # Poll so we can check API for pod IP for ITs. + + +TES_DESTENTATION_DEFAULTS: Dict[str, Any] = { + "tes_url": PARAMETER_SPECIFICATION_REQUIRED, + **COEXECUTION_DESTENTATION_DEFAULTS, +} + + +class PulsarTesJobRunner(PulsarCoexecutionJobRunner): + destination_defaults = TES_DESTENTATION_DEFAULTS + + class PulsarRESTJobRunner(PulsarJobRunner): """Flavor of Pulsar job runner with sensible defaults for RESTful usage.""" diff --git a/lib/galaxy_test/driver/driver_util.py b/lib/galaxy_test/driver/driver_util.py index 7c95fd1c452..8d0bddfdf53 100644 --- a/lib/galaxy_test/driver/driver_util.py +++ b/lib/galaxy_test/driver/driver_util.py @@ -510,6 +510,8 @@ def wait_for_http_server(host, port, prefix=None, sleep_amount=0.1, sleep_tries= prefix = f"{prefix}/" for _ in range(sleep_tries): # directly test the app, not the proxy + if port and isinstance(port, str): + port = int(port) conn = http.client.HTTPConnection(host, port) try: conn.request("GET", prefix) @@ -525,7 +527,7 @@ def wait_for_http_server(host, port, prefix=None, sleep_amount=0.1, sleep_tries= raise Exception(message) -def attempt_port(port): +def attempt_port(port: int) -> Optional[int]: sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) try: sock.bind(("", port)) @@ -535,9 +537,9 @@ def attempt_port(port): return None -def attempt_ports(port=None, set_galaxy_web_port=True): +def attempt_ports(port=None, set_galaxy_web_port=True) -> str: if port is not None: - if not attempt_port(port): + if not attempt_port(int(port)): raise Exception(f"An existing process seems bound to specified test server port [{port}]") return port else: diff --git a/lib/galaxy_test/driver/integration_util.py b/lib/galaxy_test/driver/integration_util.py index 990ae85f599..bfb718e7ba7 100644 --- a/lib/galaxy_test/driver/integration_util.py +++ b/lib/galaxy_test/driver/integration_util.py @@ -91,6 +91,13 @@ def skip_if_github_workflow(): return pytest.mark.skip("This test is skipped for Github actions.") +def skip_unless_environ(env_var): + if os.environ.get(env_var): + return _identity + + return pytest.mark.skip(f"{env_var} must be set for this test") + + class IntegrationInstance(UsesApiTestCaseMixin, UsesCeleryTasks): """Unit test case with utilities for spinning up Galaxy.""" diff --git a/mypy.ini b/mypy.ini index 87ca87e4aff..d1e7d000a0b 100644 --- a/mypy.ini +++ b/mypy.ini @@ -666,7 +666,7 @@ check_untyped_defs = False check_untyped_defs = False [mypy-integration.objectstore.test_objectstore_datatype_upload] check_untyped_defs = False -[mypy-integration.test_kubernetes_staging] +[mypy-integration.test_coexecution] check_untyped_defs = False [mypy-integration.test_interactivetools_api] check_untyped_defs = False diff --git a/packages/app/setup.cfg b/packages/app/setup.cfg index 27684524450..1639cc00577 100644 --- a/packages/app/setup.cfg +++ b/packages/app/setup.cfg @@ -57,7 +57,7 @@ install_requires = packaging paramiko!=2.9.0,!=2.9.1 pebble - pulsar-galaxy-lib>=0.14.13 + pulsar-galaxy-lib>=0.15.0.dev0 pydantic pysam PyJWT diff --git a/pyproject.toml b/pyproject.toml index fae53ce27dd..0351a35cf40 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -75,7 +75,7 @@ Parsley = "*" Paste = "*" pebble = "*" psutil = "*" -pulsar-galaxy-lib = ">=0.14.13" +pulsar-galaxy-lib = ">=0.15.0.dev0" pycryptodome = "*" pydantic = {version = "*", extras = ["email"]} PyJWT = "*" diff --git a/test/integration/test_kubernetes_staging.py b/test/integration/test_coexecution.py similarity index 61% rename from test/integration/test_kubernetes_staging.py rename to test/integration/test_coexecution.py index c868faac00a..8b87aa05c98 100644 --- a/test/integration/test_kubernetes_staging.py +++ b/test/integration/test_coexecution.py @@ -12,6 +12,7 @@ rabbitmq installed via Homebrew, and if a fixed port is set for the test. """ import os +import platform import random import string import tempfile @@ -34,9 +35,7 @@ from .test_kubernetes_runner import KubernetesDatasetPopulator from .test_local_job_cancellation import CancelsJob TOOL_DIR = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir, "tools")) -GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST = os.environ.get( - "GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST", "host.docker.internal" -) +GALAXY_TEST_INFRASTRUCTURE_HOST = os.environ.get("GALAXY_TEST_INFRASTRUCTURE_HOST", "_PLATFORM_AUTO_") AMQP_URL = integration_util.AMQP_URL GALAXY_TEST_KUBERNETES_NAMESPACE = os.environ.get("GALAXY_TEST_K8S_NAMESPACE", "default") @@ -121,10 +120,85 @@ def job_config(template_str, jobs_directory): return job_conf.name -@integration_util.skip_unless_kubernetes() +TES_CONTAINERIZED_TEMPLATE = """ +runners: + local: + load: galaxy.jobs.runners.local:LocalJobRunner + workers: 1 + pulsar_tes: + load: galaxy.jobs.runners.pulsar:PulsarTesJobRunner + amqp_url: ${amqp_url} + +execution: + default: pulsar_tes_environment + environments: + pulsar_tes_environment: + runner: pulsar_tes + tes_url: ${tes_url} + docker_enabled: true + docker_default_container_id: busybox:ubuntu-14.04 + pulsar_app_config: + message_queue_url: '${container_amqp_url}' + env: + - name: SOME_ENV_VAR + value: '42' + local_environment: + runner: local +tools: + - id: __DATA_FETCH__ + environment: local_environment +""" + + +TES_DEPENDENCY_RESOLUTION_TEMPLATE = """ +runners: + local: + load: galaxy.jobs.runners.local:LocalJobRunner + workers: 1 + pulsar_tes: + load: galaxy.jobs.runners.pulsar:PulsarTesJobRunner + amqp_url: ${amqp_url} + +execution: + default: pulsar_tes_environment + environments: + pulsar_tes_environment: + tes_url: ${tes_url} + runner: pulsar_tes + pulsar_app_config: + message_queue_url: '${container_amqp_url}' + env: + - name: SOME_ENV_VAR + value: '42' + local_environment: + runner: local +tools: + - id: __DATA_FETCH__ + environment: local_environment +""" + + +def tes_job_config(template_str, jobs_directory): + job_conf_template = string.Template(template_str) + container_amqp_url = to_infrastructure_uri(AMQP_URL) + instance_id = "".join(random.choice(string.ascii_lowercase) for i in range(8)) + tes_url = os.environ.get("FUNNEL_SERVER_TARGET") + job_conf_str = job_conf_template.substitute( + jobs_directory=jobs_directory, + tool_directory=TOOL_DIR, + instance_id=instance_id, + tes_url=tes_url, + amqp_url=AMQP_URL, + container_amqp_url=container_amqp_url, + ) + with tempfile.NamedTemporaryFile(suffix="_tes_integration_job_conf.yml", mode="w", delete=False) as job_conf: + job_conf.write(job_conf_str) + return job_conf.name + + @integration_util.skip_unless_amqp() @integration_util.skip_if_github_workflow() -class TestKubernetesStaging(BaseJobEnvironmentIntegrationTestCase, MulledJobTestCases): +class TestCoexecution(BaseJobEnvironmentIntegrationTestCase, MulledJobTestCases): def setUp(self): super().setUp() self.dataset_populator = KubernetesDatasetPopulator(self.galaxy_interactor) @@ -137,7 +211,8 @@ class TestKubernetesStaging(BaseJobEnvironmentIntegrationTestCase, MulledJobTest super().setUpClass() -class TestKubernetesStagingContainerIntegration(CancelsJob, TestKubernetesStaging): +@integration_util.skip_unless_kubernetes() +class TestKubernetesStagingContainerIntegration(CancelsJob, TestCoexecution): @classmethod def handle_galaxy_config_kwds(cls, config): config["jobs_directory"] = cls.jobs_directory @@ -189,7 +264,8 @@ class TestKubernetesStagingContainerIntegration(CancelsJob, TestKubernetesStagin return active -class TestKubernetesDependencyResolutionIntegration(TestKubernetesStaging): +@integration_util.skip_unless_kubernetes() +class TestKubernetesDependencyResolutionIntegration(TestCoexecution): @classmethod def handle_galaxy_config_kwds(cls, config): config["jobs_directory"] = cls.jobs_directory @@ -208,21 +284,66 @@ class TestKubernetesDependencyResolutionIntegration(TestKubernetesStaging): assert "0.7.15-r1140" in output +@integration_util.skip_unless_environ("FUNNEL_SERVER_TARGET") +class TestTesCoexecutionContainerIntegration(TestCoexecution): + @classmethod + def handle_galaxy_config_kwds(cls, config): + config["jobs_directory"] = cls.jobs_directory + config["file_path"] = cls.jobs_directory + config["job_config_file"] = tes_job_config(TES_CONTAINERIZED_TEMPLATE, cls.jobs_directory) + + config["default_job_shell"] = "/bin/sh" + # Disable tool dependency resolution. + config["tool_dependency_dir"] = "none" + set_infrastucture_url(config) + + +@integration_util.skip_unless_environ("FUNNEL_SERVER_TARGET") +class TestTesDependencyResolutionIntegration(TestCoexecution): + @classmethod + def handle_galaxy_config_kwds(cls, config): + config["jobs_directory"] = cls.jobs_directory + config["file_path"] = cls.jobs_directory + config["job_config_file"] = tes_job_config(TES_DEPENDENCY_RESOLUTION_TEMPLATE, cls.jobs_directory) + + config["default_job_shell"] = "/bin/sh" + # Disable tool dependency resolution. + config["tool_dependency_dir"] = "none" + set_infrastucture_url(config) + + def test_mulled_simple(self): + self.dataset_populator.run_tool("mulled_example_simple", {}, self.history_id) + self.dataset_populator.wait_for_history(self.history_id, assert_ok=True) + output = self.dataset_populator.get_history_dataset_content(self.history_id, timeout=EXTENDED_TIMEOUT) + assert "0.7.15-r1140" in output + + def set_infrastucture_url(config): - infrastructure_url = f"http://{GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST}:$GALAXY_WEB_PORT" + hostname = to_infrastructure_uri("0.0.0.0") + infrastructure_url = f"http://{hostname}:$GALAXY_WEB_PORT" config["galaxy_infrastructure_url"] = infrastructure_url -def to_infrastructure_uri(uri): +def to_infrastructure_uri(uri: str) -> str: # remap MQ or file server URI hostnames for in-container versions, this is sloppy # should actually parse the URI and rebuild with correct host - # similar code found in Pulsar integration_tests.py. + # Copied from Pulsar's integraiton tests. + infrastructure_host = os.environ.get("GALAXY_TEST_INFRASTRUCTURE_HOST") + if infrastructure_host == "_PLATFORM_AUTO_": + system = platform.system() + if system in ["Darwin", "Windows"]: + # assume Docker Desktop is installed and use its domain + infrastructure_host = "host.docker.internal" + else: + # native linux Docker sometimes sets up + infrastructure_host = "172.17.0.1" + infrastructure_uri = uri - if GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST: + if infrastructure_host: if "0.0.0.0" in infrastructure_uri: - infrastructure_uri = infrastructure_uri.replace("0.0.0.0", GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST) + infrastructure_uri = infrastructure_uri.replace("0.0.0.0", infrastructure_host) elif "localhost" in infrastructure_uri: - infrastructure_uri = infrastructure_uri.replace("localhost", GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST) + infrastructure_uri = infrastructure_uri.replace("localhost", infrastructure_host) elif "127.0.0.1" in infrastructure_uri: - infrastructure_uri = infrastructure_uri.replace("127.0.0.1", GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST) + infrastructure_uri = infrastructure_uri.replace("127.0.0.1", infrastructure_host) return infrastructure_uri diff --git a/test/integration/test_interactivetools_api.py b/test/integration/test_interactivetools_api.py index 04debf87567..1e0346b93cb 100644 --- a/test/integration/test_interactivetools_api.py +++ b/test/integration/test_interactivetools_api.py @@ -11,16 +11,16 @@ from galaxy_test.base.populators import ( wait_on, ) from galaxy_test.driver import integration_util +from .test_coexecution import ( + CONTAINERIZED_TEMPLATE, + job_config, + set_infrastucture_url, +) from .test_containerized_jobs import ( ContainerizedIntegrationTestCase, disable_dependency_resolution, DOCKERIZED_JOB_CONFIG_FILE, ) -from .test_kubernetes_staging import ( - CONTAINERIZED_TEMPLATE, - job_config, - set_infrastucture_url, -) SCRIPT_DIRECTORY = os.path.abspath(os.path.dirname(__file__)) EMBEDDED_PULSAR_JOB_CONFIG_FILE_DOCKER = os.path.join(SCRIPT_DIRECTORY, "embedded_pulsar_docker_job_conf.yml")