diff --git a/.github/workflows/integration.yaml b/.github/workflows/integration.yaml new file mode 100644 index 00000000000..3a3e5437c63 --- /dev/null +++ b/.github/workflows/integration.yaml @@ -0,0 +1,48 @@ +name: Integration +on: [push, pull_request] +env: + GALAXY_TEST_DBURI: 'postgres://postgres:postgres@localhost:5432/galaxy?client_encoding=utf8' +jobs: + + test: + name: Test + runs-on: ubuntu-18.04 + strategy: + matrix: + python-version: [3.7] + subset: ['upload_datatype', 'extended_metadata', 'kubernetes', 'not (upload_datatype or extended_metadata or kubernetes)'] + services: + postgres: + image: postgres:11 + env: + POSTGRES_USER: postgres + POSTGRES_PASSWORD: postgres + POSTGRES_DB: postgres + ports: + - 5432:5432 + steps: + - name: Setup Minikube + if: matrix.subset == 'kubernetes' + id: minikube + uses: CodingNagger/minikube-setup-action@v1.0.2 + - name: Launch Minikube + if: matrix.subset == 'kubernetes' + run: eval ${{ steps.minikube.outputs.launcher }} + - name: Check pods + if: matrix.subset == 'kubernetes' + run: | + kubectl get pods + - uses: actions/checkout@v1 + with: + fetch-depth: 1 + - uses: actions/setup-python@v1 + with: + python-version: ${{ matrix.python-version }} + - name: Cache venv dir + uses: actions/cache@v1 + id: cache-venv-galaxy + with: + path: .venv + key: cache-venv-galaxy-${{ matrix.python-version }}-${{ hashFiles('requirements.txt') }} + - name: Run tests + run: './run_tests.sh -integration test/integration -- -k "${{ matrix.subset }}"' diff --git a/lib/galaxy/config/config_manage.py b/lib/galaxy/config/config_manage.py index aacfea0c046..f88eff40452 100644 --- a/lib/galaxy/config/config_manage.py +++ b/lib/galaxy/config/config_manage.py @@ -101,6 +101,12 @@ UWSGI_OPTIONS = OrderedDict([ 'default': '/favicon.ico=static/favicon.ico', 'type': 'str', }), + ('static-safe', { + 'key': 'static-safe', + 'desc': """Allow serving images out of `client`. Most modern Galaxy interfaces bundle all of this, but some older pages still serve these via symlink, requiring this rule.""", + 'default': 'client/galaxy/images', + 'type': 'str', + }), ('master', { 'desc': """Enable the master process manager. Disabled by default for maximum compatibility with CTRL+C, but should be enabled for use with --daemon and/or production deployments.""", 'default': False, diff --git a/lib/galaxy/config/sample/galaxy.yml.sample b/lib/galaxy/config/sample/galaxy.yml.sample index eccf4c67895..9fa6e76623f 100644 --- a/lib/galaxy/config/sample/galaxy.yml.sample +++ b/lib/galaxy/config/sample/galaxy.yml.sample @@ -50,6 +50,11 @@ uwsgi: # Mapping to serve the favicon. static-map: /favicon.ico=static/favicon.ico + # Allow serving images out of `client`. Most modern Galaxy interfaces + # bundle all of this, but some older pages still serve these via + # symlink, requiring this rule. + static-safe: client/galaxy/images + # Enable the master process manager. Disabled by default for maximum # compatibility with CTRL+C, but should be enabled for use with # --daemon and/or production deployments. diff --git a/lib/galaxy/config/sample/reports.yml.sample b/lib/galaxy/config/sample/reports.yml.sample index a0d21cbc3a0..b207fd53bd7 100644 --- a/lib/galaxy/config/sample/reports.yml.sample +++ b/lib/galaxy/config/sample/reports.yml.sample @@ -32,6 +32,11 @@ uwsgi: # Mapping to serve the favicon. static-map: /favicon.ico=static/favicon.ico + # Allow serving images out of `client`. Most modern Galaxy interfaces + # bundle all of this, but some older pages still serve these via + # symlink, requiring this rule. + static-safe: client/galaxy/images + # Enable the master process manager. Disabled by default for maximum # compatibility with CTRL+C, but should be enabled for use with # --daemon and/or production deployments. diff --git a/lib/galaxy/config/sample/tool_shed.yml.sample b/lib/galaxy/config/sample/tool_shed.yml.sample index b8e5b169221..236e4b7430c 100644 --- a/lib/galaxy/config/sample/tool_shed.yml.sample +++ b/lib/galaxy/config/sample/tool_shed.yml.sample @@ -32,6 +32,11 @@ uwsgi: # Mapping to serve the favicon. static-map: /favicon.ico=static/favicon.ico + # Allow serving images out of `client`. Most modern Galaxy interfaces + # bundle all of this, but some older pages still serve these via + # symlink, requiring this rule. + static-safe: client/galaxy/images + # Enable the master process manager. Disabled by default for maximum # compatibility with CTRL+C, but should be enabled for use with # --daemon and/or production deployments. diff --git a/lib/galaxy/config/script.py b/lib/galaxy/config/script.py index f8bcb54cd66..03cfc0cd22f 100644 --- a/lib/galaxy/config/script.py +++ b/lib/galaxy/config/script.py @@ -37,6 +37,7 @@ DEFAULT_DB_CONN = 'sqlite:///./database/universe.sqlite?isolation_level=IMMEDIAT SAMPLES_PATH = os.path.abspath(os.path.join(os.path.dirname(__file__), 'sample')) GALAXY_CONFIG_TEMPLATE_FILE = os.path.join(SAMPLES_PATH, 'galaxy.yml.sample') STATIC_PATH = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, 'web', 'framework', 'static')) +CLIENT_PATH = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir, os.pardir, 'client')) MSG_CONFIG_SUMMARY = """ For help on configuring Galaxy, consult the documentation at: \n {} @@ -55,6 +56,7 @@ GALAXY_CONFIG_SUBSTITUTIONS = { ' static-map: /static/style=static/style/blue': ' static-map: /static=${static_path}/style/blue', ' static-map: /static=static': ' static-map: /static=${static_path}', ' static-map: /favicon.ico=static/favicon.ico': ' static-map: /static=${static_path}/favicon.ico', + ' static-safe: client/galaxy/images': ' ${client_path}/galaxy/images', ' virtualenv: .venv': ' #venv: .venv # not used when running installed', ' pythonpath: lib': ' #pythonpath: lib # not used when running installed', ' #config_dir: false': ' config_dir: ${config_dir}', @@ -147,6 +149,7 @@ def _handle_galaxy_yml(args, config_dir, data_dir): uwsgi_transport=uwsgi_transport, config_dir=config_dir, data_dir=data_dir, + client_dir=CLIENT_PATH, static_path=STATIC_PATH, database_connection=args.db_conn, ) diff --git a/lib/galaxy/dependencies/__init__.py b/lib/galaxy/dependencies/__init__.py index 29691a9c4a6..4d398f0f230 100644 --- a/lib/galaxy/dependencies/__init__.py +++ b/lib/galaxy/dependencies/__init__.py @@ -11,7 +11,10 @@ import pkg_resources import yaml from galaxy.containers import parse_containers_config -from galaxy.util import asbool +from galaxy.util import ( + asbool, + which, +) from galaxy.util.properties import ( find_config_file, load_app_properties @@ -145,7 +148,7 @@ class ConditionalDependencies(object): return "galaxy.jobs.runners.pbs:PBSJobRunner" in self.job_runners def check_pykube(self): - return "galaxy.jobs.runners.kubernetes:KubernetesJobRunner" in self.job_runners + return "galaxy.jobs.runners.kubernetes:KubernetesJobRunner" in self.job_runners or which('kubectl') def check_chronos_python(self): return "galaxy.jobs.runners.chronos:ChronosJobRunner" in self.job_runners @@ -194,7 +197,7 @@ class ConditionalDependencies(object): def optional(config_file=None): if not config_file: - config_file = find_config_file(['galaxy', 'universe_wsgi']) + config_file = find_config_file(['galaxy', 'universe_wsgi'], include_samples=True) if not config_file: print("galaxy.dependencies.optional: no config file found", file=sys.stderr) return [] diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 9bb04908efd..8b9d67523ca 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -2111,7 +2111,7 @@ class JobWrapper(HasResourceParameters): "connection_configuration": container.connection_configuration, }, f) - return "(python '%s'/lib/galaxy_ext/container_monitor/monitor.py &) " % exec_dir + return "(python '%s'/lib/galaxy_ext/container_monitor/monitor.py &); sleep 1 " % exec_dir @property def user(self): diff --git a/lib/galaxy/jobs/runners/cli.py b/lib/galaxy/jobs/runners/cli.py index 2da5643a6ce..6508142674b 100644 --- a/lib/galaxy/jobs/runners/cli.py +++ b/lib/galaxy/jobs/runners/cli.py @@ -161,10 +161,6 @@ class ShellJobRunner(AsynchronousJobRunner): if ajs.job_wrapper.get_state() == model.Job.states.DELETED: continue - external_metadata = not asbool(ajs.job_wrapper.job_destination.params.get("embed_metadata_in_job", DEFAULT_EMBED_METADATA_IN_JOB)) - if external_metadata: - self._handle_metadata_externally(ajs.job_wrapper, resolve_requirements=True) - log.debug("(%s/%s) job not found in batch state check" % (id_tag, external_job_id)) shell_params, job_params = self.parse_destination_params(ajs.job_destination.params) shell, job_interface = self.get_cli_plugins(shell_params, job_params) @@ -187,6 +183,9 @@ class ShellJobRunner(AsynchronousJobRunner): ajs.running = True ajs.old_state = state if state == model.Job.states.OK: + external_metadata = not asbool(ajs.job_wrapper.job_destination.params.get("embed_metadata_in_job", DEFAULT_EMBED_METADATA_IN_JOB)) + if external_metadata: + self.work_queue.put((self.handle_metadata_externally, ajs)) log.debug('(%s/%s) job execution finished, running job wrapper finish method' % (id_tag, external_job_id)) self.work_queue.put((self.finish_job, ajs)) else: @@ -194,6 +193,9 @@ class ShellJobRunner(AsynchronousJobRunner): # Replace the watch list with the updated version self.watched = new_watched + def handle_metadata_externally(self, ajs): + self._handle_metadata_externally(ajs.job_wrapper, resolve_requirements=True) + def __handle_out_of_memory(self, ajs, external_job_id): shell_params, job_params = self.parse_destination_params(ajs.job_destination.params) shell, job_interface = self.get_cli_plugins(shell_params, job_params) diff --git a/lib/galaxy/jobs/runners/util/cli/shell/rsh.py b/lib/galaxy/jobs/runners/util/cli/shell/rsh.py index 2352208229a..c6c7e1fbd1a 100644 --- a/lib/galaxy/jobs/runners/util/cli/shell/rsh.py +++ b/lib/galaxy/jobs/runners/util/cli/shell/rsh.py @@ -6,7 +6,8 @@ from pulsar.managers.util.retry import RetryActionExecutor from galaxy.util import ( smart_str, - unicodify + string_as_bool, + unicodify, ) from galaxy.util.bunch import Bunch from .local import LocalShell @@ -40,11 +41,11 @@ class RemoteShell(LocalShell): class SecureShell(RemoteShell): - SSH_NEW_KEY_STRING = 'Are you sure you want to continue connecting' def __init__(self, rsh='ssh', rcp='scp', private_key=None, port=None, strict_host_key_checking=True, **kwargs): - strict_host_key_checking = "yes" if strict_host_key_checking else "no" - options = ["-o", "StrictHostKeyChecking=%s" % strict_host_key_checking] + options = [] + if not string_as_bool(strict_host_key_checking): + options.extend(["-o", "StrictHostKeyChecking=no", "-o", "UserKnownHostsFile=/dev/null"]) options.extend(["-o", "ConnectTimeout=60"]) if private_key: options.extend(['-i', private_key]) @@ -55,13 +56,14 @@ class SecureShell(RemoteShell): class ParamikoShell(object): - def __init__(self, username, hostname, password=None, private_key=None, port=22, timeout=60, **kwargs): + def __init__(self, username, hostname, password=None, private_key=None, port=22, timeout=60, strict_host_key_checking=True, **kwargs): self.username = username self.hostname = hostname self.password = password self.private_key = private_key self.port = int(port) if port else port self.timeout = int(timeout) if timeout else timeout + self.strict_host_key_checking = string_as_bool(strict_host_key_checking) self.ssh = None self.retry_action_executor = RetryActionExecutor(max_retries=100, interval_max=300) self.connect() @@ -69,7 +71,7 @@ class ParamikoShell(object): def connect(self): log.info("Attempting establishment of new paramiko SSH channel") self.ssh = paramiko.SSHClient() - self.ssh.set_missing_host_key_policy(paramiko.AutoAddPolicy()) + self.ssh.set_missing_host_key_policy(paramiko.RejectPolicy() if self.strict_host_key_checking else paramiko.WarningPolicy()) self.ssh.connect(hostname=self.hostname, port=self.port, username=self.username, diff --git a/lib/galaxy/managers/interactivetool.py b/lib/galaxy/managers/interactivetool.py index dd740f96eb1..a6504c5c587 100644 --- a/lib/galaxy/managers/interactivetool.py +++ b/lib/galaxy/managers/interactivetool.py @@ -270,7 +270,7 @@ class InteractiveToolManager(object): entry_point = trans.sa_session.query(model.InteractiveToolEntryPoint).get(entry_point_id) if self.app.interactivetool_manager.can_access_entry_point(trans, entry_point): if entry_point.active: - return self.target_if_active(entry_point) + return self.target_if_active(trans, entry_point) elif entry_point.deleted: raise exceptions.MessageException('InteractiveTool has ended. You will have to start a new one.') else: diff --git a/lib/galaxy/tool_util/deps/mulled/mulled_build.py b/lib/galaxy/tool_util/deps/mulled/mulled_build.py index 412994c3152..6150f312d3e 100644 --- a/lib/galaxy/tool_util/deps/mulled/mulled_build.py +++ b/lib/galaxy/tool_util/deps/mulled/mulled_build.py @@ -180,7 +180,7 @@ def mull_targets( if DEST_BASE_IMAGE: dest_base_image = DEST_BASE_IMAGE else: - dest_base_image = DEFAULT_EXTENDED_BASE_IMAGE if not any_target_requires_extended_base(targets) else DEST_BASE_IMAGE + dest_base_image = DEFAULT_EXTENDED_BASE_IMAGE if any_target_requires_extended_base(targets) else DEST_BASE_IMAGE targets = list(targets) if involucro_context is None: diff --git a/lib/galaxy/util/__init__.py b/lib/galaxy/util/__init__.py index 1e488ef1c4c..1dcb20b2303 100644 --- a/lib/galaxy/util/__init__.py +++ b/lib/galaxy/util/__init__.py @@ -271,7 +271,7 @@ def xml_to_string(elem, pretty=False): except TypeError as e: # we assume this is a comment if hasattr(elem, 'text'): - return "\n" % elem.text + return u"\n" % elem.text else: raise e if xml_str and pretty: diff --git a/lib/galaxy/web/framework/webapp.py b/lib/galaxy/web/framework/webapp.py index 31dd71bd584..273eb40cc69 100644 --- a/lib/galaxy/web/framework/webapp.py +++ b/lib/galaxy/web/framework/webapp.py @@ -563,6 +563,10 @@ class GalaxyWebTransaction(base.DefaultWebTransaction, return except IndexError: pass + authnz_controller_base = url_for(controller='authnz', action='index') + if self.request.path.startswith(authnz_controller_base): + # All authnz requests pass through + return # redirect to root if the path is not in the list above if self.request.path not in allowed_paths: login_url = url_for(controller='root', action='login', redirect=self.request.path) diff --git a/lib/galaxy_ext/container_monitor/monitor.py b/lib/galaxy_ext/container_monitor/monitor.py index b1e1447f1e0..86987b81e2d 100644 --- a/lib/galaxy_ext/container_monitor/monitor.py +++ b/lib/galaxy_ext/container_monitor/monitor.py @@ -22,7 +22,7 @@ def parse_ports(container_name, connection_configuration): preexec_fn=os.setpgrp) if exit_code == 0: stdout_file.seek(0) - ports_raw = stdout_file.read() + ports_raw = stdout_file.read().decode('utf-8') return ports_raw diff --git a/lib/galaxy_test/driver/driver_util.py b/lib/galaxy_test/driver/driver_util.py index f6334b762ee..1475c8669e1 100644 --- a/lib/galaxy_test/driver/driver_util.py +++ b/lib/galaxy_test/driver/driver_util.py @@ -4,6 +4,7 @@ import fcntl import logging import os import random +import re import shutil import signal import socket @@ -62,13 +63,15 @@ INSTALLED_TOOL_PANEL_CONFIGS = [ ] REALTIME_PROXY_TEMPLATE = string.Template(r""" uwsgi: - realtime_map: $tempdir/realtime_map.sqlite - python-raw: scripts/realtime/key_type_token_mapping.py - route-host: ^([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.(realtime\.$test_host:$test_port)$ goto:realtime + http-raw-body: true + interactivetools_map: $tempdir/interactivetools_map.sqlite + python-raw: scripts/interactivetools/key_type_token_mapping.py + route-host: ^([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)-([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.(interactivetool\.$test_host:$test_port)$ goto:interactivetool route-run: goto:endendend - route-label: realtime - route-host: ^([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.(realtime\.$test_host:$test_port)$ rpcvar:TARGET_HOST rtt_key_type_token_mapper_cached $2 $1 $3 $4 $0 5 + route-label: interactivetool + route-host: ^([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)-([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.([A-Za-z0-9]+(?:-[A-Za-z0-9]+)*)\.(interactivetool\.$test_host:$test_port)$ rpcvar:TARGET_HOST rtt_key_type_token_mapper_cached $1 $3 $2 $4 $0 5 route-if-not: empty:${TARGET_HOST} httpdumb:${TARGET_HOST} + route: .* break:404 Not Found route-label: endendend """) @@ -729,8 +732,9 @@ def launch_uwsgi(kwargs, tempdir, prefix=DEFAULT_CONFIG_PREFIX, config_object=No enable_realtime_mapping = getattr(config_object, "enable_realtime_mapping", False) if enable_realtime_mapping: - config["galaxy"]["realtime_prefix"] = "realtime" - config["galaxy"]["realtime_map"] = os.path.join(tempdir, "realtime_map.sqlite") + config["galaxy"]["interactivetools_prefix"] = "interactivetool" + config["galaxy"]["interactivetools_map"] = os.path.join(tempdir, "interactivetools_map.sqlite") + config['galaxy']['interactivetools_enable'] = True yaml_config_path = os.path.join(tempdir, "galaxy.yml") with open(yaml_config_path, "w") as f: @@ -743,7 +747,7 @@ def launch_uwsgi(kwargs, tempdir, prefix=DEFAULT_CONFIG_PREFIX, config_object=No old_contents = f.read() with open(yaml_config_path, "w") as f: test_port = str(port) if port else r"[0-9]+" - test_host = host or "localhost" + test_host = re.escape(host) if host else "localhost" uwsgi_section = REALTIME_PROXY_TEMPLATE.safe_substitute(test_host=test_host, test_port=test_port, tempdir=tempdir) f.write(uwsgi_section) f.write(old_contents) diff --git a/lib/galaxy_test/driver/integration_util.py b/lib/galaxy_test/driver/integration_util.py index 37a16a33383..59e3ef17b86 100644 --- a/lib/galaxy_test/driver/integration_util.py +++ b/lib/galaxy_test/driver/integration_util.py @@ -31,7 +31,7 @@ def skip_if_jenkins(cls): def skip_unless_executable(executable): if which(executable): return _identity - return skip("PATH doesn't contain executable %s" % executable) + return pytest.mark.skip("PATH doesn't contain executable %s" % executable) def skip_unless_docker(): @@ -47,10 +47,10 @@ def k8s_config_path(): def skip_unless_fixed_port(): - if os.environ.get("GALAXY_TEST_PORT"): + if os.environ.get("GALAXY_TEST_PORT_RANDOM") != "1": return _identity - return skip("GALAXY_TEST_PORT must be set for this test.") + return pytest.mark.skip("GALAXY_TEST_PORT must be set for this test.") class IntegrationInstance(UsesApiTestCaseMixin): diff --git a/run_tests.sh b/run_tests.sh index 78ca1b2e911..77c7871043c 100755 --- a/run_tests.sh +++ b/run_tests.sh @@ -602,6 +602,7 @@ if [ -z "$skip_common_startup" ]; then export GALAXY_CONFIG_OVERRIDE_DATABASE_CONNECTION fi ./scripts/common_startup.sh $skip_venv $no_create_venv $no_replace_pip $replace_pip $skip_client_build --dev-wheels || exit 1 + unset GALAXY_CONFIG_OVERRIDE_DATABASE_CONNECTION fi . ./scripts/common_startup_functions.sh diff --git a/scripts/grt/export.py b/scripts/grt/export.py index 32dc4641ddb..a782e4c8e34 100644 --- a/scripts/grt/export.py +++ b/scripts/grt/export.py @@ -146,150 +146,146 @@ def main(argv): blacklisted_tools = config['sanitization']['tools'] annotate('export_jobs_start', 'Exporting Jobs') - handle_job = io.open(REPORT_BASE + '.jobs.tsv', 'w', encoding='utf-8') - handle_job.write('\t'.join(('id', 'tool_id', 'tool_version', 'state', 'create_time')) + '\n') - for offset_start in range(last_job_sent, end_job_id, args.batch_size): - logging.debug("Processing %s:%s", offset_start, min(end_job_id, offset_start + args.batch_size)) - for job in sa_session.query(model.Job.id, model.Job.user_id, model.Job.tool_id, model.Job.tool_version, model.Job.state, model.Job.create_time) \ - .filter(model.Job.id > offset_start) \ - .filter(model.Job.id <= min(end_job_id, offset_start + args.batch_size)) \ - .all(): - # If the tool is blacklisted, exclude everywhere - if job[2] in blacklisted_tools: - continue + with io.open(REPORT_BASE + '.jobs.tsv', 'w', encoding='utf-8') as handle_job: + handle_job.write(u'\t'.join(('id', 'tool_id', 'tool_version', 'state', 'create_time')) + '\n') + for offset_start in range(last_job_sent, end_job_id, args.batch_size): + logging.debug("Processing %s:%s", offset_start, min(end_job_id, offset_start + args.batch_size)) + for job in sa_session.query(model.Job.id, model.Job.user_id, model.Job.tool_id, model.Job.tool_version, model.Job.state, model.Job.create_time) \ + .filter(model.Job.id > offset_start) \ + .filter(model.Job.id <= min(end_job_id, offset_start + args.batch_size)) \ + .all(): + # If the tool is blacklisted, exclude everywhere + if job[2] in blacklisted_tools: + continue - try: - line = [ - str(job[0]), # id - job[2], # tool_id - job[3], # tool_version - job[4], # state - str(job[5]) # create_time - ] - cline = unicodify('\t'.join(line) + '\n') - handle_job.write(cline) - except Exception: - logging.warning("Unable to write out a 'handle_job' row. Ignoring the row.", exc_info=True) - continue - # meta counts - job_state_data[job[4]] += 1 - active_users[job[1]] += 1 - job_tool_map[job[0]] = job[2] - - handle_job.close() + try: + line = [ + str(job[0]), # id + job[2], # tool_id + job[3], # tool_version + job[4], # state + str(job[5]) # create_time + ] + cline = unicodify('\t'.join(line) + '\n') + handle_job.write(cline) + except Exception: + logging.warning("Unable to write out a 'handle_job' row. Ignoring the row.", exc_info=True) + continue + # meta counts + job_state_data[job[4]] += 1 + active_users[job[1]] += 1 + job_tool_map[job[0]] = job[2] annotate('export_jobs_end') annotate('export_datasets_start', 'Exporting Datasets') - handle_datasets = io.open(REPORT_BASE + '.datasets.tsv', 'w', encoding='utf-8') - handle_datasets.write('\t'.join(('job_id', 'dataset_id', 'extension', 'file_size', 'param_name', 'type')) + '\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)) + with io.open(REPORT_BASE + '.datasets.tsv', 'w', encoding='utf-8') as handle_datasets: + handle_datasets.write(u'\t'.join(('job_id', 'dataset_id', 'extension', 'file_size', 'param_name', 'type')) + '\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)) - # four queries: JobToInputDatasetAssociation, JobToOutputDatasetAssociation, HistoryDatasetAssociation, Dataset + # four queries: JobToInputDatasetAssociation, JobToOutputDatasetAssociation, HistoryDatasetAssociation, Dataset - job_to_input_hda_ids = sa_session.query(model.JobToInputDatasetAssociation.job_id, model.JobToInputDatasetAssociation.dataset_id, - model.JobToInputDatasetAssociation.name) \ - .filter(model.JobToInputDatasetAssociation.job_id > offset_start) \ - .filter(model.JobToInputDatasetAssociation.job_id <= min(end_job_id, offset_start + args.batch_size)) \ - .all() + job_to_input_hda_ids = sa_session.query(model.JobToInputDatasetAssociation.job_id, model.JobToInputDatasetAssociation.dataset_id, + model.JobToInputDatasetAssociation.name) \ + .filter(model.JobToInputDatasetAssociation.job_id > offset_start) \ + .filter(model.JobToInputDatasetAssociation.job_id <= min(end_job_id, offset_start + args.batch_size)) \ + .all() - job_to_output_hda_ids = sa_session.query(model.JobToOutputDatasetAssociation.job_id, model.JobToOutputDatasetAssociation.dataset_id, - model.JobToOutputDatasetAssociation.name) \ - .filter(model.JobToOutputDatasetAssociation.job_id > offset_start) \ - .filter(model.JobToOutputDatasetAssociation.job_id <= min(end_job_id, offset_start + args.batch_size)) \ - .all() + job_to_output_hda_ids = sa_session.query(model.JobToOutputDatasetAssociation.job_id, model.JobToOutputDatasetAssociation.dataset_id, + model.JobToOutputDatasetAssociation.name) \ + .filter(model.JobToOutputDatasetAssociation.job_id > offset_start) \ + .filter(model.JobToOutputDatasetAssociation.job_id <= min(end_job_id, offset_start + args.batch_size)) \ + .all() - # add type and concat - job_to_hda_ids = [[list(i), "input"] for i in job_to_input_hda_ids] + [[list(i), "output"] for i in job_to_output_hda_ids] + # add type and concat + job_to_hda_ids = [[list(i), "input"] for i in job_to_input_hda_ids] + [[list(i), "output"] for i in job_to_output_hda_ids] - # put all of the hda_ids into a list - hda_ids = [i[0][1] for i in job_to_hda_ids] + # put all of the hda_ids into a list + hda_ids = [i[0][1] for i in job_to_hda_ids] - hdas = sa_session.query(model.HistoryDatasetAssociation.id, model.HistoryDatasetAssociation.dataset_id, - model.HistoryDatasetAssociation.extension) \ - .filter(model.HistoryDatasetAssociation.id.in_(hda_ids)) \ - .all() + hdas = sa_session.query(model.HistoryDatasetAssociation.id, model.HistoryDatasetAssociation.dataset_id, + model.HistoryDatasetAssociation.extension) \ + .filter(model.HistoryDatasetAssociation.id.in_(hda_ids)) \ + .all() - # put all the dataset ids into a list - dataset_ids = [i[1] for i in hdas] + # put all the dataset ids into a list + dataset_ids = [i[1] for i in hdas] - # get the sizes of the datasets - datasets = sa_session.query(model.Dataset.id, model.Dataset.total_size) \ - .filter(model.Dataset.id.in_(dataset_ids)) \ - .all() + # get the sizes of the datasets + datasets = sa_session.query(model.Dataset.id, model.Dataset.total_size) \ + .filter(model.Dataset.id.in_(dataset_ids)) \ + .all() - # datasets to dictionay for easy search - hdas = {i[0]: i[1:] for i in hdas} - datasets = {i[0]: i[1:] for i in datasets} + # datasets to dictionay for easy search + hdas = {i[0]: i[1:] for i in hdas} + datasets = {i[0]: i[1:] for i in datasets} - for job_to_hda in job_to_hda_ids: + for job_to_hda in job_to_hda_ids: - job = job_to_hda[0] # job_id, hda_id, name - filetype = job_to_hda[1] # input|output + job = job_to_hda[0] # job_id, hda_id, name + filetype = job_to_hda[1] # input|output - # No associated job - if job[0] not in job_tool_map: - continue + # No associated job + if job[0] not in job_tool_map: + continue - # If the tool is blacklisted, exclude everywhere - if job_tool_map[job[0]] in blacklisted_tools: - continue + # If the tool is blacklisted, exclude everywhere + if job_tool_map[job[0]] in blacklisted_tools: + continue - hda_id = job[1] - if hda_id is None: - continue + hda_id = job[1] + if hda_id is None: + continue - dataset_id = hdas[hda_id][0] - if dataset_id is None: - continue + dataset_id = hdas[hda_id][0] + if dataset_id is None: + continue - try: - line = [ - str(job[0]), # Job ID - str(hda_id), # HDA ID - str(hdas[hda_id][1]), # Extension - round_to_2sd(datasets[dataset_id][0]), # File size - job[2], # Parameter name - str(filetype) # input/output - ] - cline = unicodify('\t'.join(line) + '\n') - handle_datasets.write(cline) - except Exception: - logging.warning("Unable to write out a 'handle_datasets' row. Ignoring the row.", exc_info=True) - continue - handle_datasets.close() + try: + line = [ + str(job[0]), # Job ID + str(hda_id), # HDA ID + str(hdas[hda_id][1]), # Extension + round_to_2sd(datasets[dataset_id][0]), # File size + job[2], # Parameter name + str(filetype) # input/output + ] + cline = unicodify('\t'.join(line) + '\n') + handle_datasets.write(cline) + except Exception: + logging.warning("Unable to write out a 'handle_datasets' row. Ignoring the row.", exc_info=True) + continue annotate('export_datasets_end') annotate('export_metric_num_start', 'Exporting Metrics (Numeric)') - handle_metric_num = io.open(REPORT_BASE + '.metric_num.tsv', 'w', encoding='utf-8') - handle_metric_num.write('\t'.join(('job_id', 'plugin', 'name', 'value')) + '\n') - for offset_start in range(last_job_sent, end_job_id, args.batch_size): - logging.debug("Processing %s:%s", offset_start, min(end_job_id, offset_start + args.batch_size)) - for metric in sa_session.query(model.JobMetricNumeric.job_id, model.JobMetricNumeric.plugin, model.JobMetricNumeric.metric_name, model.JobMetricNumeric.metric_value) \ - .filter(model.JobMetricNumeric.job_id > offset_start) \ - .filter(model.JobMetricNumeric.job_id <= min(end_job_id, offset_start + args.batch_size)) \ - .all(): - # No associated job - if metric[0] not in job_tool_map: - continue - # If the tool is blacklisted, exclude everywhere - if job_tool_map[metric[0]] in blacklisted_tools: - continue + with io.open(REPORT_BASE + '.metric_num.tsv', 'w', encoding='utf-8') as handle_metric_num: + handle_metric_num.write(u'\t'.join(('job_id', 'plugin', 'name', 'value')) + '\n') + for offset_start in range(last_job_sent, end_job_id, args.batch_size): + logging.debug("Processing %s:%s", offset_start, min(end_job_id, offset_start + args.batch_size)) + for metric in sa_session.query(model.JobMetricNumeric.job_id, model.JobMetricNumeric.plugin, model.JobMetricNumeric.metric_name, model.JobMetricNumeric.metric_value) \ + .filter(model.JobMetricNumeric.job_id > offset_start) \ + .filter(model.JobMetricNumeric.job_id <= min(end_job_id, offset_start + args.batch_size)) \ + .all(): + # No associated job + if metric[0] not in job_tool_map: + continue + # If the tool is blacklisted, exclude everywhere + if job_tool_map[metric[0]] in blacklisted_tools: + continue - try: - line = [ - str(metric[0]), # job id - metric[1], # plugin - metric[2], # name - str(metric[3]) # value - ] + try: + line = [ + str(metric[0]), # job id + metric[1], # plugin + metric[2], # name + str(metric[3]) # value + ] - cline = unicodify('\t'.join(line) + '\n') - handle_metric_num.write(cline) - except Exception: - logging.warning("Unable to write out a 'handle_metric_num' row. Ignoring the row.", exc_info=True) - continue - handle_metric_num.close() + cline = unicodify('\t'.join(line) + '\n') + handle_metric_num.write(cline) + except Exception: + logging.warning("Unable to write out a 'handle_metric_num' row. Ignoring the row.", exc_info=True) + continue annotate('export_metric_num_end') # Now on to outputs. diff --git a/scripts/interactivetools/key_type_token_mapping.py b/scripts/interactivetools/key_type_token_mapping.py index fa096dfbd60..df6dff6506a 100644 --- a/scripts/interactivetools/key_type_token_mapping.py +++ b/scripts/interactivetools/key_type_token_mapping.py @@ -6,7 +6,7 @@ from time import time import uwsgi -realtime_db_file = uwsgi.opt["interactivetools_map"] +realtime_db_file = uwsgi.opt["interactivetools_map"].decode('utf-8') db_conn = sqlite3.connect(realtime_db_file) DATABASE_TABLE_NAME = 'gxitproxy' @@ -51,6 +51,14 @@ class CacheList(): key_type_token_mapped_cache = CacheList() +def args_as_unicode(func): + def wrap_args(*args): + args = (arg.decode('utf-8') if isinstance(arg, bytes) else arg for arg in args) + return func(*args) + return wrap_args + + +@args_as_unicode def key_type_token_mapper_cached(key, key_type, token, route_extra, url, ttl): global key_type_token_mapped_cache cache_key = (key, key_type, token) @@ -58,11 +66,12 @@ def key_type_token_mapper_cached(key, key_type, token, route_extra, url, ttl): if entry is None: entry = key_type_token_mapper(key, key_type, token, route_extra, url) if entry is not None: - # Should we cache empt/not authorized entries, perhaps for shorter time? + # Should we cache empty/not authorized entries, perhaps for shorter time? key_type_token_mapped_cache.add_entry(cache_key, entry, ttl=float(ttl)) return entry +@args_as_unicode def key_type_token_mapper(key, key_type, token, route_extra, url): global db_conn # print 'key %s key_type %s token %s route_extra %s url %s\n' % (key, key_type, token, route_extra, url) diff --git a/test/integration/test_cli_runners.py b/test/integration/test_cli_runners.py index d09e6c07fd3..7ac42b5f333 100644 --- a/test/integration/test_cli_runners.py +++ b/test/integration/test_cli_runners.py @@ -3,6 +3,7 @@ import collections import os import string import subprocess +import sys import tempfile import unittest @@ -11,11 +12,9 @@ from galaxy_test.base.ssh_util import generate_ssh_keys from galaxy_test.driver import integration_util from .test_job_environments import BaseJobEnvironmentIntegrationTestCase - -RemoteConnection = collections.namedtuple('remote_connection', ['hostname', 'username', 'password', 'port', 'private_key', 'public_key']) +RemoteConnection = collections.namedtuple('remote_connection', ['hostname', 'username', 'port', 'private_key', 'public_key']) -@integration_util.skip_unless_docker() def start_ssh_docker(container_name, jobs_directory, port=10022, image='agaveapi/slurm'): ssh_keys = generate_ssh_keys() START_SLURM_DOCKER = ['docker', @@ -28,6 +27,7 @@ def start_ssh_docker(container_name, jobs_directory, port=10022, image='agaveapi '--name', container_name, '--rm', + '--privileged', # for torque '-v', "{jobs_directory}:{jobs_directory}".format(jobs_directory=jobs_directory), "-v", @@ -36,7 +36,9 @@ def start_ssh_docker(container_name, jobs_directory, port=10022, image='agaveapi 'nofile=2048:2048', image] subprocess.check_call(START_SLURM_DOCKER) - return RemoteConnection('localhost', 'testuser', 'testuser', port, ssh_keys.private_key_file, ssh_keys.public_key_file) + if sys.platform != 'darwin': + subprocess.check_call(['docker', 'exec', container_name, 'usermod', '-u', str(os.getuid()), 'testuser']) + return RemoteConnection('localhost', 'testuser', port, ssh_keys.private_key_file, ssh_keys.public_key_file) def stop_ssh_docker(container_name, remote_connection): @@ -58,6 +60,7 @@ def cli_job_config(remote_connection, shell_plugin='ParamikoShell', job_plugin=' $private_key $hostname $port + False False 42 @@ -72,6 +75,7 @@ def cli_job_config(remote_connection, shell_plugin='ParamikoShell', job_plugin=' return job_conf.name +@integration_util.skip_unless_docker() class BaseCliIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase): @classmethod @@ -91,7 +95,7 @@ class BaseCliIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase): super(BaseCliIntegrationTestCase, cls).tearDownClass() @classmethod - def handle_galaxy_config_kwds(cls, config, ): + def handle_galaxy_config_kwds(cls, config): config["jobs_directory"] = cls.jobs_directory config["file_path"] = cls.jobs_directory config["job_config_file"] = cli_job_config(remote_connection=cls.remote_connection, diff --git a/test/integration/test_containerized_jobs.py b/test/integration/test_containerized_jobs.py index 626d63c5e63..6d6791e8ce8 100644 --- a/test/integration/test_containerized_jobs.py +++ b/test/integration/test_containerized_jobs.py @@ -1,5 +1,6 @@ """Integration tests for running tools in Docker containers.""" +import json import os import unittest @@ -37,6 +38,11 @@ class ContainerizedIntegrationTestCase(integration_util.IntegrationTestCase): skip_if_container_type_unavailable(cls) super(ContainerizedIntegrationTestCase, cls).setUpClass() + @classmethod + def handle_galaxy_config_kwds(cls, config): + config["job_config_file"] = DOCKERIZED_JOB_CONFIG_FILE + disable_dependency_resolution(config) + def disable_dependency_resolution(config): # Disable tool dependency resolution. @@ -56,7 +62,6 @@ class DockerizedJobsIntegrationTestCase(integration_util.IntegrationTestCase, Ru job_config_file = DOCKERIZED_JOB_CONFIG_FILE build_mulled_resolver = 'build_mulled' container_type = 'docker' - default_container_home_dir = '/' @classmethod def handle_galaxy_config_kwds(cls, config): @@ -100,19 +105,24 @@ class DockerizedJobsIntegrationTestCase(integration_util.IntegrationTestCase, Ru assert job_env.pwd.endswith("/working") # Should we change env_pass_through to just always include TMP and HOME for docker? # I'm not sure, if yes this would change. - assert job_env.home == self.default_container_home_dir, job_env.home + assert not job_env.home.endswith('/home') def test_build_mulled(self): if not which('docker'): raise unittest.SkipTest("Docker not found on PATH, required for building images via involucro") resolver_type = self.build_mulled_resolver - tool_id = 'mulled_example_multi_1' - endpoint = "tools/%s/dependencies" % tool_id - data = {'id': tool_id, 'resolver_type': resolver_type} + tool_ids = ['mulled_example_multi_1'] + endpoint = "dependency_resolvers/toolbox/install" + data = {'tool_ids': json.dumps(tool_ids), 'resolver_type': resolver_type, 'container_type': self.container_type, 'include_containers': True} create_response = self._post(endpoint, data=data, admin=True) self._assert_status_code_is(create_response, 200) + create_response = self._get("dependency_resolvers/toolbox", data={'tool_ids': tool_ids, 'container_type': self.container_type, 'include_containers': True, 'index_by': 'tools'}, admin=True) response = create_response.json() - assert any([True for d in response if d['dependency_type'] == self.container_type]) + assert len(response) == 1 + status = response[0]['status'] + assert status[0]['model_class'] == 'ContainerDependency' + assert status[0]['dependency_type'] == self.container_type + assert status[0]['container_description']['identifier'].startswith('quay.io/local/mulled-v2-') class MappingContainerResolverTestCase(integration_util.IntegrationTestCase): @@ -164,5 +174,3 @@ class SingularityJobsIntegrationTestCase(DockerizedJobsIntegrationTestCase): job_config_file = SINGULARITY_JOB_CONFIG_FILE build_mulled_resolver = 'build_mulled_singularity' container_type = 'singularity' - # singularity passes $HOME by default - default_container_home_dir = os.environ.get('HOME', '/') diff --git a/test/integration/test_interactivetools_api.py b/test/integration/test_interactivetools_api.py index e214e3eee80..f3ae60d4879 100644 --- a/test/integration/test_interactivetools_api.py +++ b/test/integration/test_interactivetools_api.py @@ -47,8 +47,7 @@ class InteractiveToolsIntegrationTestCase(ContainerizedIntegrationTestCase): assert isinstance(jobs, list) assert len(jobs) == 1 job0 = jobs[0] - entry_points = self.wait_on_entry_points_active(job0["id"]) - assert len(entry_points) == 2 + entry_points = self.wait_on_entry_points_active(job0["id"], expected_num=2) entry_point0 = entry_points[0] entry_point1 = entry_points[1] target0 = self.entry_point_target(entry_point0["id"]) @@ -59,19 +58,18 @@ class InteractiveToolsIntegrationTestCase(ContainerizedIntegrationTestCase): content1 = self.wait_on_proxied_content(target1) assert content1 == "moo cow\n", content1 - assert False def wait_on_proxied_content(self, target): def get_hosted_content(): try: scheme, rest = target.split("://", 1) - prefix, host_and_port = rest.split(".realtime.") - print(rest) + prefix, host_and_port = rest.split(".interactivetool.") faked_host = rest if "/" in rest: faked_host = rest.split("/", 1)[0] - response = requests.get("%s://%s" % (scheme, host_and_port), timeout=1, headers={"Host": faked_host}) - return response.content + url = "%s://%s" % (scheme, host_and_port) + response = requests.get(url, timeout=1, headers={"Host": faked_host}) + return response.text except Exception as e: print(e) return None diff --git a/test/integration/test_kubernetes_runner.py b/test/integration/test_kubernetes_runner.py index 69966b6870c..62fe4640a82 100644 --- a/test/integration/test_kubernetes_runner.py +++ b/test/integration/test_kubernetes_runner.py @@ -1,5 +1,6 @@ """Integration tests for the Kubernetes runner.""" -# Tested on docker for mac 18.06.1-ce-mac73 using the default kubernetes setup +# Tested on docker for mac 18.06.1-ce-mac73 using the default kubernetes setup, +# also works on minikube import collections import json import os @@ -91,37 +92,40 @@ def job_config(jobs_directory): jobs-directory-claim:$jobs_directory,tool-directory-claim:$tool_directory $k8s_config_path gx-short-id + $$uid jobs-directory-claim:$jobs_directory,tool-directory-claim:$tool_directory $k8s_config_path gx-short-id 10 + $$uid jobs-directory-claim:$jobs_directory,tool-directory-claim:$tool_directory $k8s_config_path gx-short-id never + $$uid - 1.9 + 1.1 10M true busybox:ubuntu-14.04 42 - 1.9 + 1.1 10M true busybox:ubuntu-14.04 42 - 1.9 + 1.1 10M true busybox:ubuntu-14.04 @@ -182,7 +186,9 @@ class BaseKubernetesIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, M super(BaseKubernetesIntegrationTestCase, cls).tearDownClass() @classmethod - def handle_galaxy_config_kwds(cls, config, ): + def handle_galaxy_config_kwds(cls, config): + # TODO: implement metadata setting as separate job, as service or side-car + config['retry_metadata_internally'] = True config["jobs_directory"] = cls.jobs_directory config["file_path"] = cls.jobs_directory config["job_config_file"] = cls.job_config.path diff --git a/test/integration/test_scripts.py b/test/integration/test_scripts.py index cbedbd75f0e..b718f1ec6ca 100644 --- a/test/integration/test_scripts.py +++ b/test/integration/test_scripts.py @@ -116,7 +116,6 @@ class ScriptsIntegrationTestCase(integration_util.IntegrationTestCase): output = self._scripts_check_output(script, ["-c", config_file]) assert "Complete" in output - @integration_util.skip_if_jenkins def test_grt_export(self): script = "grt/export.py" self._scripts_check_argparse_help(script) @@ -134,7 +133,7 @@ class ScriptsIntegrationTestCase(integration_util.IntegrationTestCase): json_file = os.path.join(self.config_dir, json_files[0]) with open(json_file, "r") as f: export = json.load(f) - assert export["version"] == 2 + assert export["version"] == 3 def test_admin_cleanup_datasets(self): self._scripts_check_argparse_help("cleanup_datasets/admin_cleanup_datasets.py") @@ -179,13 +178,21 @@ class ScriptsIntegrationTestCase(integration_util.IntegrationTestCase): clean_env = { "PATH": os.environ.get("PATH", None), } # Don't let testing environment variables interfere with config. - return unicodify(subprocess.check_output(cmd, cwd=cwd, env=clean_env)) + try: + return unicodify(subprocess.check_output(cmd, cwd=cwd, env=clean_env)) + except Exception as e: + if isinstance(e, subprocess.CalledProcessError): + raise Exception("%s\nOutput was:\n%s" % (unicodify(e), e.output)) + raise def write_config_file(self): config_dir = self.config_dir path = os.path.join(config_dir, "galaxy.yml") self._test_driver.temp_directories.extend([config_dir]) + config = self._raw_config + # Update config dict with database_connection, which might be set through env variables + config['database_connection'] = self._app.config.database_connection with open(path, "w") as f: - yaml.dump({"galaxy": self._raw_config}, f) + yaml.dump({"galaxy": config}, f) return path diff --git a/test/unit/test_remote_shell.py b/test/unit/test_remote_shell.py index d460509c3fd..f13305ba2be 100644 --- a/test/unit/test_remote_shell.py +++ b/test/unit/test_remote_shell.py @@ -23,6 +23,7 @@ class TestCliInterface(unittest.TestCase): cls.username = 'testuser' cls.shell_params = {'username': cls.username, 'private_key': cls.ssh_keys.private_key_file, + 'strict_host_key_checking': False, 'hostname': 'localhost'} cls.cli_interface = CliInterface() diff --git a/test/unit/tools/test_watcher.py b/test/unit/tools/test_watcher.py index 352a4a8a32c..f1e0c12a21b 100644 --- a/test/unit/tools/test_watcher.py +++ b/test/unit/tools/test_watcher.py @@ -4,47 +4,47 @@ from contextlib import contextmanager from os import path from shutil import rmtree +import pytest + from galaxy.tools.toolbox import watcher from galaxy.util import bunch +@pytest.mark.skipif(not watcher.can_watch, reason="watchdog not available") def test_watcher(): - if not watcher.can_watch: - from nose.plugins.skip import SkipTest - raise SkipTest() - with __test_directory() as t: tool_path = path.join(t, "test.xml") toolbox = Toolbox() - open(tool_path, "w").write("a") + with open(tool_path, "w") as f: + f.write("a") tool_watcher = watcher.get_tool_watcher(toolbox, bunch.Bunch( watch_tools=True )) tool_watcher.start() - time.sleep(1) tool_watcher.watch_file(tool_path, "cool_tool") + time.sleep(2) assert not toolbox.was_reloaded("cool_tool") - open(tool_path, "w").write("b") + with open(tool_path, "w") as f: + f.write("b") wait_for_reload(lambda: toolbox.was_reloaded("cool_tool")) tool_watcher.shutdown() assert tool_watcher.observer is None +@pytest.mark.skipif(not watcher.can_watch, reason="watchdog not available") def test_tool_conf_watcher(): - if not watcher.can_watch: - from nose.plugins.skip import SkipTest - raise SkipTest() - callback = CallbackRecorder() conf_watcher = watcher.get_tool_conf_watcher(callback.call) conf_watcher.start() with __test_directory() as t: tool_conf_path = path.join(t, "test_conf.xml") - open(tool_conf_path, "w").write("a") + with open(tool_conf_path, "w") as f: + f.write("a") conf_watcher.watch_file(tool_conf_path) - time.sleep(1) - open(tool_conf_path, "w").write("b") + time.sleep(2) + with open(tool_conf_path, "w") as f: + f.write("b") wait_for_reload(lambda: callback.called) conf_watcher.shutdown() assert conf_watcher.thread is None