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.
This commit is contained in:
John Chilton
2022-10-16 09:56:58 +02:00
committed by mvdbeek
parent 49ecb1bc26
commit 5a3069d283
9 changed files with 193 additions and 38 deletions
@@ -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"
+36 -12
View File
@@ -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."""
+5 -3
View File
@@ -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:
@@ -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."""
+1 -1
View File
@@ -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
+1 -1
View File
@@ -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
+1 -1
View File
@@ -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 = "*"
@@ -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
@@ -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")