diff --git a/doc/source/admin/special_topics/grt.rst b/doc/source/admin/special_topics/grt.rst index cde12d76ba6..9427b6e0cf8 100644 --- a/doc/source/admin/special_topics/grt.rst +++ b/doc/source/admin/special_topics/grt.rst @@ -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 -`__. +Telescope (GRT). This can be done `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 \ - \ - \ - -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. diff --git a/scripts/grt.py b/scripts/grt.py deleted file mode 100644 index a87ea5eb1de..00000000000 --- a/scripts/grt.py +++ /dev/null @@ -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) diff --git a/scripts/grt.yml.sample b/scripts/grt.yml.sample deleted file mode 100644 index d837681622f..00000000000 --- a/scripts/grt.yml.sample +++ /dev/null @@ -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 diff --git a/scripts/grt/export.py b/scripts/grt/export.py new file mode 100644 index 00000000000..53cace28415 --- /dev/null +++ b/scripts/grt/export.py @@ -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) diff --git a/scripts/grt/grt.yml.sample b/scripts/grt/grt.yml.sample new file mode 100644 index 00000000000..810c693973c --- /dev/null +++ b/scripts/grt/grt.yml.sample @@ -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 diff --git a/scripts/grt/upload.py b/scripts/grt/upload.py new file mode 100644 index 00000000000..bc3513337b7 --- /dev/null +++ b/scripts/grt/upload.py @@ -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)