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") }';
%def>
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