Dynamic Galaxy session aware proxy for IPython, other such plugins.

This is very much an expermintal component and plugin developers shouldn't expect the API to be stable at this point. This allows plugin developers to ensure HTTP access to a single port is routes to various ports internally based on the Galaxy session of the user accessing the resource.

It is frustrating that I need to use a node-based component for this but after nights of trying I was unable to reproduce a proxy that works both dynamically and with web sockets in any Python based web framework. I tried twisted, uwsgi, tornado. My decision to use node was buttressed by the IPython team utilizing the same framework for the same reasons for their work on Juypter - https://github.com/jupyter/jupyterhub.

This is still very much a work in progress but it was languishing in my github repository and I am want to get it committed so community can hack on it and suggest (or supply) improvements.

TODO
- Replace file operations with sqlite database + write file lock (for multi-process Galaxies).
- Document configuration.
- Update galaxy-ipython to support something like proxy-prefix - configure along with the rest of this.
This commit is contained in:
John Chilton
2014-10-18 18:38:04 -04:00
parent 07c2be2bae
commit 0f0b4d226f
14 changed files with 549 additions and 0 deletions
+2
View File
@@ -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()
+5
View File
@@ -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) )
+57
View File
@@ -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
+10
View File
@@ -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
+137
View File
@@ -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
+82
View File
@@ -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()
+18
View File
@@ -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" ]
+2
View File
@@ -0,0 +1,2 @@
# A dynamic configurable reverse proxy for use within Galaxy
+41
View File
@@ -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 <n>', 'Public-facing IP of the proxy', 'localhost')
.option('--port <n>', 'Public-facing port of the proxy', parseInt)
.option('--cookie <cookiename>', 'Cookie proving authentication', 'galaxysession')
.option('--sessions <file>', '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);
+33
View File
@@ -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;
+114
View File
@@ -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;
+17
View File
@@ -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"
}
}
+19
View File
@@ -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()
+12
View File
@@ -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()