From 9803dc600cb62235d9a1abad15e5b7b4877be2aa Mon Sep 17 00:00:00 2001 From: E Rasche Date: Fri, 21 Jul 2017 15:10:41 +0000 Subject: [PATCH 01/33] More updates to GRT --- scripts/grt.py | 297 ++++++++++++++++++++++------------------- scripts/grt.yml.sample | 36 ++++- 2 files changed, 187 insertions(+), 146 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index a87ea5eb1de..2e9e92d5ce4 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -1,18 +1,21 @@ #!/usr/bin/env python -"""Script for uploading Galaxy statistics to the Galactic radio telescope. +"""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. """ from __future__ import print_function -import os -import sys -import json -import urllib2 import argparse +import gzip +import json +import os import sqlalchemy as sa +import sys +import time import yaml -import re +import logging + +from collections import defaultdict sys.path.insert(1, os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, 'lib'))) @@ -25,6 +28,14 @@ sample_config = os.path.abspath(os.path.join(os.path.dirname(__file__), 'grt.yml default_config = os.path.abspath(os.path.join(os.path.dirname(__file__), 'grt.yml')) +def dumper(obj): + try: + return obj.toJSON() + except AttributeError: + if obj.__class__.__name__ == 'Decimal': + return str(obj) + + def _init(config): if config.startswith('/'): config = os.path.abspath(config) @@ -34,6 +45,8 @@ def _init(config): properties = load_app_properties(ini_file=config) 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") return ( mapping.init( @@ -43,184 +56,188 @@ def _init(config): object_store=object_store ), object_store, - config.database_connection.split(':')[0] + config.database_connection.split(':')[0], + config ) -def _sanitize_dict(unsanitized_dict): - sanitized_dict = dict() +def kw_metrics(job): + return { + '%s_%s' % (metric.plugin, metric.metric_name): metric.metric_value + for metric in job.metrics + } - for key in unsanitized_dict: - if key == 'values' and type(unsanitized_dict[key]) is list: - sanitized_dict[key] = None + +class Sanitization: + + def __init__(self, sanitization_config): + self.sanitization_config = sanitization_config + + if 'tool_params' not in self.sanitization_config: + self.sanitization_config['tool_params'] = {} + + if '__any__' not in self.sanitization_config['tool_params']: + self.sanitization_config['tool_params']['__any__'] = [] + + def blacklisted_tree(self, path): + if path.lstrip('.') in self.sanitization_config['tool_params']['__any__']: + return True + elif self.tool_id in self.sanitization_config['tool_params']: + if path.lstrip('.') in self.sanitization_config['tool_params'][self.tool_id]: + return True + return False + + def sanitize_data(self, tool_id, data): + self.tool_id = tool_id + return self._sanitize_value(data) + + def _sanitize_dict(self, unsanitized_dict, path=""): + 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: - 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 + logging.debug("%s> Sanitizing %s = %s" % (' ' * path.count('.'), path, unsanitized_value)) + return unsanitized_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 = 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('-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:]) + args = parser.parse_args() + logging.basicConfig(level=getattr(logging, args.loglevel.upper())) - print('Loading GRT ini...') + logging.info('Loading GRT ini...') try: - with open(args.config) as f: - config_dict = yaml.load(f) + with open(args.config) as handle: + config = yaml.load(handle) except Exception: - with open(sample_config) as f: - config_dict = yaml.load(f) + with open(sample_config) as handle: + config = yaml.load(handle) - # set to 0 by default - if 'last_job_id_sent' not in config_dict: - config_dict['last_job_id_sent'] = 0 + REPORT_DIR = args.report_directory + CHECK_POINT_FILE = os.path.join(REPORT_DIR, '.checkpoint') + ARCHIVE_DIR = os.path.join(REPORT_DIR, 'archives') + METADATA_FILE = os.path.join(REPORT_DIR, 'meta.json') + REPORT_IDENTIFIER = str(time.time()) + REPORT_BASE = os.path.join(ARCHIVE_DIR, REPORT_IDENTIFIER) - 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 + 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 - print('Loading Galaxy...') - model, object_store, engine = _init(config_dict['galaxy_config']) + logging.info('Loading Galaxy...') + model, object_store, engine, gxconfig = _init(config['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'] + model.Job.table.c.id > last_job_sent ))\ .all() # Set up our arrays - active_users = [] - grt_tool_data = [] + active_users = defaultdict(int) grt_jobs_data = [] - def kw_metrics(job): - return { - '%s_%s' % (metric.plugin, metric.metric_name): metric.metric_value - for metric in job.metrics - } + san = Sanitization(config['blacklist']) # For every job for job in jobs: - if job.tool_id in config_dict['tool_blacklist']: + if job.tool_id in config['blacklist'].get('tools', []): 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)) + # If the user has run a job, they're active. + active_users[job.user_id] += 1 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) - } + logging.debug("Sanitizing %s %s" % (job.tool_id, str(params))) + job_data = ( + str(job.id), + job.tool_id, + job.tool_version, + job.update_time.strftime('%s'), + json.dumps(metrics, default=dumper), + json.dumps(san.sanitize_data(job.tool_id, params)) + ) grt_jobs_data.append(job_data) + # Remember the last job sent. 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)) + last_job_sent = jobs[-1].id else: - try: - urllib2.urlopen(config_dict['grt_url'], data=json.dumps(grt_report_data)) - except urllib2.HTTPError as htpe: - print(htpe.read()) - exit(1) + logging.info("No new jobs to report") - # 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) + # Now on to outputs. + if not os.path.exists(REPORT_DIR): + os.makedirs(REPORT_DIR) + os.makedirs(ARCHIVE_DIR) + + if os.path.exists(REPORT_DIR) and not os.path.exists(ARCHIVE_DIR): + os.makedirs(ARCHIVE_DIR) + + with open(METADATA_FILE, 'w') as handle: + json.dump(config['grt']['metadata'], handle, indent=2) + + # Now serialize the individual report data. + with open(REPORT_BASE + '.json', 'w') as handle: + json.dump({ + "version": 1, + "generated": REPORT_IDENTIFIER, + "metrics": { + }, + "users": { + "active": len(set(active_users)), + "total": sa_session.query(model.User).count(), + }, + "jobs": { + "ok": len(jobs), + }, + "tools": [ + ] + }, handle) + + with gzip.open(REPORT_BASE + '.tsv.gz', 'w') as handle: + for job in grt_jobs_data: + handle.write('\t'.join(job)) + handle.write('\n') + + # update our checkpoint + with open(CHECK_POINT_FILE, 'w') as handle: + handle.write(str(last_job_sent)) if __name__ == '__main__': diff --git a/scripts/grt.yml.sample b/scripts/grt.yml.sample index d837681622f..e6ddd298e27 100644 --- a/scripts/grt.yml.sample +++ b/scripts/grt.yml.sample @@ -1,7 +1,31 @@ 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 + +grt: + #instance_id: blah + #api_key: blah + url: https://radio-telescope.galaxyproject.org/api/v1/upload + metadata: + url: https://example.com/galaxy/ + title: NLP Galaxy + description: | + This is our Galaxy instance! + publicly_visible: True + public: False + # Optional Location + latitude: 0.00 + longitude: 0.00 + # Owners (these are the usernames of users registered with the GRT system.) + owners: + - jane.doe + + +blacklist: + # Blacklist the entire tool from appearing + tools: + - __SET_METADATA__ + - upload1 + tool_params: + __any__: + - chromInfo + edu.tamu.cpt.admin.metadataFetch: + - dbkey From 5abbf1102ed93b5d78bb126a571cef8e07a95b4e Mon Sep 17 00:00:00 2001 From: E Rasche Date: Fri, 21 Jul 2017 15:14:11 +0000 Subject: [PATCH 02/33] at loop time --- scripts/grt.py | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index 2e9e92d5ce4..f2a5f375ce6 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -157,12 +157,6 @@ def main(argv): 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 > last_job_sent - ))\ - .all() # Set up our arrays active_users = defaultdict(int) @@ -171,7 +165,12 @@ def main(argv): san = Sanitization(config['blacklist']) # For every job - for job in jobs: + for job in sa_session.query(model.Job)\ + .filter(sa.and_( + model.Job.table.c.state == "ok", + model.Job.table.c.id > last_job_sent + ))\ + .all(): if job.tool_id in config['blacklist'].get('tools', []): continue @@ -216,6 +215,7 @@ def main(argv): with open(REPORT_BASE + '.json', 'w') as handle: json.dump({ "version": 1, + "galaxy_version": gxconfig.version_major, "generated": REPORT_IDENTIFIER, "metrics": { }, From b0428352e07d444b58d428da2928608005647491 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Fri, 21 Jul 2017 15:21:44 +0000 Subject: [PATCH 03/33] Split out the submission portion --- scripts/grt-submit.py | 48 ++++++++++++++++++++++++++++++++++++++++++ scripts/grt.py | 1 + scripts/grt.yml.sample | 4 ++-- 3 files changed, 51 insertions(+), 2 deletions(-) create mode 100644 scripts/grt-submit.py diff --git a/scripts/grt-submit.py b/scripts/grt-submit.py new file mode 100644 index 00000000000..540ef5e25b2 --- /dev/null +++ b/scripts/grt-submit.py @@ -0,0 +1,48 @@ +#!/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 os +import sys +import json +import urllib2 +import argparse +import yaml +import re + +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 ini...') + 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) + + GRT_URL = config['grt']['url'] + GRT_INSTANCE_ID = config['grt']['instance_id'] + GRT_API_KEY = config['grt']['api_key'] + + # TODO: server interaction. + + +if __name__ == '__main__': + main(sys.argv) diff --git a/scripts/grt.py b/scripts/grt.py index f2a5f375ce6..f2141dc24fa 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -136,6 +136,7 @@ def main(argv): 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) diff --git a/scripts/grt.yml.sample b/scripts/grt.yml.sample index e6ddd298e27..36e847a8f55 100644 --- a/scripts/grt.yml.sample +++ b/scripts/grt.yml.sample @@ -1,8 +1,8 @@ galaxy_config: config/galaxy.ini grt: - #instance_id: blah - #api_key: blah + instance_id: blah + api_key: blah url: https://radio-telescope.galaxyproject.org/api/v1/upload metadata: url: https://example.com/galaxy/ From c14d37fbd6717a97ca700c5f4bb507444df570b9 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Fri, 21 Jul 2017 15:39:34 +0000 Subject: [PATCH 04/33] timings --- scripts/grt.py | 25 +++++++++++++++++-------- 1 file changed, 17 insertions(+), 8 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index f2141dc24fa..281b700bb2e 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -131,6 +131,8 @@ def main(argv): args = parser.parse_args() logging.basicConfig(level=getattr(logging, args.loglevel.upper())) + _times = [] + _start_time = time.time() logging.info('Loading GRT ini...') try: with open(args.config) as handle: @@ -139,6 +141,7 @@ def main(argv): logging.info('Using default GRT Configuration') with open(sample_config) as handle: config = yaml.load(handle) + _times.append(('conf_loaded', time.time() - _start_time)) REPORT_DIR = args.report_directory CHECK_POINT_FILE = os.path.join(REPORT_DIR, '.checkpoint') @@ -156,6 +159,7 @@ def main(argv): logging.info('Loading Galaxy...') model, object_store, engine, gxconfig = _init(config['galaxy_config']) sa_session = model.context.current + _times.append(('gx_conf_loaded', time.time() - _start_time)) # Fetch jobs COMPLETED with status OK that have not yet been sent. @@ -171,6 +175,7 @@ def main(argv): model.Job.table.c.state == "ok", model.Job.table.c.id > last_job_sent ))\ + .order_by(model.Job.table.c.id.asc())\ .all(): if job.tool_id in config['blacklist'].get('tools', []): continue @@ -194,10 +199,11 @@ def main(argv): json.dumps(san.sanitize_data(job.tool_id, params)) ) grt_jobs_data.append(job_data) + _times.append(('jobs_parsed', time.time() - _start_time)) # Remember the last job sent. - if len(jobs) > 0: - last_job_sent = jobs[-1].id + if len(grt_jobs_data) > 0: + last_job_sent = job.id else: logging.info("No new jobs to report") @@ -212,6 +218,13 @@ def main(argv): with open(METADATA_FILE, 'w') as handle: json.dump(config['grt']['metadata'], handle, indent=2) + _times.append(('job_meta', time.time() - _start_time)) + with gzip.open(REPORT_BASE + '.tsv.gz', 'w') as handle: + for job in grt_jobs_data: + handle.write('\t'.join(job)) + handle.write('\n') + _times.append(('job_finish', time.time() - _start_time)) + # Now serialize the individual report data. with open(REPORT_BASE + '.json', 'w') as handle: json.dump({ @@ -219,23 +232,19 @@ def main(argv): "galaxy_version": gxconfig.version_major, "generated": REPORT_IDENTIFIER, "metrics": { + "_times": _times, }, "users": { "active": len(set(active_users)), "total": sa_session.query(model.User).count(), }, "jobs": { - "ok": len(jobs), + "ok": len(grt_jobs_data), }, "tools": [ ] }, handle) - with gzip.open(REPORT_BASE + '.tsv.gz', 'w') as handle: - for job in grt_jobs_data: - handle.write('\t'.join(job)) - handle.write('\n') - # update our checkpoint with open(CHECK_POINT_FILE, 'w') as handle: handle.write(str(last_job_sent)) From 163127b4b2c513a810ece57d8deb1851123a9351 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Fri, 21 Jul 2017 15:59:25 +0000 Subject: [PATCH 05/33] Also dump toolbox --- scripts/grt.py | 15 ++++++++++----- 1 file changed, 10 insertions(+), 5 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index 281b700bb2e..0769fbd10da 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -20,6 +20,7 @@ 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 @@ -38,11 +39,11 @@ def dumper(obj): def _init(config): if config.startswith('/'): - config = os.path.abspath(config) + config_file = os.path.abspath(config) else: - config = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, config)) + config_file = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, config)) - properties = load_app_properties(ini_file=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: @@ -57,7 +58,8 @@ def _init(config): ), object_store, config.database_connection.split(':')[0], - config + config, + galaxy.app.UniverseApplication(global_conf={'__file__': config_file, 'here': os.getcwd()}), ) @@ -157,7 +159,8 @@ def main(argv): last_job_sent = -1 logging.info('Loading Galaxy...') - model, object_store, engine, gxconfig = _init(config['galaxy_config']) + model, object_store, engine, gxconfig, app = _init(config['galaxy_config']) + sa_session = model.context.current _times.append(('gx_conf_loaded', time.time() - _start_time)) @@ -242,6 +245,8 @@ def main(argv): "ok": len(grt_jobs_data), }, "tools": [ + (tool.name, tool.version, tool.tool_shed, tool.repository_id, tool.repository_name) + for tool_id, tool in app.toolbox._tools_by_id.items() ] }, handle) From 2cd2c3967a74dacbdd3d920935fe141e65675043 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Fri, 21 Jul 2017 19:06:43 +0000 Subject: [PATCH 06/33] document GRT a bit more --- scripts/grt.yml.sample | 30 +++++++++++++++++++++++++----- 1 file changed, 25 insertions(+), 5 deletions(-) diff --git a/scripts/grt.yml.sample b/scripts/grt.yml.sample index 36e847a8f55..c2570e152fb 100644 --- a/scripts/grt.yml.sample +++ b/scripts/grt.yml.sample @@ -1,20 +1,27 @@ +# Location of your galaxy config file. The radio telescope will need to +# partially initialize a copy of galaxy. galaxy_config: config/galaxy.ini grt: - instance_id: blah - api_key: blah - url: https://radio-telescope.galaxyproject.org/api/v1/upload + # Go to https://telescope.galaxyproject.org to obtain an Instance ID and API key + instance_id: None + api_key: None + # 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: https://telescope.galaxyproject.org/ + # Some metadata about your galaxy instance metadata: url: https://example.com/galaxy/ title: NLP Galaxy description: | This is our Galaxy instance! - publicly_visible: True + # Is the instance open to the public? public: False # Optional Location latitude: 0.00 longitude: 0.00 - # Owners (these are the usernames of users registered with the GRT system.) + # Owners (these are the usernames of users registered with the galactic-radio-telescope system.) owners: - jane.doe @@ -24,8 +31,21 @@ blacklist: 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: + # If you wish to remove a specific parameter from *all* tools, you can + # do this by setting the name under `__all__` __any__: - chromInfo + # Or to blacklist under a specific tool, just specify the ID edu.tamu.cpt.admin.metadataFetch: - 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 From 11ac796f2df407ff2f547324842333591a543daa Mon Sep 17 00:00:00 2001 From: E Rasche Date: Fri, 21 Jul 2017 19:27:44 +0000 Subject: [PATCH 07/33] Include hash for transport --- scripts/grt.py | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/scripts/grt.py b/scripts/grt.py index 0769fbd10da..dbe3569366a 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -10,6 +10,7 @@ import gzip import json import os import sqlalchemy as sa +import subprocess import sys import time import yaml @@ -227,6 +228,10 @@ def main(argv): handle.write('\t'.join(job)) handle.write('\n') _times.append(('job_finish', time.time() - _start_time)) + sha = subprocess.check_output(['sha256sum', REPORT_BASE + '.tsv.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: @@ -234,6 +239,7 @@ def main(argv): "version": 1, "galaxy_version": gxconfig.version_major, "generated": REPORT_IDENTIFIER, + "report_hash": "sha256:" + sha, "metrics": { "_times": _times, }, @@ -245,7 +251,7 @@ def main(argv): "ok": len(grt_jobs_data), }, "tools": [ - (tool.name, tool.version, tool.tool_shed, tool.repository_id, tool.repository_name) + (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() ] }, handle) From d453231226543face2eee1c7d652029450b056fd Mon Sep 17 00:00:00 2001 From: E Rasche Date: Mon, 31 Jul 2017 16:30:31 +0200 Subject: [PATCH 08/33] Various optimisations --- scripts/grt.py | 150 ++++++++++++++++++++++++++++++----------- scripts/grt.yml.sample | 4 ++ 2 files changed, 115 insertions(+), 39 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index dbe3569366a..98b1706cd5b 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -15,6 +15,7 @@ import sys import time import yaml import logging +# logging.getLogger('sqlalchemy.engine').setLevel(logging.INFO) from collections import defaultdict @@ -38,7 +39,7 @@ def dumper(obj): return str(obj) -def _init(config): +def _init(config, need_app=False): if config.startswith('/'): config_file = os.path.abspath(config) else: @@ -50,6 +51,11 @@ def _init(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, @@ -60,7 +66,7 @@ def _init(config): object_store, config.database_connection.split(':')[0], config, - galaxy.app.UniverseApplication(global_conf={'__file__': config_file, 'here': os.getcwd()}), + app ) @@ -73,8 +79,12 @@ def kw_metrics(job): class Sanitization: - def __init__(self, sanitization_config): + 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'] = {} @@ -94,7 +104,37 @@ class Sanitization: self.tool_id = tool_id return self._sanitize_value(data) + 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) @@ -132,7 +172,7 @@ def main(argv): help="Set the logging level", default='warning') args = parser.parse_args() - logging.basicConfig(level=getattr(logging, args.loglevel.upper())) + logging.getLogger().setLevel(getattr(logging, args.loglevel.upper())) _times = [] _start_time = time.time() @@ -160,9 +200,10 @@ def main(argv): last_job_sent = -1 logging.info('Loading Galaxy...') - model, object_store, engine, gxconfig, app = _init(config['galaxy_config']) + model, object_store, engine, gxconfig, app = _init(config['galaxy_config'], need_app=config['grt']['metadata']['share_toolbox']) sa_session = model.context.current + logging.info('Configuration Loaded') _times.append(('gx_conf_loaded', time.time() - _start_time)) # Fetch jobs COMPLETED with status OK that have not yet been sent. @@ -171,39 +212,64 @@ def main(argv): active_users = defaultdict(int) grt_jobs_data = [] - san = Sanitization(config['blacklist']) + logging.info('Building Sanitizer') + _times.append(('san_init', time.time() - _start_time)) + san = Sanitization(config['blacklist'], model, sa_session) + _times.append(('san_fin', time.time() - _start_time)) - # For every job - for job in sa_session.query(model.Job)\ - .filter(sa.and_( - model.Job.table.c.state == "ok", - model.Job.table.c.id > last_job_sent - ))\ - .order_by(model.Job.table.c.id.asc())\ - .all(): - if job.tool_id in config['blacklist'].get('tools', []): - continue + logging.info('Loading Jobs') + _times.append(('job_init', time.time() - _start_time)) - # If the user has run a job, they're active. - active_users[job.user_id] += 1 + # Batch the database queries to improve the performance. + # Get the maximum value. + max_job_id = sa_session.query(model.Job.id) \ + .filter(model.Job.id > last_job_sent) \ + .order_by(model.Job.id.desc()) \ + .first()[0] - metrics = kw_metrics(job) + for selection_start in range(last_job_sent, max_job_id + 1, 1000): + logging.info("Processing %s - %s", selection_start, selection_start + 1000) + _processing_times = [] + # For every job + for job in sa_session.query(model.Job)\ + .filter(sa.and_( + model.Job.table.c.state == "ok", + model.Job.table.c.id > selection_start, + model.Job.table.c.id < selection_start + 1000 + ))\ + .order_by(model.Job.table.c.id.asc())\ + .all(): + if job.id % 100 == 0: + logging.info(str(job.id)) - params = job.raw_param_dict() - for key in params: - params[key] = json.loads(params[key]) + _job_start_time = time.time() + if job.tool_id in config['blacklist'].get('tools', []): + continue - logging.debug("Sanitizing %s %s" % (job.tool_id, str(params))) - job_data = ( - str(job.id), - job.tool_id, - job.tool_version, - job.update_time.strftime('%s'), - json.dumps(metrics, default=dumper), - json.dumps(san.sanitize_data(job.tool_id, params)) - ) - grt_jobs_data.append(job_data) - _times.append(('jobs_parsed', time.time() - _start_time)) + # If the user has run a job, they're active. + active_users[job.user_id] += 1 + + metrics = kw_metrics(job) + + params = job.raw_param_dict() + for key in params: + params[key] = json.loads(params[key]) + + logging.debug("Sanitizing %s %s" % (job.tool_id, str(params))) + job_data = ( + str(job.id), + job.tool_id, + job.tool_version, + job.update_time.strftime('%s'), + json.dumps(metrics, default=dumper), + json.dumps(san.sanitize_data(job.tool_id, params)) + ) + grt_jobs_data.append(job_data) + _processing_times.append(time.time() - _job_start_time) + logging.info('Min %s', min(_processing_times)) + logging.info('Max %s', max(_processing_times)) + logging.info('Avg %s', sum(t / len(_processing_times) for t in _processing_times)) + _times.append(('job_fin', time.time() - _start_time)) # Remember the last job sent. if len(grt_jobs_data) > 0: @@ -219,10 +285,11 @@ def main(argv): if os.path.exists(REPORT_DIR) and not os.path.exists(ARCHIVE_DIR): os.makedirs(ARCHIVE_DIR) + _times.append(('job_meta_start', time.time() - _start_time)) with open(METADATA_FILE, 'w') as handle: json.dump(config['grt']['metadata'], handle, indent=2) + _times.append(('job_meta_end', time.time() - _start_time)) - _times.append(('job_meta', time.time() - _start_time)) with gzip.open(REPORT_BASE + '.tsv.gz', 'w') as handle: for job in grt_jobs_data: handle.write('\t'.join(job)) @@ -235,6 +302,14 @@ def main(argv): # Now serialize the individual report data. with open(REPORT_BASE + '.json', 'w') as handle: + if config['grt']['metadata']['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, @@ -245,15 +320,12 @@ def main(argv): }, "users": { "active": len(set(active_users)), - "total": sa_session.query(model.User).count(), + "total": sa_session.query(model.User.id).count(), }, "jobs": { "ok": len(grt_jobs_data), }, - "tools": [ - (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() - ] + "tools": toolbox }, handle) # update our checkpoint diff --git a/scripts/grt.yml.sample b/scripts/grt.yml.sample index c2570e152fb..fa89a32bafe 100644 --- a/scripts/grt.yml.sample +++ b/scripts/grt.yml.sample @@ -24,6 +24,10 @@ grt: # Owners (these are the usernames of users registered with the galactic-radio-telescope system.) owners: - jane.doe + # 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. + share_toolbox: False blacklist: From a0292a46c7f966d4f1c70a5cc19320f8255d64b1 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Tue, 1 Aug 2017 10:56:03 +0200 Subject: [PATCH 09/33] Revert submission of instance metadata --- scripts/grt-submit.py | 15 +++++++++++---- scripts/grt.py | 10 ++-------- scripts/grt.yml.sample | 28 +++++++--------------------- 3 files changed, 20 insertions(+), 33 deletions(-) diff --git a/scripts/grt-submit.py b/scripts/grt-submit.py index 540ef5e25b2..b1d3a8fc342 100644 --- a/scripts/grt-submit.py +++ b/scripts/grt-submit.py @@ -5,13 +5,13 @@ See doc/source/admin/grt.rst for more detailed usage information. """ from __future__ import print_function +import argparse import os import sys -import json -import urllib2 -import argparse +import time import yaml -import re +import logging +import urllib2 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')) @@ -37,9 +37,16 @@ def main(argv): with open(sample_config) as handle: config = yaml.load(handle) + REPORT_DIR = args.report_directory + CHECK_POINT_FILE = os.path.join(REPORT_DIR, '.checkpoint') + ARCHIVE_DIR = os.path.join(REPORT_DIR, 'archives') + METADATA_FILE = os.path.join(REPORT_DIR, 'meta.json') + REPORT_BASE = os.path.join(ARCHIVE_DIR, REPORT_IDENTIFIER) + 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. # TODO: server interaction. diff --git a/scripts/grt.py b/scripts/grt.py index 98b1706cd5b..147fde47836 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -189,7 +189,6 @@ def main(argv): REPORT_DIR = args.report_directory CHECK_POINT_FILE = os.path.join(REPORT_DIR, '.checkpoint') ARCHIVE_DIR = os.path.join(REPORT_DIR, 'archives') - METADATA_FILE = os.path.join(REPORT_DIR, 'meta.json') REPORT_IDENTIFIER = str(time.time()) REPORT_BASE = os.path.join(ARCHIVE_DIR, REPORT_IDENTIFIER) @@ -200,7 +199,7 @@ def main(argv): last_job_sent = -1 logging.info('Loading Galaxy...') - model, object_store, engine, gxconfig, app = _init(config['galaxy_config'], need_app=config['grt']['metadata']['share_toolbox']) + model, object_store, engine, gxconfig, app = _init(config['galaxy_config'], need_app=config['grt']['share_toolbox']) sa_session = model.context.current logging.info('Configuration Loaded') @@ -285,11 +284,6 @@ def main(argv): if os.path.exists(REPORT_DIR) and not os.path.exists(ARCHIVE_DIR): os.makedirs(ARCHIVE_DIR) - _times.append(('job_meta_start', time.time() - _start_time)) - with open(METADATA_FILE, 'w') as handle: - json.dump(config['grt']['metadata'], handle, indent=2) - _times.append(('job_meta_end', time.time() - _start_time)) - with gzip.open(REPORT_BASE + '.tsv.gz', 'w') as handle: for job in grt_jobs_data: handle.write('\t'.join(job)) @@ -302,7 +296,7 @@ def main(argv): # Now serialize the individual report data. with open(REPORT_BASE + '.json', 'w') as handle: - if config['grt']['metadata']['share_toolbox']: + 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() diff --git a/scripts/grt.yml.sample b/scripts/grt.yml.sample index fa89a32bafe..54ee5d61a73 100644 --- a/scripts/grt.yml.sample +++ b/scripts/grt.yml.sample @@ -4,30 +4,16 @@ galaxy_config: config/galaxy.ini grt: # Go to https://telescope.galaxyproject.org to obtain an Instance ID and API key - instance_id: None - api_key: None + instance_id: 5f126e1d-d216-493e-a67a-3d6f20a0f1ac + api_key: 9e43e2b9-b54a-418c-9712-36e79f4d9183 # 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: https://telescope.galaxyproject.org/ - # Some metadata about your galaxy instance - metadata: - url: https://example.com/galaxy/ - title: NLP Galaxy - description: | - This is our Galaxy instance! - # Is the instance open to the public? - public: False - # Optional Location - latitude: 0.00 - longitude: 0.00 - # Owners (these are the usernames of users registered with the galactic-radio-telescope system.) - owners: - - jane.doe - # 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. - share_toolbox: False + 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. + share_toolbox: False blacklist: From 93eac6dfcf8d2517381c0cff3b10429ae97aacad Mon Sep 17 00:00:00 2001 From: E Rasche Date: Tue, 1 Aug 2017 11:50:38 +0200 Subject: [PATCH 10/33] remove redundant code --- scripts/grt.py | 7 +------ 1 file changed, 1 insertion(+), 6 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index 147fde47836..c25abe7fd5c 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -188,9 +188,8 @@ def main(argv): REPORT_DIR = args.report_directory CHECK_POINT_FILE = os.path.join(REPORT_DIR, '.checkpoint') - ARCHIVE_DIR = os.path.join(REPORT_DIR, 'archives') REPORT_IDENTIFIER = str(time.time()) - REPORT_BASE = os.path.join(ARCHIVE_DIR, REPORT_IDENTIFIER) + REPORT_BASE = os.path.join(REPORT_DIR, REPORT_IDENTIFIER) if os.path.exists(CHECK_POINT_FILE): with open(CHECK_POINT_FILE, 'r') as handle: @@ -279,10 +278,6 @@ def main(argv): # Now on to outputs. if not os.path.exists(REPORT_DIR): os.makedirs(REPORT_DIR) - os.makedirs(ARCHIVE_DIR) - - if os.path.exists(REPORT_DIR) and not os.path.exists(ARCHIVE_DIR): - os.makedirs(ARCHIVE_DIR) with gzip.open(REPORT_BASE + '.tsv.gz', 'w') as handle: for job in grt_jobs_data: From 50156d8e7be4cb7693cd3e92ccd5df682b4799a6 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Tue, 1 Aug 2017 12:03:18 +0200 Subject: [PATCH 11/33] import from grt --- scripts/grt-submit.py | 31 +++++++++++++++++++++++-------- 1 file changed, 23 insertions(+), 8 deletions(-) diff --git a/scripts/grt-submit.py b/scripts/grt-submit.py index b1d3a8fc342..8081926ed10 100644 --- a/scripts/grt-submit.py +++ b/scripts/grt-submit.py @@ -11,7 +11,7 @@ import sys import time import yaml import logging -import urllib2 +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')) @@ -38,17 +38,32 @@ def main(argv): config = yaml.load(handle) REPORT_DIR = args.report_directory - CHECK_POINT_FILE = os.path.join(REPORT_DIR, '.checkpoint') - ARCHIVE_DIR = os.path.join(REPORT_DIR, 'archives') - METADATA_FILE = os.path.join(REPORT_DIR, 'meta.json') - REPORT_BASE = os.path.join(ARCHIVE_DIR, REPORT_IDENTIFIER) - 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. - # TODO: server interaction. + # 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 not report_id in remote_reports: + print("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) + print(r.json()) if __name__ == '__main__': From df33d4f7080bd75a5d2e25ec20469ebc2ac54574 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Tue, 1 Aug 2017 13:50:21 +0200 Subject: [PATCH 12/33] Super optimised GRT --- scripts/grt.py | 150 +++++++++++++++++++++++++++++-------------------- 1 file changed, 90 insertions(+), 60 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index c25abe7fd5c..acc29aa98e8 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -7,6 +7,7 @@ from __future__ import print_function import argparse import gzip +import tarfile import json import os import sqlalchemy as sa @@ -170,13 +171,19 @@ def main(argv): default=default_config) parser.add_argument("-l", "--loglevel", choices=['debug', 'info', 'warning', 'error', 'critical'], help="Set the logging level", default='warning') + # parser.add_argument("-m", "--max-batch", type=int, min=1, help="The maximum number of records to be exported in a single invocation of GRT.") args = parser.parse_args() logging.getLogger().setLevel(getattr(logging, args.loglevel.upper())) _times = [] _start_time = time.time() - logging.info('Loading GRT ini...') + 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 ini...') try: with open(args.config) as handle: config = yaml.load(handle) @@ -184,7 +191,7 @@ def main(argv): logging.info('Using default GRT Configuration') with open(sample_config) as handle: config = yaml.load(handle) - _times.append(('conf_loaded', time.time() - _start_time)) + annotate('init_end') REPORT_DIR = args.report_directory CHECK_POINT_FILE = os.path.join(REPORT_DIR, '.checkpoint') @@ -197,12 +204,10 @@ def main(argv): else: last_job_sent = -1 - logging.info('Loading Galaxy...') + annotate('galaxy_init', 'Loading Galaxy...') model, object_store, engine, gxconfig, app = _init(config['galaxy_config'], need_app=config['grt']['share_toolbox']) - sa_session = model.context.current - logging.info('Configuration Loaded') - _times.append(('gx_conf_loaded', time.time() - _start_time)) + annotate('galaxy_end') # Fetch jobs COMPLETED with status OK that have not yet been sent. @@ -210,64 +215,91 @@ def main(argv): active_users = defaultdict(int) grt_jobs_data = [] - logging.info('Building Sanitizer') - _times.append(('san_init', time.time() - _start_time)) + annotate('san_init', 'Building Sanitizer') san = Sanitization(config['blacklist'], model, sa_session) - _times.append(('san_fin', time.time() - _start_time)) + annotate('san_end') - logging.info('Loading Jobs') - _times.append(('job_init', time.time() - _start_time)) + if not os.path.exists(REPORT_DIR): + os.makedirs(REPORT_DIR) - # Batch the database queries to improve the performance. - # Get the maximum value. - max_job_id = sa_session.query(model.Job.id) \ - .filter(model.Job.id > last_job_sent) \ + + # 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] + annotate('endpoint_end') - for selection_start in range(last_job_sent, max_job_id + 1, 1000): - logging.info("Processing %s - %s", selection_start, selection_start + 1000) - _processing_times = [] - # For every job - for job in sa_session.query(model.Job)\ - .filter(sa.and_( - model.Job.table.c.state == "ok", - model.Job.table.c.id > selection_start, - model.Job.table.c.id < selection_start + 1000 - ))\ - .order_by(model.Job.table.c.id.asc())\ - .all(): - if job.id % 100 == 0: - logging.info(str(job.id)) + annotate('export_jobs_start', 'Exporting Jobs') + handle_job = open(REPORT_BASE + '.jobs.tsv', 'w') + for job in sa_session.query(model.Job.id, model.Job.tool_id, model.Job.state) \ + .filter(model.Job.id > last_job_sent) \ + .filter(model.Job.id <= end_job_id) \ + .all(): - _job_start_time = time.time() - if job.tool_id in config['blacklist'].get('tools', []): - continue + handle_job.write(str(job[0])) + handle_job.write('\t') + handle_job.write(job[1]) + handle_job.write('\n') + handle_job.close() + annotate('export_jobs_end') - # If the user has run a job, they're active. - active_users[job.user_id] += 1 - metrics = kw_metrics(job) + annotate('export_metric_num_start', 'Exporting Metrics (Numeric)') + handle_metric_num = open(REPORT_BASE + '.metric_num.tsv', 'w') + 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 > last_job_sent) \ + .filter(model.JobMetricNumeric.job_id <= end_job_id) \ + .all(): - params = job.raw_param_dict() - for key in params: - params[key] = json.loads(params[key]) + 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') - logging.debug("Sanitizing %s %s" % (job.tool_id, str(params))) - job_data = ( - str(job.id), - job.tool_id, - job.tool_version, - job.update_time.strftime('%s'), - json.dumps(metrics, default=dumper), - json.dumps(san.sanitize_data(job.tool_id, params)) - ) - grt_jobs_data.append(job_data) - _processing_times.append(time.time() - _job_start_time) - logging.info('Min %s', min(_processing_times)) - logging.info('Max %s', max(_processing_times)) - logging.info('Avg %s', sum(t / len(_processing_times) for t in _processing_times)) - _times.append(('job_fin', time.time() - _start_time)) + + annotate('export_metric_txt_start', 'Exporting Metrics (Text)') + handle_metric_txt = open(REPORT_BASE + '.metric_txt.tsv', 'w') + for metric in sa_session.query(model.JobMetricText.job_id, model.JobMetricText.plugin, model.JobMetricText.metric_name, model.JobMetricText.metric_value) \ + .filter(model.JobMetricText.job_id > last_job_sent) \ + .filter(model.JobMetricText.job_id <= end_job_id) \ + .all(): + + handle_metric_txt.write(str(metric[0])) + handle_metric_txt.write('\t') + handle_metric_txt.write(metric[1]) + handle_metric_txt.write('\t') + handle_metric_txt.write(metric[2]) + handle_metric_txt.write('\t') + handle_metric_txt.write(metric[3]) + handle_metric_txt.write('\n') + handle_metric_txt.close() + annotate('export_metric_txt_end') + + + # job, metric_text, metric_num, params + + annotate('export_params_start', 'Export Job Parameters') + handle_params = open(REPORT_BASE + '.params.tsv', 'w') + for param in sa_session.query(model.JobParameter.job_id, model.JobParameter.name, model.JobParameter.value) \ + .filter(model.JobParameter.job_id > last_job_sent) \ + .filter(model.JobParameter.job_id <= end_job_id) \ + .all(): + + handle_params.write(str(param[0])) + handle_params.write('\t') + handle_params.write(param[1]) + handle_params.write('\t') + handle_params.write(param[2]) + handle_params.write('\n') + handle_params.close() + annotate('export_params_end') # Remember the last job sent. if len(grt_jobs_data) > 0: @@ -276,15 +308,13 @@ def main(argv): logging.info("No new jobs to report") # Now on to outputs. - if not os.path.exists(REPORT_DIR): - os.makedirs(REPORT_DIR) - with gzip.open(REPORT_BASE + '.tsv.gz', 'w') as handle: - for job in grt_jobs_data: - handle.write('\t'.join(job)) - handle.write('\n') + with tarfile.open(REPORT_BASE + '.tar.gz', 'w:gz') as handle: + for name in ('jobs', 'metric_num', 'metric_txt', 'params'): + handle.add(REPORT_BASE + '.' + name + '.tsv') + _times.append(('job_finish', time.time() - _start_time)) - sha = subprocess.check_output(['sha256sum', REPORT_BASE + '.tsv.gz']) + 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(' ')] From 07faeaa5ce2645be8e2ce9bfc61f1b8d8658505a Mon Sep 17 00:00:00 2001 From: E Rasche Date: Tue, 1 Aug 2017 14:35:50 +0200 Subject: [PATCH 13/33] Batched sql queries, sanitization --- scripts/grt.py | 146 +++++++++++++++++++++++++++---------------------- 1 file changed, 81 insertions(+), 65 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index acc29aa98e8..acbec5f33bb 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -6,7 +6,6 @@ See doc/source/admin/grt.rst for more detailed usage information. from __future__ import print_function import argparse -import gzip import tarfile import json import os @@ -171,13 +170,15 @@ def main(argv): default=default_config) parser.add_argument("-l", "--loglevel", choices=['debug', 'info', 'warning', 'error', 'critical'], help="Set the logging level", default='warning') - # parser.add_argument("-m", "--max-batch", type=int, min=1, help="The maximum number of records to be exported in a single invocation of GRT.") + parser.add_argument("-b", "--batch-size", type=int, default=1000, + help="Batch size for sql queries") 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) @@ -213,7 +214,7 @@ def main(argv): # Set up our arrays active_users = defaultdict(int) - grt_jobs_data = [] + job_state_data = defaultdict(int) annotate('san_init', 'Building Sanitizer') san = Sanitization(config['blacklist'], model, sa_session) @@ -222,97 +223,114 @@ def main(argv): 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] - annotate('endpoint_end') + 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 = {} annotate('export_jobs_start', 'Exporting Jobs') handle_job = open(REPORT_BASE + '.jobs.tsv', 'w') - for job in sa_session.query(model.Job.id, model.Job.tool_id, model.Job.state) \ - .filter(model.Job.id > last_job_sent) \ - .filter(model.Job.id <= end_job_id) \ - .all(): + 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.tool_id, model.Job.state, model.Job.user_id) \ + .filter(model.Job.id > offset_start) \ + .filter(model.Job.id <= min(end_job_id, offset_start + args.batch_size)) \ + .all(): + + handle_job.write(str(job[0])) + handle_job.write('\t') + handle_job.write(job[1]) + handle_job.write('\t') + handle_job.write(job[2]) + handle_job.write('\n') + # meta counts + job_state_data[job[2]] += 1 + active_users[job[3]] += 1 + job_tool_map[job[0]] = job[1] - handle_job.write(str(job[0])) - handle_job.write('\t') - handle_job.write(job[1]) - handle_job.write('\n') 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') - 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 > last_job_sent) \ - .filter(model.JobMetricNumeric.job_id <= end_job_id) \ - .all(): + 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(): - 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.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_metric_txt_start', 'Exporting Metrics (Text)') handle_metric_txt = open(REPORT_BASE + '.metric_txt.tsv', 'w') - for metric in sa_session.query(model.JobMetricText.job_id, model.JobMetricText.plugin, model.JobMetricText.metric_name, model.JobMetricText.metric_value) \ - .filter(model.JobMetricText.job_id > last_job_sent) \ - .filter(model.JobMetricText.job_id <= end_job_id) \ - .all(): + 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.JobMetricText.job_id, model.JobMetricText.plugin, model.JobMetricText.metric_name, model.JobMetricText.metric_value) \ + .filter(model.JobMetricText.job_id > offset_start) \ + .filter(model.JobMetricText.job_id <= min(end_job_id, offset_start + args.batch_size)) \ + .all(): - handle_metric_txt.write(str(metric[0])) - handle_metric_txt.write('\t') - handle_metric_txt.write(metric[1]) - handle_metric_txt.write('\t') - handle_metric_txt.write(metric[2]) - handle_metric_txt.write('\t') - handle_metric_txt.write(metric[3]) - handle_metric_txt.write('\n') + handle_metric_txt.write(str(metric[0])) + handle_metric_txt.write('\t') + handle_metric_txt.write(metric[1]) + handle_metric_txt.write('\t') + handle_metric_txt.write(metric[2]) + handle_metric_txt.write('\t') + handle_metric_txt.write(metric[3]) + handle_metric_txt.write('\n') handle_metric_txt.close() annotate('export_metric_txt_end') - - # job, metric_text, metric_num, params - annotate('export_params_start', 'Export Job Parameters') handle_params = open(REPORT_BASE + '.params.tsv', 'w') - for param in sa_session.query(model.JobParameter.job_id, model.JobParameter.name, model.JobParameter.value) \ - .filter(model.JobParameter.job_id > last_job_sent) \ - .filter(model.JobParameter.job_id <= end_job_id) \ - .all(): + 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(): - handle_params.write(str(param[0])) - handle_params.write('\t') - handle_params.write(param[1]) - handle_params.write('\t') - handle_params.write(param[2]) - handle_params.write('\n') + unsanitized = {param[1]: json.loads(param[2])} + sanitized = san.sanitize_data(job_tool_map[param[0]], unsanitized) + + 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') - # Remember the last job sent. - if len(grt_jobs_data) > 0: - last_job_sent = job.id - else: - logging.info("No new jobs to report") - # Now on to outputs. - with tarfile.open(REPORT_BASE + '.tar.gz', 'w:gz') as handle: for name in ('jobs', 'metric_num', 'metric_txt', 'params'): handle.add(REPORT_BASE + '.' + name + '.tsv') + for name in ('jobs', 'metric_num', 'metric_txt', '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)) @@ -338,18 +356,16 @@ def main(argv): "_times": _times, }, "users": { - "active": len(set(active_users)), + "active": len(active_users.keys()), "total": sa_session.query(model.User.id).count(), }, - "jobs": { - "ok": len(grt_jobs_data), - }, + "jobs": job_state_data, "tools": toolbox }, handle) - # update our checkpoint + # Write our checkpoint file so we know where to start next time. with open(CHECK_POINT_FILE, 'w') as handle: - handle.write(str(last_job_sent)) + handle.write(str(end_job_id)) if __name__ == '__main__': From 94a8a9a841c4352c8dccfd42844f4fab2c97db7c Mon Sep 17 00:00:00 2001 From: E Rasche Date: Tue, 1 Aug 2017 15:10:36 +0200 Subject: [PATCH 14/33] default to sharing toolbox --- scripts/grt.py | 4 +++- scripts/grt.yml.sample | 2 +- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index acbec5f33bb..af2e5a27ed8 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -52,7 +52,7 @@ def _init(config, need_app=False): 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()}), + app = galaxy.app.UniverseApplication(global_conf={'__file__': config_file, 'here': os.getcwd()}) else: app = None @@ -207,6 +207,8 @@ def main(argv): 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') diff --git a/scripts/grt.yml.sample b/scripts/grt.yml.sample index 54ee5d61a73..7ed7ca9e8db 100644 --- a/scripts/grt.yml.sample +++ b/scripts/grt.yml.sample @@ -13,7 +13,7 @@ grt: # 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. - share_toolbox: False + share_toolbox: True blacklist: From cf4dc28c5fa762c93dcc2fd51d42a7bb684f7ecb Mon Sep 17 00:00:00 2001 From: E Rasche Date: Tue, 1 Aug 2017 15:15:17 +0200 Subject: [PATCH 15/33] Disable sanitization by default --- scripts/grt.py | 10 +++++++--- scripts/grt.yml.sample | 7 ++++++- 2 files changed, 13 insertions(+), 4 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index af2e5a27ed8..1be4f080c95 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -198,6 +198,7 @@ def main(argv): CHECK_POINT_FILE = os.path.join(REPORT_DIR, '.checkpoint') REPORT_IDENTIFIER = str(time.time()) REPORT_BASE = os.path.join(REPORT_DIR, REPORT_IDENTIFIER) + SANITIZATION_ENABLED = config['sanitization']['enabled'] if os.path.exists(CHECK_POINT_FILE): with open(CHECK_POINT_FILE, 'r') as handle: @@ -219,7 +220,7 @@ def main(argv): job_state_data = defaultdict(int) annotate('san_init', 'Building Sanitizer') - san = Sanitization(config['blacklist'], model, sa_session) + san = Sanitization(config['sanitization'], model, sa_session) annotate('san_end') if not os.path.exists(REPORT_DIR): @@ -313,8 +314,11 @@ def main(argv): .filter(model.JobParameter.job_id <= min(end_job_id, offset_start + args.batch_size)) \ .all(): - unsanitized = {param[1]: json.loads(param[2])} - sanitized = san.sanitize_data(job_tool_map[param[0]], unsanitized) + if SANITIZATION_ENABLED: + unsanitized = {param[1]: json.loads(param[2])} + sanitized = san.sanitize_data(job_tool_map[param[0]], unsanitized) + else: + sanitized = param[2] handle_params.write(str(param[0])) handle_params.write('\t') diff --git a/scripts/grt.yml.sample b/scripts/grt.yml.sample index 7ed7ca9e8db..622fb155b46 100644 --- a/scripts/grt.yml.sample +++ b/scripts/grt.yml.sample @@ -16,7 +16,12 @@ grt: share_toolbox: True -blacklist: +sanitization: + # This defaults to disabled as it has a serious performance impact and may + # not be necessary for your instance. Without sanitization we see + # performance on the order of 4k parameters parsed per second. With + # sanitizatoin on, it averages to 100 parameters per second. + enabled: False # Blacklist the entire tool from appearing tools: - __SET_METADATA__ From cefa47c6986429ba37949f49b849e624d226a62e Mon Sep 17 00:00:00 2001 From: E Rasche Date: Tue, 1 Aug 2017 15:38:48 +0200 Subject: [PATCH 16/33] remove keys from testing --- scripts/grt.yml.sample | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/scripts/grt.yml.sample b/scripts/grt.yml.sample index 622fb155b46..d8273cdfe76 100644 --- a/scripts/grt.yml.sample +++ b/scripts/grt.yml.sample @@ -4,8 +4,8 @@ galaxy_config: config/galaxy.ini grt: # Go to https://telescope.galaxyproject.org to obtain an Instance ID and API key - instance_id: 5f126e1d-d216-493e-a67a-3d6f20a0f1ac - api_key: 9e43e2b9-b54a-418c-9712-36e79f4d9183 + 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. From 6caded21251539c721907c281c8cd83b4797391e Mon Sep 17 00:00:00 2001 From: E Rasche Date: Tue, 1 Aug 2017 15:40:34 +0200 Subject: [PATCH 17/33] unused function --- scripts/grt.py | 8 -------- 1 file changed, 8 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index 1be4f080c95..ab9ef0ed0a3 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -31,14 +31,6 @@ sample_config = os.path.abspath(os.path.join(os.path.dirname(__file__), 'grt.yml default_config = os.path.abspath(os.path.join(os.path.dirname(__file__), 'grt.yml')) -def dumper(obj): - try: - return obj.toJSON() - except AttributeError: - if obj.__class__.__name__ == 'Decimal': - return str(obj) - - def _init(config, need_app=False): if config.startswith('/'): config_file = os.path.abspath(config) From 5e6d1cba409f082a3f2e7748d6eb1fa8010c4470 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Wed, 2 Aug 2017 18:12:39 +0200 Subject: [PATCH 18/33] Some simplification which should significantly speed up parsing --- scripts/grt.py | 40 +++++++++++++++++++++++++--------------- scripts/grt.yml.sample | 11 +---------- 2 files changed, 26 insertions(+), 25 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index ab9ef0ed0a3..dda82715181 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -92,9 +92,19 @@ class Sanitization: return True return False - def sanitize_data(self, tool_id, data): + 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' + + # If it isn't in tool_params, return quickly without parsing. + if tool_id not in self.sanitization_config['tool_params']: + return value + + # Slow path. + unsanitized = {key: json.loads(value)} self.tool_id = tool_id - return self._sanitize_value(data) + return json.dumps(self._sanitize_value(unsanitized)) def _file_dict(self, data): key = '{src}-{id}'.format(**data) @@ -190,7 +200,6 @@ def main(argv): CHECK_POINT_FILE = os.path.join(REPORT_DIR, '.checkpoint') REPORT_IDENTIFIER = str(time.time()) REPORT_BASE = os.path.join(REPORT_DIR, REPORT_IDENTIFIER) - SANITIZATION_ENABLED = config['sanitization']['enabled'] if os.path.exists(CHECK_POINT_FILE): with open(CHECK_POINT_FILE, 'r') as handle: @@ -238,21 +247,26 @@ def main(argv): handle_job = open(REPORT_BASE + '.jobs.tsv', 'w') 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.tool_id, model.Job.state, model.Job.user_id) \ + 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(): - handle_job.write(str(job[0])) + # TODO: blacklisted + handle_job.write(str(job[0])) # id handle_job.write('\t') - handle_job.write(job[1]) + handle_job.write(job[2]) # tool_id handle_job.write('\t') - handle_job.write(job[2]) + handle_job.write(job[3]) # tool_version + handle_job.write('\t') + handle_job.write(job[4]) # state + handle_job.write('\t') + handle_job.write(job[5]) # create_time handle_job.write('\n') # meta counts - job_state_data[job[2]] += 1 - active_users[job[3]] += 1 - job_tool_map[job[0]] = job[1] + 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') @@ -306,11 +320,7 @@ def main(argv): .filter(model.JobParameter.job_id <= min(end_job_id, offset_start + args.batch_size)) \ .all(): - if SANITIZATION_ENABLED: - unsanitized = {param[1]: json.loads(param[2])} - sanitized = san.sanitize_data(job_tool_map[param[0]], unsanitized) - else: - sanitized = param[2] + sanitized = san.sanitize_data(job_tool_map[param[0]], param[i], param[2]) handle_params.write(str(param[0])) handle_params.write('\t') diff --git a/scripts/grt.yml.sample b/scripts/grt.yml.sample index d8273cdfe76..62402b5b0f1 100644 --- a/scripts/grt.yml.sample +++ b/scripts/grt.yml.sample @@ -17,11 +17,6 @@ grt: sanitization: - # This defaults to disabled as it has a serious performance impact and may - # not be necessary for your instance. Without sanitization we see - # performance on the order of 4k parameters parsed per second. With - # sanitizatoin on, it averages to 100 parameters per second. - enabled: False # Blacklist the entire tool from appearing tools: - __SET_METADATA__ @@ -29,12 +24,8 @@ sanitization: # Or you can blacklist individual parameters from being submitted, e.g. if # you have API keys as a tool parameter. tool_params: - # If you wish to remove a specific parameter from *all* tools, you can - # do this by setting the name under `__all__` - __any__: - - chromInfo # Or to blacklist under a specific tool, just specify the ID - edu.tamu.cpt.admin.metadataFetch: + 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 From 08aeea3e92158a45a3e88b5da5e8f3b31821e6ba Mon Sep 17 00:00:00 2001 From: E Rasche Date: Fri, 4 Aug 2017 13:05:43 +0200 Subject: [PATCH 19/33] add headers --- scripts/grt.py | 7 +++++-- scripts/grt.yml.sample | 4 +++- 2 files changed, 8 insertions(+), 3 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index dda82715181..898ec574fbb 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -245,6 +245,7 @@ def main(argv): 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) \ @@ -261,7 +262,7 @@ def main(argv): handle_job.write('\t') handle_job.write(job[4]) # state handle_job.write('\t') - handle_job.write(job[5]) # create_time + handle_job.write(str(job[5])) # create_time handle_job.write('\n') # meta counts job_state_data[job[4]] += 1 @@ -273,6 +274,7 @@ def main(argv): 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) \ @@ -313,6 +315,7 @@ def main(argv): 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) \ @@ -320,7 +323,7 @@ def main(argv): .filter(model.JobParameter.job_id <= min(end_job_id, offset_start + args.batch_size)) \ .all(): - sanitized = san.sanitize_data(job_tool_map[param[0]], param[i], param[2]) + sanitized = san.sanitize_data(job_tool_map[param[0]], param[1], param[2]) handle_params.write(str(param[0])) handle_params.write('\t') diff --git a/scripts/grt.yml.sample b/scripts/grt.yml.sample index 62402b5b0f1..bc2cc212aaf 100644 --- a/scripts/grt.yml.sample +++ b/scripts/grt.yml.sample @@ -12,7 +12,9 @@ grt: 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. + # 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. share_toolbox: True From 058c4a74b83897ba36ec548b6be8162976a549c7 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Fri, 4 Aug 2017 13:05:48 +0200 Subject: [PATCH 20/33] remove text metrics --- scripts/grt.py | 24 ++---------------------- 1 file changed, 2 insertions(+), 22 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index 898ec574fbb..9ba71cc8eb7 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -293,26 +293,6 @@ def main(argv): handle_metric_num.close() annotate('export_metric_num_end') - annotate('export_metric_txt_start', 'Exporting Metrics (Text)') - handle_metric_txt = open(REPORT_BASE + '.metric_txt.tsv', 'w') - 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.JobMetricText.job_id, model.JobMetricText.plugin, model.JobMetricText.metric_name, model.JobMetricText.metric_value) \ - .filter(model.JobMetricText.job_id > offset_start) \ - .filter(model.JobMetricText.job_id <= min(end_job_id, offset_start + args.batch_size)) \ - .all(): - - handle_metric_txt.write(str(metric[0])) - handle_metric_txt.write('\t') - handle_metric_txt.write(metric[1]) - handle_metric_txt.write('\t') - handle_metric_txt.write(metric[2]) - handle_metric_txt.write('\t') - handle_metric_txt.write(metric[3]) - handle_metric_txt.write('\n') - handle_metric_txt.close() - annotate('export_metric_txt_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') @@ -336,10 +316,10 @@ def main(argv): # Now on to outputs. with tarfile.open(REPORT_BASE + '.tar.gz', 'w:gz') as handle: - for name in ('jobs', 'metric_num', 'metric_txt', 'params'): + for name in ('jobs', 'metric_num', 'params'): handle.add(REPORT_BASE + '.' + name + '.tsv') - for name in ('jobs', 'metric_num', 'metric_txt', 'params'): + for name in ('jobs', 'metric_num', 'params'): os.unlink(REPORT_BASE + '.' + name + '.tsv') _times.append(('job_finish', time.time() - _start_time)) From 16f17efe497e39bcd5fc1f493b7e5bb810a437b2 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Fri, 4 Aug 2017 13:07:19 +0200 Subject: [PATCH 21/33] remove remenants of __any__ --- scripts/grt.py | 7 +------ 1 file changed, 1 insertion(+), 6 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index 9ba71cc8eb7..9c8181bcb3a 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -81,13 +81,8 @@ class Sanitization: if 'tool_params' not in self.sanitization_config: self.sanitization_config['tool_params'] = {} - if '__any__' not in self.sanitization_config['tool_params']: - self.sanitization_config['tool_params']['__any__'] = [] - def blacklisted_tree(self, path): - if path.lstrip('.') in self.sanitization_config['tool_params']['__any__']: - return True - elif self.tool_id in self.sanitization_config['tool_params']: + if self.tool_id in self.sanitization_config['tool_params']: if path.lstrip('.') in self.sanitization_config['tool_params'][self.tool_id]: return True return False From 70f4bc957e9a29b85f41e35cf8530122c8284356 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Fri, 4 Aug 2017 13:45:01 +0200 Subject: [PATCH 22/33] update documentation and fix blacklisting bug --- doc/source/admin/special_topics/grt.rst | 116 +++++++++++++++++++----- scripts/grt.py | 24 ++++- 2 files changed, 115 insertions(+), 25 deletions(-) diff --git a/doc/source/admin/special_topics/grt.rst b/doc/source/admin/special_topics/grt.rst index cde12d76ba6..ddf5d7aaf31 100644 --- a/doc/source/admin/special_topics/grt.rst +++ b/doc/source/admin/special_topics/grt.rst @@ -11,31 +11,105 @@ 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.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 faesible. -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.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.py -l debug + + +``grt.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/``. (This script could really just be a +set of SQL queries, but it has been written in python to be database agnostic.) +Once the files have been exported, they are put in a compresesd 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-submit.py + +``scripts/grt-submit.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 index 9c8181bcb3a..0511a886c39 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -82,20 +82,29 @@ class Sanitization: self.sanitization_config['tool_params'] = {} def blacklisted_tree(self, path): - if self.tool_id in self.sanitization_config['tool_params']: - if path.lstrip('.') in self.sanitization_config['tool_params'][self.tool_id]: - return True + 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, return quickly without parsing. + # 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 @@ -237,6 +246,7 @@ def main(argv): # 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') @@ -276,6 +286,9 @@ def main(argv): .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') @@ -297,6 +310,9 @@ def main(argv): .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]) From 5ed50dba1c4478ab58b9c96e1077652675fe9b2f Mon Sep 17 00:00:00 2001 From: E Rasche Date: Fri, 4 Aug 2017 16:03:54 +0200 Subject: [PATCH 23/33] Allow sending subsets --- scripts/grt.py | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/scripts/grt.py b/scripts/grt.py index 0511a886c39..3a8978cc2e5 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -178,6 +178,8 @@ def main(argv): 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())) @@ -244,6 +246,11 @@ def main(argv): # So we can just quit now. sys.exit(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 + # Unfortunately we have to keep this mapping for the sanitizer to work properly. job_tool_map = {} blacklisted_tools = config['sanitization']['tools'] From 7deff206226ce90b68e2e306146a1ea59c745073 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Fri, 4 Aug 2017 16:17:24 +0200 Subject: [PATCH 24/33] move up to fix print statement --- scripts/grt.py | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/scripts/grt.py b/scripts/grt.py index 3a8978cc2e5..af0f81e899d 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -238,6 +238,12 @@ def main(argv): 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. @@ -246,10 +252,6 @@ def main(argv): # So we can just quit now. sys.exit(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 # Unfortunately we have to keep this mapping for the sanitizer to work properly. job_tool_map = {} From b381a90fcb2b4e169a133f3ff0667b4bf93375df Mon Sep 17 00:00:00 2001 From: E Rasche Date: Fri, 4 Aug 2017 17:00:47 +0200 Subject: [PATCH 25/33] flake8 --- scripts/grt-submit.py | 3 +-- scripts/grt.py | 15 +++++++-------- 2 files changed, 8 insertions(+), 10 deletions(-) diff --git a/scripts/grt-submit.py b/scripts/grt-submit.py index 8081926ed10..92d4a28a691 100644 --- a/scripts/grt-submit.py +++ b/scripts/grt-submit.py @@ -8,7 +8,6 @@ from __future__ import print_function import argparse import os import sys -import time import yaml import logging import requests @@ -53,7 +52,7 @@ def main(argv): # 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 not report_id in remote_reports: + if report_id not in remote_reports: print("Uploading %s" % report_id) files = { 'meta': open(os.path.join(sys.argv[1], report_id + '.json'), 'rb'), diff --git a/scripts/grt.py b/scripts/grt.py index af0f81e899d..86235bf3515 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -117,8 +117,8 @@ class Sanitization: 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() + .filter_by(id=data['id']) \ + .first() if dataset and dataset[1]: data['size'] = int(dataset[1]) else: @@ -132,7 +132,6 @@ class Sanitization: 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 \ @@ -268,15 +267,15 @@ def main(argv): .all(): # TODO: blacklisted - handle_job.write(str(job[0])) # id + handle_job.write(str(job[0])) # id handle_job.write('\t') - handle_job.write(job[2]) # tool_id + handle_job.write(job[2]) # tool_id handle_job.write('\t') - handle_job.write(job[3]) # tool_version + handle_job.write(job[3]) # tool_version handle_job.write('\t') - handle_job.write(job[4]) # state + handle_job.write(job[4]) # state handle_job.write('\t') - handle_job.write(str(job[5])) # create_time + handle_job.write(str(job[5])) # create_time handle_job.write('\n') # meta counts job_state_data[job[4]] += 1 From d3b271761f4d2f998bea3a290e2b35389f968fe9 Mon Sep 17 00:00:00 2001 From: Eric Rasche Date: Fri, 4 Aug 2017 17:18:09 +0200 Subject: [PATCH 26/33] remove a blank line --- scripts/grt.py | 1 - 1 file changed, 1 deletion(-) diff --git a/scripts/grt.py b/scripts/grt.py index 86235bf3515..c16ccb24d94 100644 --- a/scripts/grt.py +++ b/scripts/grt.py @@ -251,7 +251,6 @@ def main(argv): # 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'] From e8b367bde3d429c991d223c3d8df76399d6cff81 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Mon, 7 Aug 2017 10:44:54 +0200 Subject: [PATCH 27/33] move to subdir --- scripts/{ => grt}/grt-submit.py | 0 scripts/{ => grt}/grt.py | 0 scripts/{ => grt}/grt.yml.sample | 0 3 files changed, 0 insertions(+), 0 deletions(-) rename scripts/{ => grt}/grt-submit.py (100%) rename scripts/{ => grt}/grt.py (100%) rename scripts/{ => grt}/grt.yml.sample (100%) diff --git a/scripts/grt-submit.py b/scripts/grt/grt-submit.py similarity index 100% rename from scripts/grt-submit.py rename to scripts/grt/grt-submit.py diff --git a/scripts/grt.py b/scripts/grt/grt.py similarity index 100% rename from scripts/grt.py rename to scripts/grt/grt.py diff --git a/scripts/grt.yml.sample b/scripts/grt/grt.yml.sample similarity index 100% rename from scripts/grt.yml.sample rename to scripts/grt/grt.yml.sample From 0ed3571f0d6285d5415eeffc0a8a8f97724aab77 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Mon, 7 Aug 2017 10:48:10 +0200 Subject: [PATCH 28/33] Fix various logging statements --- scripts/grt/grt-submit.py | 8 ++++---- scripts/grt/grt.py | 7 ++----- 2 files changed, 6 insertions(+), 9 deletions(-) diff --git a/scripts/grt/grt-submit.py b/scripts/grt/grt-submit.py index 92d4a28a691..bc3513337b7 100644 --- a/scripts/grt/grt-submit.py +++ b/scripts/grt/grt-submit.py @@ -27,12 +27,12 @@ def main(argv): help="Set the logging level", default='warning') args = parser.parse_args() - logging.info('Loading GRT ini...') + logging.info('Loading GRT configuration...') try: with open(args.config) as handle: config = yaml.load(handle) except Exception: - logging.info('Using default GRT Configuration') + logging.info('Using default GRT configuration') with open(sample_config) as handle: config = yaml.load(handle) @@ -53,7 +53,7 @@ def main(argv): 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: - print("Uploading %s" % report_id) + 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') @@ -62,7 +62,7 @@ def main(argv): 'identifier': report_id } r = requests.post(GRT_URL + 'api/v2/upload', files=files, headers=headers, data=data) - print(r.json()) + logging.info("Uploaded successfully", report_id) if __name__ == '__main__': diff --git a/scripts/grt/grt.py b/scripts/grt/grt.py index 86235bf3515..f92c7568c06 100644 --- a/scripts/grt/grt.py +++ b/scripts/grt/grt.py @@ -3,8 +3,6 @@ See doc/source/admin/grt.rst for more detailed usage information. """ -from __future__ import print_function - import argparse import tarfile import json @@ -15,7 +13,6 @@ import sys import time import yaml import logging -# logging.getLogger('sqlalchemy.engine').setLevel(logging.INFO) from collections import defaultdict @@ -191,12 +188,12 @@ def main(argv): logging.info(human_label) _times.append((label, time.time() - _start_time)) - annotate('init_start', 'Loading GRT ini...') + 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') + logging.info('Using default GRT configuration') with open(sample_config) as handle: config = yaml.load(handle) annotate('init_end') From 7e94cf20b0bd05e73c50976cb1f5196e776ebb1e Mon Sep 17 00:00:00 2001 From: E Rasche Date: Mon, 7 Aug 2017 10:48:16 +0200 Subject: [PATCH 29/33] fix todo --- scripts/grt/grt.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/scripts/grt/grt.py b/scripts/grt/grt.py index f92c7568c06..8533a5cd899 100644 --- a/scripts/grt/grt.py +++ b/scripts/grt/grt.py @@ -262,8 +262,10 @@ def main(argv): .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 - # TODO: blacklisted handle_job.write(str(job[0])) # id handle_job.write('\t') handle_job.write(job[2]) # tool_id From 1043ee5fd4e8eb559795e16c63f735dbbe9924a6 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Mon, 7 Aug 2017 10:49:09 +0200 Subject: [PATCH 30/33] update docs to reflect new locations --- doc/source/admin/special_topics/grt.rst | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/doc/source/admin/special_topics/grt.rst b/doc/source/admin/special_topics/grt.rst index ddf5d7aaf31..01fc1ef90a3 100644 --- a/doc/source/admin/special_topics/grt.rst +++ b/doc/source/admin/special_topics/grt.rst @@ -18,7 +18,7 @@ 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 +an API key which are used to run ``scripts/grt/grt.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 @@ -32,7 +32,7 @@ Data Privacy 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.yml.sample`` file allows you to build up +built-in sanitization. ``scripts/grt/grt.yml.sample`` file allows you to build up sanitization for the job logs. .. code-block:: yaml @@ -72,7 +72,7 @@ Data Collection Process .. code-block:: console - cd $GALAXY; python scripts/grt.py -l debug + cd $GALAXY; python scripts/grt/grt.py -l debug ``grt.py`` connects to your galaxy database and makes queries against the @@ -99,9 +99,9 @@ Data Submission .. code-block:: console - cd $GALAXY; python scripts/grt-submit.py + cd $GALAXY; python scripts/grt/grt-submit.py -``scripts/grt-submit.py`` is a script which will submit your data to the +``scripts/grt/grt-submit.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. From ddcddecb4a3f049bc35fa7a6832339dc9f66bad9 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Mon, 7 Aug 2017 10:49:56 +0200 Subject: [PATCH 31/33] rename scripts --- doc/source/admin/special_topics/grt.rst | 10 +++++----- scripts/grt/{grt.py => export.py} | 0 scripts/grt/{grt-submit.py => upload.py} | 0 3 files changed, 5 insertions(+), 5 deletions(-) rename scripts/grt/{grt.py => export.py} (100%) rename scripts/grt/{grt-submit.py => upload.py} (100%) diff --git a/doc/source/admin/special_topics/grt.rst b/doc/source/admin/special_topics/grt.rst index 01fc1ef90a3..20292c1d80a 100644 --- a/doc/source/admin/special_topics/grt.rst +++ b/doc/source/admin/special_topics/grt.rst @@ -18,7 +18,7 @@ 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/grt.py``. The tool itself is very +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 @@ -72,10 +72,10 @@ Data Collection Process .. code-block:: console - cd $GALAXY; python scripts/grt/grt.py -l debug + cd $GALAXY; python scripts/grt/export.py -l debug -``grt.py`` connects to your galaxy database and makes queries against the +``export.py`` connects to your galaxy database and makes queries against the database for three primary tables: - job @@ -99,9 +99,9 @@ Data Submission .. code-block:: console - cd $GALAXY; python scripts/grt/grt-submit.py + cd $GALAXY; python scripts/grt/upload.py -``scripts/grt/grt-submit.py`` is a script which will submit your data to the +``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. diff --git a/scripts/grt/grt.py b/scripts/grt/export.py similarity index 100% rename from scripts/grt/grt.py rename to scripts/grt/export.py diff --git a/scripts/grt/grt-submit.py b/scripts/grt/upload.py similarity index 100% rename from scripts/grt/grt-submit.py rename to scripts/grt/upload.py From a6fcb95a15134e9ef42cd39663ebe70dfed37fd1 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Mon, 7 Aug 2017 10:52:39 +0200 Subject: [PATCH 32/33] default false with expl --- scripts/grt/grt.yml.sample | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/scripts/grt/grt.yml.sample b/scripts/grt/grt.yml.sample index bc2cc212aaf..810c693973c 100644 --- a/scripts/grt/grt.yml.sample +++ b/scripts/grt/grt.yml.sample @@ -14,8 +14,8 @@ grt: # 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. - share_toolbox: True + # complete list. WARNING: This may not work on production instances, currently. + share_toolbox: False sanitization: From 89887689b479e34e3d049dd689a5f8c8e3f47fc2 Mon Sep 17 00:00:00 2001 From: E Rasche Date: Mon, 7 Aug 2017 11:46:24 +0200 Subject: [PATCH 33/33] correct some spelling errors --- doc/source/admin/special_topics/grt.rst | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/doc/source/admin/special_topics/grt.rst b/doc/source/admin/special_topics/grt.rst index 20292c1d80a..9427b6e0cf8 100644 --- a/doc/source/admin/special_topics/grt.rst +++ b/doc/source/admin/special_topics/grt.rst @@ -5,7 +5,7 @@ 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 ------------ @@ -24,7 +24,7 @@ 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 faesible. +is feasible. Data Privacy ------------ @@ -83,11 +83,11 @@ database for three primary tables: - job_metric_numeric these are exported with very little processing, as tabular files to the GRT -reports directory, ``$GALAXY/reports/``. (This script could really just be a -set of SQL queries, but it has been written in python to be database agnostic.) -Once the files have been exported, they are put in a compresesd archive, and -some metadata about the export process is written to a json file with the same -name as the report archive. +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.