From 276d8dce229d1ddb7b312c77aee19857caebb58e Mon Sep 17 00:00:00 2001 From: Dannon Baker Date: Wed, 30 Sep 2020 16:26:49 -0400 Subject: [PATCH 1/5] Pin pytest < 6.1 for now (backport) --- packages/test.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/test.sh b/packages/test.sh index 2f7a62af826..c7f702fb6ea 100755 --- a/packages/test.sh +++ b/packages/test.sh @@ -14,7 +14,7 @@ TEST_ENV_DIR=${TEST_ENV_DIR:-$(mktemp -d -t gxpkgtestenvXXXXXX)} virtualenv -p "$TEST_PYTHON" "$TEST_ENV_DIR" . "${TEST_ENV_DIR}/bin/activate" -pip install pytest +pip install "pytest<6.1" # ensure ordered by dependency dag PACKAGE_DIRS=( From 9d311d4557a90e7991a6817b25d6c8c76b47bff2 Mon Sep 17 00:00:00 2001 From: Nuwan Goonasekera <2070605+nuwang@users.noreply.github.com> Date: Mon, 28 Sep 2020 11:42:15 +0530 Subject: [PATCH 2/5] Prevent endless cleanup loops when duplicate k8s job or k8s job not found --- lib/galaxy/jobs/runners/kubernetes.py | 25 ++++--------------------- 1 file changed, 4 insertions(+), 21 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 362fd99666d..9fcec554f81 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -2,7 +2,6 @@ Offload jobs to a Kubernetes cluster. """ -import errno import logging import math import os @@ -465,30 +464,14 @@ class KubernetesJobRunner(AsynchronousJobRunner): # there is no job responding to this job_id, it is either lost or something happened. log.error("No Jobs are available under expected selector app=%s", job_state.job_id) self.mark_as_failed(job_state) - try: - with open(job_state.error_file, 'w') as error_file: - error_file.write("No Kubernetes Jobs are available under expected selector app=%s\n" % job_state.job_id) - except EnvironmentError as e: - # Python 2/3 compatible handling of FileNotFoundError - if e.errno == errno.ENOENT: - log.error("Job directory already cleaned up. Assuming already handled for selector app=%s", job_state.job_id) - else: - raise - return job_state + # job is no longer viable - remove from watched jobs + return None else: # there is more than one job associated to the expected unique job id used as selector. log.error("More than one Kubernetes Job associated to job id '%s'", job_state.job_id) self.mark_as_failed(job_state) - try: - with open(job_state.error_file, 'w') as error_file: - error_file.write("More than one Kubernetes Job associated with job id '%s'\n" % job_state.job_id) - except EnvironmentError as e: - # Python 2/3 compatible handling of FileNotFoundError - if e.errno == errno.ENOENT: - log.error("Job directory already cleaned up. Assuming already handled for selector app=%s", job_state.job_id) - else: - raise - return job_state + # job is no longer viable - remove from watched jobs + return None def _handle_job_failure(self, job, job_state): # Figure out why job has failed From 5d1b54cf6cfa895ae18eae008f9744b3b0e438e5 Mon Sep 17 00:00:00 2001 From: Nuwan Goonasekera <2070605+nuwang@users.noreply.github.com> Date: Fri, 15 May 2020 18:17:26 +0530 Subject: [PATCH 3/5] Added test to check that jobs fail when deleted directly via k8s api --- test/integration/test_kubernetes_runner.py | 51 ++++++++++++++++++++++ 1 file changed, 51 insertions(+) diff --git a/test/integration/test_kubernetes_runner.py b/test/integration/test_kubernetes_runner.py index 09b325b92b0..d6625661e15 100644 --- a/test/integration/test_kubernetes_runner.py +++ b/test/integration/test_kubernetes_runner.py @@ -264,6 +264,57 @@ class BaseKubernetesIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, M subprocess.check_output(['kubectl', 'get', 'job', external_id, '-o', 'json'], stderr=subprocess.STDOUT) assert "not found" in unicodify(excinfo.value.output) + @skip_without_tool('cat_data_and_sleep') + def test_external_job_delete(self): + with self.dataset_populator.test_history() as history_id: + hda1 = self.dataset_populator.new_dataset(history_id, content="1 2 3") + running_inputs = { + "input1": {"src": "hda", "id": hda1["id"]}, + "sleep_time": 240, + } + running_response = self.dataset_populator.run_tool( + "cat_data_and_sleep", + running_inputs, + history_id, + assert_ok=False, + ).json() + job_dict = running_response["jobs"][0] + + app = self._app + sa_session = app.model.context.current + external_id = None + state = False + + job = sa_session.query(app.model.Job).filter_by(tool_id="cat_data_and_sleep").one() + # Not checking the state here allows the change from queued to running to overwrite + # the change from queued to deleted_new in the API thread - this is a problem because + # the job will still run. See issue https://github.com/galaxyproject/galaxy/issues/4960. + max_tries = 60 + while max_tries > 0 and external_id is None or state != app.model.Job.states.RUNNING: + sa_session.refresh(job) + assert not job.finished + external_id = job.job_runner_external_id + state = job.state + time.sleep(1) + max_tries -= 1 + + output = unicodify(subprocess.check_output(['kubectl', 'get', 'job', external_id, '-o', 'json'])) + status = json.loads(output) + assert status['status']['active'] == 1 + + output = unicodify(subprocess.check_output(['kubectl', 'delete', 'job', external_id, '-o', 'name'])) + assert 'job.batch/%s' % external_id in output + + max_tries = 60 + state = False + while max_tries > 0 and external_id is None or state != app.model.Job.states.ERROR: + sa_session.refresh(job) + state = job.state + time.sleep(1) + max_tries -= 1 + + assert state == app.model.Job.states.ERROR + @skip_without_tool('job_properties') def test_exit_code_127(self): inputs = { From ada9ee648d2864b6feb3c7fde71438d3a21b0592 Mon Sep 17 00:00:00 2001 From: Nuwan Goonasekera <2070605+nuwang@users.noreply.github.com> Date: Fri, 15 May 2020 19:02:19 +0530 Subject: [PATCH 4/5] Minor refactoring of k8s tests --- test/integration/test_kubernetes_runner.py | 59 ++++++++-------------- 1 file changed, 21 insertions(+), 38 deletions(-) diff --git a/test/integration/test_kubernetes_runner.py b/test/integration/test_kubernetes_runner.py index d6625661e15..8354a7c94d8 100644 --- a/test/integration/test_kubernetes_runner.py +++ b/test/integration/test_kubernetes_runner.py @@ -216,6 +216,17 @@ class BaseKubernetesIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, M job_env = self._run_and_get_environment_properties() assert job_env.some_env == '42' + @staticmethod + def _wait_for_external_state(sa_session, job, expected): + # Not checking the state here allows the change from queued to running to overwrite + # the change from queued to deleted_new in the API thread - this is a problem because + # the job will still run. See issue https://github.com/galaxyproject/galaxy/issues/4960. + max_tries = 60 + while max_tries > 0 and job.job_runner_external_id is None or job.state != expected: + sa_session.refresh(job) + time.sleep(1) + max_tries -= 1 + @skip_without_tool('cat_data_and_sleep') def test_kill_process(self): with self.dataset_populator.test_history() as history_id: @@ -234,22 +245,12 @@ class BaseKubernetesIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, M app = self._app sa_session = app.model.context.current - external_id = None - state = False - job = sa_session.query(app.model.Job).filter_by(tool_id="cat_data_and_sleep").one() - # Not checking the state here allows the change from queued to running to overwrite - # the change from queued to deleted_new in the API thread - this is a problem because - # the job will still run. See issue https://github.com/galaxyproject/galaxy/issues/4960. - max_tries = 60 - while max_tries > 0 and external_id is None or state != app.model.Job.states.RUNNING: - sa_session.refresh(job) - assert not job.finished - external_id = job.job_runner_external_id - state = job.state - time.sleep(1) - max_tries -= 1 + self._wait_for_external_state(sa_session, job, app.model.Job.states.RUNNING) + assert not job.finished + + external_id = job.job_runner_external_id output = unicodify(subprocess.check_output(['kubectl', 'get', 'job', external_id, '-o', 'json'])) status = json.loads(output) assert status['status']['active'] == 1 @@ -277,27 +278,15 @@ class BaseKubernetesIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, M running_inputs, history_id, assert_ok=False, - ).json() - job_dict = running_response["jobs"][0] + ) app = self._app sa_session = app.model.context.current - external_id = None - state = False - job = sa_session.query(app.model.Job).filter_by(tool_id="cat_data_and_sleep").one() - # Not checking the state here allows the change from queued to running to overwrite - # the change from queued to deleted_new in the API thread - this is a problem because - # the job will still run. See issue https://github.com/galaxyproject/galaxy/issues/4960. - max_tries = 60 - while max_tries > 0 and external_id is None or state != app.model.Job.states.RUNNING: - sa_session.refresh(job) - assert not job.finished - external_id = job.job_runner_external_id - state = job.state - time.sleep(1) - max_tries -= 1 + self._wait_for_external_state(sa_session, job, app.model.Job.states.RUNNING) + + external_id = job.job_runner_external_id output = unicodify(subprocess.check_output(['kubectl', 'get', 'job', external_id, '-o', 'json'])) status = json.loads(output) assert status['status']['active'] == 1 @@ -305,15 +294,9 @@ class BaseKubernetesIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, M output = unicodify(subprocess.check_output(['kubectl', 'delete', 'job', external_id, '-o', 'name'])) assert 'job.batch/%s' % external_id in output - max_tries = 60 - state = False - while max_tries > 0 and external_id is None or state != app.model.Job.states.ERROR: - sa_session.refresh(job) - state = job.state - time.sleep(1) - max_tries -= 1 + self._wait_for_external_state(sa_session, job, app.model.Job.states.ERROR) - assert state == app.model.Job.states.ERROR + assert job.state == app.model.Job.states.ERROR @skip_without_tool('job_properties') def test_exit_code_127(self): From e0d221a97b512f9e61b8fb1f52d1fb5831acfd42 Mon Sep 17 00:00:00 2001 From: Nuwan Goonasekera <2070605+nuwang@users.noreply.github.com> Date: Fri, 15 May 2020 19:17:57 +0530 Subject: [PATCH 5/5] Also check for expected output in k8s delete test --- test/integration/test_kubernetes_runner.py | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/test/integration/test_kubernetes_runner.py b/test/integration/test_kubernetes_runner.py index 8354a7c94d8..dada94cdfb3 100644 --- a/test/integration/test_kubernetes_runner.py +++ b/test/integration/test_kubernetes_runner.py @@ -245,7 +245,7 @@ class BaseKubernetesIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, M app = self._app sa_session = app.model.context.current - job = sa_session.query(app.model.Job).filter_by(tool_id="cat_data_and_sleep").one() + job = sa_session.query(app.model.Job).get(app.security.decode_id(job_dict["id"])) self._wait_for_external_state(sa_session, job, app.model.Job.states.RUNNING) assert not job.finished @@ -279,10 +279,11 @@ class BaseKubernetesIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, M history_id, assert_ok=False, ) + job_dict = running_response.json()["jobs"][0] app = self._app sa_session = app.model.context.current - job = sa_session.query(app.model.Job).filter_by(tool_id="cat_data_and_sleep").one() + job = sa_session.query(app.model.Job).get(app.security.decode_id(job_dict["id"])) self._wait_for_external_state(sa_session, job, app.model.Job.states.RUNNING) @@ -294,9 +295,11 @@ class BaseKubernetesIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, M output = unicodify(subprocess.check_output(['kubectl', 'delete', 'job', external_id, '-o', 'name'])) assert 'job.batch/%s' % external_id in output - self._wait_for_external_state(sa_session, job, app.model.Job.states.ERROR) + result = self.dataset_populator.wait_for_tool_run(run_response=running_response, history_id=history_id, + assert_ok=False).json() + details = self.dataset_populator.get_job_details(result['jobs'][0]['id'], full=True).json() - assert job.state == app.model.Job.states.ERROR + assert details['state'] == app.model.Job.states.ERROR, details @skip_without_tool('job_properties') def test_exit_code_127(self):