Merge pull request #7494 from blankenberg/realtimetools

Add InteractiveTools.
This commit is contained in:
John Chilton
2019-08-21 20:13:55 -04:00
committed by GitHub
62 changed files with 2287 additions and 17 deletions
@@ -0,0 +1,65 @@
<template>
<div class="infomessagelarge">
<p v-if="entryPoints.length==0" >
Waiting for InteractiveTool result view(s) to become available.
</p>
<p v-else-if="entryPoints.length==1" >
<span v-if="entryPoints[0].active" >
There is an InteractiveTool result view available, <a :href="entryPoints[0].target">click here to display</a>.
</span>
<span v-else>
There is an InteractiveTool result view available, waiting for view to become active...
</span>
</p>
<p v-else>
There are multiple InteractiveTool result views available:
<ul>
<li v-for="entryPoint of entryPoints" v-bind:key="entryPoint.id" >
{{ entryPoint.name }}
<span v-if="entryPoint.active">
(<a :href="entryPoints[0].target">click here to display</a>)
</span>
<span v-else>
(waiting to become active...)
</span>
</li>
</ul>
</p>
You may also access all active InteractiveTools from the User menu.
</div>
</template>
<script>
import { clearPolling, pollUntilActive } from "mvc/entrypoints/poll";
export default {
props: {
jobId: {
type: String,
required: true
}
},
data() {
return {
entryPoints: [],
};
},
created: function() {
this.pollEntryPoints();
},
beforeDestroy: function(){
clearPolling();
},
methods: {
pollEntryPoints: function() {
const onUpdate = (entryPoints) => {
this.entryPoints = entryPoints;
}
const onError = (e) => {
console.error(e);
}
pollUntilActive(onUpdate, onError, {"job_id": this.jobId})
}
}
};
</script>
@@ -22,6 +22,7 @@ import UserPreferences from "mvc/user/user-preferences";
import CustomBuilds from "mvc/user/user-custom-builds";
import Tours from "mvc/tours";
import GridView from "mvc/grid/grid-view";
import EntryPointGridView from "mvc/entrypoints/view";
import GridShared from "mvc/grid/grid-shared";
import Workflows from "mvc/workflow/workflow";
import WorkflowImport from "components/WorkflowImport.vue";
@@ -74,7 +75,8 @@ export const getAnalysisRouter = Galaxy =>
"(/)datasets(/)list(/)": "show_datasets",
"(/)custom_builds": "show_custom_builds",
"(/)datasets/edit": "show_dataset_edit_attributes",
"(/)datasets/error": "show_dataset_error"
"(/)datasets/error": "show_dataset_error",
"(/)realtime_entry_points(/)list": "show_realtime_list"
},
require_login: ["show_user", "show_user_form", "show_workflows", "show_cloud_auth"],
@@ -111,6 +113,15 @@ export const getAnalysisRouter = Galaxy =>
this.page.display(new FormWrapper.View(_.extend(model.get(form_id), { active_tab: "user" })));
},
show_realtime_list: function() {
this.page.display(
new EntryPointGridView({
url_base: `${getAppRoot()}realtime/list`,
active_tab: "analysis"
})
);
},
show_cloud_auth: function() {
this._display_vue_helper(CloudAuth);
},
+8
View File
@@ -309,6 +309,14 @@ const Collection = Backbone.Collection.extend({
target: "__use_router__"
});
}
if (Galaxy.config.interactivetools_enable) {
userTab.menu[userTab.menu.length - 1].divider = true;
userTab.menu.push({
title: _l("Active InteractiveTools"),
url: "realtime_entry_points/list",
target: "__use_router__"
});
}
}
this.add(userTab);
return new $.Deferred().resolve().promise();
@@ -0,0 +1,33 @@
import axios from "axios";
import { getAppRoot } from "onload/loadConfig";
let interval;
export const clearPolling = () => {
clearInterval(interval);
}
export const pollUntilActive = (onUpdate, onError, params) => {
clearPolling();
const url = getAppRoot() + `api/entry_points`;
console.log(params);
axios
.get(url, {params: params})
.then(response => {
const entryPoints = [];
let allReady = true;
response.data.forEach((entryPoint, i) => {
entryPoints.push(entryPoint);
if(! entryPoint.active) {
allReady = false;
}
});
onUpdate(entryPoints)
if(! allReady || entryPoints.length == 0) {
interval = setInterval(() => {
pollUntilActive(onUpdate, onError, params);
}, 3000);
}
})
.catch(onError);
}
@@ -0,0 +1,34 @@
import $ from "jquery";
import GridView from "mvc/grid/grid-view";
import { clearPolling, pollUntilActive } from "mvc/entrypoints/poll";
export default GridView.extend({
init_grid_elements: function() {
GridView.prototype.init_grid_elements.call(this);
const activated = {};
const onUpdate = (entryPoints) => {
entryPoints.forEach((entryPoint) => {
const entryPointId = entryPoint.id;
if (entryPoint.active && ! activated[entryPointId]) {
const $link = $(`.entry-point-link[entry_point_id='${entryPointId}']`)
if ($link.length > 0) {
$link.attr("href", entryPoint["target"]);
activated[entryPointId] = true;
}
}
});
}
const onError = (e) => {
console.error(e);
}
pollUntilActive(onUpdate, onError, {"running": true});
},
remove: function() {
// Your processing code here
clearPolling();
GridView.prototype.remove.apply(this, arguments);
}
});
@@ -11,6 +11,8 @@ import Ui from "mvc/ui/ui-misc";
import Modal from "mvc/ui/ui-modal";
import ToolFormBase from "mvc/tool/tool-form-base";
import Webhooks from "mvc/webhooks";
import Vue from "vue";
import ToolEntryPoints from "components/ToolEntryPoints/ToolEntryPoints";
const View = Backbone.View.extend({
initialize: function(options) {
@@ -230,6 +232,19 @@ const View = Backbone.View.extend({
success: response => {
callback && callback();
this.$el.children().hide();
if (response.produces_entry_points) {
for (const job of response.jobs) {
const toolEntryPointsInstance = Vue.extend(ToolEntryPoints);
const vm = document.createElement("div");
this.$el.append(vm);
const instance = new toolEntryPointsInstance({
propsData: {
jobId: job.id
}
});
instance.$mount(vm);
}
}
this.$el.append(this._templateSuccess(response, job_def));
this.$el.parent().scrollTop(0);
// Show Webhook if job is running
+36
View File
@@ -0,0 +1,36 @@
uwsgi:
http: localhost:8080
threads: 8
http-raw-body: True
offload-threads: 8
master: true
module: galaxy.webapps.galaxy.buildapp:uwsgi_app()
interactivetools_map: database/interactivetools_map.sqlite
python-raw: scripts/interactivetools/key_type_token_mapping.py
route-host: ^([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.(interactivetool\.localhost:8080)$ goto:interactivetool
route-run: goto:endendend
route-label: interactivetool
route-host: ^([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.(interactivetool\.localhost:8080)$ rpcvar:TARGET_HOST rtt_key_type_token_mapper_cached $2 $1 $3 $4 $0 5
route-if-not: empty:${TARGET_HOST} httpdumb:${TARGET_HOST}
route: .* break:404 Not Found
route-label: endendend
galaxy:
interactivetools_enable: true
# outputs_to_working_directory will provide you with a better level of isolation. It is highly recommended to set
# this parameter with InteractiveTools.
outputs_to_working_directory: true
interactivetools_prefix: interactivetool
interactivetools_map: database/interactivetools_map.sqlite
# If you develop InteractiveTools locally and do not have a full FQDN you can
# use an arbritrary one, e.g. 'my-hostname' here, if you set this hostname in your
# job_conf.xml as well (see the corresponding comment).
# galaxy_infrastructure_url: http://my-hostname:8080
+37
View File
@@ -0,0 +1,37 @@
<?xml version="1.0"?>
<!-- A sample job config for RealTimeTools using local runner. -->
<job_conf>
<plugins>
<plugin id="local" type="runner" load="galaxy.jobs.runners.local:LocalJobRunner" workers="4"/>
</plugins>
<destinations default="docker_dispatch">
<destination id="local" runner="local"/>
<destination id="docker_local" runner="local">
<param id="docker_enabled">true</param>
<!-- If you have not set 'outputs_to_working_directory: true' in galaxy.yml you can remove the docker_volumes setting. -->
<param id="docker_volumes">$galaxy_root:ro,$tool_directory:ro,$job_directory:rw,$working_directory:rw,$default_file_path:ro</param>
<param id="docker_sudo">false</param>
<param id="docker_net">bridge</param>
<param id="docker_auto_rm">true</param>
<param id="require_container">true</param>
<param id="docker_set_user"></param>
<!-- InteractiveTools do need real hostnames or URLs to work - simply specifying IPs will not work.
If you develop interactive tools on your 'localhost' and don't have a proper domain name
you need to tell all Docker containers a hostname where Galaxy is running.
This can be done via the --add-host parameter during the `docker run` command.
'my-hostname' here is an arbritrary hostname that matches the IP address of your
Galaxy host. Make sure this hostname ('my-hostname') is also set in your galaxy.yml file, e.g.
`galaxy_infrastructure_url: http://my-hostname:8080`. -->
<!--param id="docker_run_extra_arguments">--add-host my-hostname:10.1.2.3</param-->
</destination>
<destination id="docker_dispatch" runner="dynamic">
<param id="type">docker_dispatch</param>
<param id="docker_destination_id">docker_local</param>
<param id="default_destination_id">local</param>
</destination>
</destinations>
</job_conf>
+10
View File
@@ -1236,6 +1236,16 @@
:Type: bool
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
``interactivetools_enable``
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
:Description:
Set this to true to enable InteractiveTools.
:Default: ``false``
:Type: bool
~~~~~~~~~~~~~~~~~~~~~~~~~~
``visualizations_visible``
~~~~~~~~~~~~~~~~~~~~~~~~~~
@@ -0,0 +1,125 @@
Galaxy RealTimeTools
=====================================
A Galaxy RealTimeTool allows launching a container-backed Galaxy Tool
and enabling a Galaxy User to gain access to content inside in real-time.
How Galaxy RealTimeTools work
-----------------------------
A RealTimeTool is defined in the same familiar way as standard Galaxy Tools,
but are specified with **tool_type="realtime"**, and providing additional entry point
information:
.. code-block:: xml
<realtime>
<entry_point name="Display name">
<port>80</port>
<url><![CDATA[optional/path/can/be/${templated}]]></url>
</entry_point>
</realtime>
**Note** that name, port, and url are each able to be templated from the RealTimeTool's parameter dictionary.
Some important benefits of using Galaxy RealTimeTools
-----------------------------------------------------
- You can have and access **any number of RealTimeTools at a time** (admin configurable)
- If you accidentally close the **RealTimeTool browser window**, you can **regain access** by selecting from a **list of active RealTimeTools**
- A single **RealTimeTool** can **grant access** to **multiple running applications, servers, and interfaces**
- **RealTimeTools** can be **added to** Galaxy **Workflows**
- **Native, out-of-the box support for RealTimeTools** in Galaxy via uWSGI proxying
- **RealTimeTools** are **bonafide Galaxy Tools**; just specify **tool_type as "realtime"** and list the ports you want to expose
- **RealTimeTools** can be **added** to and **installed from the ToolShed**.
- **R Shiny apps**, **Javascript-based VNC** access to desktop environments, **genome-browsers-in-a-box**, **interactive notebook environments**, etc, are all possible with **RealTimeTools**
Server-side configuration of Galaxy RealTimeTools
-------------------------------------------------
The **galaxy.yml** file will need to be populated as seen in **config/galaxy.yml.realtime**.
In the **uwsgi:** section:
.. code-block:: yaml
http-raw-body: true
# master: true
realtime_map: database/realtime_map.sqlite
python-raw: scripts/realtime/key_type_token_mapping.py
route-host: ^([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.(realtime\.localhost:8080)$ goto:realtime
route-run: goto:endendend
route-label: realtime
route-host: ^([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.(realtime\.localhost:8080)$ rpcvar:TARGET_HOST rtt_key_type_token_mapper_cached $2 $1 $3 $4 $0 5
route-if-not: empty:${TARGET_HOST} httpdumb:${TARGET_HOST}
route-label: endendend
In the **galaxy:** section:
.. code-block:: yaml
realtime_prefix: realtime
The admin should modify the **route-host**s and **realtime_prefix** to match their preferred configuration.
An example **job_conf.xml** file as seen in **config/galaxy.yml.realtime**:
.. code-block:: xml
<job_conf>
<plugins>
<plugin id="local" type="runner" load="galaxy.jobs.runners.local:LocalJobRunner" workers="4"/>
</plugins>
<destinations default="docker_dispatch">
<destination id="local" runner="local"/>
<destination id="docker_local" runner="local">
<param id="docker_enabled">true</param>
<param id="docker_volumes">$galaxy_root:ro,$tool_directory:ro,$job_directory:rw,$working_directory:rw,$default_file_path:ro</param>
<param id="docker_sudo">false</param>
<param id="docker_net">bridge</param>
<param id="docker_auto_rm">true</param>
<param id="require_container">true</param>
</destination>
<destination id="docker_dispatch" runner="dynamic">
<param id="type">docker_dispatch</param>
<param id="docker_destination_id">docker_local</param>
<param id="default_destination_id">local</param>
</destination>
</destinations>
</job_conf>
Alternatively to the local job runner, RealTimeTools have been enabled for the condor job runner, e.g.:
.. code-block:: xml
<destination id="condor" runner="condor">
<param id="docker_enabled">true</param>
<param id="docker_sudo">false</param>
</destination>
**Note on resource consumption:** Keep in mind that Distributed Resource Management (DRM) / cluster systems may have a maximum runtime configured for jobs. From the Galaxy point of view, such a container could run as long as the user desires, this may not be advisable and an admin may want to restrict the runtime of RealTimeTools *(and jobs in general)*. However, if the job is killed by the DRM, the user is not informed beforehand and data in the container could be discarded.
Two **example test RealTimeTools** have been defined, and can be added to the **config/tool_conf.xml**:
.. code-block:: xml
<tool file="../test/functional/tools/realtimetool_juypter_notebook.xml" />
<tool file="../test/functional/tools/realtimetool_cellxgene.xml" />
A few words on the condor integration
-------------------------------------
Galaxy needs to be able to stop a container gracefully. This is not a problem with the local job runner, where we assume that Docker is either running on the same host. However, if you are using production scale DRM, like condor, then your job is running
somewhere on your cluster and you can not easily **docker stop** your container. For the condor integration we are using a great
condor feature and commandline utility called **condor_ssh_to_job**. This tool (assuming your condor setup is configured correctly) will bring us directly to the host in question and we can execute the **docker stop** command. Galaxy will simply run **condor_ssh_to_job <condor_job_id> docker stop <container_name>** to stop the container gracefully.
+3
View File
@@ -17,6 +17,7 @@ from galaxy.managers.collections import DatasetCollectionManager
from galaxy.managers.folders import FolderManager
from galaxy.managers.histories import HistoryManager
from galaxy.managers.libraries import LibraryManager
from galaxy.managers.realtime import RealTimeManager
from galaxy.managers.tools import DynamicToolManager
from galaxy.model.database_heartbeat import DatabaseHeartbeat
from galaxy.model.tags import GalaxyTagHandler
@@ -222,6 +223,8 @@ class UniverseApplication(config.ConfiguresGalaxyMixin):
containers_conf=self.config.containers_conf
)
self.realtime_manager = RealTimeManager(self)
# Configure handling of signals
handlers = {}
if self.heartbeat:
+6
View File
@@ -708,6 +708,7 @@ class GalaxyAppConfiguration(BaseAppConfiguration):
self.fluent_log = string_as_bool(kwargs.get('fluent_log', False))
self.fluent_host = kwargs.get('fluent_host', 'localhost')
self.fluent_port = int(kwargs.get('fluent_port', 24224))
# directory where the visualization registry searches for plugins
self.visualization_plugins_directory = kwargs.get(
'visualization_plugins_directory', 'config/plugins/visualizations')
@@ -734,6 +735,11 @@ class GalaxyAppConfiguration(BaseAppConfiguration):
self.dynamic_proxy_golang_docker_address = kwargs.get("dynamic_proxy_golang_docker_address", "unix:///var/run/docker.sock")
self.dynamic_proxy_golang_api_key = kwargs.get("dynamic_proxy_golang_api_key", None)
# InteractiveTools propagator mapping file
self.realtime_map = self.resolve_path(kwargs.get("interactivetools_map", "database/interactivetools_map.sqlite"))
self.realtime_prefix = kwargs.get("interactivetools_prefix", "interactivetool")
self.interactivetools_enable = string_as_bool(kwargs.get('interactivetools_enable', False))
# Default chunk size for chunkable datatypes -- 64k
self.display_chunk_size = int(kwargs.get('display_chunk_size', 65536))
@@ -423,6 +423,9 @@ galaxy:
# scenarios than the watchdog default.
#watch_tool_data_dir: 'false'
# Enable InteractiveTools by setting to true
#interactivetools_enable: false
# File containing old-style genome builds
#builds_file_path: tool-data/shared/ucsc/builds.txt
+56
View File
@@ -4,6 +4,7 @@ Support for running a tool in Galaxy via an internal job management system
import copy
import datetime
import errno
import json
import logging
import os
import pwd
@@ -1054,6 +1055,7 @@ class JobWrapper(HasResourceParameters):
# if the server was stopped and restarted before the job finished
job.command_line = unicodify(self.command_line)
job.dependencies = self.tool.dependencies
self.realtimetools = getattr(tool_evaluator, 'realtimetools', None)
self.sa_session.add(job)
self.sa_session.flush()
# Return list of all extra files
@@ -1963,6 +1965,54 @@ class JobWrapper(HasResourceParameters):
command = "%s; %s" % (dependency_shell_commands, command)
return command
def check_for_entry_points(self, check_already_configured=True):
if not self.tool.produces_entry_points:
return True
job = self.get_job()
if check_already_configured and job.all_entry_points_configured:
return True
working_directory = self.working_directory
container_runtime_path = os.path.join(working_directory, "container_runtime.json")
if os.path.exists(container_runtime_path):
with open(container_runtime_path, "r") as f:
try:
container_runtime = json.load(f)
except ValueError:
# File exists, but is not fully populated yet
return False
log.debug("found container runtime %s" % container_runtime)
self.app.realtime_manager.configure_entry_points(job, container_runtime)
return True
def container_monitor_command(self, container, **kwds):
if not container or not self.tool.produces_entry_points:
return None
from os.path import abspath
from os import getcwd
exec_dir = kwds.get('exec_dir', abspath(getcwd()))
work_dir = self.working_directory
container_config = os.path.join(work_dir, "container_config.json")
# TODO: implement callback via URL callback for configuring ports instead
# of fs polling.
# if self.app.config.galaxy_infrastructure_url_set:
# # TODO: settable from job destination...
# infrastructure_url = self.app.config.galaxy_infrastructure_url
# else:
# infrastructure_url = "host.docker.internal"
with open(container_config, "w") as f:
json.dump({
"container_name": container.container_name,
"container_type": container.container_type,
"connection_configuration": container.connection_configuration,
}, f)
return "(python '%s'/lib/galaxy_ext/container_monitor/monitor.py &) " % exec_dir
@property
def user(self):
job = self.get_job()
@@ -2077,6 +2127,12 @@ class JobWrapper(HasResourceParameters):
for dataset in job.output_datasets:
self.app.error_reports.default_error_plugin.submit_report(dataset, job, tool, user_submission=False)
def set_container(self, container):
if container:
cont = model.JobContainerAssociation(job=self.get_job(), container_type=container.container_type, container_name=container.container_name, container_info=container.container_info)
self.sa_session.add(cont)
self.sa_session.flush()
class TaskWrapper(JobWrapper):
"""
+4
View File
@@ -103,6 +103,10 @@ def build_command(
# xref https://github.com/galaxyproject/galaxy/issues/3289
commands_builder.prepend_command("rm -rf working; mkdir -p working; cd working")
container_monitor_command = job_wrapper.container_monitor_command(container)
if container_monitor_command:
commands_builder.prepend_command(container_monitor_command)
if include_work_dir_outputs:
__handle_work_dir_outputs(commands_builder, job_wrapper, runner, remote_command_params)
+6 -2
View File
@@ -419,7 +419,8 @@ class BaseJobRunner(object):
compute_tmp_directory = job_wrapper.tmp_directory()
tool = job_wrapper.tool
tool_info = ToolInfo(tool.containers, tool.requirements, tool.requires_galaxy_python_environment, tool.docker_env_pass_through)
guest_ports = [ep.get('port') for ep in getattr(job_wrapper, 'realtimetools', [])]
tool_info = ToolInfo(tool.containers, tool.requirements, tool.requires_galaxy_python_environment, tool.docker_env_pass_through, guest_ports=guest_ports)
job_info = JobInfo(
compute_working_directory,
compute_tool_directory,
@@ -429,11 +430,14 @@ class BaseJobRunner(object):
)
destination_info = job_wrapper.job_destination.params
return self.app.container_finder.find_container(
container = self.app.container_finder.find_container(
tool_info,
destination_info,
job_info
)
if container:
job_wrapper.set_container(container)
return container
def _handle_runner_state(self, runner_state, job_state):
try:
+74 -3
View File
@@ -12,6 +12,7 @@ it at this time.
"""
import logging
import os
import subprocess
from galaxy import model
from galaxy.jobs.runners import (
@@ -174,6 +175,7 @@ class CondorJobRunner(AsynchronousJobRunner):
for cjs in self.watched:
job_id = cjs.job_id
galaxy_id_tag = cjs.job_wrapper.get_id_tag()
log.debug("### (%s/%s) whats-up" % (galaxy_id_tag, job_id))
try:
if os.stat(cjs.user_log).st_size == cjs.user_log_size:
new_watched.append(cjs)
@@ -190,6 +192,11 @@ class CondorJobRunner(AsynchronousJobRunner):
cjs.fail_message = "Cluster could not complete job"
self.work_queue.put((self.fail_job, cjs))
continue
if job_running:
# If running, check for entry points...
cjs.job_wrapper.job_wrapper.check_for_entry_points()
if job_running and not cjs.running:
log.debug("(%s/%s) job is now running" % (galaxy_id_tag, job_id))
cjs.job_wrapper.change_state(model.Job.states.RUNNING)
@@ -219,9 +226,41 @@ class CondorJobRunner(AsynchronousJobRunner):
"""Attempts to delete a job from the DRM queue"""
job = job_wrapper.get_job()
external_id = job.job_runner_external_id
failure_message = condor_stop(external_id)
if failure_message:
log.debug("(%s). Failed to stop condor %s" % (external_id, failure_message))
galaxy_id_tag = job_wrapper.get_id_tag()
if job.container:
try:
log.info("stop_job(): %s: trying to stop container .... (%s)" % (job.id, external_id))
# self.watched = [cjs for cjs in self.watched if cjs.job_id != external_id]
new_watch_list = list()
cjs = None
for tcjs in self.watched:
if tcjs.job_id != external_id:
new_watch_list.append(tcjs)
else:
cjs = tcjs
break
self.watched = new_watch_list
self._stop_container(job_wrapper)
# self.watched.append(cjs)
if cjs.job_wrapper.get_state() != model.Job.states.DELETED:
external_metadata = not asbool(cjs.job_wrapper.job_destination.params.get("embed_metadata_in_job", True))
if external_metadata:
self._handle_metadata_externally(cjs.job_wrapper, resolve_requirements=True)
log.debug("(%s/%s) job has completed" % (galaxy_id_tag, external_id))
self.work_queue.put((self.finish_job, cjs))
except Exception as e:
log.warning("stop_job(): %s: trying to stop container failed. (%s)" % (job.id, e))
try:
self._kill_container(job_wrapper)
except Exception as e:
log.warning("stop_job(): %s: trying to kill container failed. (%s)" % (job.id, e))
failure_message = condor_stop(external_id)
if failure_message:
log.debug("(%s). Failed to stop condor %s" % (external_id, failure_message))
else:
failure_message = condor_stop(external_id)
if failure_message:
log.debug("(%s). Failed to stop condor %s" % (external_id, failure_message))
def recover(self, job, job_wrapper):
"""Recovers jobs stuck in the queued/running state when Galaxy started"""
@@ -246,3 +285,35 @@ class CondorJobRunner(AsynchronousJobRunner):
log.debug("(%s/%s) is still in DRM queued state, adding to the DRM queue" % (job.id, job.job_runner_external_id))
cjs.running = False
self.monitor_queue.put(cjs)
def _stop_container(self, job_wrapper):
return self._run_container_command(job_wrapper, 'stop')
def _kill_container(self, job_wrapper):
return self._run_container_command(job_wrapper, 'kill')
def _run_container_command(self, job_wrapper, command):
job = job_wrapper.get_job()
external_id = job.job_runner_external_id
if job:
cont = job.container
if cont:
if cont.container_type == 'docker':
return self._run_command(cont.container_info['commands'][command], external_id)[0]
def _run_command(self, command, external_job_id):
command = 'condor_ssh_to_job %s %s' % (external_job_id, command)
p = subprocess.Popen(command, stdout=subprocess.PIPE, stderr=subprocess.PIPE, shell=True, close_fds=True, preexec_fn=os.setpgrp)
stdout, stderr = p.communicate()
exit_code = p.returncode
ret = None
if exit_code == 0:
ret = stdout.strip()
else:
log.debug(stderr)
# exit_code = subprocess.call(command,
# shell=True,
# preexec_fn=os.setpgrp)
log.debug('_run_command(%s) exit code (%s) and failure: %s', command, exit_code, stderr)
return (exit_code, ret)
+12
View File
@@ -108,6 +108,8 @@ class LocalJobRunner(BaseJobRunner):
job_wrapper.set_job_destination(job_wrapper.job_destination, proc.pid)
job_wrapper.change_state(model.Job.states.RUNNING)
self._handle_container(job_wrapper, proc)
terminated = self.__poll_if_needed(proc, job_wrapper, job_id)
if terminated:
return
@@ -220,6 +222,16 @@ class LocalJobRunner(BaseJobRunner):
os.killpg(proc.pid, 9)
return proc.wait() # reap
def _handle_container(self, job_wrapper, proc):
if not job_wrapper.tool.produces_entry_points:
return
while proc.poll() is None:
if job_wrapper.check_for_entry_points(check_already_configured=False):
return
sleep(0.5)
def __poll_if_needed(self, proc, job_wrapper, job_id):
# Only poll if needed (i.e. job limits are set)
if not job_wrapper.has_limits():
+1
View File
@@ -79,6 +79,7 @@ class ConfigSerializer(base.ModelSerializer):
'require_login' : _defaults_to(None),
'inactivity_box_content' : _defaults_to(None),
'visualizations_visible' : _defaults_to(True),
'interactivetools_enable' : _defaults_to(False),
'message_box_content' : _defaults_to(None),
'message_box_visible' : _defaults_to(False),
'message_box_class' : _defaults_to('info'),
+295
View File
@@ -0,0 +1,295 @@
import logging
import sqlite3
from six import string_types
from sqlalchemy import or_
from galaxy import (
exceptions,
model
)
from galaxy.util.filelock import FileLock
log = logging.getLogger(__name__)
DATABASE_TABLE_NAME = 'gxrtproxy'
class RealtimeSqlite(object):
def __init__(self, sqlite_filename, encode_id):
self.sqlite_filename = sqlite_filename
self.encode_id = encode_id
def get(self, key, key_type):
with FileLock(self.sqlite_filename):
conn = sqlite3.connect(self.sqlite_filename)
try:
c = conn.cursor()
select = '''SELECT token, host, port, info
FROM %s
WHERE key=? and key_type=?''' % (DATABASE_TABLE_NAME)
c.execute(select, (key, key_type,))
try:
token, host, port, info = c.fetchone()
except TypeError:
log.warning('get(): invalid key: %s key_type %s', key, key_type)
return None
return dict(
key=key,
key_type=key_type,
token=token,
host=host,
port=port,
info=info)
finally:
conn.close()
def save(self, key, key_type, token, host, port, info=None):
"""
Writeout a key, key_type, token, value store that is can be used for coordinating
with external resources.
"""
assert key, ValueError("A non-zero length key is required.")
assert key_type, ValueError("A non-zero length key_type is required.")
assert token, ValueError("A non-zero length token is required.")
with FileLock(self.sqlite_filename):
conn = sqlite3.connect(self.sqlite_filename)
try:
c = conn.cursor()
try:
# Create table
c.execute('''CREATE TABLE %s
(key text,
key_type text,
token text,
host text,
port integer,
info text,
PRIMARY KEY (key, key_type)
)''' % (DATABASE_TABLE_NAME))
except Exception:
pass
delete = '''DELETE FROM %s WHERE key=? and key_type=?''' % (DATABASE_TABLE_NAME)
c.execute(delete, (key, key_type,))
insert = '''INSERT INTO %s
(key, key_type, token, host, port, info)
VALUES (?, ?, ?, ?, ?, ?)''' % (DATABASE_TABLE_NAME)
c.execute(insert,
(key,
key_type,
token,
host,
port,
info,
))
conn.commit()
finally:
conn.close()
def remove(self, **kwd):
"""
Remove entry from a key, key_type, token, value store that is can be used for coordinating
with external resources. Remove entries that match all provided key=values
"""
assert kwd, ValueError("You must provide some values to key upon")
delete = 'DELETE FROM %s WHERE' % (DATABASE_TABLE_NAME)
value_list = []
for i, (key, value) in enumerate(kwd.items()):
if i != 0:
delete += ' and'
delete += ' %s=?' % (key)
value_list.append(value)
with FileLock(self.sqlite_filename):
conn = sqlite3.connect(self.sqlite_filename)
try:
c = conn.cursor()
try:
# Delete entry
# NB: This does not invalidate in-memory caches used by uwsgi (if any)
c.execute(delete, tuple(value_list))
except Exception as e:
log.debug('Error removing entry (%s): %s', delete, e)
conn.commit()
finally:
conn.close()
def save_entry_point(self, entry_point):
"""Convenience method to easily save an entry_point.
"""
return self.save(self.encode_id(entry_point.id), entry_point.__class__.__name__.lower(), entry_point.token, entry_point.host, entry_point.port, None)
def remove_entry_point(self, entry_point):
"""Convenience method to easily remove an entry_point.
"""
return self.remove(key=self.encode_id(entry_point.id), key_type=entry_point.__class__.__name__.lower())
def remove_entry_points0(self, entry_points):
"""Convenience method to easily remove entry_points.
"""
rval = []
for entry_point in entry_points:
rval.append(self.remove_entry_point(entry_point))
return rval
def remove_realtime0(self, rtt):
"""Convenience method to easily remove a RealTimeTool.
"""
for ep in rtt.entry_points:
self.remove_entry_point(ep)
class RealTimeManager(object):
"""
Manager for dealing with RealTimeTools
"""
def __init__(self, app):
self.app = app
self.model = app.model
self.security = app.security
self.sa_session = app.model.context
self.job_manager = app.job_manager
self.propagator = RealtimeSqlite(app.config.realtime_map, app.security.encode_id)
def create_entry_points(self, job, tool, entry_points=None, flush=True):
entry_points = entry_points or tool.ports
for entry in entry_points:
ep = self.model.InteractiveToolEntryPoint(job=job, tool_port=entry['port'], entry_url=entry['url'], name=entry['name'])
self.sa_session.add(ep)
if flush:
self.sa_session.flush()
def configure_entry_point(self, job, tool_port=None, host=None, port=None, protocol=None):
return self.configure_entry_points(job, {tool_port: dict(tool_port=tool_port, host=host, port=port, protocol=protocol)})
def configure_entry_points(self, job, ports_dict):
# There can be multiple entry points that reference the same tool port (could have different entry URLs)
configured = []
not_configured = []
for ep in job.realtimetool_entry_points:
port_dict = ports_dict.get(str(ep.tool_port), None)
if port_dict is None:
log.error("Did not find port to assign to InteractiveToolEntryPoint by tool port: %s.", ep.tool_port)
not_configured.append(ep)
else:
ep.host = port_dict['host']
ep.port = port_dict['port']
ep.protocol = port_dict['protocol']
ep.configured = True
self.sa_session.add(ep)
self.save_entry_point(ep)
configured.append(ep)
if configured:
self.sa_session.flush()
return dict(not_configured=not_configured, configured=configured)
def save_entry_point(self, entry_point):
"""
Writeout a key, key_type, token, value store that is used validating access
"""
self.propagator.save_entry_point(entry_point)
def create_realtime(self, job, tool, entry_points):
# create from initial job
if job and tool:
self.create_entry_points(job, tool, entry_points)
else:
log.warning('Called RealTimeManager.create_realtime, but job (%s) or tool (%s) is None', job, tool)
def get_nonterminal_for_user_by_trans(self, trans):
if trans.user:
jobs = trans.sa_session.query(trans.app.model.Job).filter(trans.app.model.Job.user == trans.user)
else:
jobs = trans.sa_session.query(trans.app.model.Job).filter(trans.app.model.Job.session_id == trans.get_galaxy_session().id)
def build_and_apply_filters(query, objects, filter_func):
if objects is not None:
if isinstance(objects, string_types):
query = query.filter(filter_func(objects))
elif isinstance(objects, list):
t = []
for obj in objects:
t.append(filter_func(obj))
query = query.filter(or_(*t))
return query
jobs = build_and_apply_filters(jobs, trans.app.model.Job.non_ready_states, lambda s: trans.app.model.Job.state == s)
return trans.sa_session.query(trans.app.model.InteractiveToolEntryPoint).filter(trans.app.model.InteractiveToolEntryPoint.job_id.in_([job.id for job in jobs]))
def can_access_job(self, trans, job):
if job:
if trans.user is None:
galaxy_session = trans.get_galaxy_session()
if galaxy_session is None or job.session_id != galaxy_session.id:
return False
elif job.user != trans.user:
return False
else:
return False
return True
def can_access_entry_point(self, trans, entry_point):
if entry_point:
return self.can_access_job(trans, entry_point.job)
return False
def can_access_entry_points(self, trans, entry_points):
for ep in entry_points:
if not self.can_access_entry_point(trans, ep):
return False
return True
def stop(self, trans, entry_point):
try:
self.remove_entry_point(entry_point)
job = entry_point.job
if not job.finished:
log.debug('Stopping Job: %s for InteractiveToolEntryPoint: %s', job, entry_point)
job.mark_deleted(trans.app.config.track_jobs_in_database)
# This self.job_manager.stop(job) does nothing without changing job.state, manually or e.g. with .mark_deleted()
self.job_manager.stop(job)
trans.sa_session.add(job)
trans.sa_session.flush()
except Exception as e:
log.debug('Unable to stop job for InteractiveToolEntryPoint (%s): %s', entry_point, e)
return False
return True
def remove_realtime0(self, realtime):
return self.propagator.remove_realtime(realtime)
def remove_entry_points(self, entry_points):
if entry_points:
for entry_point in entry_points:
self.remove_entry_point(entry_point, flush=False)
self.sa_session.flush()
def remove_entry_point(self, entry_point, flush=True):
entry_point.deleted = True
self.sa_session.add(entry_point)
if flush:
self.sa_session.flush()
self.propagator.remove_entry_point(entry_point)
def target_if_active(self, trans, entry_point):
if entry_point.active and not entry_point.deleted:
request_host = trans.request.host
rval = '%s//%s.%s.%s.%s.%s/' % (trans.request.host_url.split('//', 1)[0], entry_point.__class__.__name__.lower(), trans.security.encode_id(entry_point.id),
entry_point.token, self.app.config.realtime_prefix, request_host)
if entry_point.entry_url:
rval = '%s/%s' % (rval.rstrip('/'), entry_point.entry_url.lstrip('/'))
return rval
def access_entry_point_target(self, trans, entry_point_id):
entry_point = trans.sa_session.query(model.InteractiveToolEntryPoint).get(entry_point_id)
if self.app.realtime_manager.can_access_entry_point(trans, entry_point):
if entry_point.active:
return self.target_if_active(entry_point)
elif entry_point.deleted:
raise exceptions.MessageException('RealTimeTool has ended. You will have to start a new one.')
else:
raise exceptions.MessageException('RealTimeTool is not active. If you recently launched this tool it may not be ready yet, please wait a moment and refresh this page.')
raise exceptions.ItemAccessibilityException("You do not have access to this RealTimeTool entry point.")
+44
View File
@@ -942,6 +942,14 @@ class Job(JobLike, UsesCreateAndUpdateTime, Dictifiable, RepresentById):
def add_post_job_action(self, pja):
self.post_job_actions.append(PostJobActionAssociation(pja, self))
@property
def all_entry_points_configured(self):
# consider an actual DB attribute for this.
all_configured = True
for ep in self.realtimetool_entry_points:
all_configured = ep.configured and all_configured
return all_configured
def set_state(self, state):
"""
Save state history
@@ -1474,6 +1482,42 @@ class JobImportHistoryArchive(RepresentById):
self.archive_dir = archive_dir
class JobContainerAssociation(RepresentById):
def __init__(self, job=None, container_type=None, container_name=None, container_info=None):
self.job = job
self.container_type = container_type
self.container_name = container_name
self.container_info = container_info or {}
class InteractiveToolEntryPoint(Dictifiable, RepresentById):
dict_collection_visible_keys = ['id', 'name', 'active']
dict_element_visible_keys = ['id', 'name', 'active']
def __init__(self, job=None, name=None, token=None, tool_port=None, host=None, port=None, protocol=None,
entry_url=None, info=None, configured=False, deleted=False):
self.job = job
self.name = name
if not token:
token = uuid4().hex
self.token = token
self.tool_port = tool_port
self.host = host
self.port = port
self.protocol = protocol
self.entry_url = entry_url
self.info = info or {}
self.configured = configured
self.deleted = deleted
@property
def active(self):
if self.configured and not self.deleted:
# FIXME: don't included queued?
return not self.job.finished
return False
class GenomeIndexToolData(RepresentById):
def __init__(self, job=None, params=None, dataset=None, deferred_job=None,
transfer_job=None, fasta_path=None, created_time=None, modified_time=None,
+35
View File
@@ -792,6 +792,33 @@ model.GenomeIndexToolData.table = Table(
Column("indexer", String(64)),
Column("user_id", Integer, ForeignKey("galaxy_user.id"), index=True))
model.InteractiveToolEntryPoint.table = Table(
"interactivetool_entry_point", metadata,
Column("id", Integer, primary_key=True),
Column("job_id", Integer, ForeignKey("job.id"), index=True),
Column("name", TEXT),
Column("token", TEXT),
Column("tool_port", Integer),
Column("host", TEXT),
Column("port", Integer),
Column("protocol", TEXT),
Column("entry_url", TEXT),
Column("info", JSONType, nullable=True),
Column("configured", Boolean, default=False),
Column("deleted", Boolean, default=False),
Column("created_time", DateTime, default=now),
Column("modified_time", DateTime, default=now, onupdate=now))
model.JobContainerAssociation.table = Table(
"job_container_association", metadata,
Column("id", Integer, primary_key=True),
Column("job_id", Integer, ForeignKey("job.id"), index=True),
Column("container_type", TEXT),
Column("container_name", TEXT),
Column("container_info", JSONType, nullable=True),
Column("created_time", DateTime, default=now),
Column("modified_time", DateTime, default=now, onupdate=now))
model.Task.table = Table(
"task", metadata,
Column("id", Integer, primary_key=True),
@@ -2236,6 +2263,14 @@ mapper(model.GenomeIndexToolData, model.GenomeIndexToolData.table, properties=di
transfer=relation(model.TransferJob, backref='transfer_job')
))
mapper(model.InteractiveToolEntryPoint, model.InteractiveToolEntryPoint.table, properties=dict(
job=relation(model.Job, backref=backref('realtimetool_entry_points', uselist=True), uselist=False)
))
mapper(model.JobContainerAssociation, model.JobContainerAssociation.table, properties=dict(
job=relation(model.Job, backref=backref('container', uselist=False), uselist=False)
))
mapper(model.PostJobAction, model.PostJobAction.table, properties=dict(
workflow_step=relation(model.WorkflowStep,
backref='post_job_actions',
@@ -0,0 +1,72 @@
"""
Migration script to add new tables for InteractiveTools.
"""
from __future__ import print_function
import logging
from sqlalchemy import Boolean, Column, DateTime, ForeignKey, Integer, MetaData, Table, TEXT
from galaxy.model.custom_types import JSONType
from galaxy.model.orm.now import now
log = logging.getLogger(__name__)
metadata = MetaData()
interactivetool_entry_point = Table(
"interactivetool_entry_point", metadata,
Column("id", Integer, primary_key=True),
Column("job_id", Integer, ForeignKey("job.id"), index=True),
Column("name", TEXT),
Column("token", TEXT),
Column("tool_port", Integer),
Column("host", TEXT),
Column("port", Integer),
Column("protocol", TEXT),
Column("entry_url", TEXT),
Column("info", JSONType, nullable=True),
Column("configured", Boolean, default=False),
Column("deleted", Boolean, default=False),
Column("created_time", DateTime, default=now),
Column("modified_time", DateTime, default=now, onupdate=now))
job_container_association = Table(
"job_container_association", metadata,
Column("id", Integer, primary_key=True),
Column("job_id", Integer, ForeignKey("job.id"), index=True),
Column("container_type", TEXT),
Column("container_name", TEXT),
Column("container_info", JSONType, nullable=True),
Column("created_time", DateTime, default=now),
Column("modified_time", DateTime, default=now, onupdate=now))
def upgrade(migrate_engine):
print(__doc__)
metadata.bind = migrate_engine
metadata.reflect()
try:
job_container_association.create()
except Exception:
log.exception("Failed to create job_container_association table")
try:
interactivetool_entry_point.create()
except Exception:
log.exception("Failed to create interactivetool_entry_point table")
def downgrade(migrate_engine):
metadata.bind = migrate_engine
metadata.reflect()
try:
job_container_association.drop()
except Exception:
log.exception("Failed to drop job_container_association table")
try:
interactivetool_entry_point.drop()
except Exception:
log.exception("Failed to drop interactivetool_entry_point table")
+27 -3
View File
@@ -4,6 +4,8 @@ from abc import (
ABCMeta,
abstractmethod
)
from logging import getLogger
from uuid import uuid4
import six
@@ -21,6 +23,8 @@ from .requirements import (
DEFAULT_CONTAINER_SHELL,
)
log = getLogger(__name__)
DOCKER_CONTAINER_TYPE = "docker"
SINGULARITY_CONTAINER_TYPE = "singularity"
@@ -57,13 +61,15 @@ EOF
@six.add_metaclass(ABCMeta)
class Container(object):
def __init__(self, container_id, app_info, tool_info, destination_info, job_info, container_description):
def __init__(self, container_id, app_info, tool_info, destination_info, job_info, container_description, container_name=None):
self.container_id = container_id
self.app_info = app_info
self.tool_info = tool_info
self.destination_info = destination_info
self.job_info = job_info
self.container_description = container_description
self.container_name = container_name or uuid4().hex
self.container_info = {}
def prop(self, name, default):
destination_name = "docker_%s" % name
@@ -217,6 +223,10 @@ class DockerContainer(Container, HasDockerLikeVolumes):
)
return docker_host_props
@property
def connection_configuration(self):
return self.docker_host_props
def build_pull_command(self):
return docker_util.build_pull_command(self.container_id, **self.docker_host_props)
@@ -263,13 +273,21 @@ class DockerContainer(Container, HasDockerLikeVolumes):
volumes_from=volumes_from,
env_directives=env_directives,
working_directory=working_directory,
net=self.prop("net", "none"), # By default, docker instance has networking disabled
net=self.prop("net", None), # By default, docker instance has networking disabled
auto_rm=asbool(self.prop("auto_rm", docker_util.DEFAULT_AUTO_REMOVE)),
set_user=self.prop("set_user", docker_util.DEFAULT_SET_USER),
run_extra_arguments=self.prop("run_extra_arguments", docker_util.DEFAULT_RUN_EXTRA_ARGUMENTS),
guest_ports=self.tool_info.guest_ports,
container_name=self.container_name,
**docker_host_props
)
return "%s\n%s" % (cache_command, run_command)
kill_command = docker_util.build_docker_simple_command("kill", container_name=self.container_name, **docker_host_props)
return """
_on_exit() {
%s
}
trap _on_exit 0
%s\n%s""" % (kill_command, cache_command, run_command)
def __cache_from_file_command(self, cached_image_file, docker_host_props):
images_cmd = docker_util.build_docker_images_command(truncate=False, **docker_host_props)
@@ -312,6 +330,10 @@ class SingularityContainer(Container, HasDockerLikeVolumes):
sudo_cmd=self.prop("sudo_cmd", singularity_util.DEFAULT_SUDO_COMMAND),
)
@property
def connection_configuration(self):
return self.get_singularity_target_kwds()
def build_mulled_singularity_pull_command(self, cache_directory, namespace="biocontainers"):
return singularity_util.pull_mulled_singularity_command(
docker_image_identifier=self.container_id,
@@ -349,6 +371,8 @@ class SingularityContainer(Container, HasDockerLikeVolumes):
env=env,
working_directory=working_directory,
run_extra_arguments=self.prop("run_extra_arguments", singularity_util.DEFAULT_RUN_EXTRA_ARGUMENTS),
guest_ports=self.tool_info.guest_ports,
container_name=self.container_name,
**self.get_singularity_target_kwds()
)
return run_command
+2 -1
View File
@@ -36,7 +36,7 @@ class ToolInfo(object):
# variables they can consume (e.g. JVM options, license keys, etc..)
# and add these to env_path_through
def __init__(self, container_descriptions=None, requirements=None, requires_galaxy_python_environment=False, env_pass_through=["GALAXY_SLOTS"]):
def __init__(self, container_descriptions=None, requirements=None, requires_galaxy_python_environment=False, env_pass_through=["GALAXY_SLOTS"], guest_ports=None):
if container_descriptions is None:
container_descriptions = []
if requirements is None:
@@ -45,6 +45,7 @@ class ToolInfo(object):
self.requirements = requirements
self.requires_galaxy_python_environment = requires_galaxy_python_environment
self.env_pass_through = env_pass_through
self.guest_ports = guest_ports
class JobInfo(object):
+46
View File
@@ -82,6 +82,24 @@ def build_docker_load_command(**kwds):
return command_shell("load", [])
def build_docker_simple_command(
command,
docker_cmd=DEFAULT_DOCKER_COMMAND,
sudo=DEFAULT_SUDO,
sudo_cmd=DEFAULT_SUDO_COMMAND,
container_name=None,
**kwd
):
command_parts = _docker_prefix(
docker_cmd=docker_cmd,
sudo=sudo,
sudo_cmd=sudo_cmd,
)
command_parts.append(command)
command_parts.append(container_name or '{CONTAINER_NAME}')
return " ".join(command_parts)
def build_docker_run_command(
container_command,
image,
@@ -102,6 +120,8 @@ def build_docker_run_command(
auto_rm=DEFAULT_AUTO_REMOVE,
set_user=DEFAULT_SET_USER,
host=DEFAULT_HOST,
guest_ports=False,
container_name=None
):
command_parts = _docker_prefix(
docker_cmd=docker_cmd,
@@ -118,6 +138,16 @@ def build_docker_run_command(
# e.g. -e "GALAXY_SLOTS=$GALAXY_SLOTS"
# These are environment variable expansions so we don't quote these.
command_parts.extend(["-e", env_directive])
if guest_ports is True:
# When is True, expose all ports
command_parts.append("-P")
elif guest_ports:
if not isinstance(guest_ports, list):
guest_ports = [guest_ports]
for guest_port in guest_ports:
command_parts.extend(["-p", guest_port])
if container_name:
command_parts.extend(["--name", container_name])
for volume in volumes:
# These are environment variable expansions so we don't quote these.
volume_str = str(volume)
@@ -192,3 +222,19 @@ def _docker_prefix(
if host:
command_parts.extend(["-H", host])
return command_parts
def parse_port_text(port_text):
ports = None
if port_text is not None:
ports = {}
for line in port_text.strip().split('\n'):
if " -> " not in line:
raise Exception("Cannot parse host and port from line [%s]" % line)
tool, host = line.split(" -> ", 1)
hostname, port = host.split(':')
port = int(port)
tool_p, tool_prot = tool.split("/")
tool_p = int(tool_p)
ports[tool_p] = dict(tool_port=tool_p, host=hostname, port=port, protocol=tool_prot)
return ports
@@ -40,6 +40,8 @@ def build_singularity_run_command(
run_extra_arguments=DEFAULT_RUN_EXTRA_ARGUMENTS,
sudo=DEFAULT_SUDO,
sudo_cmd=DEFAULT_SUDO_COMMAND,
guest_ports=False,
container_name=None
):
command_parts = []
# http://singularity.lbl.gov/docs-environment-metadata
+5
View File
@@ -122,6 +122,11 @@ class ToolSource(object):
adjacent to the tool).
"""
@abstractmethod
def parse_realtime(self):
""" Return RealTimeTool entry point templates to expose.
"""
def parse_redirect_url_params_elem(self):
""" Return an XML element describing redirect_url_params.
+19 -1
View File
@@ -201,6 +201,25 @@ class XmlToolSource(ToolSource):
return ParallelismInfo(parallelism)
return parallelism_info
def parse_realtime(self):
realtimetool_el = self.root.find("entry_points")
rtt = []
if realtimetool_el is None:
return rtt
for ep_el in realtimetool_el.findall("entry_point"):
port = ep_el.find("port")
assert port is not None, ValueError('A port is required for InteractiveTools')
port = port.text.strip()
url = ep_el.find("url")
if url is not None:
url = url.text.strip()
name = ep_el.get('name', None)
if name:
name = name.strip()
requires_domain = string_as_bool(ep_el.attrib.get("requires_domain", False))
rtt.append(dict(port=port, url=url, name=name, requires_domain=requires_domain))
return rtt
def parse_hidden(self):
hidden = xml_text(self.root, "hidden")
if hidden:
@@ -221,7 +240,6 @@ class XmlToolSource(ToolSource):
for option_elem in root.findall("options"):
if key in option_elem.attrib:
return string_as_bool(option_elem.get(key))
return default
@property
+3
View File
@@ -172,6 +172,9 @@ class YamlToolSource(ToolSource):
def parse_profile(self):
return self.root_dict.get("profile", "16.04")
def parse_realtime(self):
return self.root_dict.get("entry_points", [])
def parse_python_template_version(self):
python_template_version = self.root_dict.get("python_template_version", None)
if python_template_version is not None:
+70
View File
@@ -49,6 +49,7 @@ A ``data_source`` tool contains a few more relevant attributes.
<xs:element name="edam_operations" type="EdamOperations" minOccurs="0"/>
<xs:element name="xrefs" type="xrefs" minOccurs="0" />
<xs:element name="requirements" type="Requirements" minOccurs="0"/>
<xs:element name="entry_points" type="EntryPoints" minOccurs="0" maxOccurs="1" />
<xs:element name="description" type="xs:string" minOccurs="0">
<xs:annotation gxdocs:best_practices="tool-descriptions">
<xs:documentation xml:lang="en"><![CDATA[The value is displayed in
@@ -225,6 +226,74 @@ For more information, see https://planemo.readthedocs.io/en/latest/writing_advan
</xs:sequence>
</xs:complexType>
<xs:complexType name="EntryPoints">
<xs:annotation>
<xs:documentation xml:lang="en"><![CDATA[
This is a container tag set for the ``entry_point`` tag that contains ``port`` and ``url`` tags
described in greater detail below. ``entry_point``s describe InteractiveTool entry points
to a tool.
]]></xs:documentation>
</xs:annotation>
<xs:sequence>
<xs:element name="entry_point" type="EntryPoint" minOccurs="0" maxOccurs="unbounded"/>
</xs:sequence>
</xs:complexType>
<xs:complexType name="EntryPoint" mixed="true">
<xs:annotation>
<xs:documentation xml:lang="en"><![CDATA[
This tag set is contained within the ``<entry_point>`` tag set. Access to entry point
ports and urls are included in this tag set. These are used by InteractiveTools
to provide access to graphical tools in real-time.
```xml
<entry_points>
<entry_point name="Example name">
<port>80</port>
<url>landing/${template_enabled}/index.html</url>
</entry_point>
</entry_points>
```
]]></xs:documentation>
</xs:annotation>
<xs:sequence>
<xs:element name="port" type="EntryPointPort" minOccurs="1" maxOccurs="1"/>
<xs:element name="url" type="EntryPointURL" minOccurs="0" maxOccurs="1"/>
</xs:sequence>
<xs:attribute name="name" type="xs:string" use="required">
<xs:annotation>
<xs:documentation xml:lang="en">This value defines the name of the entry point.</xs:documentation>
</xs:annotation>
</xs:attribute>
<xs:attribute name="requires_domain" type="PermissiveBoolean" default="false">
<xs:annotation>
<xs:documentation xml:lang="en">This value declares if domain-based proxying is required. Default is False. Currently only works when True.</xs:documentation>
</xs:annotation>
</xs:attribute>
</xs:complexType>
<xs:complexType name="EntryPointPort" mixed="true">
<xs:annotation>
<xs:documentation xml:lang="en"><![CDATA[
This tag set is contained within the ``<entry_point>`` tag set. It contains the entry port.
]]></xs:documentation>
</xs:annotation>
</xs:complexType>
<xs:complexType name="EntryPointURL" mixed="true">
<xs:annotation>
<xs:documentation xml:lang="en"><![CDATA[
This tag set is contained within the ``<entry_point>`` tag set. It contains the entry URL.
]]></xs:documentation>
</xs:annotation>
</xs:complexType>
<xs:complexType name="ToolAction">
<xs:annotation>
<xs:documentation xml:lang="en">Describe the backend Python action to execute for this Galaxy tool.</xs:documentation>
@@ -5359,6 +5428,7 @@ and ``bibtex`` are the only supported options.</xs:documentation>
<xs:restriction base="xs:string">
<xs:enumeration value="data_source"/>
<xs:enumeration value="manage_data"/>
<xs:enumeration value="interactive"/>
<xs:enumeration value="expression"/>
</xs:restriction>
</xs:simpleType>
+37 -3
View File
@@ -404,6 +404,7 @@ class Tool(Dictifiable):
tool_type = 'default'
requires_setting_metadata = True
produces_entry_points = False
default_tool_action = DefaultToolAction
dict_collection_visible_keys = ['id', 'name', 'version', 'description', 'labels']
@@ -536,7 +537,9 @@ class Tool(Dictifiable):
"""Indicates this tool's runtime requires Galaxy's Python environment."""
# All special tool types (data source, history import/export, etc...)
# seem to require Galaxy's Python.
if self.tool_type not in ["default", "manage_data"]:
# FIXME: the (instantiated) tool class should emit this behavior, and not
# use inspection by string check
if self.tool_type not in ["default", "manage_data", "interactive"]:
return True
if self.tool_type == "manage_data" and self.profile < 18.09:
@@ -812,6 +815,7 @@ class Tool(Dictifiable):
self.__parse_trackster_conf(tool_source)
# Record macro paths so we can reload a tool if any of its macro has changes
self._macro_paths = tool_source.macro_paths()
self.ports = tool_source.parse_realtime()
def __parse_legacy_features(self, tool_source):
self.code_namespace = dict()
@@ -1918,7 +1922,8 @@ class Tool(Dictifiable):
tool_dict['panel_section_id'], tool_dict['panel_section_name'] = self.get_panel_section()
tool_class = self.__class__
regular_form = tool_class == Tool or isinstance(self, DatabaseOperationTool)
# FIXME: the Tool class should declare directly, instead of ad hoc inspection
regular_form = tool_class == Tool or isinstance(self, (DatabaseOperationTool, RealTimeTool))
tool_dict["form_style"] = "regular" if regular_form else "special"
return tool_dict
@@ -2440,6 +2445,35 @@ class ImportHistoryTool(Tool):
tool_type = 'import_history'
class RealTimeTool(Tool):
tool_type = 'interactive'
produces_entry_points = True
def __init__(self, config_file, tool_source, app, **kwd):
assert app.config.interactivetools_enable, ValueError('Trying to load an InteractiveTool, but InteractiveTools are not enabled.')
super(RealTimeTool, self).__init__(config_file, tool_source, app, **kwd)
for port in self.ports:
assert port.get('requires_domain', None), ValueError('InteractiveTools currently only work when requires_domain is True for each entry_point.')
def __remove_realtime_by_job(self, job):
if job:
eps = job.realtimetool_entry_points
log.debug('__remove_realtime_by_job: %s', eps)
self.app.realtime_manager.remove_entry_points(eps)
else:
log.warning("Could not determine job to stop InteractiveTool: %s", job)
def exec_after_process(self, app, inp_data, out_data, param_dict, job=None, **kwds):
# run original exec_after_process
super(RealTimeTool, self).exec_after_process(app, inp_data, out_data, param_dict, job=job, **kwds)
self.__remove_realtime_by_job(job)
def job_failed(self, job_wrapper, message, exception=False):
super(RealTimeTool, self).job_failed(job_wrapper, message, exception=exception)
job = job_wrapper.sa_session.query(model.Job).get(job_wrapper.job_id)
self.__remove_realtime_by_job(job)
class DataManagerTool(OutputParameterJSONTool):
tool_type = 'manage_data'
default_tool_action = DataManagerToolAction
@@ -3050,7 +3084,7 @@ class FilterFromFileTool(DatabaseOperationTool):
# Populate tool_type to ToolClass mappings
tool_types = {}
for tool_class in [Tool, SetMetadataTool, OutputParameterJSONTool, ExpressionTool,
for tool_class in [Tool, SetMetadataTool, OutputParameterJSONTool, ExpressionTool, RealTimeTool,
DataManagerTool, DataSourceTool, AsyncDataSourceTool,
UnzipCollectionTool, ZipCollectionTool, MergeCollectionTool, RelabelFromFileTool, FilterFromFileTool,
BuildListCollectionTool, ExtractDatasetCollectionTool,
+26
View File
@@ -149,6 +149,8 @@ class ToolEvaluator(object):
self.__sanitize_param_dict(param_dict)
# Parameters added after this line are not sanitized
self.__populate_non_job_params(param_dict)
# Populate and store templated RealTimeTools values
self.__populate_realtimetools(param_dict)
# Return the dictionary of parameters
return param_dict
@@ -404,6 +406,30 @@ class ToolEvaluator(object):
# the paths rewritten.
self.__walk_inputs(self.tool.inputs, param_dict, rewrite_unstructured_paths)
def __populate_realtimetools(self, param_dict):
"""
Populate RealTimeTools templated values.
"""
rtt = []
for ep in getattr(self.tool, 'ports', []):
ep_dict = {}
for key in 'port', 'name', 'url':
val = ep.get(key, None)
if val is not None:
val = fill_template(val, context=param_dict, python_template_version=self.tool.python_template_version)
clean_val = []
for line in val.split('\n'):
clean_val.append(line.strip())
val = '\n'.join(clean_val)
val = val.replace("\n", " ").replace("\r", " ").strip()
ep_dict[key] = val
rtt.append(ep_dict)
self.realtimetools = rtt
rtt_man = getattr(self.app, "realtime_manager", None)
if rtt_man:
rtt_man.create_realtime(self.job, self.tool, rtt)
return rtt
def __sanitize_param_dict(self, param_dict):
"""
Sanitize all values that will be substituted on the command line, with the exception of ToolParameterValueWrappers,
@@ -0,0 +1,77 @@
""" API for asynchronous job running mechanisms can use to fetch or put files
related to running and queued jobs.
"""
import logging
from galaxy import exceptions, util
from galaxy.managers.realtime import RealTimeManager
from galaxy.web import expose_api_anonymous_and_sessionless
from galaxy.webapps.base.controller import BaseAPIController
log = logging.getLogger(__name__)
class ToolEntryPointsAPIController(BaseAPIController):
def __init__(self, app):
self.app = app
self.realtime_manager = RealTimeManager(app)
@expose_api_anonymous_and_sessionless
def index(self, trans, running=False, job_id=None, **kwd):
"""
* GET /api/entry_points
Returns tool entry point information. Currently passing a job_id
parameter is required, as this becomes more general that won't be
needed.
:type job_id: string
:param job_id: Encoded job id
:type running: boolean
:param running: filter to only include running job entry points.
:rtype: list
:returns: list of entry point dictionaries.
"""
running = util.asbool(running)
if job_id is None and not running:
raise exceptions.RequestParameterInvalidException("Currently this API must passed a job id or running=true")
if job_id is not None and running:
raise exceptions.RequestParameterInvalidException("Currently this API must passed only a job id or running=true")
if job_id is not None:
job = trans.sa_session.query(trans.app.model.Job).get(self.decode_id(job_id))
if not self.realtime_manager.can_access_job(trans, job):
raise exceptions.ItemAccessibilityException()
entry_points = job.realtimetool_entry_points
if running:
entry_points = self.realtime_manager.get_nonterminal_for_user_by_trans(trans)
rval = []
for entry_point in entry_points:
as_dict = self.encode_all_ids(trans, entry_point.to_dict(), True)
target = self.realtime_manager.target_if_active(trans, entry_point)
if target:
as_dict["target"] = target
rval.append(as_dict)
return rval
@expose_api_anonymous_and_sessionless
def access_entry_point(self, trans, id, **kwd):
"""
* GET /api/entry_points/{id}/access
Return the URL target described by the entry point.
:type id: string
:param id: Encoded entry point id
:rtype: dictionary
:returns: dictionary containing target for realtime entry point
"""
# Because of auto id encoding needed for link from grid, the item.id keyword must be 'id'
if not id:
raise exceptions.RequestParameterMissingException("Must supply entry point ID.")
entry_point_id = self.decode_id(id)
return {"target": self.realtime_manager.access_entry_point_target(trans, entry_point_id)}
+1 -1
View File
@@ -530,7 +530,7 @@ class ToolsController(BaseAPIController, UsesVisualizationMixin):
# TODO: check for errors and ensure that output dataset(s) are available.
output_datasets = vars.get('out_data', [])
rval = {'outputs': [], 'output_collections': [], 'jobs': [], 'implicit_collections': []}
rval['produces_entry_points'] = tool.produces_entry_points
job_errors = vars.get('job_errors', [])
if job_errors:
# If we are here - some jobs were successfully executed but some failed.
+4
View File
@@ -151,6 +151,7 @@ def app_factory(global_conf, load_app_kwds={}, **kwargs):
webapp.add_client_route('/workflows/run')
webapp.add_client_route('/workflows/import')
webapp.add_client_route('/custom_builds')
webapp.add_client_route('/realtime_entry_points/list')
# ==== Done
# Indicate that all configuration settings have been provided
@@ -362,6 +363,9 @@ def populate_api_routes(webapp, app):
webapp.mapper.resource('tool', 'tools', path_prefix='/api')
webapp.mapper.resource('dynamic_tools', 'dynamic_tools', path_prefix='/api')
webapp.mapper.connect('/api/entry_points', action='index', controller="tool_entry_points")
webapp.mapper.connect('/api/entry_points/{id:.+?}/access', action='access_entry_point', controller="tool_entry_points")
webapp.mapper.connect('/api/dependency_resolvers/clean', action="clean", controller="tool_dependencies", conditions=dict(method=["POST"]))
webapp.mapper.connect('/api/dependency_resolvers/dependency', action="manager_dependency", controller="tool_dependencies", conditions=dict(method=["GET"]))
webapp.mapper.connect('/api/dependency_resolvers/dependency', action="install_dependency", controller="tool_dependencies", conditions=dict(method=["POST"]))
@@ -928,6 +928,13 @@ mapping:
Galaxy server, but contains a redirect to a third-party server, tricking a
Galaxy user to access said site.
interactivetools_enable:
type: bool
default: false
required: false
desc: |
Enable InteractiveTools.
visualizations_visible:
type: bool
default: true
@@ -0,0 +1,106 @@
"""
Provides web interaction with RealTimeTools
"""
import logging
from galaxy import (
model,
web
)
from galaxy.web.framework.helpers import (
grids,
time_ago,
)
from galaxy.webapps.base.controller import (
BaseUIController,
)
log = logging.getLogger(__name__)
class JobStatusColumn(grids.StateColumn):
def get_value(self, trans, grid, item):
return super(JobStatusColumn, self).get_value(trans, grid, item.job)
class EntryPointLinkColumn(grids.GridColumn):
def get_value(self, trans, grid, item):
return '<a class="entry-point-link" entry_point_id="%s">%s</a>' % (trans.security.encode_id(item.id), item.name)
class RealTimeToolEntryPointListGrid(grids.Grid):
use_panels = True
title = "Available InteractiveTools"
model_class = model.InteractiveToolEntryPoint
default_filter = {"name": "All"}
default_sort_key = "-update_time"
columns = [
EntryPointLinkColumn("Name", filterable="advanced"),
JobStatusColumn("Job Info", key="job_state", model_class=model.Job),
grids.GridColumn("Created", key="created_time", format=time_ago),
grids.GridColumn("Last Updated", key="modified_time", format=time_ago),
]
columns.append(
grids.MulticolFilterColumn(
"Search",
cols_to_filter=[columns[0]],
key="free-text-search", visible=False, filterable="standard"
)
)
operations = [
grids.GridOperation("Stop", condition=(lambda item: item.active), async_compatible=False),
]
def build_initial_query(self, trans, **kwargs):
# Get list of user's active RealTimeTools
return trans.app.realtime_manager.get_nonterminal_for_user_by_trans(trans)
class RealTime(BaseUIController):
entry_point_grid = RealTimeToolEntryPointListGrid()
@web.expose_api_anonymous
def list(self, trans, **kwargs):
"""List all available realtimetools"""
if not trans.app.config.interactivetools_enable:
raise web.httpexceptions.HTTPNotFound()
operation = kwargs.get('operation', None)
message = None
status = None
if operation:
eps = []
ids = kwargs.get('id', None)
if ids:
if not isinstance(ids, list):
ids = [ids]
for entry_point_id in ids:
entry_point_id = self.decode_id(entry_point_id)
entry_point = trans.sa_session.query(trans.app.model.InteractiveToolEntryPoint).get(entry_point_id)
if trans.app.realtime_manager.can_access_entry_point(trans, entry_point):
eps.append(entry_point)
if eps:
failed = []
succeeded = []
jobs = []
if operation == 'stop':
for ep in eps:
if ep.job not in jobs:
stopped = trans.app.realtime_manager.stop(trans, ep)
if stopped:
succeeded.append(ep)
jobs.append(ep.job)
else:
failed.append(ep)
else:
succeeded.append(ep)
if failed:
message = 'Unable to stop %i InteractiveTools.' % (len(failed))
status = 'error'
if succeeded:
message = 'Stopped %i InteractiveTools.' % (len(succeeded))
status = 'ok'
if message and status:
kwargs['message'] = message
kwargs['status'] = status
return self.entry_point_grid(trans, **kwargs)
@@ -61,7 +61,8 @@ class ToolRunner(BaseUIController):
redirect=redirect))
if not tool.allow_user_access(trans.user):
return __tool_404__()
if tool.tool_type == 'default':
# FIXME: Tool class should define behavior
if tool.tool_type in ['default', 'realtime']:
return trans.response.send_redirect(url_for(controller='root', tool_id=tool_id))
# execute tool without displaying form (used for datasource tools)
@@ -0,0 +1,48 @@
import json
import os
import subprocess
import sys
import tempfile
# insert *this* galaxy before all others on sys.path
sys.path.insert(1, os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir)))
from galaxy.tool_util.deps import docker_util
def main():
with open("container_config.json", "r") as f:
container_config = json.load(f)
container_type = container_config["container_type"]
container_name = container_config["container_name"]
connection_configuration = container_config["connection_configuration"]
if container_type != "docker":
raise Exception("Monitoring container type [%s], not yet implemented." % container_type)
ports_raw = None
try:
while True:
ports_command = docker_util.build_docker_simple_command("port", container_name=container_name, **connection_configuration)
with tempfile.TemporaryFile() as stdout_file:
exit_code = subprocess.call(ports_command,
shell=True,
stdout=stdout_file,
preexec_fn=os.setpgrp)
if exit_code == 0:
stdout_file.seek(0)
ports_raw = stdout_file.read()
break
if ports_raw is not None:
with open("container_runtime.json", "w") as f:
json.dump(docker_util.parse_port_text(ports_raw), f)
else:
raise Exception("Failed to recover ports...")
except Exception as e:
with open("exception.txt", "w") as f:
f.write(str(e))
if __name__ == "__main__":
main()
@@ -0,0 +1,86 @@
import itertools
import sqlite3
from threading import RLock
from time import time
import uwsgi
realtime_db_file = uwsgi.opt["interactivetools_map"]
db_conn = sqlite3.connect(realtime_db_file)
class CacheEntry():
def __init__(self, key, value, ttl=20):
self.key = key
self.value = value
self.expires_at = time() + ttl
self._expired = False
def expired(self):
if self._expired is False:
return (self.expires_at < time())
else:
return self._expired
class CacheList():
def __init__(self):
self.entries = []
self.lock = RLock()
def add_entry(self, key, value, ttl=20):
with self.lock:
self.entries.append(CacheEntry(key, value, ttl))
def read_entries(self):
with self.lock:
self.entries = list(itertools.dropwhile(lambda x: x.expired(), self.entries))
return self.entries
def get_entry_value(self, key, default=None):
entries = self.read_entries()
for entry in entries:
if entry.key == key:
return entry.value
return default
key_type_token_mapped_cache = CacheList()
def key_type_token_mapper_cached(key, key_type, token, route_extra, url, ttl):
global key_type_token_mapped_cache
cache_key = (key, key_type, token)
entry = key_type_token_mapped_cache.get_entry_value(cache_key, None)
if entry is None:
entry = key_type_token_mapper(key, key_type, token, route_extra, url)
if entry is not None:
# Should we cache empt/not authorized entries, perhaps for shorter time?
key_type_token_mapped_cache.add_entry(cache_key, entry, ttl=float(ttl))
return entry
def key_type_token_mapper(key, key_type, token, route_extra, url):
global db_conn
# print 'key %s key_type %s token %s route_extra %s url %s\n' % (key, key_type, token, route_extra, url)
if key and key_type and token:
# sqlite3.ProgrammingError: SQLite objects created in a thread can only be used in that same thread. The object was created in thread id x and this is thread id y.
# So try upto 2 times
for i in range(2):
# Order by rowid gives us the last row added
try:
row = db_conn.execute("select host, port from gxrtproxy where key=? and key_type=? and token=? order by rowid desc limit 1", (key, key_type, token)).fetchone()
if row:
rval = '%s:%s' % (tuple(row))
return rval.encode()
break
except sqlite3.ProgrammingError:
db_conn = sqlite3.connect(realtime_db_file)
continue
break
return None
uwsgi.register_rpc('rtt_key_type_token_mapper', key_type_token_mapper)
uwsgi.register_rpc('rtt_key_type_token_mapper_cached', key_type_token_mapper_cached)
+29 -1
View File
@@ -19,6 +19,7 @@ import nose.config
import nose.core
import nose.loader
import nose.plugins.manager
import yaml
from paste import httpserver
from six.moves import (
http_client,
@@ -59,6 +60,17 @@ MIGRATED_TOOL_PANEL_CONFIG = 'config/migrated_tools_conf.xml'
INSTALLED_TOOL_PANEL_CONFIGS = [
os.environ.get('GALAXY_TEST_SHED_TOOL_CONF', 'config/shed_tool_conf.xml')
]
REALTIME_PROXY_TEMPLATE = string.Template(r"""
uwsgi:
realtime_map: $tempdir/realtime_map.sqlite
python-raw: scripts/realtime/key_type_token_mapping.py
route-host: ^([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.(realtime\.$test_host:$test_port)$ goto:realtime
route-run: goto:endendend
route-label: realtime
route-host: ^([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.(realtime\.$test_host:$test_port)$ rpcvar:TARGET_HOST rtt_key_type_token_mapper_cached $2 $1 $3 $4 $0 5
route-if-not: empty:${TARGET_HOST} httpdumb:${TARGET_HOST}
route-label: endendend
""")
DEFAULT_LOCALES = "en"
@@ -717,11 +729,27 @@ def launch_uwsgi(kwargs, tempdir, prefix=DEFAULT_CONFIG_PREFIX, config_object=No
config = {}
config["galaxy"] = kwargs.copy()
enable_realtime_mapping = getattr(config_object, "enable_realtime_mapping", False)
if enable_realtime_mapping:
config["galaxy"]["realtime_prefix"] = "realtime"
config["galaxy"]["realtime_map"] = os.path.join(tempdir, "realtime_map.sqlite")
yaml_config_path = os.path.join(tempdir, "galaxy.yml")
with open(yaml_config_path, "w") as f:
import yaml
yaml.dump(config, f)
if enable_realtime_mapping:
# Avoid YAML.dump configuration since uwsgi doesn't like real YAML :( -
# though maybe it would work?
with open(yaml_config_path, "r") as f:
old_contents = f.read()
with open(yaml_config_path, "w") as f:
test_port = str(port) if port else r"[0-9]+"
test_host = host or "localhost"
uwsgi_section = REALTIME_PROXY_TEMPLATE.safe_substitute(test_host=test_host, test_port=test_port, tempdir=tempdir)
f.write(uwsgi_section)
f.write(old_contents)
def attempt_port_bind(port):
uwsgi_command = [
"uwsgi",
@@ -0,0 +1,53 @@
{
"cells": [
{
"cell_type": "markdown",
"metadata": {},
"source": [
"# Welcome to the interactive Galaxy IPython Notebook."
]
},
{
"cell_type": "markdown",
"metadata": {},
"source": [
"You can access your data via the dataset number. Using a Python kernel, you can access dataset number 42 with ``handle = open(get(42), 'r')``.\n",
"To save data, write your data to a file, and then call ``put('filename.txt')``. The dataset will then be available in your galaxy history.\n<br>",
"When using a non-Python kernel, ``get`` and ``put`` are available as command-line tools, which can be accessed using system calls in R, Julia, and Ruby. For example, to read dataset number 42 into R, you can write ```handle <- file(system('get -i 42', intern = TRUE))```.\n",
"To save data in R, write the data to a file and then call ``system('put -p filename.txt')``.\n",
"Notebooks can be saved to Galaxy by clicking the large green button at the top right of the IPython interface.<br>\n",
"More help and informations can be found on the project [website](https://github.com/bgruening/docker-jupyter-notebook)."
]
},
{
"cell_type": "code",
"execution_count": 1,
"metadata": {
"collapsed": false
},
"outputs": [],
"source": []
}
],
"metadata": {
"kernelspec": {
"display_name": "Python 2",
"language": "python",
"name": "python2"
},
"language_info": {
"codemirror_mode": {
"name": "ipython",
"version": 2
},
"file_extension": ".py",
"mimetype": "text/x-python",
"name": "python",
"nbconvert_exporter": "python",
"pygments_lexer": "ipython2",
"version": "2.7.10"
}
},
"nbformat": 4,
"nbformat_minor": 0
}
@@ -0,0 +1,59 @@
<tool id="interactive_tool_askomics" tool_type="interactive" name="AskOmics" version="0.1">
<description>AskOmics, a visual SPARQL query builder</description>
<requirements>
<container type="docker">quay.io/askomics/askomics-ie:17.12_g19.09</container>
</requirements>
<entry_points>
<entry_point name="AskOmics instance on $infile.display_name" requires_domain="True">
<port>6543</port>
<url>/login_api_gie?key=abcd</url>
</entry_point>
</entry_points>
<environment_variables>
<environment_variable name="GALAXY_URL">${__app__.config.galaxy_infrastructure_url}</environment_variable> <!-- FIXME: Warning: The use of __app__ is deprecated and will break backward compatibility in the near future -->
<environment_variable name="API_KEY" strip="True">
#if $__user__:
#for $api_key in $__user__.api_keys:
${api_key.key}
#break
#end for
#end if
</environment_variable>
</environment_variables>
<command><![CDATA[
#import re
## ToDo: the key could be generated randomly
export ASKO_load_url='http://localhost:6543' &&
export ASKOMICS_API_KEY='abcd' &&
export ASKO_files_dir='/tmp/askomics-ie' &&
export VIRT_Parameters_NumberOfBuffers='10000' &&
export VIRT_Parameters_MaxDirtyBuffers='6000' &&
#set link_name = re.sub('[^\w_]', '_', $infile.element_identifier)
#if $infile.ext == 'tabular':
#set link_name = $link_name + '.tsv'
#elif $infile.ext == 'interval':
#set link_name = $link_name + '.bed'
#else:
#set link_name = $link_name + '.' + $infile.ext
#end if
mkdir -p /import &&
ln -s '$infile' '/import/$link_name' &&
bash /start.sh
]]>
</command>
<inputs>
<param name="infile" type="data" format="tabular,gff,gff3,bed,interval" label="A datasets with genomic coordinates"/>
</inputs>
<outputs>
<data name="outfile" format="txt" />
</outputs>
<tests>
</tests>
<help>
AskOmics is a visual SPARQL query interface supporting both intuitive data integration and
querying while shielding the user from most of the technical difficulties underlying RDF and SPARQL.
</help>
</tool>
@@ -0,0 +1,44 @@
<tool id="interactive_tool_bam_iobio" tool_type="interactive" name="BAM (iobio) Visualisation" version="0.1">
<requirements>
<container type="docker">qiaoy/iobio-bundle.bam-iobio:1.0-ondemand</container>
</requirements>
<entry_points>
<entry_point name="BAM io.bio visualisation of $infile.display_name" requires_domain="True">
<port>80</port>
<url><![CDATA[/?bam=http://localhost/tmp/bamfile.bam&region=1]]></url>
</entry_point>
</entry_points>
<command><![CDATA[
## ToDo: websocket could not be found
## WebSocket connection to 'ws://localhost/bamreaddepther/' failed: Error in connection establishment: net::ERR_CONNECTION_REFUSED
#set $PUB_HOSTNAME = 'localhost'
#set $PUB_HTTP_PORT = '80'
cd /var/www/html &&
sed -i "s@\"wss://services.iobio.io/samtools/\"@((window.location.protocol === \"https:\") ? \"wss://\" : \"ws://\") + window.location.host + \"/samtools/\"@" js/bam.iobio.js/bam.iobio.js &&
sed -i "s@\"wss://services.iobio.io/bamreaddepther/\"@((window.location.protocol === \"https:\") ? \"wss://\" : \"ws://\") + window.location.host + \"/bamreaddepther/\"@" js/bam.iobio.js/bam.iobio.js &&
sed -i "s@\"wss://services.iobio.io/bamstatsalive/\"@((window.location.protocol === \"https:\") ? \"wss://\" : \"ws://\") + window.location.host + \"/bamstatsalive/\"@" js/bam.iobio.js/bam.iobio.js &&
sed -i "s@\"wss://services.iobio.io/samheader/\"@((window.location.protocol === \"https:\") ? \"wss://\" : \"ws://\") + window.location.host + \"/samheader/\"@" js/bam.iobio.js/bam.iobio.js &&
cp '${infile}' /input/bamfile.bam &&
cp '${infile.metadata.bam_index}' /input/bamfile.bam.bai &&
mkdir /var/log/supervisor/ &&
head -n -2 /etc/supervisor.d/app.conf > /tmp/app.conf &&
mv /tmp/app.conf /etc/supervisor.d/app.conf &&
/usr/bin/supervisord -c /etc/supervisord.conf
]]>
</command>
<inputs>
<param name="infile" type="data" format="bam" label="BAM file"/>
</inputs>
<outputs>
<data name="outfile" format="txt" />
</outputs>
<tests>
</tests>
<help>
BAM iobio visualisation.
</help>
</tool>
@@ -0,0 +1,29 @@
<tool id="interactive_tool_cellxgene" tool_type="interactive" name="Interactive CellXgene Environment" version="0.1">
<requirements>
<container type="docker">quay.io/galaxy/cellxgene-galaxy-ie:ie2</container>
</requirements>
<entry_points>
<entry_point name="Cellxgene Single Cell Visualisation on $infile.display_name" requires_domain="True">
<port>80</port>
</entry_point>
</entry_points>
<command><![CDATA[
#import re
#set $fancy_name = '/tmp/galaxy_cellxgene_' + re.sub('[^\w\-_]', '_', $infile.element_identifier) + '.h5ad'
cp '${infile}' '${fancy_name}'
&&
cellxgene launch --host 0.0.0.0 --port 80 '${fancy_name}'
]]>
</command>
<inputs>
<param name="infile" type="data" format="h5ad" label="Concatenate Dataset"/>
</inputs>
<outputs>
<data name="out_file1" format="txt" />
</outputs>
<tests>
</tests>
<help>
Interactive tool for visualising AnnData.
</help>
</tool>
@@ -0,0 +1,63 @@
<tool id="interactive_tool_ethercalc" tool_type="interactive" name="EtherCalc" version="0.1">
<requirements>
<container type="docker">shiltemann/ethercalc-galaxy-ie:17.05</container>
</requirements>
<entry_points>
<entry_point name="EtherCalc instance on $infile.display_name" requires_domain="True">
<port>8000</port>
</entry_point>
</entry_points>
<command><![CDATA[
forever start `which ethercalc` --cors &
bash $wait_script
&&
curl --include --request POST --header "Content-Type: application/json" --data-binary "{ \"room\": \"galaxy\", \"snapshot\": \"...\"}" 'http://localhost:8000/_'
&&
## put loading message into worksheet in case of large files
echo -e "Loading your,dataset\nThis may take, some time for, large files.." > loading.txt
&&
curl --include --request PUT --header "Content-Type: text/csv" --data-binary @loading.txt http://localhost:8000/_/galaxy
&&
## remove dump file so this doesnt appear in audit trail
rm /dump.json
&&
## load dataset into worksheet
curl --include --request PUT --header "Content-Type: text/csv" --data-binary @$infile http://localhost:8000/_/galaxy
&&
tail -f /etc/hosts
]]>
</command>
<configfiles>
<configfile name="wait_script">
<![CDATA[
## make sure ethercalc has finished starting up
STATUS=\$(curl --include 'http://localhost:8000/_/galaxy' 2>&1)
while [[ \${STATUS} =~ "refused" ]]
do
echo "waiting for ethercalc: \$STATUS \n"
STATUS=\$(curl --include 'http://localhost:8000/_/galaxy' 2>&1)
sleep 2
done
]]>
</configfile>
</configfiles>
<inputs>
<param name="infile" type="data" format="tabular,csv,tsv" label="Some tabular dataset"/>
</inputs>
<outputs>
<data name="outfile" format="txt" />
</outputs>
<tests>
</tests>
<help>
EtherCalc is a web spreadsheet.
https://ethercalc.net
</help>
</tool>
@@ -0,0 +1,33 @@
<tool id="interactive_tool_hicbrowser" tool_type="interactive" name="HiCBrowser" version="0.1">
<requirements>
<container type="docker">bgruening/hicbrowser</container>
</requirements>
<entry_points>
<entry_point name="HiCBrowser on $infile.display_name" requires_domain="True">
<port>80</port>
</entry_point>
</entry_points>
<command><![CDATA[
#set $PROXY_PREFIX = ""
sed -i "s|PROXY_PREFIX|${PROXY_PREFIX}|" /etc/nginx/conf.d/nginx.conf &&
rm /data/* -rf &&
cp '$infile' /data/data_pack.zip &&
cd /data &&
unzip data_pack.zip &&
supervisord &&
tail -f /var/log/supervisor/*
]]>
</command>
<inputs>
<param name="infile" type="data" format="zip" label="HiCExplorer archive"/>
</inputs>
<outputs>
<data name="outfile" format="txt" />
</outputs>
<tests>
</tests>
<help>
Visualising HiC data with HiCBrowser.
</help>
</tool>
@@ -0,0 +1,100 @@
<tool id="interactive_tool_jupyter_notebook" tool_type="interactive" name="Interactive Jupyter Notebook" version="0.1">
<requirements>
<container type="docker">quay.io/bgruening/docker-jupyter-notebook:ie2</container>
</requirements>
<entry_points>
<entry_point name="Jupyter Interactive Tool" requires_domain="True">
<port>8888</port>
<url>ipython/tree</url>
</entry_point>
</entry_points>
<environment_variables>
<environment_variable name="HISTORY_ID" strip="True">${__app__.security.encode_id($jupyter_notebook.history_id)}</environment_variable> <!-- FIXME: Warning: The use of __app__ is deprecated and will break backward compatibility in the near future -->
<environment_variable name="REMOTE_HOST">${__app__.config.galaxy_infrastructure_url}</environment_variable> <!-- FIXME: Warning: The use of __app__ is deprecated and will break backward compatibility in the near future -->
<environment_variable name="GALAXY_WEB_PORT">8080</environment_variable>
<environment_variable name="GALAXY_URL">${__app__.config.galaxy_infrastructure_url}</environment_variable> <!-- FIXME: Warning: The use of __app__ is deprecated and will break backward compatibility in the near future -->
<environment_variable name="API_KEY" strip="True">
#if $__user__:
#for $api_key in $__user__.api_keys:
${api_key.key}
#break
#end for
#end if
</environment_variable> <!-- FIXME: We should have a better way to get user's API key -->
</environment_variables>
<command detect_errors="aggressive"><![CDATA[
#import re
export GALAXY_WORKING_DIR=`pwd` &&
mkdir -p ./jupyter/outputs/ &&
mkdir -p ./jupyter/data &&
#set $cleaned_name = re.sub('[^\w\-\.]', '_', str($input.element_identifier))
ln -sf '$input' './jupyter/data/${cleaned_name}' &&
## change into the directory where the notebooks are located
cd ./jupyter/ &&
export PATH=/home/jovyan/.local/bin:\$PATH &&
#if $mode.mode_select == 'scratch':
## copy default notebook
cp '$__tool_directory__/default_notebook.ipynb' ./ipython_galaxy_notebook.ipynb &&
jupyter trust ./ipython_galaxy_notebook.ipynb &&
jupyter lab --no-browser --NotebookApp.shutdown_button=True &&
cp ./ipython_galaxy_notebook.ipynb '$jupyter_notebook'
#else:
#set $cleaned_name = re.sub('[^\w\-\.]', '_', str($input.element_identifier))
cp '$mode.ipynb' ./${cleaned_name}.ipynb &&
jupyter trust ./${cleaned_name}.ipynb &&
#if $mode.run_it:
jupyter nbconvert --to notebook --execute --output ./ipython_galaxy_notebook.ipynb --allow-errors ./*.ipynb &&
#else:
jupyter lab --no-browser --NotebookApp.shutdown_button=True &&
#end if
cp ./ipython_galaxy_notebook.ipynb '$jupyter_notebook'
#end if
]]>
</command>
<inputs>
<conditional name="mode">
<param name="mode_select" type="select" label="Do you already have a notebook?" help="If not, no problem we will provide you with a default one.">
<option value="scratch">Start with a fresh notebook</option>
<option value="previous">Load a previous notebook</option>
</param>
<when value="scratch"/>
<when value="previous">
<param name="ipynb" type="data" format="ipynb" label="IPython Notebook"/>
<param name="run_it" type="boolean" truevalue="true" falsevalue="false" label="Execute notebook and return a new one."
help="This option is useful in workflows when you just want to execute a notebook and not dive into the webfrontend."/>
</when>
</conditional>
<param name="input" type="data" optional="true" label="Include data into the environment"/>
</inputs>
<outputs>
<data name="jupyter_notebook" format="ipynb" label="Executed Notebook"></data>
</outputs>
<tests>
<test expect_num_outputs="1">
<param name="mode" value="previous" />
<param name="ipynb" value="test.ipynb" />
<param name="run_it" value="true" />
<output name="jupyter_notebook" file="test.ipynb" ftype="ipynb"/>
</test>
</tests>
<help>
The Jupyter Notebook is an open-source web application that allows you to create and share documents that contain live code, equations,
visualizations and narrative text. Uses include: data cleaning and transformation, numerical simulation, statistical modeling, data visualization,
machine learning, and much more.
Galaxy offers you to use Jupyter Notebooks directly in Galaxy accessing and interacting with Galaxy datasets as you like. A very common use-case is to
do the heavy lifting and data reduction steps in Galaxy and the plotting and more `interactive` part on smaller datasets in Jupyter.
You can start with a new Jupyter notebook from scratch or load an already existing one, e.g. from your collegue and execute it on your dataset.
If you have a defined input dataset you can even execute a Jupyter notebook in a workflow, given that the notebook is writing the output back to the history.
You can import data into the notebook via a predefined `get()` function and write results back to Galaxy with a `put()` function.
</help>
</tool>
@@ -0,0 +1,41 @@
<tool id="interactive_tool_neo4j" tool_type="interactive" name="Neo4j (Graph Database)" version="0.1">
<requirements>
<container type="docker">quay.io/sanbi-sa/neo_ie:3.1.9</container>
</requirements>
<entry_points>
<entry_point name="Neo4j on $infile.display_name" requires_domain="True">
<port>80</port>
</entry_point>
<!-- I get some web-socket error that 7687 is not reachable.
but docker does not expose 7687?
-->
</entry_points>
<environment_variables>
<environment_variable name="USER_GID">2345</environment_variable>
<environment_variable name="USER_UID">2345</environment_variable>
<environment_variable name="MONITOR_TRAFFIC">false</environment_variable>
</environment_variables>
<command><![CDATA[
##export USER_UID=2345 &&
##export USER_GID=2345 &&
##export MONITOR_TRAFFIC=false &&
mkdir -p /data/neo4jdb &&
unzip '$infile' -d /data/neo4jdb &&
ls -l /data/neo4jdb &&
sed -i.bak 's@/monitor_traffic.*@@g' /docker-entrypoint.sh &&
cd /opt/neo4j &&
bash /docker-entrypoint.sh neo4j
]]>
</command>
<inputs>
<param name="infile" type="data" format="zip" label="Neo4J data store"/>
</inputs>
<outputs>
<data name="outfile" format="txt" />
</outputs>
<tests>
</tests>
<help>
Neo4j is a highly scalable, robust native graph database.
</help>
</tool>
@@ -0,0 +1,37 @@
<tool id="interactive_tool_pinch" tool_type="interactive" name="Phinch Visualisation" version="0.1">
<requirements>
<container type="docker">shiltemann/docker-phinch-galaxy:16.04</container>
</requirements>
<entry_points>
<entry_point name="Phinch Visualisation of $infile.display_name" requires_domain="True">
<port>80</port>
</entry_point>
</entry_points>
<command><![CDATA[
## ToDo nginx, proxy.conf etc can be removed from the container
#import os
#set $name = os.path.splitext(str($infile.display_name))[0]
## in case someone names the data testdata.
rm /home/Phinch/data/testdata.biom | true &&
ln -s '$infile' /home/Phinch/data/${name}.biom &&
cd /home/Phinch/data &&
sed -i "s/'REPLACE_ME'/'${name}.biom'/g" /home/Phinch/scripts/readFile.js &&
## keep it running
cd /home/Phinch &&
php -S 0.0.0.0:80 2>&1 > /var/log/phinch.log
]]>
</command>
<inputs>
<param name="infile" type="data" format="biom1" label="Biom1 dataset"/>
</inputs>
<outputs>
<data name="outfile" format="txt" />
</outputs>
<tests>
</tests>
<help>
Interactive tool for visualising Biom data.
</help>
</tool>
@@ -0,0 +1,24 @@
<tool id="interactive_tool_simple" name="realtimetool_simple" tool_type="interactive" version="0.1">
<requirements>
<container type="docker">galaxy/test-http-example:0.1</container>
</requirements>
<entry_points>
<entry_point name="Simple" requires_domain="True">
<port>7000</port>
<url>/</url>
</entry_point>
</entry_points>
<command detect_errors="exit_code"><![CDATA[
cd /;
python -m SimpleHTTPServer 7000
]]>
</command>
<inputs>
</inputs>
<outputs>
</outputs>
<tests>
</tests>
<help>
</help>
</tool>
@@ -0,0 +1,29 @@
<tool id="interactive_tool_two_entry_points" name="realtimetool_two_entry_points" tool_type="interactive" version="0.1">
<requirements>
<container type="docker">galaxy/test-http-example:0.1</container>
</requirements>
<entry_points>
<entry_point name="Server 1" requires_domain="True">
<port>7000</port>
<url>/</url>
</entry_point>
<entry_point name="Server 2" requires_domain="True">
<port>7001</port>
<url>/</url>
</entry_point>
</entry_points>
<command detect_errors="exit_code"><![CDATA[
cd /;
python -m SimpleHTTPServer 7000 &
python -m SimpleHTTPServer 7001
]]>
</command>
<inputs>
</inputs>
<outputs>
</outputs>
<tests>
</tests>
<help>
</help>
</tool>
@@ -0,0 +1,4 @@
FROM python:2.7
EXPOSE 7000
EXPOSE 7001
RUN echo 'moo cow' > /index.html
@@ -0,0 +1,3 @@
all:
docker build . -t galaxy/test-http-example:0.1
@@ -0,0 +1,63 @@
<tool id="interactive_tool_ethercalc" tool_type="interactive" name="EtherCalc" version="0.1">
<requirements>
<container type="docker">shiltemann/ethercalc-galaxy-ie:17.05</container>
</requirements>
<entry_points>
<entry_point name="ethercalc" requires_domain="True">
<port>8000</port>
</entry_point>
</entry_points>
<command><![CDATA[
forever start `which ethercalc` --cors &
## make sure ethercalc has finished starting up
STATUS=$(curl --include 'http://localhost:8000/_/galaxy' 2>&1)
while [[ ${STATUS} =~ "refused" ]]
do
echo "waiting for ethercalc: $STATUS \n"
STATUS=$(curl --include 'http://localhost:8000/_/galaxy' 2>&1)
sleep 2
done
&&
curl --include \
--request POST \
--header "Content-Type: application/json" \
--data-binary "{ \"room\": \"galaxy\", \"snapshot\": \"...\"}" \
'http://localhost:8000/_'
&&
## put loading message into worksheet in case of large files
echo -e "Loading your,dataset\nThis may take, some time for, large files.." > loading.txt
curl --include --request PUT --header "Content-Type: text/csv" --data-binary @loading.txt http://localhost:8000/_/galaxy
&&
## remove dump file so this doesnt appear in audit trail
rm /dump.json
&&
## load dataset into worksheet
curl --include --request PUT \
--header "Content-Type: text/csv" \
--data-binary $infile http://localhost:8000/_/galaxy
&&
tail -f /etc/hosts
]]>
</command>
<inputs>
<param name="infile" type="data" format="tabular,csv,tsv" label="Some tabular dataset"/>
</inputs>
<outputs>
<data name="out_file1" format="txt" />
</outputs>
<tests>
</tests>
<help>
EtherCalc is a web spreadsheet.
https://ethercalc.net
</help>
</tool>
@@ -164,6 +164,9 @@
<tool file="multiple_versions_v01.xml" />
<tool file="multiple_versions_v02.xml" />
<tool file="realtimetool_simple.xml" />
<tool file="realtimetool_two_entry_points.xml" />
<!-- Tools interesting only for building up test workflows. -->
<!-- Next three tools demonstrate concatenating multiple datasets
@@ -31,6 +31,15 @@ class MulledJobTestCases(object):
assert "0.7.15-r1140" in output
class ContainerizedIntegrationTestCase(integration_util.IntegrationTestCase):
@classmethod
def setUpClass(cls):
if not which(cls.container_type):
raise unittest.SkipTest("Executable '%s' not found on PATH" % cls.container_type)
super(ContainerizedIntegrationTestCase, cls).setUpClass()
class DockerizedJobsIntegrationTestCase(integration_util.IntegrationTestCase, RunsEnvironmentJobs, MulledJobTestCases):
framework_tool_and_types = True
+104
View File
@@ -0,0 +1,104 @@
"""Integration tests for realtime tools."""
import os
import requests
from base import api_asserts
from base.populators import (
DatasetPopulator,
wait_on,
)
from .test_containerized_jobs import ContainerizedIntegrationTestCase
SCRIPT_DIRECTORY = os.path.abspath(os.path.dirname(__file__))
class RealtimeToolsIntegrationTestCase(ContainerizedIntegrationTestCase):
framework_tool_and_types = True
container_type = "docker"
require_uwsgi = True
enable_realtime_mapping = True
def setUp(self):
super(RealtimeToolsIntegrationTestCase, self).setUp()
self.dataset_populator = DatasetPopulator(self.galaxy_interactor)
self.history_id = self.dataset_populator.new_history()
def test_simple_execution(self):
response_dict = self.dataset_populator.run_tool("realtimetool_simple", {}, self.history_id, assert_ok=True)
assert "jobs" in response_dict, response_dict
jobs = response_dict["jobs"]
assert isinstance(jobs, list)
assert len(jobs) == 1
job0 = jobs[0]
entry_points = self.wait_on_entry_points_active(job0["id"])
assert len(entry_points) == 1
entry_point0 = entry_points[0]
target = self.entry_point_target(entry_point0["id"])
content = self.wait_on_proxied_content(target)
assert content == "moo cow\n", content
def test_multi_server_realtime_tool(self):
response_dict = self.dataset_populator.run_tool("realtimetool_two_entry_points", {}, self.history_id, assert_ok=True)
assert "jobs" in response_dict, response_dict
jobs = response_dict["jobs"]
assert isinstance(jobs, list)
assert len(jobs) == 1
job0 = jobs[0]
entry_points = self.wait_on_entry_points_active(job0["id"])
assert len(entry_points) == 2
entry_point0 = entry_points[0]
entry_point1 = entry_points[1]
target0 = self.entry_point_target(entry_point0["id"])
target1 = self.entry_point_target(entry_point1["id"])
assert target0 != target1
content0 = self.wait_on_proxied_content(target0)
assert content0 == "moo cow\n", content0
content1 = self.wait_on_proxied_content(target1)
assert content1 == "moo cow\n", content1
assert False
def wait_on_proxied_content(self, target):
def get_hosted_content():
try:
scheme, rest = target.split("://", 1)
prefix, host_and_port = rest.split(".realtime.")
print(rest)
faked_host = rest
if "/" in rest:
faked_host = rest.split("/", 1)[0]
response = requests.get("%s://%s" % (scheme, host_and_port), timeout=1, headers={"Host": faked_host})
return response.content
except Exception as e:
print(e)
return None
content = wait_on(get_hosted_content, "realtime hosted content at %s" % target)
return content
def entry_point_target(self, entry_point_id):
entry_point_access_response = self._get("entry_points/%s/access" % entry_point_id)
api_asserts.assert_status_code_is(entry_point_access_response, 200)
access_json = entry_point_access_response.json()
api_asserts.assert_has_key(access_json, "target")
return access_json["target"]
def wait_on_entry_points_active(self, job_id, expected_num=1):
def active_entry_points():
entry_points = self.entry_points_for_job(job_id)
if len(entry_points) != expected_num:
return None
elif any([not e["active"] for e in entry_points]):
return None
else:
return entry_points
return wait_on(active_entry_points, "entry points to become active")
def entry_points_for_job(self, job_id):
entry_points_response = self._get("entry_points?job_id=%s" % job_id)
api_asserts.assert_status_code_is(entry_points_response, 200)
return entry_points_response.json()
+3
View File
@@ -183,6 +183,9 @@ class MockJobWrapper(object):
def get_command_line(self):
return self.command_line
def container_monitor_command(self, *args, **kwds):
return None
@property
def requires_setting_metadata(self):
return self.metadata_line is not None
+3
View File
@@ -176,6 +176,9 @@ class MockJobWrapper(object):
def get_command_line(self):
return self.command_line
def container_monitor_command(self, *args, **kwds):
return None
def get_id_tag(self):
return "1"