mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge pull request #4376 from erasche/grt-update-2
Galactic Radio Telescope Update
This commit is contained in:
@@ -5,37 +5,111 @@ This is an opt-in service which Galaxy admins can configure to contribute their
|
||||
job run data back to the community. We hope that by collecting this information
|
||||
we can build accurate models of tool CPU/memory/time requirements. In turn,
|
||||
admins will be able to use this analyzed data to optimize their job
|
||||
distribution across highly heterogenous clusters.
|
||||
distribution across highly heterogeneous clusters.
|
||||
|
||||
Registration
|
||||
------------
|
||||
|
||||
You will need to register your Galaxy instance with the Galactic Radio
|
||||
Telescope (GRT). This can be done `https://radio-telescope.galaxyproject.org
|
||||
<https://radio-telescope.galaxyproject.org>`__.
|
||||
Telescope (GRT). This can be done `https://telescope.galaxyproject.org
|
||||
<https://telescope.galaxyproject.org>`__.
|
||||
|
||||
Submitting Data
|
||||
---------------
|
||||
About the Script
|
||||
----------------
|
||||
|
||||
Once you've registered your Galaxy instance, you'll receive an instance ID and
|
||||
an API key which are used to run ``scripts/grt.py``. The tool itself is very simple
|
||||
to run. It collects the last 7 days (by default) of data from your Galaxy
|
||||
server, and sends them to the GRT for processing and display. Additionally
|
||||
it collects the total number of users, and the number of users who ran
|
||||
jobs in the last N days.
|
||||
an API key which are used to run ``scripts/grt/export.py``. The tool itself is very
|
||||
simple to run. GRT will run and produce a directory of reports that can be
|
||||
synced with the GRT server. Every time it is run, GRT only processes the list
|
||||
of jobs that were run since the last time it was run. On first run, GRT will
|
||||
attempt to export all job data for your instance which may be very slow
|
||||
depending on your instance size. We have attempted to optimize this as much as
|
||||
is feasible.
|
||||
|
||||
Running the tool is simple:
|
||||
Data Privacy
|
||||
------------
|
||||
|
||||
.. code-block:: shell
|
||||
All data submitted to the GRT will be released into the public domain. If there
|
||||
are certain tools you do not want included, or certain parameters you wish to
|
||||
hide (e.g. because they contain API keys), then you can take advantage of the
|
||||
built-in sanitization. ``scripts/grt/grt.yml.sample`` file allows you to build up
|
||||
sanitization for the job logs.
|
||||
|
||||
python scripts/grt.py \
|
||||
<INSTANCE_UUID> \
|
||||
<API_KEY> \
|
||||
-c config/galaxy.ini \
|
||||
--grt-url https://radio-telescope.galaxyproject.org/api/v1/upload/
|
||||
--days 7
|
||||
.. code-block:: yaml
|
||||
|
||||
The only required parameters are the instance ID and API key. As you can see in
|
||||
the example command, the GRT URL is configurable. If you do not wish to
|
||||
participate in the public version of this experiment you can host your own
|
||||
radio telescope to collect Galactic information.
|
||||
sanitization:
|
||||
# Blacklist the entire tool from appearing
|
||||
tools:
|
||||
- __SET_METADATA__
|
||||
- upload1
|
||||
# Or you can blacklist individual parameters from being submitted, e.g. if
|
||||
# you have API keys as a tool parameter.
|
||||
tool_params:
|
||||
# Or to blacklist under a specific tool, just specify the ID
|
||||
some_tool_id:
|
||||
- dbkey
|
||||
# If you need to specify a parameter multiple levels deep, you can
|
||||
# do that as well. Currently we only support blacklisting via the
|
||||
# full path, rather than just a path component. So everything under
|
||||
# `path.to.parameter` will be blacklisted.
|
||||
- path.to.parameter
|
||||
# However you could not do "parameter" and have everything under
|
||||
# `path.to.parameter` be removed.
|
||||
# Repeats are rendered as an *, e.g.: repeat_name.*.values
|
||||
|
||||
To blacklist the results from specific tools appearing in results, just add the
|
||||
tool ID under the ``tools`` list.
|
||||
|
||||
Blacklisting tool parameters is more complex. In a key under the ``tool_params`` key,
|
||||
supply a list of parameters you wish to blacklist. *NB: This will slow down
|
||||
processing of records associated with that tool.* Selecting keys is done
|
||||
identically to writing test cases, except if you have a repeat element, just
|
||||
replace the location of the numeric identifier with ``*``, e.g.
|
||||
``repeat_name.*.some_subkey``
|
||||
|
||||
Data Collection Process
|
||||
-----------------------
|
||||
|
||||
.. code-block:: console
|
||||
|
||||
cd $GALAXY; python scripts/grt/export.py -l debug
|
||||
|
||||
|
||||
``export.py`` connects to your galaxy database and makes queries against the
|
||||
database for three primary tables:
|
||||
|
||||
- job
|
||||
- job_parameter
|
||||
- job_metric_numeric
|
||||
|
||||
these are exported with very little processing, as tabular files to the GRT
|
||||
reports directory, ``$GALAXY/reports/``. We only collect new job data that we
|
||||
have not seen since the previous run. The last-seen job ID is stored in
|
||||
``$GALAXY/reports/.checkpoint``. Once the files have been exported, they are
|
||||
put in a compressed archive, and some metadata about the export process is
|
||||
written to a json file with the same name as the report archive.
|
||||
|
||||
You may wish to inspect these files to be sure that you're comfortable with the
|
||||
information being sent.
|
||||
|
||||
Once you're happy with the data, you can submit it with the GRT submission tool...
|
||||
|
||||
Data Submission
|
||||
---------------
|
||||
|
||||
.. code-block:: console
|
||||
|
||||
cd $GALAXY; python scripts/grt/upload.py
|
||||
|
||||
``scripts/grt/upload.py`` is a script which will submit your data to the
|
||||
configured GRT server. You must first be registered with the server which will
|
||||
also walk you through the setup process.
|
||||
|
||||
With your reports, submitting them is very simple. The script will login to the
|
||||
server and determine which reports the server does not have yet. Then it will
|
||||
begin uploading those.
|
||||
|
||||
For administrators with firewalled galaxies and no internet access, if you are
|
||||
able to exfiltrate your files to somewhere with internet, then you can still
|
||||
take advantage of GRT. Alternatively you can deploy GRT on your own
|
||||
infrastructure if you don't want to share your job logs.
|
||||
|
||||
-227
@@ -1,227 +0,0 @@
|
||||
#!/usr/bin/env python
|
||||
"""Script for uploading Galaxy statistics to the Galactic radio telescope.
|
||||
|
||||
See doc/source/admin/grt.rst for more detailed usage information.
|
||||
"""
|
||||
from __future__ import print_function
|
||||
|
||||
import os
|
||||
import sys
|
||||
import json
|
||||
import urllib2
|
||||
import argparse
|
||||
import sqlalchemy as sa
|
||||
import yaml
|
||||
import re
|
||||
|
||||
sys.path.insert(1, os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, 'lib')))
|
||||
|
||||
from galaxy.util.properties import load_app_properties
|
||||
import galaxy.config
|
||||
from galaxy.objectstore import build_object_store_from_config
|
||||
from galaxy.model import mapping
|
||||
|
||||
sample_config = os.path.abspath(os.path.join(os.path.dirname(__file__), 'grt.yml.sample'))
|
||||
default_config = os.path.abspath(os.path.join(os.path.dirname(__file__), 'grt.yml'))
|
||||
|
||||
|
||||
def _init(config):
|
||||
if config.startswith('/'):
|
||||
config = os.path.abspath(config)
|
||||
else:
|
||||
config = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, config))
|
||||
|
||||
properties = load_app_properties(ini_file=config)
|
||||
config = galaxy.config.Configuration(**properties)
|
||||
object_store = build_object_store_from_config(config)
|
||||
|
||||
return (
|
||||
mapping.init(
|
||||
config.file_path,
|
||||
config.database_connection,
|
||||
create_tables=False,
|
||||
object_store=object_store
|
||||
),
|
||||
object_store,
|
||||
config.database_connection.split(':')[0]
|
||||
)
|
||||
|
||||
|
||||
def _sanitize_dict(unsanitized_dict):
|
||||
sanitized_dict = dict()
|
||||
|
||||
for key in unsanitized_dict:
|
||||
if key == 'values' and type(unsanitized_dict[key]) is list:
|
||||
sanitized_dict[key] = None
|
||||
else:
|
||||
sanitized_dict[key] = _sanitize_value(unsanitized_dict[key])
|
||||
|
||||
if sanitized_dict[key] is None:
|
||||
del sanitized_dict[key]
|
||||
|
||||
if len(sanitized_dict) == 0:
|
||||
return None
|
||||
else:
|
||||
return sanitized_dict
|
||||
|
||||
|
||||
def _sanitize_list(unsanitized_list):
|
||||
sanitized_list = list()
|
||||
|
||||
for key in range(len(unsanitized_list)):
|
||||
sanitized_value = _sanitize_value(unsanitized_list[key])
|
||||
if not None:
|
||||
sanitized_list.append(sanitized_value)
|
||||
|
||||
if len(sanitized_list) == 0:
|
||||
return None
|
||||
else:
|
||||
return sanitized_list
|
||||
|
||||
|
||||
def _sanitize_value(unsanitized_value):
|
||||
sanitized_value = None
|
||||
|
||||
fp_regex = re.compile('^(\/[^\/]+)+$')
|
||||
|
||||
if type(unsanitized_value) is dict:
|
||||
sanitized_value = _sanitize_dict(unsanitized_value)
|
||||
elif type(unsanitized_value) is list:
|
||||
sanitized_value = _sanitize_list(unsanitized_value)
|
||||
else:
|
||||
if fp_regex.match(str(unsanitized_value)):
|
||||
sanitized_value = None
|
||||
else:
|
||||
sanitized_value = unsanitized_value
|
||||
|
||||
return sanitized_value
|
||||
|
||||
|
||||
def main(argv):
|
||||
"""Entry point for GRT statistics collection."""
|
||||
parser = argparse.ArgumentParser()
|
||||
parser.add_argument('instance_id', help='Galactic Radio Telescope Instance ID')
|
||||
parser.add_argument('api_key', help='Galactic Radio Telescope API Key')
|
||||
|
||||
parser.add_argument('-c', '--config', dest='config', help='Path to GRT config file (scripts/grt.ini)', default=default_config)
|
||||
parser.add_argument('--dry-run', dest='dryrun', help='Dry run (show data to be sent, but do not send)', action='store_true', default=False)
|
||||
parser.add_argument('--grt-url', dest='grt_url', help='GRT Server (You can run your own!)')
|
||||
args = parser.parse_args(argv[1:])
|
||||
|
||||
print('Loading GRT ini...')
|
||||
try:
|
||||
with open(args.config) as f:
|
||||
config_dict = yaml.load(f)
|
||||
except Exception:
|
||||
with open(sample_config) as f:
|
||||
config_dict = yaml.load(f)
|
||||
|
||||
# set to 0 by default
|
||||
if 'last_job_id_sent' not in config_dict:
|
||||
config_dict['last_job_id_sent'] = 0
|
||||
|
||||
if args.instance_id:
|
||||
config_dict['instance_id'] = args.instance_id
|
||||
if args.api_key:
|
||||
config_dict['api_key'] = args.api_key
|
||||
if args.grt_url:
|
||||
config_dict['grt_url'] = args.grt_url
|
||||
|
||||
print('Loading Galaxy...')
|
||||
model, object_store, engine = _init(config_dict['galaxy_config'])
|
||||
sa_session = model.context.current
|
||||
|
||||
# Fetch jobs COMPLETED with status OK that have not yet been sent.
|
||||
jobs = sa_session.query(model.Job)\
|
||||
.filter(sa.and_(
|
||||
model.Job.table.c.state == "ok",
|
||||
model.Job.table.c.id > config_dict['last_job_id_sent']
|
||||
))\
|
||||
.all()
|
||||
|
||||
# Set up our arrays
|
||||
active_users = []
|
||||
grt_tool_data = []
|
||||
grt_jobs_data = []
|
||||
|
||||
def kw_metrics(job):
|
||||
return {
|
||||
'%s_%s' % (metric.plugin, metric.metric_name): metric.metric_value
|
||||
for metric in job.metrics
|
||||
}
|
||||
|
||||
# For every job
|
||||
for job in jobs:
|
||||
if job.tool_id in config_dict['tool_blacklist']:
|
||||
continue
|
||||
|
||||
# Append an active user, we'll reduce at the end
|
||||
active_users.append(job.user_id)
|
||||
|
||||
# Find the tool in our normalized tool table.
|
||||
if (job.tool_id, job.tool_version) not in grt_tool_data:
|
||||
grt_tool_idx = len(grt_tool_data)
|
||||
grt_tool_data.append((job.tool_id, job.tool_version))
|
||||
else:
|
||||
grt_tool_idx = grt_tool_data.index((job.tool_id, job.tool_version))
|
||||
|
||||
metrics = kw_metrics(job)
|
||||
|
||||
wanted_metrics = ('core_galaxy_slots', 'core_runtime_seconds')
|
||||
|
||||
grt_metrics = {
|
||||
k: int(metrics.get(k, 0))
|
||||
for k in wanted_metrics
|
||||
}
|
||||
|
||||
params = job.raw_param_dict()
|
||||
for key in params:
|
||||
params[key] = json.loads(params[key])
|
||||
|
||||
job_data = {
|
||||
'tool': grt_tool_idx,
|
||||
'date': job.update_time.strftime('%s'),
|
||||
'metrics': grt_metrics,
|
||||
'params': _sanitize_dict(params)
|
||||
}
|
||||
grt_jobs_data.append(job_data)
|
||||
|
||||
if len(jobs) > 0:
|
||||
config_dict['last_job_id_sent'] = jobs[-1].id
|
||||
|
||||
grt_report_data = {
|
||||
'meta': {
|
||||
'version': 1,
|
||||
'instance_uuid': config_dict['instance_id'],
|
||||
'instance_api_key': config_dict['api_key'],
|
||||
# We do not record ANYTHING about your users other than count.
|
||||
'active_users': len(set(active_users)),
|
||||
'total_users': sa_session.query(model.User).count(),
|
||||
'recent_jobs': len(jobs),
|
||||
},
|
||||
'tools': [
|
||||
{
|
||||
'tool_id': a,
|
||||
'tool_version': b,
|
||||
}
|
||||
for (a, b) in grt_tool_data
|
||||
],
|
||||
'jobs': grt_jobs_data,
|
||||
}
|
||||
|
||||
if args.dryrun:
|
||||
print(json.dumps(grt_report_data, indent=2))
|
||||
else:
|
||||
try:
|
||||
urllib2.urlopen(config_dict['grt_url'], data=json.dumps(grt_report_data))
|
||||
except urllib2.HTTPError as htpe:
|
||||
print(htpe.read())
|
||||
exit(1)
|
||||
|
||||
# Update grt.ini with last id of job (prevent duplicates from being sent)
|
||||
with open(args.config, 'w') as f:
|
||||
yaml.dump(config_dict, f, default_flow_style=False)
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
main(sys.argv)
|
||||
@@ -1,7 +0,0 @@
|
||||
galaxy_config: config/galaxy.ini
|
||||
#instance_id: blah
|
||||
#api_key: blah
|
||||
grt_url: https://radio-telescope.galaxyproject.org/api/v1/upload
|
||||
tool_blacklist:
|
||||
- __SET_METADATA__
|
||||
- upload1
|
||||
@@ -0,0 +1,380 @@
|
||||
#!/usr/bin/env python
|
||||
"""Script for parsing Galaxy job information in preparation for submission to the Galactic radio telescope.
|
||||
|
||||
See doc/source/admin/grt.rst for more detailed usage information.
|
||||
"""
|
||||
import argparse
|
||||
import tarfile
|
||||
import json
|
||||
import os
|
||||
import sqlalchemy as sa
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
import yaml
|
||||
import logging
|
||||
|
||||
from collections import defaultdict
|
||||
|
||||
sys.path.insert(1, os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, 'lib')))
|
||||
|
||||
from galaxy.util.properties import load_app_properties
|
||||
import galaxy
|
||||
import galaxy.config
|
||||
from galaxy.objectstore import build_object_store_from_config
|
||||
from galaxy.model import mapping
|
||||
|
||||
sample_config = os.path.abspath(os.path.join(os.path.dirname(__file__), 'grt.yml.sample'))
|
||||
default_config = os.path.abspath(os.path.join(os.path.dirname(__file__), 'grt.yml'))
|
||||
|
||||
|
||||
def _init(config, need_app=False):
|
||||
if config.startswith('/'):
|
||||
config_file = os.path.abspath(config)
|
||||
else:
|
||||
config_file = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, config))
|
||||
|
||||
properties = load_app_properties(ini_file=config_file)
|
||||
config = galaxy.config.Configuration(**properties)
|
||||
object_store = build_object_store_from_config(config)
|
||||
if not config.database_connection:
|
||||
logging.warning("The database connection is empty. If you are using the default value, please uncomment that in your galaxy.ini")
|
||||
|
||||
if need_app:
|
||||
app = galaxy.app.UniverseApplication(global_conf={'__file__': config_file, 'here': os.getcwd()})
|
||||
else:
|
||||
app = None
|
||||
|
||||
return (
|
||||
mapping.init(
|
||||
config.file_path,
|
||||
config.database_connection,
|
||||
create_tables=False,
|
||||
object_store=object_store
|
||||
),
|
||||
object_store,
|
||||
config.database_connection.split(':')[0],
|
||||
config,
|
||||
app
|
||||
)
|
||||
|
||||
|
||||
def kw_metrics(job):
|
||||
return {
|
||||
'%s_%s' % (metric.plugin, metric.metric_name): metric.metric_value
|
||||
for metric in job.metrics
|
||||
}
|
||||
|
||||
|
||||
class Sanitization:
|
||||
|
||||
def __init__(self, sanitization_config, model, sa_session):
|
||||
self.sanitization_config = sanitization_config
|
||||
# SA Stuff
|
||||
self.model = model
|
||||
self.sa_session = sa_session
|
||||
self.filesize_cache = {}
|
||||
|
||||
if 'tool_params' not in self.sanitization_config:
|
||||
self.sanitization_config['tool_params'] = {}
|
||||
|
||||
def blacklisted_tree(self, path):
|
||||
if path.lstrip('.') in self.sanitization_config['tool_params'][self.tool_id]:
|
||||
return True
|
||||
return False
|
||||
|
||||
def sanitize_data(self, tool_id, key, value):
|
||||
# If the tool is blacklisted, skip it.
|
||||
if tool_id in self.sanitization_config['tools']:
|
||||
return 'null'
|
||||
# Thus, all tools below here are not blacklisted at the top level.
|
||||
|
||||
# If it isn't in tool_params, there are no keys being sanitized for
|
||||
# this tool so we can return quickly without parsing.
|
||||
if tool_id not in self.sanitization_config['tool_params']:
|
||||
return value
|
||||
|
||||
# If the key is listed precisely (not a sub-tree), we can also return slightly more quickly.
|
||||
if key in self.sanitization_config['tool_params'][tool_id]:
|
||||
return 'null'
|
||||
|
||||
# If the key isn't a prefix for any of the keys being sanitized, then this is safe.
|
||||
if not any(san_key.startswith(key) for san_key in self.sanitization_config['tool_params'][tool_id]):
|
||||
return value
|
||||
|
||||
# Slow path.
|
||||
unsanitized = {key: json.loads(value)}
|
||||
self.tool_id = tool_id
|
||||
return json.dumps(self._sanitize_value(unsanitized))
|
||||
|
||||
def _file_dict(self, data):
|
||||
key = '{src}-{id}'.format(**data)
|
||||
if key in self.filesize_cache:
|
||||
return self.filesize_cache[data]
|
||||
if data['src'] == 'hda':
|
||||
try:
|
||||
dataset = self.sa_session.query(self.model.Dataset.id, self.model.Dataset.total_size) \
|
||||
.filter_by(id=data['id']) \
|
||||
.first()
|
||||
if dataset and dataset[1]:
|
||||
data['size'] = int(dataset[1])
|
||||
else:
|
||||
data['size'] = None
|
||||
except sa.orm.exc.NoResultFound:
|
||||
data['size'] = None
|
||||
|
||||
# Push to cache for later.
|
||||
self.filesize_cache[data['id']] = data
|
||||
return data
|
||||
else:
|
||||
raise Exception("Cannot handle {src} yet".format(data))
|
||||
|
||||
def _sanitize_dict(self, unsanitized_dict, path=""):
|
||||
# if it is a file dictionary, handle specially.
|
||||
if len(unsanitized_dict.keys()) == 2 and \
|
||||
'id' in unsanitized_dict and \
|
||||
'src' in unsanitized_dict and \
|
||||
unsanitized_dict['src'] in ('hda', 'ldda'):
|
||||
return self._file_dict(unsanitized_dict)
|
||||
|
||||
return {
|
||||
k: self._sanitize_value(v, path=path + '.' + k)
|
||||
for (k, v)
|
||||
in unsanitized_dict.items()
|
||||
}
|
||||
|
||||
def _sanitize_list(self, unsanitized_list, path=""):
|
||||
return [
|
||||
self._sanitize_value(v, path=path + '.*')
|
||||
for v in unsanitized_list
|
||||
]
|
||||
|
||||
def _sanitize_value(self, unsanitized_value, path=""):
|
||||
logging.debug("%sSAN %s" % (' ' * path.count('.'), unsanitized_value))
|
||||
if self.blacklisted_tree(path):
|
||||
logging.debug("%sSAN ***REDACTED***" % (' ' * path.count('.')))
|
||||
return None
|
||||
|
||||
if type(unsanitized_value) is dict:
|
||||
return self._sanitize_dict(unsanitized_value, path=path)
|
||||
elif type(unsanitized_value) is list:
|
||||
return self._sanitize_list(unsanitized_value, path=path)
|
||||
else:
|
||||
logging.debug("%s> Sanitizing %s = %s" % (' ' * path.count('.'), path, unsanitized_value))
|
||||
return unsanitized_value
|
||||
|
||||
|
||||
def main(argv):
|
||||
parser = argparse.ArgumentParser(formatter_class=argparse.ArgumentDefaultsHelpFormatter)
|
||||
parser.add_argument('-r', '--report-directory', help='Directory to store reports in',
|
||||
default=os.path.abspath(os.path.join('.', 'reports')))
|
||||
parser.add_argument('-c', '--config', help='Path to GRT config file',
|
||||
default=default_config)
|
||||
parser.add_argument("-l", "--loglevel", choices=['debug', 'info', 'warning', 'error', 'critical'],
|
||||
help="Set the logging level", default='warning')
|
||||
parser.add_argument("-b", "--batch-size", type=int, default=1000,
|
||||
help="Batch size for sql queries")
|
||||
parser.add_argument("-m", "--max-records", type=int, default=0,
|
||||
help="Maximum number of records to include in a single report. This option should ONLY be used when reporting historical data. Setting this may require running GRT multiple times to capture all historical logs.")
|
||||
|
||||
args = parser.parse_args()
|
||||
logging.getLogger().setLevel(getattr(logging, args.loglevel.upper()))
|
||||
|
||||
_times = []
|
||||
_start_time = time.time()
|
||||
|
||||
def annotate(label, human_label=None):
|
||||
if human_label:
|
||||
logging.info(human_label)
|
||||
_times.append((label, time.time() - _start_time))
|
||||
|
||||
annotate('init_start', 'Loading GRT configuration...')
|
||||
try:
|
||||
with open(args.config) as handle:
|
||||
config = yaml.load(handle)
|
||||
except Exception:
|
||||
logging.info('Using default GRT configuration')
|
||||
with open(sample_config) as handle:
|
||||
config = yaml.load(handle)
|
||||
annotate('init_end')
|
||||
|
||||
REPORT_DIR = args.report_directory
|
||||
CHECK_POINT_FILE = os.path.join(REPORT_DIR, '.checkpoint')
|
||||
REPORT_IDENTIFIER = str(time.time())
|
||||
REPORT_BASE = os.path.join(REPORT_DIR, REPORT_IDENTIFIER)
|
||||
|
||||
if os.path.exists(CHECK_POINT_FILE):
|
||||
with open(CHECK_POINT_FILE, 'r') as handle:
|
||||
last_job_sent = int(handle.read())
|
||||
else:
|
||||
last_job_sent = -1
|
||||
|
||||
annotate('galaxy_init', 'Loading Galaxy...')
|
||||
model, object_store, engine, gxconfig, app = _init(config['galaxy_config'], need_app=config['grt']['share_toolbox'])
|
||||
# Galaxy overrides our logging level.
|
||||
logging.getLogger().setLevel(getattr(logging, args.loglevel.upper()))
|
||||
sa_session = model.context.current
|
||||
annotate('galaxy_end')
|
||||
|
||||
# Fetch jobs COMPLETED with status OK that have not yet been sent.
|
||||
|
||||
# Set up our arrays
|
||||
active_users = defaultdict(int)
|
||||
job_state_data = defaultdict(int)
|
||||
|
||||
annotate('san_init', 'Building Sanitizer')
|
||||
san = Sanitization(config['sanitization'], model, sa_session)
|
||||
annotate('san_end')
|
||||
|
||||
if not os.path.exists(REPORT_DIR):
|
||||
os.makedirs(REPORT_DIR)
|
||||
|
||||
# Pick an end point so our queries can return uniform data.
|
||||
annotate('endpoint_start', 'Identifying a safe endpoint for SQL queries')
|
||||
end_job_id = sa_session.query(model.Job.id) \
|
||||
.order_by(model.Job.id.desc()) \
|
||||
.first()[0]
|
||||
|
||||
# Allow users to only report N records at once.
|
||||
if args.max_records > 0:
|
||||
if end_job_id - last_job_sent > args.max_records:
|
||||
end_job_id = last_job_sent + args.max_records
|
||||
|
||||
annotate('endpoint_end', 'Processing jobs (%s, %s]' % (last_job_sent, end_job_id))
|
||||
|
||||
# Remember the last job sent.
|
||||
if end_job_id == last_job_sent:
|
||||
logging.info("No new jobs to report")
|
||||
# So we can just quit now.
|
||||
sys.exit(0)
|
||||
|
||||
# Unfortunately we have to keep this mapping for the sanitizer to work properly.
|
||||
job_tool_map = {}
|
||||
blacklisted_tools = config['sanitization']['tools']
|
||||
|
||||
annotate('export_jobs_start', 'Exporting Jobs')
|
||||
handle_job = open(REPORT_BASE + '.jobs.tsv', 'w')
|
||||
handle_job.write('\t'.join(('id', 'tool_id', 'tool_version', 'state', 'create_time')) + '\n')
|
||||
for offset_start in range(last_job_sent, end_job_id, args.batch_size):
|
||||
logging.debug("Processing %s:%s", offset_start, min(end_job_id, offset_start + args.batch_size))
|
||||
for job in sa_session.query(model.Job.id, model.Job.user_id, model.Job.tool_id, model.Job.tool_version, model.Job.state, model.Job.create_time) \
|
||||
.filter(model.Job.id > offset_start) \
|
||||
.filter(model.Job.id <= min(end_job_id, offset_start + args.batch_size)) \
|
||||
.all():
|
||||
# If the tool is blacklisted, exclude everywhere
|
||||
if job[2] in blacklisted_tools:
|
||||
continue
|
||||
|
||||
handle_job.write(str(job[0])) # id
|
||||
handle_job.write('\t')
|
||||
handle_job.write(job[2]) # tool_id
|
||||
handle_job.write('\t')
|
||||
handle_job.write(job[3]) # tool_version
|
||||
handle_job.write('\t')
|
||||
handle_job.write(job[4]) # state
|
||||
handle_job.write('\t')
|
||||
handle_job.write(str(job[5])) # create_time
|
||||
handle_job.write('\n')
|
||||
# meta counts
|
||||
job_state_data[job[4]] += 1
|
||||
active_users[job[1]] += 1
|
||||
job_tool_map[job[0]] = job[2]
|
||||
|
||||
handle_job.close()
|
||||
annotate('export_jobs_end')
|
||||
|
||||
annotate('export_metric_num_start', 'Exporting Metrics (Numeric)')
|
||||
handle_metric_num = open(REPORT_BASE + '.metric_num.tsv', 'w')
|
||||
handle_metric_num.write('\t'.join(('job_id', 'plugin', 'name', 'value')) + '\n')
|
||||
for offset_start in range(last_job_sent, end_job_id, args.batch_size):
|
||||
logging.debug("Processing %s:%s", offset_start, min(end_job_id, offset_start + args.batch_size))
|
||||
for metric in sa_session.query(model.JobMetricNumeric.job_id, model.JobMetricNumeric.plugin, model.JobMetricNumeric.metric_name, model.JobMetricNumeric.metric_value) \
|
||||
.filter(model.JobMetricNumeric.job_id > offset_start) \
|
||||
.filter(model.JobMetricNumeric.job_id <= min(end_job_id, offset_start + args.batch_size)) \
|
||||
.all():
|
||||
# If the tool is blacklisted, exclude everywhere
|
||||
if job_tool_map[metric[0]] in blacklisted_tools:
|
||||
continue
|
||||
|
||||
handle_metric_num.write(str(metric[0]))
|
||||
handle_metric_num.write('\t')
|
||||
handle_metric_num.write(metric[1])
|
||||
handle_metric_num.write('\t')
|
||||
handle_metric_num.write(metric[2])
|
||||
handle_metric_num.write('\t')
|
||||
handle_metric_num.write(str(metric[3]))
|
||||
handle_metric_num.write('\n')
|
||||
handle_metric_num.close()
|
||||
annotate('export_metric_num_end')
|
||||
|
||||
annotate('export_params_start', 'Export Job Parameters')
|
||||
handle_params = open(REPORT_BASE + '.params.tsv', 'w')
|
||||
handle_params.write('\t'.join(('job_id', 'name', 'value')) + '\n')
|
||||
for offset_start in range(last_job_sent, end_job_id, args.batch_size):
|
||||
logging.debug("Processing %s:%s", offset_start, min(end_job_id, offset_start + args.batch_size))
|
||||
for param in sa_session.query(model.JobParameter.job_id, model.JobParameter.name, model.JobParameter.value) \
|
||||
.filter(model.JobParameter.job_id > offset_start) \
|
||||
.filter(model.JobParameter.job_id <= min(end_job_id, offset_start + args.batch_size)) \
|
||||
.all():
|
||||
# If the tool is blacklisted, exclude everywhere
|
||||
if job_tool_map[param[0]] in blacklisted_tools:
|
||||
continue
|
||||
|
||||
sanitized = san.sanitize_data(job_tool_map[param[0]], param[1], param[2])
|
||||
|
||||
handle_params.write(str(param[0]))
|
||||
handle_params.write('\t')
|
||||
handle_params.write(param[1])
|
||||
handle_params.write('\t')
|
||||
handle_params.write(json.dumps(sanitized))
|
||||
handle_params.write('\n')
|
||||
handle_params.close()
|
||||
annotate('export_params_end')
|
||||
|
||||
# Now on to outputs.
|
||||
with tarfile.open(REPORT_BASE + '.tar.gz', 'w:gz') as handle:
|
||||
for name in ('jobs', 'metric_num', 'params'):
|
||||
handle.add(REPORT_BASE + '.' + name + '.tsv')
|
||||
|
||||
for name in ('jobs', 'metric_num', 'params'):
|
||||
os.unlink(REPORT_BASE + '.' + name + '.tsv')
|
||||
|
||||
_times.append(('job_finish', time.time() - _start_time))
|
||||
sha = subprocess.check_output(['sha256sum', REPORT_BASE + '.tar.gz'])
|
||||
_times.append(('hash_finish', time.time() - _start_time))
|
||||
# Strip out to space
|
||||
sha = sha[0:sha.index(' ')]
|
||||
|
||||
# Now serialize the individual report data.
|
||||
with open(REPORT_BASE + '.json', 'w') as handle:
|
||||
if config['grt']['share_toolbox']:
|
||||
toolbox = [
|
||||
(tool.id, tool.name, tool.version, tool.tool_shed, tool.repository_id, tool.repository_name)
|
||||
for tool_id, tool in app.toolbox._tools_by_id.items()
|
||||
]
|
||||
else:
|
||||
toolbox = None
|
||||
|
||||
json.dump({
|
||||
"version": 1,
|
||||
"galaxy_version": gxconfig.version_major,
|
||||
"generated": REPORT_IDENTIFIER,
|
||||
"report_hash": "sha256:" + sha,
|
||||
"metrics": {
|
||||
"_times": _times,
|
||||
},
|
||||
"users": {
|
||||
"active": len(active_users.keys()),
|
||||
"total": sa_session.query(model.User.id).count(),
|
||||
},
|
||||
"jobs": job_state_data,
|
||||
"tools": toolbox
|
||||
}, handle)
|
||||
|
||||
# Write our checkpoint file so we know where to start next time.
|
||||
with open(CHECK_POINT_FILE, 'w') as handle:
|
||||
handle.write(str(end_job_id))
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
main(sys.argv)
|
||||
@@ -0,0 +1,39 @@
|
||||
# Location of your galaxy config file. The radio telescope will need to
|
||||
# partially initialize a copy of galaxy.
|
||||
galaxy_config: config/galaxy.ini
|
||||
|
||||
grt:
|
||||
# Go to https://telescope.galaxyproject.org to obtain an Instance ID and API key
|
||||
instance_id:
|
||||
api_key:
|
||||
# Galaxy Project offers a public galactic-radio-telescope instance, however
|
||||
# you are free to run your own if you need. We would love it if you were
|
||||
# willing and able to contribute your data publicly.
|
||||
url: http://localhost:8080
|
||||
# Are you willing to share your toolbox? I.e. what tools are installed.
|
||||
# If your instance is public, this can help us direct users to your
|
||||
# instance. If you elect to not send your entire toolbox, we will still
|
||||
# (necessarily) receive a list of tools that have been run, just not the
|
||||
# complete list. WARNING: This may not work on production instances, currently.
|
||||
share_toolbox: False
|
||||
|
||||
|
||||
sanitization:
|
||||
# Blacklist the entire tool from appearing
|
||||
tools:
|
||||
- __SET_METADATA__
|
||||
- upload1
|
||||
# Or you can blacklist individual parameters from being submitted, e.g. if
|
||||
# you have API keys as a tool parameter.
|
||||
tool_params:
|
||||
# Or to blacklist under a specific tool, just specify the ID
|
||||
some_tool_id:
|
||||
- dbkey
|
||||
# If you need to specify a parameter multiple levels deep, you can
|
||||
# do that as well. Currently we only support blacklisting via the
|
||||
# full path, rather than just a path component. So everything under
|
||||
# `path.to.parameter` will be blacklisted.
|
||||
- path.to.parameter
|
||||
# However you could not do "parameter" and have everything under
|
||||
# `path.to.parameter` be removed.
|
||||
# Repeats are rendered as an *, e.g.: repeat_name.*.values
|
||||
@@ -0,0 +1,69 @@
|
||||
#!/usr/bin/env python
|
||||
"""Script for submitting Galaxy job information to the Galactic radio telescope.
|
||||
|
||||
See doc/source/admin/grt.rst for more detailed usage information.
|
||||
"""
|
||||
from __future__ import print_function
|
||||
|
||||
import argparse
|
||||
import os
|
||||
import sys
|
||||
import yaml
|
||||
import logging
|
||||
import requests
|
||||
|
||||
sample_config = os.path.abspath(os.path.join(os.path.dirname(__file__), 'grt.yml.sample'))
|
||||
default_config = os.path.abspath(os.path.join(os.path.dirname(__file__), 'grt.yml'))
|
||||
|
||||
|
||||
def main(argv):
|
||||
parser = argparse.ArgumentParser()
|
||||
parser = argparse.ArgumentParser(formatter_class=argparse.ArgumentDefaultsHelpFormatter)
|
||||
parser.add_argument('-r', '--report-directory', help='Directory in which reports are stored',
|
||||
default=os.path.abspath(os.path.join('.', 'reports')))
|
||||
parser.add_argument('-c', '--config', help='Path to GRT config file',
|
||||
default=default_config)
|
||||
parser.add_argument("-l", "--loglevel", choices=['debug', 'info', 'warning', 'error', 'critical'],
|
||||
help="Set the logging level", default='warning')
|
||||
args = parser.parse_args()
|
||||
|
||||
logging.info('Loading GRT configuration...')
|
||||
try:
|
||||
with open(args.config) as handle:
|
||||
config = yaml.load(handle)
|
||||
except Exception:
|
||||
logging.info('Using default GRT configuration')
|
||||
with open(sample_config) as handle:
|
||||
config = yaml.load(handle)
|
||||
|
||||
REPORT_DIR = args.report_directory
|
||||
GRT_URL = config['grt']['url']
|
||||
GRT_INSTANCE_ID = config['grt']['instance_id']
|
||||
GRT_API_KEY = config['grt']['api_key']
|
||||
|
||||
# Contact the server and check auth details.
|
||||
headers = {
|
||||
'AUTHORIZATION': GRT_INSTANCE_ID + ':' + GRT_API_KEY
|
||||
}
|
||||
r = requests.post(GRT_URL + 'api/whoami', headers=headers)
|
||||
data = r.json()
|
||||
# we get back some information about which reports had previously been uploaded.
|
||||
remote_reports = data['uploaded_reports']
|
||||
# so now we can know which to send.
|
||||
local_reports = [x.strip('.json') for x in os.listdir(REPORT_DIR) if x.endswith('.json')]
|
||||
for report_id in local_reports:
|
||||
if report_id not in remote_reports:
|
||||
logging.info("Uploading %s", report_id)
|
||||
files = {
|
||||
'meta': open(os.path.join(sys.argv[1], report_id + '.json'), 'rb'),
|
||||
'data': open(os.path.join(sys.argv[1], report_id + '.tsv.gz'), 'rb')
|
||||
}
|
||||
data = {
|
||||
'identifier': report_id
|
||||
}
|
||||
r = requests.post(GRT_URL + 'api/v2/upload', files=files, headers=headers, data=data)
|
||||
logging.info("Uploaded successfully", report_id)
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
main(sys.argv)
|
||||
Reference in New Issue
Block a user