From 08f66a924b8a46d673fc6237dec3a4e8fcb34f21 Mon Sep 17 00:00:00 2001 From: Nate Coraor Date: Tue, 28 Feb 2017 17:09:56 -0500 Subject: [PATCH] GIE container waiting and swarm management. 1. Create a new API endpoint and client spin to wait for GIE containers to become ready. IEs are not required to use it, but usage does require a slight change to their `load_notebook()` functions. 2. Add a Docker swarm manager daemon that runs automatically to add nodes when services are waiting, clean up exited services, and remove idle nodes. Shuts down automatically when there's nothing to do. The format of the GIE proxy session database has changed, if GIEs are running when you restart Galaxy, they will be interrupted. If you run the GIE proxy separate from Galaxy, you should restart it after upgrading. --- .../galaxy.interactive_environments.js | 55 +- config/galaxy.ini.sample | 5 + .../bam_iobio/static/js/bam_iobio.js | 10 +- .../bam_iobio/templates/bam_iobio.mako | 4 +- .../common/templates/ie.mako | 1 + .../jupyter/static/js/jupyter.js | 10 +- .../jupyter/templates/jupyter.mako | 4 +- .../neo/static/js/neo.js | 11 +- .../neo/templates/neo.mako | 4 +- .../phinch/static/js/phinch.js | 6 +- .../phinch/templates/phinch.mako | 4 +- .../rstudio/static/js/rstudio.js | 100 ++-- .../rstudio/templates/rstudio.mako | 4 +- config/swarm_manager_conf.yml.sample | 84 +++ cron/clean_docker_swarm_mode_services.sh | 24 - doc/source/admin/interactive_environments.rst | 36 +- lib/galaxy/config.py | 1 + lib/galaxy/container/__init__.py | 0 lib/galaxy/container/docker_swarm.py | 552 ++++++++++++++++++ .../web/base/interactive_environments.py | 8 +- lib/galaxy/web/proxy/__init__.py | 113 +++- lib/galaxy/web/proxy/js/lib/mapper.js | 14 +- .../controllers/interactive_environments.py | 46 ++ .../galaxy.interactive_environments.js.map | 2 +- .../galaxy.interactive_environments.js | 2 +- 25 files changed, 943 insertions(+), 157 deletions(-) create mode 100644 config/swarm_manager_conf.yml.sample delete mode 100644 cron/clean_docker_swarm_mode_services.sh create mode 100644 lib/galaxy/container/__init__.py create mode 100644 lib/galaxy/container/docker_swarm.py create mode 100644 lib/galaxy/webapps/galaxy/controllers/interactive_environments.py diff --git a/client/galaxy/scripts/galaxy.interactive_environments.js b/client/galaxy/scripts/galaxy.interactive_environments.js index eb2ed98837a..180011d5906 100644 --- a/client/galaxy/scripts/galaxy.interactive_environments.js +++ b/client/galaxy/scripts/galaxy.interactive_environments.js @@ -17,6 +17,59 @@ function display_spinner(){ $('#main').append(''); } +/** + * Check a URL for a boolean true/false and call a callback when done. + */ +function load_when_ready(url, success_callback){ + var request_count = 0; + var timeout_time = 1000; + var timeout_time_max = 15000; + var timeout_time_step = 1000; + var timeout = function(){ + $.ajax({ + url: url, + xhrFields: { + withCredentials: true + }, + type: "GET", + timeout: 500, + dataType: "json", + success: function(data){ + if(data == true){ + console.log("Galaxy reports IE container ready, returning"); + clear_main_area(); + toastr.clear(); + success_callback(); + }else if(data == false){ + if(request_count == 0){ + display_spinner(); + toastr.info( + "Galaxy is launching a container in which to run this interactive environment. Please wait...", + {'closeButton': true, 'tapToDismiss': false} + ); + } + request_count++; + if(timeout_time < timeout_time_max){ + timeout_time += timeout_time_step; + } + console.log("Readiness request " + request_count + " sleeping " + timeout_time / 1000 + "s"); + window.setTimeout(timeout, timeout_time) + }else{ + clear_main_area(); + toastr.clear(); + toastr.error( + "Galaxy failed to launch a container in which to run this interactive environment, contact your administrator.", + "Error", + {'closeButton': true, 'tapToDismiss': false} + ); + } + } + }); + } + window.setTimeout(timeout, timeout_time); +} + + /** * Test availability of a URL, and call a callback when done. @@ -43,7 +96,7 @@ function test_ie_availability(url, success_callback){ }, error: function(jqxhr, status, error){ request_count++; - console.log("Request " + request_count); + console.log("Availability request " + request_count); if(request_count > 30){ clearInterval(interval); clear_main_area(); diff --git a/config/galaxy.ini.sample b/config/galaxy.ini.sample index 15f849cd585..c201ffc8818 100644 --- a/config/galaxy.ini.sample +++ b/config/galaxy.ini.sample @@ -365,6 +365,11 @@ paste.app_factory = galaxy.web.buildapp:app_factory # the plugin's ini config file. #interactive_environment_swarm_mode = False +# Galaxy can run a "swarm manager" service that will monitor utilization of the +# swarm and provision/deprovision worker nodes as necessary. The service has +# its own configuration file. +#swarm_manager_config_file = config/swarm_manager_conf.yml + # Interactive tour directory: where to store interactive tour definition files. # Galaxy ships with several basic interface tours enabled, though a different # directory with custom tours can be specified here. The path is relative to the diff --git a/config/plugins/interactive_environments/bam_iobio/static/js/bam_iobio.js b/config/plugins/interactive_environments/bam_iobio/static/js/bam_iobio.js index 395904e329b..7b019f63630 100644 --- a/config/plugins/interactive_environments/bam_iobio/static/js/bam_iobio.js +++ b/config/plugins/interactive_environments/bam_iobio/static/js/bam_iobio.js @@ -22,12 +22,10 @@ function message_failed_connection(){ * */ function load_notebook(notebook_access_url){ - $( document ).ready(function() { - // Test notebook_login_url for accessibility, executing the login+load function whenever - // we've successfully connected to the IE. - test_ie_availability(notebook_access_url, function(){ - _handle_notebook_loading(notebook_access_url); - }); + // Test notebook_login_url for accessibility, executing the login+load function whenever + // we've successfully connected to the IE. + test_ie_availability(notebook_access_url, function(){ + _handle_notebook_loading(notebook_access_url); }); } diff --git a/config/plugins/interactive_environments/bam_iobio/templates/bam_iobio.mako b/config/plugins/interactive_environments/bam_iobio/templates/bam_iobio.mako index 8f18630074e..97eb4b78d25 100644 --- a/config/plugins/interactive_environments/bam_iobio/templates/bam_iobio.mako +++ b/config/plugins/interactive_environments/bam_iobio/templates/bam_iobio.mako @@ -46,7 +46,9 @@ root = h.url_for( '/' ) var startup = function(){ // Load notebook requirejs(['interactive_environments', 'plugin/bam_iobio'], function(){ - load_notebook(notebook_access_url); + load_when_ready(ie_readiness_url, function(){ + load_notebook(notebook_access_url); + }); }); }; diff --git a/config/plugins/interactive_environments/common/templates/ie.mako b/config/plugins/interactive_environments/common/templates/ie.mako index fac6366e50d..338ccc7e263 100644 --- a/config/plugins/interactive_environments/common/templates/ie.mako +++ b/config/plugins/interactive_environments/common/templates/ie.mako @@ -9,6 +9,7 @@ ie_password = '${ ie_request.notebook_pw }'; var galaxy_root = '${ ie_request.attr.root }'; var app_root = '${ ie_request.attr.app_root }'; +var ie_readiness_url = '${ ie_request.url_template("${PROXY_PREFIX}/interactive_environments/ready") }'; diff --git a/config/plugins/interactive_environments/jupyter/static/js/jupyter.js b/config/plugins/interactive_environments/jupyter/static/js/jupyter.js index 3b0ca53e51e..0e717c796b0 100644 --- a/config/plugins/interactive_environments/jupyter/static/js/jupyter.js +++ b/config/plugins/interactive_environments/jupyter/static/js/jupyter.js @@ -34,12 +34,10 @@ function message_no_auth(){ * */ function load_notebook(password, notebook_login_url, notebook_access_url){ - $( document ).ready(function() { - // Test notebook_login_url for accessibility, executing the login+load function whenever - // we've successfully connected to the IE. - test_ie_availability(notebook_login_url, function(){ - _handle_notebook_loading(password, notebook_login_url, notebook_access_url); - }); + // Test notebook_login_url for accessibility, executing the login+load function whenever + // we've successfully connected to the IE. + test_ie_availability(notebook_login_url, function(){ + _handle_notebook_loading(password, notebook_login_url, notebook_access_url); }); } diff --git a/config/plugins/interactive_environments/jupyter/templates/jupyter.mako b/config/plugins/interactive_environments/jupyter/templates/jupyter.mako index 69bf37fc454..8f8490f5367 100644 --- a/config/plugins/interactive_environments/jupyter/templates/jupyter.mako +++ b/config/plugins/interactive_environments/jupyter/templates/jupyter.mako @@ -60,7 +60,9 @@ requirejs(['interactive_environments', 'plugin/jupyter'], function(){ // Load notebook requirejs(['interactive_environments', 'plugin/jupyter'], function(){ - load_notebook(ie_password, notebook_login_url, notebook_access_url); + load_when_ready(ie_readiness_url, function(){ + load_notebook(ie_password, notebook_login_url, notebook_access_url); + }); }); diff --git a/config/plugins/interactive_environments/neo/static/js/neo.js b/config/plugins/interactive_environments/neo/static/js/neo.js index 068d68b6515..f1af1254819 100644 --- a/config/plugins/interactive_environments/neo/static/js/neo.js +++ b/config/plugins/interactive_environments/neo/static/js/neo.js @@ -1,12 +1,9 @@ // Load an interactive environment (IE) from a remote URL // @param {String} notebook_access_url: the URL embeded in the page and loaded function load_notebook(notebook_access_url){ - // When the page has completely loaded... - $( document ).ready(function() { - // Test if we can access the GIE, and if so, execute the function - // to load the GIE for the user. - test_ie_availability(notebook_access_url, function(){ - append_notebook(notebook_access_url); - }); + // Test if we can access the GIE, and if so, execute the function + // to load the GIE for the user. + test_ie_availability(notebook_access_url, function(){ + append_notebook(notebook_access_url); }); } diff --git a/config/plugins/interactive_environments/neo/templates/neo.mako b/config/plugins/interactive_environments/neo/templates/neo.mako index b241b9dae87..ab080fde489 100644 --- a/config/plugins/interactive_environments/neo/templates/neo.mako +++ b/config/plugins/interactive_environments/neo/templates/neo.mako @@ -33,7 +33,9 @@ requirejs(['interactive_environments', 'plugin/neo'], function () { - load_notebook(url); + load_when_ready(ie_readiness_url, function(){ + load_notebook(url); + }); }); diff --git a/config/plugins/interactive_environments/phinch/static/js/phinch.js b/config/plugins/interactive_environments/phinch/static/js/phinch.js index 1b4c993e554..e750943f96d 100644 --- a/config/plugins/interactive_environments/phinch/static/js/phinch.js +++ b/config/plugins/interactive_environments/phinch/static/js/phinch.js @@ -1,8 +1,6 @@ function load_notebook(url){ - $( document ).ready(function() { - test_ie_availability(url, function(){ - append_notebook(url) - }); + test_ie_availability(url, function(){ + append_notebook(url) }); } diff --git a/config/plugins/interactive_environments/phinch/templates/phinch.mako b/config/plugins/interactive_environments/phinch/templates/phinch.mako index c446bf7d6e5..54cd9aa5503 100644 --- a/config/plugins/interactive_environments/phinch/templates/phinch.mako +++ b/config/plugins/interactive_environments/phinch/templates/phinch.mako @@ -36,7 +36,9 @@ requirejs(['interactive_environments', 'plugin/phinch'], function(){ // Load notebook requirejs(['interactive_environments', 'plugin/phinch'], function(){ - load_notebook(url); + load_when_ready(ie_readiness_url, function(){ + load_notebook(url); + }); }); diff --git a/config/plugins/interactive_environments/rstudio/static/js/rstudio.js b/config/plugins/interactive_environments/rstudio/static/js/rstudio.js index cc0b62ee7a1..38764d84567 100644 --- a/config/plugins/interactive_environments/rstudio/static/js/rstudio.js +++ b/config/plugins/interactive_environments/rstudio/static/js/rstudio.js @@ -15,58 +15,56 @@ function message_failed_connection(){ * */ function load_notebook(notebook_login_url, notebook_access_url, notebook_pubkey_url, username){ - $( document ).ready(function() { - // Test notebook_login_url for accessibility, executing the login+load function whenever - // we've successfully connected to the IE. - test_ie_availability(notebook_pubkey_url, function(){ - var payload = username + "\n" + ie_password; - $.ajax({ - type: 'GET', - url: notebook_pubkey_url, - xhrFields: { + // Test notebook_login_url for accessibility, executing the login+load function whenever + // we've successfully connected to the IE. + test_ie_availability(notebook_pubkey_url, function(){ + var payload = username + "\n" + ie_password; + $.ajax({ + type: 'GET', + url: notebook_pubkey_url, + xhrFields: { + withCredentials: true + }, + success: function(response_text){ + var chunks = response_text.split(':', 2); + var exp = chunks[0]; + var mod = chunks[1]; + console.log("Found " + exp +" and " + mod); + var rsa = new RSAKey(); + rsa.setPublic(mod, exp); + console.log("Encrypting '" + username + "', '" + ie_password + "'"); + var enc_hex = rsa.encrypt(payload); + var encrypted = hex2b64(enc_hex); + console.log("E: " + encrypted); + + // Now we can login + $.ajax({ + type: "POST", + // to the Login URL + url: notebook_login_url, + // With our password + data: { + 'v': encrypted, + 'persist': 1, + 'clientPath': '/rstudio/auth-sign-in', + 'appUri': '', + }, + contentType: "application/x-www-form-urlencoded", + xhrFields: { withCredentials: true - }, - success: function(response_text){ - var chunks = response_text.split(':', 2); - var exp = chunks[0]; - var mod = chunks[1]; - console.log("Found " + exp +" and " + mod); - var rsa = new RSAKey(); - rsa.setPublic(mod, exp); - console.log("Encrypting '" + username + "', '" + ie_password + "'"); - var enc_hex = rsa.encrypt(payload); - var encrypted = hex2b64(enc_hex); - console.log("E: " + encrypted); - - // Now we can login - $.ajax({ - type: "POST", - // to the Login URL - url: notebook_login_url, - // With our password - data: { - 'v': encrypted, - 'persist': 1, - 'clientPath': '/rstudio/auth-sign-in', - 'appUri': '', - }, - contentType: "application/x-www-form-urlencoded", - xhrFields: { - withCredentials: true - }, - // If that is successful, load the notebook - success: function(){ - append_notebook(notebook_access_url); - }, - error: function(jqxhr, status, error){ - message_failed_connection(); - // Do we want to try and load the notebook anyway? Just in case? - append_notebook(notebook_access_url); - } - }); - } - }); - + }, + // If that is successful, load the notebook + success: function(){ + append_notebook(notebook_access_url); + }, + error: function(jqxhr, status, error){ + message_failed_connection(); + // Do we want to try and load the notebook anyway? Just in case? + append_notebook(notebook_access_url); + } + }); + } }); + }); } diff --git a/config/plugins/interactive_environments/rstudio/templates/rstudio.mako b/config/plugins/interactive_environments/rstudio/templates/rstudio.mako index fdb99cde32d..c9b472c6f98 100644 --- a/config/plugins/interactive_environments/rstudio/templates/rstudio.mako +++ b/config/plugins/interactive_environments/rstudio/templates/rstudio.mako @@ -64,7 +64,9 @@ requirejs([ 'crypto/base64', 'plugin/rstudio' ], function(){ - load_notebook(notebook_login_url, notebook_access_url, notebook_pubkey_url, "${ USERNAME }"); + load_when_ready(ie_readiness_url, function(){ + load_notebook(notebook_login_url, notebook_access_url, notebook_pubkey_url, "${ USERNAME }"); + }); });
diff --git a/config/swarm_manager_conf.yml.sample b/config/swarm_manager_conf.yml.sample new file mode 100644 index 00000000000..31c5632519c --- /dev/null +++ b/config/swarm_manager_conf.yml.sample @@ -0,0 +1,84 @@ +--- +# Galaxy docker swarm manager configuration file +# +# To configure the location of this file, use the `swarm_manager_config_file` +# setting in galaxy.ini + +# When the swarm manager daemonizes, it writes a pid file so that only one +# manager will run at a time. This is the path to that pid file. +# {xdg_data_home} will be templated automatically and defaults to +# ~/.local/share as per the XDG specification +#pid_file: '{xdg_data_home}/galaxy_swarm_manager.pid' + +# Program output will be written to the log +#log_file: '{xdg_data_home}/galaxy_swarm_manager.log' + +# As with GIE plugins, you can modify the base docker command ({docker_args} +# must be present and will be filled in with the docker subcommand and +# arguments) +#command: 'docker {docker_args}' + +# Managed services should be started with this string at the beginning of their +# name. It should match the value of CONTAINER_NAME_PREFIX in +# lib/galaxy/web/base/interactive_environments.py, so you should not change +# this unless you change both. +#service_prefix: galaxy_gie_ + +# Limits: +# +# - max_waiting_services: number of services that should be waiting of each +# "CPU class" (number of CPUs requested e.g. with --reserve-cpu) before +# attempting to spawn a node +# - max_wait_time: number of seconds a service should be waiting before +# attempting to spawn a node +# - max_node_idle_time: number of seconds a node should be idle before +# terminating it +# - max_node_counts: a dictionary controlling the maximum number of nodes of +# each CPU class that the swarm manager will attempt to spawn, e.g.: +# max_node_counts: +# 1: 10 # spawn up to 10 x 1-CPU nodes +# 2: 3 # spawn up to 3 x 2-CPU nodes +# 4: 1 # spawn up to 1 x 4-CPU nodes +#max_waiting_services: 0 +#max_wait_time: 5 +#max_node_idle_time: 120 +#max_node_counts: {} + +# If set, only manage nodes whose swarm hostnames begin with this prefix. +# Otherwise, attempt to manage all nodes +#node_prefix: null + +# Amount of time to wait for a spawning node to appear in `docker node ls` +# before considering it failed +#spawn_wait_time: 30 + +# Command to run to spawn new nodes. This command should join the node to the +# swarm. Can include template variables: +# - {cpu_class}: CPU class as explained above (the value of --reserve-cpu) +# - {cpus_needed}: Total number of CPUs of the given class needed to run the +# waiting services +# If this command does not block until the node is joined to the swarm, make +# sure it at least completes that step in `spawn_wait_time` once it returns +# control. This command should return a space-separated list of nodes. If the +# nodes have a different number of CPUs than the class that they were started +# for, you can include that number after a colon (e.g. `node1:4`). +#spawn_command: /bin/true + +# Command to run to destroy idle nodes. Can include template variables: +# - {nodes}: Space-separated list of node names to destroy +# This command should block until at least the point at which any nodes being +# deallocated no longer appear in `docker node ls`. +#destroy_command: /bin/true + +# Command to run if either of the above commands failed (e.g. to notify an +# administrator). Can include template variables: +# - {failed_command}: Command line of the command that failed +#command_failure_command: /bin/true + +# Number of times to retry spawn/destroy commands before considering them to +# have failed, and seconds to wait between retries +#command_retries: 0 +#command_retry_wait: 10 + +# Stop the swarm manager daemon when there are no services or nodes to manage +#terminate_when_idle: True diff --git a/cron/clean_docker_swarm_mode_services.sh b/cron/clean_docker_swarm_mode_services.sh deleted file mode 100644 index e3aacd400b7..00000000000 --- a/cron/clean_docker_swarm_mode_services.sh +++ /dev/null @@ -1,24 +0,0 @@ -#!/bin/bash -# -# When running Galaxy Interactive Environments on a swarm (using Docker Engine -# swarm mode), GIE sessions that have ended will leave behind "shut down" -# Docker services. These must be removed, which you can do with this script. -# -# Note that this is dependent on the specific ordering and output format of -# `docker service ls` and `docker service ps`. As of the time of writing -# (Docker version 1.13.1) these formats cannot be controlled as can be done -# with `docker ps` and the `--format` option, so be careful when upgrading -# Docker releases. - - -CONTAINER_NAME_PREFIX='galaxy_gie_' -LOG_PATH=${1:-'/tmp/galaxy_gie_service_clean.log'} - - -{ - echo "Running cleanup at $(date)" - for service_name in $(docker service ls | awk "\$2 ~ /^$CONTAINER_NAME_PREFIX/ {print \$2}"); do - docker service ps --no-trunc $service_name | tail -1 | awk '{print $6}' | grep -q '^Shutdown$' && docker service rm $service_name; - done - echo "Done" -} >>$LOG_PATH diff --git a/doc/source/admin/interactive_environments.rst b/doc/source/admin/interactive_environments.rst index 42cf91ca012..f2f0cb30926 100644 --- a/doc/source/admin/interactive_environments.rst +++ b/doc/source/admin/interactive_environments.rst @@ -283,11 +283,11 @@ Galaxy supports both Docker Engine swarm mode and the legacy Docker Swarm system. Legacy Docker Swarm is supported without any special configuration, because the containers are still run with ``docker run`` as before. To support Docker Engine swarm mode, additional configuration is required. Begin by -editing your GIE config plugin's ini configuration file (e.g. ``jupyter.ini``) -and set the ``docker_connect_port`` and ``swarm_mode options`` in addition to -any other relevant options. Unless you are using a non-standard Docker image, -the correct value for ``docker_connect_port`` should be suggested to you in the -sample configuration file: +editing your GIE plugin's ini configuration file (e.g. ``jupyter.ini``) and set +the ``docker_connect_port`` and ``swarm_mode options`` in addition to any other +relevant options. Unless you are using a non-standard Docker image, the correct +value for ``docker_connect_port`` should be suggested to you in the sample +configuration file: .. code-block:: ini @@ -295,6 +295,12 @@ sample configuration file: docker_connect_port = 8888 swarm_mode = True +You can also enable swarm mode for *all* GIE plugins by setting +``interactive_environment_swarm_mode`` in ``galaxy.ini`` to ``True``. If using +this setting, you must still set ``docker_connect_port`` in each GIE plugin's +ini configuration file. The ``swarm_mode`` setting in individual GIE plugin +config files will override the value set in ``galaxy.ini``. + Note that your Galaxy server does not need to be a member of the swarm itself. It can use the method outlined above in the `Docker on Another Host`_ section to connect as a client to a Docker daemon acting as a swarm mode manager. @@ -303,18 +309,10 @@ Once configured, you should see that your GIE containers are started and run as services, which you can inspect using the ``docker service ls`` command and other ``docker service`` subcommands. -**Docker services are not cleaned up by Galaxy**. To clean them up, we have -provided a script that can be run from cron which will locate "shut down" -services (GIE containers which have stopped themselves) at -`cron/clean_docker_swarm_mode_services.sh -`__ -in the Galaxy source. This script can be run from cron with a crontab entry -like this example which runs every 15 minutes: +**Galaxy swarm manager** -.. code-block:: bash - - */15 * * * * bash /path/to/clean_docker_swarm_mode_services.sh /path/to/galaxy/log/dir/clean_docker_swarm_mode_services.log - -This entry would be suitable to be run as ``root`` on a swarm mode manager. You -could also run it as the Galaxy user on the Galaxy server (with modifications -to set the correct daemon socket, if running remotely). +Galaxy will start a "swarm manager" process when the first swarm mode GIE is +launched. You can control this daemon with the config file +``config/swarm_mode_manager.yml``. Consult the sample configuration at +``config/swarm_mode_manager.yml.sample`` for syntax. It will automatically shut +down when no services or nodes remain to be managed. diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index fd6e11a3117..c965906517b 100644 --- a/lib/galaxy/config.py +++ b/lib/galaxy/config.py @@ -52,6 +52,7 @@ PATH_DEFAULTS = dict( workflow_schedulers_config_file=['config/workflow_schedulers_conf.xml', 'config/workflow_schedulers_conf.xml.sample'], modules_mapping_files=['config/environment_modules_mapping.yml', 'config/environment_modules_mapping.yml.sample'], local_conda_mapping_file=['config/local_conda_mapping.yml', 'config/local_conda_mapping.yml.sample'], + swarm_manager_config_file=['config/swarm_manager_conf.yml', 'config/swarm_manager_conf.yml.sample'], ) PATH_LIST_DEFAULTS = dict( diff --git a/lib/galaxy/container/__init__.py b/lib/galaxy/container/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/lib/galaxy/container/docker_swarm.py b/lib/galaxy/container/docker_swarm.py new file mode 100644 index 00000000000..b7b67de6e09 --- /dev/null +++ b/lib/galaxy/container/docker_swarm.py @@ -0,0 +1,552 @@ +""" +Docker Swarm mode management +""" +import argparse +import errno +import json +import logging +import os +import subprocess +import sys +import time + +try: + import daemon + import daemon.pidfile + import lockfile +except ImportError: + daemon = None +import yaml + +try: + import galaxy # noqa: F401 this is a test import +except ImportError: + sys.path.insert(0, os.path.abspath(os.path.join( + os.path.dirname(__file__), + os.pardir, + os.pardir))) + +from galaxy.config import ( + configure_logging, + find_path, + find_root, +) +from galaxy.util.properties import find_config_file, load_app_properties + + +DESCRIPTION = "Daemon to manage a Docker Swarm (running in Docker Swarm mode)." +SWARM_MANAGER_CONF_DEFAULTS = { + 'pid_file': '{xdg_data_home}/galaxy_swarm_manager.pid', + 'log_file': '{xdg_data_home}/galaxy_swarm_manager.log', + 'command': 'docker {docker_args}', + 'service_prefix': 'galaxy_gie_', + 'max_waiting_services': 0, + 'max_wait_time': 5, + 'max_node_counts': {}, # max number of nodes per class to spawn + 'max_node_idle_time': 120, + 'node_prefix': None, + 'spawn_wait_time': 30, + 'spawn_command': '/bin/true', + 'destroy_command': '/bin/true', + 'command_failure_command': '/bin/true', + 'command_retries': 0, + 'command_retry_wait': 10, + 'terminate_when_idle': True, +} +OK_NODE_STATE = 'ready-active' +NODE_CPU_CLASS_LABEL = '_galaxy_cpu_class' +log = lambda *x: None # noqa: E731 + + +# TODO: pass around instances or at least namedtuples rather than these +# arbitrary dictionaries + + +class DockerInterface(object): + + def __init__(self, swarm_manager_conf): + self.swarm_manager_conf = swarm_manager_conf + self.command = swarm_manager_conf['command'] + self.service_prefix = swarm_manager_conf['service_prefix'] + + def _run_docker(self, docker_args): + raw_cmd = self.command.format(docker_args=docker_args) + p = subprocess.Popen(raw_cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, close_fds=True, shell=True) + stdout, stderr = p.communicate() + if p.returncode != 0: + log.error("%s\n%s" % (stdout, stderr)) + return None + else: + return stdout + + def _parse_docker_column_output(self, output): + """Many docker commands do not provide an option to format the output + or output in a machine-readily-parseable format (e.g. json). In order + to deal with such output and hopefully stay compatible with future + column order changes, key returned rows based on column headers. + + An assumption is made that a single space in the header row does not + separate columns - column names can have spaces in them, and columns + are separated by at least 2 spaces. This seems to be true as of Docker + 1.13.1. + """ + parsed = [] + output = output.splitlines() + header = output[0] + colstarts = [0] + colidx = 0 + spacect = 0 + if not output: + return parsed + for i, c in enumerate(header): + if c != ' ' and spacect > 1: + colidx += 1 + colstarts.append(i) + spacect = 0 + elif c == ' ': + spacect += 1 + colstarts.append(None) + colheadings = [] + for i in range(0, len(colstarts) - 1): + colheadings.append(header[colstarts[i]:colstarts[i + 1]].strip()) + for line in output[1:]: + row = {} + for i, key in enumerate(colheadings): + row[key] = line[colstarts[i]:colstarts[i + 1]].strip() + parsed.append(row) + return parsed + + def _service_inspect(self, service_id): + return self._run_docker(docker_args='service inspect {service_id}'.format(service_id=service_id)) + + def _get_reserved_cpu_count(self, service_id=None, inspect_output=None): + assert service_id or inspect_output, "Either `service_id` or `inspect_output` is required" + if not inspect_output: + inspect_output = self._service_inspect(service_id) + try: + return json.loads(inspect_output)[0]['Spec']['Resources']['Reservations']['NanoCPUs'] / 1000000000 + except KeyError: + return 1 + + def _node_inspect(self, node_name): + return self._run_docker(docker_args='node inspect {node_name}'.format(node_name=node_name)) + + def _get_node_cpu_class(self, node_name=None, inspect_output=None): + assert node_name or inspect_output, "Either `node_name` or `inspect_output` is required" + if not inspect_output: + inspect_output = self._node_inspect(node_name) + try: + return json.loads(inspect_output)[0]['Spec']['Labels'][NODE_CPU_CLASS_LABEL] + except KeyError: + return None + + def waiting_services_by_cpu_class(self): + rval = {} + service_ids = self._services_in_state('Running', 'Pending') + for service_id in service_ids: + cpu_class = self._get_reserved_cpu_count(service_id=service_id) + if cpu_class not in rval: + rval[cpu_class] = [] + rval[cpu_class].append(service_id) + return rval + + def active_nodes_by_cpu_class(self): + rval = {} + ls_output = self._run_docker(docker_args='node ls') + for node in self._parse_docker_column_output(ls_output): + cpu_class = self._get_node_cpu_class(node_name=node['HOSTNAME']) + if cpu_class: + cpu_class = int(cpu_class) + rval[cpu_class] = rval.get(cpu_class, 0) + 1 + return rval + + def completed_services(self): + return self._services_in_state('Shutdown', 'Complete') + + def _services_in_state(self, desired, current): + service_ids = [] + for service_detail in self._service_details(): + if service_detail['DESIRED STATE'] == desired and service_detail['CURRENT STATE'].startswith(current): + service_ids.append(service_detail['ID']) + return service_ids + + def _service_details(self): + ls_output = self._run_docker(docker_args='service ls') + for service in self._parse_docker_column_output(ls_output): + if not service['NAME'].startswith(self.service_prefix): + continue + ps_output = self._run_docker(docker_args='service ps --no-trunc {service_id}'.format( + service_id=service['ID'])) + service_details = self._parse_docker_column_output(ps_output)[0] + for col in ('ID', 'NAME'): + service_details['PROCESS ' + col] = service_details[col] + service_details[col] = service[col] + yield service_details + + def clean_services(self): + cleaned_services = [] + services = self.completed_services() + if services: + cleaned_services = self._run_docker(docker_args='service rm {service_ids}'.format( + service_ids=' '.join(services))).splitlines() + return cleaned_services + + def node_states(self): + nodes = {} + ls_output = self._run_docker(docker_args='node ls') + for node in self._parse_docker_column_output(ls_output): + nodes[node['HOSTNAME']] = { + 'state': ('%s-%s' % (node['STATUS'], node['AVAILABILITY'])).lower(), + 'manager': True if node['MANAGER STATUS'] else False, + } + return nodes + + def node_job_count(self, node_name): + ps_output = self._run_docker(docker_args='node ps --no-trunc {node_name}'.format( + node_name=node_name)) + jobs = filter(lambda x: x['NAME'].startswith(self.service_prefix), self._parse_docker_column_output(ps_output)) + return len(jobs) + + def ensure_node_cpu_class(self, node, cpu_class): + cur_cpu_class = self._get_node_cpu_class(node_name=node) + if str(cpu_class) != cur_cpu_class: + log.info("setting node '%s' cpu class from '%s' to '%s'", node, cur_cpu_class, cpu_class) + self._run_docker(docker_args='node update --label-add {label_name}={label_val} {node}'.format( + label_name=NODE_CPU_CLASS_LABEL, + label_val=cpu_class, + node=node)) + else: + log.debug("node '%s' cpu class is '%s'", node, cur_cpu_class) + + def node_cpu_class(self, node): + return self._get_node_cpu_class(node_name=node) + + +class SwarmManager(object): + + def __init__(self, conf): + self.conf = conf + self.docker_interface = DockerInterface(conf) + self.state = SwarmState(conf) + self.spawn_wait_time = conf['spawn_wait_time'] + self.spawn_command = conf['spawn_command'] + self.destroy_command = conf['destroy_command'] + self.command_retries = conf['command_retries'] + self.command_retry_wait = conf['command_retry_wait'] + self.node_prefix = conf['node_prefix'] + self.terminate_when_idle = conf['terminate_when_idle'] + + def run(self): + while True: + node_states = None + self._spawn_if_waiting() + self._check_for_new_nodes(node_states=node_states) + self._destroy_if_surplus(node_states=node_states) + self._clean_services() + self._terminate_if_idle() + time.sleep(1) + + def _run_command(self, command, command_retries=None, **kwargs): + stdout = None + attempt = 0 + if not command_retries: + command_retries = self.command_retries + raw_cmd = command.format(**kwargs) + log.debug('running command: %s', raw_cmd) + while not stdout and attempt < command_retries + 1: + attempt += 1 + p = subprocess.Popen(raw_cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, close_fds=True, shell=True) + stdout, stderr = p.communicate() + if p.returncode != 0: + msg = "error running '%s'" % raw_cmd + if attempt < command_retries + 1: + msg += ', waiting %s seconds' % self.command_retry_wait + time.sleep(self.command_retry_wait) + log.warning(msg + "\nstdout: %s\nstderr: %s\n", stdout, stderr) + else: + msg += ' (final attempt)' + log.error(msg + "\nstdout: %s\nstderr: %s\n", stdout, stderr) + self._run_command(self.conf['command_failure_command'].format(failed_command=raw_cmd), command_retries=0) + stdout = None + else: + stdout = stdout.strip() + return stdout + + def _spawn_if_waiting(self): + waiting = self.docker_interface.waiting_services_by_cpu_class() + active = self.docker_interface.active_nodes_by_cpu_class() + cpus_needed = self.state.need_nodes(waiting, active) + for cpu_class in cpus_needed.keys(): + log.info("requesting node(s) for services requesting %d CPUs total (%d CPUs each): %s", cpus_needed[cpu_class], cpu_class, ' '.join(waiting[cpu_class])) + command = '{spawn_command}'.format(spawn_command=self.spawn_command).format( + cpu_class=cpu_class, + cpus_needed=cpus_needed[cpu_class]) + new_nodes = self._run_command(command) + if not new_nodes: + log.warning('spawn_command returned no new nodes, cannot manage nodes') + else: + log.info("node allocator will spawn: %s", new_nodes) + self.state.nodes_requested(cpu_class, new_nodes.split(), waiting[cpu_class]) + self.state.mark_services_handled(waiting[cpu_class]) + + def _check_for_new_nodes(self, node_states=None): + for node_name, elapsed, cpu_class, node in self.state.spawning_nodes(): + if not node_states: + node_states = self.docker_interface.node_states() + if node_name not in node_states: + if elapsed > self.spawn_wait_time: + log.warning("spawning node '%s' not found in `docker node ls` and spawn_wait_time exceeded! %d seconds have elapsed", node_name, elapsed) + self._run_command(self.conf['command_failure_command'].format(failed_command='wait_for_spawning_node %s' % node_name), command_retries=0) + self.mark_spawning_node_timeout(node_name) + elif node_states[node_name]['state'] == OK_NODE_STATE: + self.docker_interface.ensure_node_cpu_class(node_name, cpu_class) + self.state.mark_spawning_node_ready(node_name) + log.info("spawning node '%s' is ready!", node_name) + elif node_states[node_name]['state'] != node['state']: + log.info("spawning node '%s' state changed from '%s' to '%s'", node_name, node['state'], node_states[node_name]['state']) + self.docker_interface.ensure_node_cpu_class(node_name, cpu_class) + self.state.mark_spawning_node_state(node_name, node_states[node_name]['state']) + elif elapsed > self.spawn_wait_time: + log.warning("spawning node '%s' state is '%s' after %s seconds", node_name, node_states[node_name]['state'], elapsed) + + def _destroy_if_surplus(self, node_states=None): + destroy_nodes = [] + if not node_states: + node_states = self.docker_interface.node_states() + for node_name, node_state in node_states.items(): + if self._node_ready_for_destruction(node_name, node_state): + destroy_nodes.append(node_name) + if destroy_nodes: + command = '{destroy_command}'.format(destroy_command=self.destroy_command).format( + nodes=' '.join(destroy_nodes)) + destroyed_nodes = self._run_command(command) + if not destroyed_nodes: + log.warning('destroy_command returned no destroyed nodes') + else: + log.info("destroyed nodes: %s", destroyed_nodes) + + def _node_is_managed(self, node_name, node_state): + return (not self.node_prefix or node_name.startswith(self.node_prefix)) and not node_state['manager'] + + def _node_ready_for_destruction(self, node_name, node_state): + ready = False + if (self._node_is_managed(node_name, node_state) and + node_state['state'] == 'ready-active'): + if self.docker_interface.node_job_count(node_name) == 0: + self.state.mark_node_idle(node_name) + ready = self.state.is_destruction_time(node_name) + else: + self.state.clear_node_idle(node_name) + return ready + + def _clean_services(self): + cleaned_services = self.docker_interface.clean_services() + if cleaned_services: + self.state.clean_services(cleaned_services) + log.info("cleaned services: %s", ', '.join(cleaned_services)) + + def _terminate_if_idle(self): + if not self.terminate_when_idle: + return + node_states = self.docker_interface.node_states() + for node_name, node_state in node_states.items(): + if self._node_is_managed(node_name, node_state): + return # nonterminated managed nodes remain + elif self.docker_interface.node_job_count(node_name) > 0: + return # unmanaged nodes are running a galaxy service + # FIXME: there's a race condition here + if self.docker_interface.waiting_services_by_cpu_class(): + return # waiting jobs remain + log.info('nothing to manage, shutting down') + sys.exit(0) + + +class SwarmState(object): + + def __init__(self, conf): + self._handled_services = set() + self._waiting_since = {} + self._spawning_nodes = {} + self._surplus_nodes = {} + self.max_waiting_services = conf['max_waiting_services'] + self.max_wait_time = conf['max_wait_time'] + self.max_node_idle_time = conf['max_node_idle_time'] + self.max_node_counts = conf['max_node_counts'] + + def need_nodes(self, waiting_services, active_nodes): + rval = {} + need_cpus = self._needed_cpu_counts(waiting_services) + spawning_cpus = self._spawning_cpu_counts() + for cpu_class in need_cpus.keys(): + cpus_needed = need_cpus[cpu_class] - spawning_cpus.get(cpu_class, 0) + if cpus_needed > 0 and active_nodes.get(cpu_class, 0) < self.max_node_counts.get(cpu_class, sys.maxint): + rval[cpu_class] = cpus_needed + return rval + + def _needed_cpu_counts(self, waiting_services): + """Given a count of services waiting of each cpu class, return the + count of nodes needed of each node type if the maximum wait times and + waiting service count thresholds have been reached. + """ + rval = {} + new_waiting_since = {} + for cpu_class in waiting_services.keys(): + new_waiting_since[cpu_class] = self._waiting_since.get(cpu_class, time.time()) + # filter out any services that have already been handled + unhandled_waiting_services = [ s for s in waiting_services[cpu_class] if s not in self._handled_services ] + if (len(unhandled_waiting_services) > self.max_waiting_services and + time.time() - new_waiting_since[cpu_class] > self.max_wait_time): + # need waiting[cpu_class] nodes of this class + rval[cpu_class] = len(unhandled_waiting_services) + # drop any cpu_classes from waiting_since that are no longer waiting + self._waiting_since = new_waiting_since + return rval + + def spawning_nodes(self): + now = time.time() + for cpu_class in self._spawning_nodes.keys(): + for node_name in self._spawning_nodes[cpu_class].keys(): + node = self._spawning_nodes[cpu_class][node_name] + yield (node_name, now - node['time_requested'], cpu_class, node) + + def _spawning_cpu_counts(self): + rval = {} + for cpu_class in self._spawning_nodes.keys(): + rval[cpu_class] = sum([ v['cpu_count'] for k, v in self._spawning_nodes[cpu_class].items() ]) + return rval + + def nodes_requested(self, cpu_class, nodes, services): + if cpu_class not in self._spawning_nodes: + self._spawning_nodes[cpu_class] = {} + for node in nodes: + node_name = node.split(':')[0] + try: + cpu_count = node.split(':')[1] + except IndexError: + cpu_count = cpu_class + self._spawning_nodes[cpu_class][node_name] = { + 'state': 'requested', + 'cpu_count': cpu_count, + 'time_requested': time.time(), + } + + def mark_services_handled(self, services): + self._handled_services.update(services) + + def mark_spawning_node_ready(self, node_name): + self._delete_spawning_node(node_name) + + def mark_spawning_node_timeout(self, node_name): + self._delete_spawning_node(node_name) + + def _delete_spawning_node(self, node_name): + for cpu_class in self._spawning_nodes.keys(): + if node_name in self._spawning_nodes[cpu_class]: + del self._spawning_nodes[cpu_class][node_name] + + def mark_spawning_node_state(self, node_name, state): + for cpu_class in self._spawning_nodes.keys(): + if node_name in self._spawning_nodes[cpu_class]: + self._spawning_nodes[cpu_class][node_name]['state'] = state + + def is_destruction_time(self, node_name): + now = time.time() + return now - self._surplus_nodes.get(node_name, now) > self.max_node_idle_time + + def mark_node_idle(self, node_name): + if node_name not in self._surplus_nodes: + self._surplus_nodes[node_name] = time.time() + + def clear_node_idle(self, node_name): + if node_name in self._surplus_nodes: + del self._surplus_nodes[node_name] + + def clean_services(self, services): + self._handled_services.difference_update(services) + + +def main(argv=None, fork=False): + if not daemon: + log.warning('The daemon module is required to use the swarm manager, install it with `pip install python-daemon`') + return + if argv is None: + argv = sys.argv[1:] + if fork: + p = subprocess.Popen([sys.executable, __file__] + argv) + p.wait() + else: + args = _arg_parser().parse_args(argv) + kwargs = _app_properties(args) + _run_swarm_manager(kwargs, args) + + +def _app_properties(args): + galaxy_config_file = find_config_file("config/galaxy.ini", "universe_wsgi.ini", args.galaxy_config_file) + app_properties = load_app_properties(ini_file=galaxy_config_file) + return app_properties + + +def _arg_parser(): + parser = argparse.ArgumentParser(description=DESCRIPTION) + parser.add_argument("-c", "--galaxy-config-file", default=None) + return parser + + +def _run_swarm_manager(kwargs, args): + configure_logging(kwargs) + global log + log = logging.getLogger(__name__) + + root = find_root(kwargs) + swarm_manager_config_file = find_path(kwargs, "swarm_manager_config_file", root) + swarm_manager_conf = _parse_swarm_manager_conf(swarm_manager_config_file) + try: + os.makedirs(os.path.dirname(swarm_manager_conf['pid_file'])) + except (IOError, OSError) as exc: + if exc.errno != errno.EEXIST: + raise + log.debug("daemonizing, logs will be written to '%s'", swarm_manager_conf['log_file']) + pidfile = daemon.pidfile.PIDLockFile(swarm_manager_conf['pid_file']) + with open(swarm_manager_conf['log_file'], 'a') as logfh: + try: + with daemon.DaemonContext( + pidfile=pidfile, + stdout=logfh, + stderr=logfh, + ): + _swarm_manager(swarm_manager_conf) + except lockfile.AlreadyLocked: + log.debug("attempt to daemonize with swarm manager already running ignored") + + +def _load_xdg_environment(): + return dict( + data_home=os.path.expanduser(os.environ.get('XDG_DATA_HOME', '~/.local/share')), + ) + + +def _parse_swarm_manager_conf(swarm_manager_config_file): + conf = SWARM_MANAGER_CONF_DEFAULTS.copy() + xdg_env = _load_xdg_environment() + try: + with open(swarm_manager_config_file) as fh: + conf.update(yaml.load(fh)) + except (OSError, IOError) as exc: + if exc.errno == errno.ENOENT: + log.warning("config file '%s' does not exist, running with default config", swarm_manager_config_file) + else: + raise + for opt in ('pid_file', 'log_file'): + conf[opt] = conf[opt].format(xdg_data_home=xdg_env['data_home']) + return conf + + +def _swarm_manager(conf): + swarm_manager = SwarmManager(conf) + log.debug("swarm manager loaded, running...") + swarm_manager.run() + + +if __name__ == '__main__': + __name__ = 'swarm_manager' + main() diff --git a/lib/galaxy/web/base/interactive_environments.py b/lib/galaxy/web/base/interactive_environments.py index c3a7d782d43..e43b775c072 100644 --- a/lib/galaxy/web/base/interactive_environments.py +++ b/lib/galaxy/web/base/interactive_environments.py @@ -11,6 +11,7 @@ from subprocess import Popen, PIPE from galaxy.util import string_as_bool_or_none from galaxy.util.bunch import Bunch +from galaxy.container import docker_swarm from galaxy import web, model from galaxy.managers import api_keys from galaxy.tools.deps.docker_util import DockerVolume @@ -370,6 +371,9 @@ class InteractiveEnvironmentRequest(object): log.debug( "Container host: %s", self.attr.docker_hostname ) host_port = None + if self.attr.swarm_mode: + docker_swarm.main(argv=['-c', self.trans.app.config.config_file], fork=True) + if len(port_mappings) > 1: if self.attr.docker_connect_port is not None: for _service, _host_ip, _host_port in port_mappings: @@ -393,7 +397,9 @@ class InteractiveEnvironmentRequest(object): port=host_port, proxy_prefix=self.attr.proxy_prefix, route_name=self.attr.viz_id, - container_ids=[container_id], + container_ids=[container_id] if not self.attr.swarm_mode else [], + service_ids=[container_id] if self.attr.swarm_mode else [], + docker_command=self.attr.viz_config.get("docker", "command"), ) # These variables then become available for use in templating URLs self.attr.proxy_url = self.attr.proxy_request[ 'proxy_url' ] diff --git a/lib/galaxy/web/proxy/__init__.py b/lib/galaxy/web/proxy/__init__.py index 8db25a730b6..3efb82f8156 100644 --- a/lib/galaxy/web/proxy/__init__.py +++ b/lib/galaxy/web/proxy/__init__.py @@ -1,6 +1,7 @@ import logging import os import json +from collections import namedtuple from galaxy.util.filelock import FileLock from galaxy.util import sockets @@ -46,7 +47,7 @@ class ProxyManager(object): def shutdown( self ): self.lazy_process.shutdown() - def setup_proxy( self, trans, host=DEFAULT_PROXY_TO_HOST, port=None, proxy_prefix="", route_name="", container_ids=None ): + def setup_proxy( self, trans, host=DEFAULT_PROXY_TO_HOST, port=None, proxy_prefix="", route_name="", container_ids=None, service_ids=None, docker_command=None ): if self.manage_dynamic_proxy: log.info("Attempting to start dynamic proxy process") log.debug("Cmd: " + ' '.join(self.lazy_process.command_and_args)) @@ -54,6 +55,8 @@ class ProxyManager(object): if container_ids is None: container_ids = [] + if service_ids is None: + service_ids = [] authentication = AuthenticationToken(trans) proxy_requests = ProxyRequests(host=host, port=port) @@ -61,7 +64,9 @@ class ProxyManager(object): authentication, proxy_requests, '/%s' % route_name, - container_ids + container_ids, + service_ids, + docker_command, ) # TODO: These shouldn't need to be request.host and request.scheme - # though they are reasonable defaults. @@ -79,6 +84,10 @@ class ProxyManager(object): 'proxied_host': proxy_requests.host, } + def query_proxy( self, trans ): + authentication = AuthenticationToken(trans) + return self.proxy_ipc.fetch_requests(authentication) + def __setup_lazy_process( self, config ): launcher = self.proxy_launcher() command = launcher.launch_proxy_command(config) @@ -174,7 +183,10 @@ def proxy_ipc(config): class ProxyIpc(object): - def handle_requests(self, authentication, proxy_requests, route_name, container_ids): + def handle_requests(self, authentication, proxy_requests, route_name, container_ids, service_ids, docker_command): + raise NotImplementedError() + + def fetch_requests(self, authentication, key): raise NotImplementedError() @@ -183,50 +195,103 @@ class JsonFileProxyIpc(object): def __init__(self, proxy_session_map): self.proxy_session_map = proxy_session_map - def handle_requests(self, authentication, proxy_requests, route_name, container_ids): - key = "%s:%s" % ( proxy_requests.host, proxy_requests.port ) - secure_id = authentication.cookie_value + def handle_requests(self, authentication, proxy_requests, route_name, container_ids, service_ids, docker_command): + key = authentication.cookie_value with FileLock( self.proxy_session_map ): if not os.path.exists( self.proxy_session_map ): open( self.proxy_session_map, "w" ).write( "{}" ) json_data = open( self.proxy_session_map, "r" ).read() session_map = json.loads( json_data ) - to_remove = [] - for k, value in session_map.items(): - if value == secure_id: - to_remove.append( k ) - for k in to_remove: - del session_map[ k ] - session_map[ key ] = secure_id + session_map[ key ] = { + 'host': proxy_requests.host, + 'port': proxy_requests.port, + 'container_ids': container_ids, + 'service_ids': service_ids, + 'docker_command': docker_command, + } new_json_data = json.dumps( session_map ) open( self.proxy_session_map, "w" ).write( new_json_data ) + def fetch_requests(self, authentication): + key = authentication.cookie_value + try: + with open(self.proxy_session_map) as fh: + session_map = json.load(fh) + m = session_map[key] + return ProxyMapping( + host=m['host'], + port=m['port'], + container_ids=m['container_ids'], + service_ids=m['service_ids'], + docker_command=m['docker_command'], + ) + except (TypeError, KeyError): + log.warning('fetch_requests(): invalid key: %s', key) + return None + class SqliteProxyIpc(object): def __init__(self, proxy_session_map): self.proxy_session_map = proxy_session_map - def handle_requests(self, authentication, proxy_requests, route_name, container_ids): - key = "%s:%s" % ( proxy_requests.host, proxy_requests.port ) - secure_id = authentication.cookie_value + def handle_requests(self, authentication, proxy_requests, route_name, container_ids, service_ids, docker_command): + key = authentication.cookie_value with FileLock( self.proxy_session_map ): conn = sqlite.connect(self.proxy_session_map) try: c = conn.cursor() try: # Create table - c.execute('''CREATE TABLE gxproxy - (key text PRIMARY_KEY, secret text)''') + c.execute('''CREATE TABLE gxproxy2 + (key text PRIMARY KEY, + host text, + port integer, + container_ids text, + service_ids text, + docker_command text)''') except Exception: pass - insert_tmpl = '''INSERT INTO gxproxy (key, secret) VALUES ('%s', '%s');''' - insert = insert_tmpl % (key, secure_id) - c.execute(insert) + delete = '''DELETE FROM gxproxy2 WHERE key=?''' + c.execute(delete, (key,)) + insert = '''INSERT INTO gxproxy2 + (key, host, port, container_ids, service_ids, docker_command) + VALUES (?, ?, ?, ?, ?, ?)''' + c.execute(insert, + (key, + proxy_requests.host, + proxy_requests.port, + json.dumps(container_ids), + json.dumps(service_ids), + docker_command)) conn.commit() finally: conn.close() + def fetch_requests(self, authentication): + key = authentication.cookie_value + with FileLock( self.proxy_session_map): + conn = sqlite.connect(self.proxy_session_map) + try: + c = conn.cursor() + select = '''SELECT host, port, container_ids, service_ids, docker_command + FROM gxproxy2 + WHERE key=?''' + c.execute(select, (key,)) + try: + host, port, container_ids, service_ids, docker_command = c.fetchone() + except TypeError: + log.warning('fetch_requests(): invalid key: %s', key) + return None + return ProxyMapping( + host=host, + port=port, + container_ids=json.loads(container_ids), + service_ids=json.loads(service_ids), + docker_command=docker_command) + finally: + conn.close() + class RestGolangProxyIpc(object): @@ -234,7 +299,7 @@ class RestGolangProxyIpc(object): self.config = config self.api_url = 'http://127.0.0.1:%s/api?api_key=%s' % (self.config.dynamic_proxy_bind_port, self.config.dynamic_proxy_golang_api_key) - def handle_requests(self, authentication, proxy_requests, route_name, container_ids, sleep=1): + def handle_requests(self, authentication, proxy_requests, route_name, container_ids, service_ids, docker_command, sleep=1): """Make a POST request to the GO proxy to register a route """ values = { @@ -260,4 +325,8 @@ class RestGolangProxyIpc(object): self.handle_requests(authentication, proxy_requests, route_name, container_ids, sleep=sleep + 1) pass + +ProxyMapping = namedtuple('ProxyMapping', ['host', 'port', 'container_ids', 'service_ids', 'docker_command']) + + # TODO: MQ diven proxy? diff --git a/lib/galaxy/web/proxy/js/lib/mapper.js b/lib/galaxy/web/proxy/js/lib/mapper.js index 20f0f14e262..54fb097078e 100644 --- a/lib/galaxy/web/proxy/js/lib/mapper.js +++ b/lib/galaxy/web/proxy/js/lib/mapper.js @@ -15,9 +15,7 @@ var updateFromJson = function(path, map) { var keyToSession = JSON.parse(content); var newSessions = {}; for(var key in keyToSession) { - var hostAndPort = key.split(":"); - // 'host': hostAndPort[0], - newSessions[keyToSession[key]] = {'target': {'host': hostAndPort[0], 'port': parseInt(hostAndPort[1])}}; + newSessions[key] = {'target': {'host': keyToSession[key]['host'], 'port': parseInt(keyToSession[key]['port'])}}; } for(var oldSession in map) { if(!(oldSession in newSessions)) { @@ -32,12 +30,10 @@ var updateFromJson = function(path, map) { var updateFromSqlite = function(path, map) { var newSessions = {}; var loadSessions = function() { - db.each("SELECT key, secret FROM gxproxy", function(err, row) { + db.each("SELECT key, host, port FROM gxproxy2", function(err, row) { var key = row['key']; - var secret = row['secret']; - var hostAndPort = key.split(":"); - var target = {'host': hostAndPort[0], 'port': parseInt(hostAndPort[1])}; - newSessions[secret] = {'target': target}; + var target = {'host': row['host'], 'port': parseInt(row['port'])}; + newSessions[key] = {'target': target}; }, finish); }; @@ -75,4 +71,4 @@ var mapFor = function(path) { return map; }; -exports.mapFor = mapFor; \ No newline at end of file +exports.mapFor = mapFor; diff --git a/lib/galaxy/webapps/galaxy/controllers/interactive_environments.py b/lib/galaxy/webapps/galaxy/controllers/interactive_environments.py new file mode 100644 index 00000000000..36d986fcdfe --- /dev/null +++ b/lib/galaxy/webapps/galaxy/controllers/interactive_environments.py @@ -0,0 +1,46 @@ +""" +API check for whether the current session's interactive environment launch is ready +""" +from subprocess import Popen, PIPE + +from galaxy.web import expose, json +from galaxy.web.base.controller import BaseUIController + +import logging +log = logging.getLogger(__name__) + + +class InteractiveEnvironmentsController(BaseUIController): + + @expose + @json + def ready(self, trans, **kwd): + """ + GET /interactive_environments/ready/ + + Queries the GIE proxy IPC to determine whether the current user's session's GIE launch is ready + + :returns: ``true`` if ready else ``false`` + :rtype: boolean + """ + proxy_map = self.app.proxy_manager.query_proxy(trans) + if proxy_map.container_ids: + command = proxy_map.docker_command.format( + docker_args='ps --format {{{{.Status}}}} --filter id={container_id}'.format( + container_id=proxy_map.container_ids[0])) + match_col = 0 + match_str = 'Up' + elif proxy_map.service_ids: + command = proxy_map.docker_command.format( + docker_args='service ps --no-trunc {service_id}'.format(service_id=proxy_map.service_ids[0])) + match_col = 5 + match_str = 'Running' + else: + raise Exception('Proxy map has neither container ids nor service ids stored') + p = Popen(command, stdout=PIPE, stderr=PIPE, close_fds=True, shell=True) + stdout, stderr = p.communicate() + if p.returncode != 0: + log.error( "%s\n%s" % (stdout, stderr) ) + return None + else: + return stdout.splitlines()[-1].strip().split()[match_col].startswith(match_str) diff --git a/static/maps/galaxy.interactive_environments.js.map b/static/maps/galaxy.interactive_environments.js.map index ee5a8e61608..75370a01be8 100644 --- a/static/maps/galaxy.interactive_environments.js.map +++ b/static/maps/galaxy.interactive_environments.js.map @@ -1 +1 @@ -{"version":3,"file":"galaxy.interactive_environments.js","sources":["../src/galaxy.interactive_environments.js"],"names":["append_notebook","url","clear_main_area","$","append","remove","children","display_spinner","galaxy_root","test_ie_availability","success_callback","request_count","interval","setInterval","ajax","xhrFields","withCredentials","type","timeout","success","console","log","clearInterval","error","toastr","closeButton","timeOut","tapToDismiss"],"mappings":"AAIA,QAASA,iBAAgBC,GACrBC,kBACAC,EAAE,SAASC,OAAO,uHAAwHH,EAAK,eAInJ,QAASC,mBACLC,EAAE,YAAYE,SACdF,EAAE,SAASG,WAAWD,SAG1B,QAASE,mBACDJ,EAAE,SAASC,OAAO,0BAA4BI,YAAc,wGAWpE,QAASC,sBAAqBR,EAAKS,GAC/B,GAAIC,GAAgB,CACpBJ,mBACAK,SAAWC,YAAY,WACnBV,EAAEW,MACEb,IAAKA,EACLc,WACIC,iBAAiB,GAErBC,KAAM,MACNC,QAAS,IACTC,QAAS,WACLC,QAAQC,IAAI,8BACZC,cAAcV,UACdF,KAEJa,MAAO,WACHZ,IACAS,QAAQC,IAAI,WAAaV,GACtBA,EAAgB,KACfW,cAAcV,UACdV,kBACAsB,OAAOD,MACH,sDACA,SACCE,aAAe,EAAMC,QAAW,IAAOC,cAAgB,SAKzE"} \ No newline at end of file +{"version":3,"file":"galaxy.interactive_environments.js","sources":["../src/galaxy.interactive_environments.js"],"names":["append_notebook","url","clear_main_area","$","append","remove","children","display_spinner","galaxy_root","load_when_ready","success_callback","request_count","timeout_time","timeout_time_max","timeout_time_step","timeout","ajax","xhrFields","withCredentials","type","dataType","success","data","console","log","toastr","clear","info","closeButton","tapToDismiss","window","setTimeout","error","test_ie_availability","interval","setInterval","clearInterval","timeOut"],"mappings":"AAIA,QAASA,iBAAgBC,GACrBC,kBACAC,EAAE,SAASC,OAAO,uHAAwHH,EAAK,eAInJ,QAASC,mBACLC,EAAE,YAAYE,SACdF,EAAE,SAASG,WAAWD,SAG1B,QAASE,mBACDJ,EAAE,SAASC,OAAO,0BAA4BI,YAAc,wGAMpE,QAASC,iBAAgBR,EAAKS,GAC1B,GAAIC,GAAgB,EAChBC,EAAe,IACfC,EAAmB,KACnBC,EAAoB,IACpBC,EAAU,WACVZ,EAAEa,MACEf,IAAKA,EACLgB,WACIC,iBAAiB,GAErBC,KAAM,MACNJ,QAAS,IACTK,SAAU,OACVC,QAAS,SAASC,GACH,GAARA,GACCC,QAAQC,IAAI,gDACZtB,kBACAuB,OAAOC,QACPhB,KACa,GAARY,GACe,GAAjBX,IACCJ,kBACAkB,OAAOE,KACH,gGACCC,aAAe,EAAMC,cAAgB,KAG9ClB,IACkBE,EAAfD,IACCA,GAAgBE,GAEpBS,QAAQC,IAAI,qBAAuBb,EAAgB,aAAeC,EAAe,IAAO,KACxFkB,OAAOC,WAAWhB,EAASH,KAE3BV,kBACAuB,OAAOC,QACPD,OAAOO,MACH,gHACA,SACCJ,aAAe,EAAMC,cAAgB,QAM1DC,QAAOC,WAAWhB,EAASH,GAY/B,QAASqB,sBAAqBhC,EAAKS,GAC/B,GAAIC,GAAgB,CACpBJ,mBACA2B,SAAWC,YAAY,WACnBhC,EAAEa,MACEf,IAAKA,EACLgB,WACIC,iBAAiB,GAErBC,KAAM,MACNJ,QAAS,IACTM,QAAS,WACLE,QAAQC,IAAI,8BACZY,cAAcF,UACdxB,KAEJsB,MAAO,WACHrB,IACAY,QAAQC,IAAI,wBAA0Bb,GACnCA,EAAgB,KACfyB,cAAcF,UACdhC,kBACAuB,OAAOO,MACH,sDACA,SACCJ,aAAe,EAAMS,QAAW,IAAOR,cAAgB,SAKzE"} \ No newline at end of file diff --git a/static/scripts/galaxy.interactive_environments.js b/static/scripts/galaxy.interactive_environments.js index c859c44e03d..4769bb67db0 100644 --- a/static/scripts/galaxy.interactive_environments.js +++ b/static/scripts/galaxy.interactive_environments.js @@ -1,2 +1,2 @@ -function append_notebook(a){clear_main_area(),$("#main").append('')}function clear_main_area(){$("#spinner").remove(),$("#main").children().remove()}function display_spinner(){$("#main").append('')}function test_ie_availability(a,b){var c=0;display_spinner(),interval=setInterval(function(){$.ajax({url:a,xhrFields:{withCredentials:!0},type:"GET",timeout:500,success:function(){console.log("Connected to IE, returning"),clearInterval(interval),b()},error:function(){c++,console.log("Request "+c),c>30&&(clearInterval(interval),clear_main_area(),toastr.error("Could not connect to IE, contact your administrator","Error",{closeButton:!0,timeOut:2e4,tapToDismiss:!1}))}})},1e3)} +function append_notebook(a){clear_main_area(),$("#main").append('')}function clear_main_area(){$("#spinner").remove(),$("#main").children().remove()}function display_spinner(){$("#main").append('')}function load_when_ready(a,b){var c=0,d=1e3,e=15e3,f=1e3,g=function(){$.ajax({url:a,xhrFields:{withCredentials:!0},type:"GET",timeout:500,dataType:"json",success:function(a){1==a?(console.log("Galaxy reports IE container ready, returning"),clear_main_area(),toastr.clear(),b()):0==a?(0==c&&(display_spinner(),toastr.info("Galaxy is launching a container in which to run this interactive environment. Please wait...",{closeButton:!0,tapToDismiss:!1})),c++,e>d&&(d+=f),console.log("Readiness request "+c+" sleeping "+d/1e3+"s"),window.setTimeout(g,d)):(clear_main_area(),toastr.clear(),toastr.error("Galaxy failed to launch a container in which to run this interactive environment, contact your administrator.","Error",{closeButton:!0,tapToDismiss:!1}))}})};window.setTimeout(g,d)}function test_ie_availability(a,b){var c=0;display_spinner(),interval=setInterval(function(){$.ajax({url:a,xhrFields:{withCredentials:!0},type:"GET",timeout:500,success:function(){console.log("Connected to IE, returning"),clearInterval(interval),b()},error:function(){c++,console.log("Availability request "+c),c>30&&(clearInterval(interval),clear_main_area(),toastr.error("Could not connect to IE, contact your administrator","Error",{closeButton:!0,timeOut:2e4,tapToDismiss:!1}))}})},1e3)} //# sourceMappingURL=../maps/galaxy.interactive_environments.js.map \ No newline at end of file