From 3ba46709fd04a81b78d1a7bceb08ac8089237c3f Mon Sep 17 00:00:00 2001 From: Sasan Ahmadi Date: Wed, 6 Apr 2022 15:57:08 -0700 Subject: [PATCH 1/4] bugfix for when polling for the created job, if fail to get job info it should not fail the task, instead it should continue polling until reaches the max allowed polling tries --- .../jenkins/operators/jenkins_job_trigger.py | 30 ++++++++++++------- 1 file changed, 19 insertions(+), 11 deletions(-) diff --git a/airflow/providers/jenkins/operators/jenkins_job_trigger.py b/airflow/providers/jenkins/operators/jenkins_job_trigger.py index 3f3bafa27bbbf..e821e227d7310 100644 --- a/airflow/providers/jenkins/operators/jenkins_job_trigger.py +++ b/airflow/providers/jenkins/operators/jenkins_job_trigger.py @@ -153,17 +153,25 @@ def poll_job_in_queue(self, location: str, jenkins_server: Jenkins) -> int: # once it will be available in python-jenkins (v > 0.4.15) self.log.info('Polling jenkins queue at the url %s', location) while try_count < self.max_try_before_job_appears: - location_answer = jenkins_request_with_headers( - jenkins_server, Request(method='POST', url=location) - ) - if location_answer is not None: - json_response = json.loads(location_answer['body']) - if 'executable' in json_response and 'number' in json_response['executable']: - build_number = json_response['executable']['number'] - self.log.info('Job executed on Jenkins side with the build number %s', build_number) - return build_number - try_count += 1 - time.sleep(self.sleep_time) + try: + location_answer = jenkins_request_with_headers( + jenkins_server, Request(method='POST', url=location) + ) + if location_answer is not None: + json_response = json.loads(location_answer['body']) + if 'executable' in json_response and json_response['executable'] is not None and 'number' in json_response['executable']: + build_number = json_response['executable']['number'] + self.log.info('Job executed on Jenkins side with the build number %s', build_number) + return build_number + try_count += 1 + time.sleep(self.sleep_time) + # we don't want to fail the operator, this will continue to poll + # until max_try_before_job_appears reached + except (HTTPError, JenkinsException) as ex: + self.log.info(f'polling failed, retry polling. Failure reason: {ex}') + try_count += 1 + time.sleep(self.sleep_time) + continue raise AirflowException( "The job hasn't been executed after polling " f"the queue {self.max_try_before_job_appears} times" ) From e4a61e84b34edd13f9fd991f9d7d06f33a732f74 Mon Sep 17 00:00:00 2001 From: Sasan Ahmadi Date: Wed, 6 Apr 2022 18:23:45 -0700 Subject: [PATCH 2/4] correct static code problem --- airflow/providers/jenkins/operators/jenkins_job_trigger.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/airflow/providers/jenkins/operators/jenkins_job_trigger.py b/airflow/providers/jenkins/operators/jenkins_job_trigger.py index c028193ee920f..3acb19b0c5636 100644 --- a/airflow/providers/jenkins/operators/jenkins_job_trigger.py +++ b/airflow/providers/jenkins/operators/jenkins_job_trigger.py @@ -169,7 +169,8 @@ def poll_job_in_queue(self, location: str, jenkins_server: Jenkins) -> int: return build_number try_count += 1 time.sleep(self.sleep_time) - # we don't want to fail the operator, this will continue to poll + + # we don't want to fail the operator, this will continue to poll # until max_try_before_job_appears reached except (HTTPError, JenkinsException) as ex: self.log.info(f'polling failed, retry polling. Failure reason: {ex}') From f8235b1597c0751d57d776e10a2ca573063c7323 Mon Sep 17 00:00:00 2001 From: Sasan Ahmadi Date: Wed, 6 Apr 2022 20:48:02 -0700 Subject: [PATCH 3/4] reduce the size of enclosing code in try except try except accounts for only the risky part of the code --- .../jenkins/operators/jenkins_job_trigger.py | 27 ++++++++++--------- 1 file changed, 14 insertions(+), 13 deletions(-) diff --git a/airflow/providers/jenkins/operators/jenkins_job_trigger.py b/airflow/providers/jenkins/operators/jenkins_job_trigger.py index 3acb19b0c5636..a8f29442cd946 100644 --- a/airflow/providers/jenkins/operators/jenkins_job_trigger.py +++ b/airflow/providers/jenkins/operators/jenkins_job_trigger.py @@ -157,19 +157,6 @@ def poll_job_in_queue(self, location: str, jenkins_server: Jenkins) -> int: location_answer = jenkins_request_with_headers( jenkins_server, Request(method='POST', url=location) ) - if location_answer is not None: - json_response = json.loads(location_answer['body']) - if ( - 'executable' in json_response - and json_response['executable'] is not None - and 'number' in json_response['executable'] - ): - build_number = json_response['executable']['number'] - self.log.info('Job executed on Jenkins side with the build number %s', build_number) - return build_number - try_count += 1 - time.sleep(self.sleep_time) - # we don't want to fail the operator, this will continue to poll # until max_try_before_job_appears reached except (HTTPError, JenkinsException) as ex: @@ -177,6 +164,20 @@ def poll_job_in_queue(self, location: str, jenkins_server: Jenkins) -> int: try_count += 1 time.sleep(self.sleep_time) continue + + if location_answer is not None: + json_response = json.loads(location_answer['body']) + if ( + 'executable' in json_response + and json_response['executable'] is not None + and 'number' in json_response['executable'] + ): + build_number = json_response['executable']['number'] + self.log.info('Job executed on Jenkins side with the build number %s', build_number) + return build_number + try_count += 1 + time.sleep(self.sleep_time) + raise AirflowException( "The job hasn't been executed after polling " f"the queue {self.max_try_before_job_appears} times" ) From c1b92ec61b61481b69f80c3be5383003114caa6e Mon Sep 17 00:00:00 2001 From: Sasan Ahmadi Date: Thu, 7 Apr 2022 00:18:10 -0700 Subject: [PATCH 4/4] corrected for code review --- airflow/providers/jenkins/operators/jenkins_job_trigger.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/airflow/providers/jenkins/operators/jenkins_job_trigger.py b/airflow/providers/jenkins/operators/jenkins_job_trigger.py index a8f29442cd946..b7dcb25913b27 100644 --- a/airflow/providers/jenkins/operators/jenkins_job_trigger.py +++ b/airflow/providers/jenkins/operators/jenkins_job_trigger.py @@ -159,8 +159,8 @@ def poll_job_in_queue(self, location: str, jenkins_server: Jenkins) -> int: ) # we don't want to fail the operator, this will continue to poll # until max_try_before_job_appears reached - except (HTTPError, JenkinsException) as ex: - self.log.info(f'polling failed, retry polling. Failure reason: {ex}') + except (HTTPError, JenkinsException): + self.log.warning('polling failed, retrying', exc_info=True) try_count += 1 time.sleep(self.sleep_time) continue @@ -179,7 +179,7 @@ def poll_job_in_queue(self, location: str, jenkins_server: Jenkins) -> int: time.sleep(self.sleep_time) raise AirflowException( - "The job hasn't been executed after polling " f"the queue {self.max_try_before_job_appears} times" + f"The job hasn't been executed after polling the queue {self.max_try_before_job_appears} times" ) def get_hook(self) -> JenkinsHook: