diff --git a/lib/galaxy/app.py b/lib/galaxy/app.py index fcac7526590..9cda4879f70 100644 --- a/lib/galaxy/app.py +++ b/lib/galaxy/app.py @@ -17,6 +17,7 @@ from galaxy.openid.providers import OpenIDProviders from galaxy.tools.data_manager.manager import DataManagers from galaxy.jobs import metrics as job_metrics from galaxy.web.base import pluginframework +from galaxy.web.proxy import ProxyManager from galaxy.queue_worker import GalaxyQueueWorker from tool_shed.galaxy_install import update_repository_manager @@ -138,6 +139,7 @@ class UniverseApplication( object, config.ConfiguresGalaxyMixin ): # FIXME: These are exposed directly for backward compatibility self.job_queue = self.job_manager.job_queue self.job_stop_queue = self.job_manager.job_stop_queue + self.proxy_manager = ProxyManager( self.config ) # Initialize the external service types self.external_service_types = external_service_types.ExternalServiceTypesCollection( self.config.external_service_type_config_file, self.config.external_service_type_path, self ) self.model.engine.dispose() diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index bfc2ec0e43a..649c9a76165 100644 --- a/lib/galaxy/config.py +++ b/lib/galaxy/config.py @@ -357,6 +357,11 @@ class Configuration( object ): # directory where the visualization/registry searches for plugins self.visualization_plugins_directory = kwargs.get( 'visualization_plugins_directory', 'config/plugins/visualizations' ) + self.proxy_session_map = self.resolve_path( kwargs.get( "dynamic_proxy_session_map", "database/session_map.json" ) ) + self.manage_dynamic_proxy = string_as_bool( kwargs.get( "dynamic_proxy_manage", "True" ) ) # Set to false if being launched externally + self.dynamic_proxy_bind_port = int( kwargs.get( "dynamic_proxy_bind_port", "8800" ) ) + self.dynamic_proxy_bind_ip = kwargs.get( "dynamic_proxy_bind_ip", "0.0.0.0" ) + # Default chunk size for chunkable datatypes -- 64k self.display_chunk_size = int( kwargs.get( 'display_chunk_size', 65536) ) diff --git a/lib/galaxy/util/lazy_process.py b/lib/galaxy/util/lazy_process.py new file mode 100644 index 00000000000..e2e90ff8f5e --- /dev/null +++ b/lib/galaxy/util/lazy_process.py @@ -0,0 +1,57 @@ +import threading +import time +import subprocess + + +class LazyProcess( object ): + """ Abstraction describing a command line launching a service - probably + as needed as functionality is accessed in Galaxy. + """ + + def __init__( self, command_and_args ): + self.command_and_args = command_and_args + self.thread_lock = threading.Lock() + self.allow_process_request = True + self.process = None + + def start_process( self ): + with self.thread_lock: + if self.allow_process_request: + self.allow_process_request = False + t = threading.Thread(target=self.__start) + t.daemon = True + t.start() + + def __start(self): + with self.thread_lock: + self.process = subprocess.Popen( self.command_and_args, close_fds=True ) + + def shutdown( self ): + with self.thread_lock: + self.allow_process_request = False + if self.running: + self.process.terminate() + time.sleep(.01) + if self.running: + self.process.kill() + + @property + def running( self ): + return self.process and not self.process.poll() + + +class NoOpLazyProcess( object ): + """ LazyProcess abstraction meant to describe potentially optional + services, in those cases where one is not configured or valid, this + class can be used in place of LazyProcess. + """ + + def start_process( self ): + return + + def shutdown( self ): + return + + @property + def running( self ): + return False diff --git a/lib/galaxy/util/sockets.py b/lib/galaxy/util/sockets.py new file mode 100644 index 00000000000..d01732fb764 --- /dev/null +++ b/lib/galaxy/util/sockets.py @@ -0,0 +1,10 @@ +import socket + + +def unused_port(): + # TODO: Allow ranges (though then need to guess and check)... + s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + s.bind(('localhost', 0)) + addr, port = s.getsockname() + s.close() + return port diff --git a/lib/galaxy/web/proxy/__init__.py b/lib/galaxy/web/proxy/__init__.py new file mode 100644 index 00000000000..08a371013a5 --- /dev/null +++ b/lib/galaxy/web/proxy/__init__.py @@ -0,0 +1,137 @@ +import logging +import os +import json + +from .filelock import FileLock +from galaxy.util import sockets +from galaxy.util.lazy_process import LazyProcess, NoOpLazyProcess + +log = logging.getLogger( __name__ ) + + +DEFAULT_PROXY_TO_HOST = "localhost" +SECURE_COOKIE = "galaxysession" + + +class ProxyManager(object): + + def __init__( self, config ): + for option in [ "manage_dynamic_proxy", "dynamic_proxy_bind_port", "dynamic_proxy_bind_ip" ]: + setattr( self, option, getattr( config, option ) ) + self.launch_by = "node" # TODO: Support docker + if self.manage_dynamic_proxy: + self.lazy_process = self.__setup_lazy_process( config ) + else: + self.lazy_process = NoOpLazyProcess() + self.proxy_ipc = proxy_ipc(config) + + def shutdown( self ): + self.lazy_process.shutdown() + + def setup_proxy( self, trans, host=DEFAULT_PROXY_TO_HOST, port=None ): + if self.manage_dynamic_proxy: + log.info("Attempting to start dynamic proxy process") + self.lazy_process.start_process() + + authentication = AuthenticationToken(trans) + proxy_requests = ProxyRequests(host=host, port=port) + self.proxy_ipc.handle_requests(authentication, proxy_requests) + # TODO: These shouldn't need to be request.host and request.scheme - + # though they are reasonable defaults. + host = trans.request.host + if ':' in host: + host = host[0:host.index(':')] + scheme = trans.request.scheme + proxy_url = '%s://%s:%d' % (scheme, host, self.dynamic_proxy_bind_port) + return { + 'proxy_url': proxy_url, + 'proxied_port': proxy_requests.port, + 'proxied_host': proxy_requests.host, + } + + def __setup_lazy_process( self, config ): + launcher = proxy_launcher(self) + command = launcher.launch_proxy_command(config) + return LazyProcess(command) + + +def proxy_launcher(config): + return NodeProxyLauncher() + + +class ProxyLauncher(object): + + def launch_proxy_command(self, config): + raise NotImplementedError() + + +class NodeProxyLauncher(object): + + def launch_proxy_command(self, config): + args = [ + "--sessions", config.proxy_session_map, + "--ip", config.dynamic_proxy_bind_ip, + "--port", str(config.dynamic_proxy_bind_port), + ] + parent_directory = os.path.dirname( __file__ ) + path_to_application = os.path.join( parent_directory, "js", "lib", "main.js" ) + command = [ path_to_application ] + args + return command + + +class AuthenticationToken(object): + + def __init__(self, trans): + self.cookie_name = SECURE_COOKIE + self.cookie_value = trans.get_cookie( self.cookie_name ) + + +class ProxyRequests(object): + + def __init__(self, host=None, port=None): + if host is None: + host = DEFAULT_PROXY_TO_HOST + if port is None: + port = sockets.unused_port() + log.info("Obtained unused port %d" % port) + self.host = host + self.port = port + + +def proxy_ipc(config): + proxy_session_map = config.proxy_session_map + return JsonFileProxyIpc(proxy_session_map) + + +class ProxyIpc(object): + + def handle_requests(self, cookie, host, port): + raise NotImplementedError() + + +class JsonFileProxyIpc(object): + + def __init__(self, proxy_session_map): + self.proxy_session_map = proxy_session_map + + def handle_requests(self, authentication, proxy_requests): + key = "%s:%s" % ( proxy_requests.host, proxy_requests.port ) + secure_id = 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 + new_json_data = json.dumps( session_map ) + open( self.proxy_session_map, "w" ).write( new_json_data ) + + +# TODO: sqlitefileipc +# TODO: RESTful API driven proxy diff --git a/lib/galaxy/web/proxy/filelock.py b/lib/galaxy/web/proxy/filelock.py new file mode 100644 index 00000000000..6b51e3f5347 --- /dev/null +++ b/lib/galaxy/web/proxy/filelock.py @@ -0,0 +1,82 @@ +""" Code obtained from https://github.com/dmfrey/FileLock + +See full license at: + +https://github.com/dmfrey/FileLock/blob/master/LICENSE.txt + +""" +import os +import time +import errno + + +class FileLockException(Exception): + pass + + +class FileLock(object): + """ A file locking mechanism that has context-manager support so + you can use it in a with statement. This should be relatively cross + compatible as it doesn't rely on msvcrt or fcntl for the locking. + """ + + def __init__(self, file_name, timeout=10, delay=.05): + """ Prepare the file locker. Specify the file to lock and optionally + the maximum timeout and the delay between each attempt to lock. + """ + self.is_locked = False + full_path = os.path.abspath(file_name) + self.lockfile = "%s.lock" % full_path + self.file_name = full_path + self.timeout = timeout + self.delay = delay + + def acquire(self): + """ Acquire the lock, if possible. If the lock is in use, it check again + every `wait` seconds. It does this until it either gets the lock or + exceeds `timeout` number of seconds, in which case it throws + an exception. + """ + start_time = time.time() + while True: + try: + self.fd = os.open(self.lockfile, os.O_CREAT | os.O_EXCL | os.O_RDWR) + break + except OSError as e: + if e.errno != errno.EEXIST: + raise + if (time.time() - start_time) >= self.timeout: + raise FileLockException("Timeout occured.") + time.sleep(self.delay) + self.is_locked = True + + def release(self): + """ Get rid of the lock by deleting the lockfile. + When working in a `with` statement, this gets automatically + called at the end. + """ + if self.is_locked: + os.close(self.fd) + os.unlink(self.lockfile) + self.is_locked = False + + def __enter__(self): + """ Activated when used in the with statement. + Should automatically acquire a lock to be used in the with block. + """ + if not self.is_locked: + self.acquire() + return self + + def __exit__(self, type, value, traceback): + """ Activated at the end of the with statement. + It automatically releases the lock if it isn't locked. + """ + if self.is_locked: + self.release() + + def __del__(self): + """ Make sure that the FileLock instance doesn't leave a lockfile + lying around. + """ + self.release() diff --git a/lib/galaxy/web/proxy/js/Dockerfile b/lib/galaxy/web/proxy/js/Dockerfile new file mode 100644 index 00000000000..05e057d1419 --- /dev/null +++ b/lib/galaxy/web/proxy/js/Dockerfile @@ -0,0 +1,18 @@ +# Have not yet gotten this to work - goal was to launch the prox in a Docker container. +# Networking is a bit tricky though - could not get the child proxy to talk to the child +# IPython container. + + +# sudo docker build --no-cache=true -t gxproxy . +# sudo docker run --net host -v /home/john/workspace/galaxy-central/database:/var/gxproxy -p 8800:8800 -t gxproxy lib/main.js --sessions /var/gxproxy/session_map.json --ip 0.0.0.0 --port 8800 + +FROM node:0.11.13 + +RUN mkdir -p /usr/src/gxproxy +WORKDIR /usr/src/gxproxy + +ADD package.json /usr/src/gxproxy/ +RUN npm install +ADD . /usr/src/gxproxy + +CMD [ "lib/main.js" ] diff --git a/lib/galaxy/web/proxy/js/README.md b/lib/galaxy/web/proxy/js/README.md new file mode 100644 index 00000000000..f397ef876ad --- /dev/null +++ b/lib/galaxy/web/proxy/js/README.md @@ -0,0 +1,2 @@ +# A dynamic configurable reverse proxy for use within Galaxy + diff --git a/lib/galaxy/web/proxy/js/lib/main.js b/lib/galaxy/web/proxy/js/lib/main.js new file mode 100755 index 00000000000..3aa3a43542c --- /dev/null +++ b/lib/galaxy/web/proxy/js/lib/main.js @@ -0,0 +1,41 @@ +#!/usr/bin/env node +/* +Inspiration taken from + https://github.com/jupyter/multiuser-server/blob/master/multiuser/js/main.js +*/ +var fs = require('fs'); +var args = require('commander'); + +package_info = require('../package') + +args + .version(package_info) + .option('--ip ', 'Public-facing IP of the proxy', 'localhost') + .option('--port ', 'Public-facing port of the proxy', parseInt) + .option('--cookie ', 'Cookie proving authentication', 'galaxysession') + .option('--sessions ', 'Routes file to monitor') + .option('--verbose') + +args.parse(process.argv); + +var DynamicProxy = require('./proxy.js').DynamicProxy; +var mapFor = require('./mapper.js').mapFor; + +var sessions = mapFor(args.sessions); + +var dynamic_proxy_options = { + sessionCookie: args['cookie'], + sessionMap: sessions, + verbose: args.verbose +} + +var dynamic_proxy = new DynamicProxy(dynamic_proxy_options); + +var listen = {}; +listen.port = args.port || 8000; +listen.ip = args.ip; + +if(args.verbose) { + console.log("Listening on " + listen.ip + ":" + listen.port); +} +dynamic_proxy.proxy_server.listen(listen.port, listen.ip); diff --git a/lib/galaxy/web/proxy/js/lib/mapper.js b/lib/galaxy/web/proxy/js/lib/mapper.js new file mode 100644 index 00000000000..1e2b8120ce0 --- /dev/null +++ b/lib/galaxy/web/proxy/js/lib/mapper.js @@ -0,0 +1,33 @@ +var fs = require('fs'); + + +var mapFor = function(path) { + var map = {}; + + var loadMap = function() { + var content = fs.readFileSync(path, 'utf8'); + 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])}}; + } + for(var oldSession in map) { + if(!(oldSession in newSessions)) { + delete map[ oldSession ]; + } + } + for(var newSession in newSessions) { + map[newSession] = newSessions[newSession]; + } + } + + console.log("Watching path " + path); + loadMap(); + fs.watch(path, loadMap); + + return map; +} + +exports.mapFor = mapFor; \ No newline at end of file diff --git a/lib/galaxy/web/proxy/js/lib/proxy.js b/lib/galaxy/web/proxy/js/lib/proxy.js new file mode 100644 index 00000000000..1c77b212a3a --- /dev/null +++ b/lib/galaxy/web/proxy/js/lib/proxy.js @@ -0,0 +1,114 @@ +var http = require('http'), + httpProxy = require('http-proxy'); + +var bound = function (that, method) { + // bind a method, to ensure `this=that` when it is called + // because prototype languages are bad + return function () { + method.apply(that, arguments); + }; +}; + +var DynamicProxy = function(options) { + var dynamicProxy = this; + this.sessionCookie = options.sessionCookie; + this.sessionMap = options.sessionMap; + + var log_errors = function(handler) { + return function (req, res) { + try { + return handler.apply(dynamicProxy, arguments); + } catch (e) { + console.log("Error in handler for " + req.method + ' ' + req.url + ': ', e); + } + }; + }; + + var proxy = this.proxy = httpProxy.createProxyServer({ + ws : true, + }); + + this.proxy_server = http.createServer( + log_errors(dynamicProxy.handleProxyRequest) + ); + this.proxy_server.on('upgrade', bound(this, this.handleWs)); +}; + +DynamicProxy.prototype.rewriteRequest = function(request) { +} + +DynamicProxy.prototype.targetForRequest = function(request) { + // return proxy target for a given url + var session = this.findSession(request); + for (var mappedSession in this.sessionMap) { + if(session == mappedSession) { + return this.sessionMap[session].target; + } + } + + return null; +}; + +DynamicProxy.prototype.findSession = function(request) { + var sessionCookie = this.sessionCookie; + rc = request.headers.cookie; + if(!rc) { + return null; + } + var cookies = rc.split(';'); + for(var cookieIndex in cookies) { + var cookie = cookies[cookieIndex]; + var parts = cookie.split('='); + var partName = parts.shift().trim(); + if(partName == sessionCookie) { + return unescape(parts.join('=')) + } + } + + return null; +}; + +DynamicProxy.prototype.handleProxyRequest = function(req, res) { + var target = this.targetForRequest(req); + console.log("PROXY " + req.method + " " + req.url + " to " + target); + var origin = req.headers.origin; + this.rewriteRequest(req); + res.oldWriteHead = res.writeHead; + res.writeHead = function(statusCode, headers) { + res.setHeader('Access-Control-Allow-Origin', origin); + res.setHeader('Access-Control-Allow-Credentials', 'true'); + res.oldWriteHead(statusCode, headers); + } + this.proxy.web(req, res, { + target: target + }, function (e) { + console.log("Proxy error: ", e); + res.writeHead(502); + res.write("Proxy target missing"); + res.end(); + }); +}; + +DynamicProxy.prototype.handleWs = function(req, res, head) { + // no local route found, time to proxy + var target = this.targetForRequest(req); + console.log("PROXY WS " + req.url + " to " + req.url); + var origin = req.headers.origin; + this.rewriteRequest(req); + res.oldWriteHead = res.writeHead; + res.writeHead = function(statusCode, headers) { + res.setHeader('Access-Control-Allow-Origin', origin); + res.setHeader('Access-Control-Allow-Credentials', 'true'); + res.oldWriteHead(statusCode, headers); + } + this.proxy.ws(req, res, head, { + target: target + }, function (e) { + console.log("Proxy error: ", e); + res.writeHead(502); + res.write("Proxy target missing"); + res.end(); + }); +}; + +exports.DynamicProxy = DynamicProxy; diff --git a/lib/galaxy/web/proxy/js/package.json b/lib/galaxy/web/proxy/js/package.json new file mode 100644 index 00000000000..40323beab01 --- /dev/null +++ b/lib/galaxy/web/proxy/js/package.json @@ -0,0 +1,17 @@ +{ + "name": "galaxy-proxy", + "version": "0.0.1", + "description": "A dynamic reverse proxy for use within Galaxy", + "main": "index.js", + "author": "John Chilton", + "license": "AFL v3", + "readmeFilename": "README.md", + "repository": { + "type": "mercurial", + "url": "https://bitbucket.org/galaxy/galaxy-central" + }, + "dependencies": { + "http-proxy": "~1.1", + "commander": "~2.2" + } +} diff --git a/test/unit/test_lazy_process.py b/test/unit/test_lazy_process.py new file mode 100644 index 00000000000..6554c84ed87 --- /dev/null +++ b/test/unit/test_lazy_process.py @@ -0,0 +1,19 @@ +import os +import tempfile +import time + +from galaxy.util.lazy_process import LazyProcess + + +def test_lazy_process(): + t = tempfile.NamedTemporaryFile() + os.remove(t.name) + lazy_process = LazyProcess(["bash", "-c", "touch %s; sleep 100" % t.name]) + assert not os.path.exists(t.name) + lazy_process.start_process() + time.sleep(.02) + assert lazy_process.process.poll() is None + assert os.path.exists(t.name) + lazy_process.shutdown() + time.sleep(.02) + assert lazy_process.process.poll() diff --git a/test/unit/test_sockets.py b/test/unit/test_sockets.py new file mode 100644 index 00000000000..0c9fc5bdb8c --- /dev/null +++ b/test/unit/test_sockets.py @@ -0,0 +1,12 @@ +import socket +from galaxy.util import sockets + + +def test_unused_free_port_unconstrained(): + port = sockets.unused_port() + + s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + # would throw exception if port was not free. + s.bind(('localhost', port)) + s.close() +