Merge pull request #1421 from alphagov/celery-fix

celery.task.retry exc param should be a throwable.
This commit is contained in:
Leo Hemsted
2017-11-28 10:19:14 +00:00
committed by GitHub
3 changed files with 30 additions and 19 deletions
+2 -2
View File
@@ -13,6 +13,6 @@ def process_ses_results(self, response):
errors = process_ses_response(response) errors = process_ses_response(response)
if errors: if errors:
current_app.logger.error(errors) current_app.logger.error(errors)
except Exception: except Exception as exc:
current_app.logger.exception('Error processing SES results') current_app.logger.exception('Error processing SES results')
self.retry(queue=QueueNames.RETRY, exc="SES responses processed with error") self.retry(queue=QueueNames.RETRY)
+9 -8
View File
@@ -18,6 +18,8 @@ from requests import (
RequestException RequestException
) )
from sqlalchemy.exc import SQLAlchemyError from sqlalchemy.exc import SQLAlchemyError
from botocore.exceptions import ClientError as BotoClientError
from app import ( from app import (
create_uuid, create_uuid,
create_random_identifier, create_random_identifier,
@@ -323,12 +325,13 @@ def build_dvla_file(self, job_id):
) )
dao_update_job_status(job_id, JOB_STATUS_READY_TO_SEND) dao_update_job_status(job_id, JOB_STATUS_READY_TO_SEND)
else: else:
current_app.logger.info("All notifications for job {} are not persisted".format(job_id)) msg = "All notifications for job {} are not persisted".format(job_id)
self.retry(queue=QueueNames.RETRY, exc="All notifications for job {} are not persisted".format(job_id)) current_app.logger.info(msg)
except Exception as e: self.retry(queue=QueueNames.RETRY)
# ? should this retry? # specifically don't catch celery.retry errors
except (SQLAlchemyError, BotoClientError):
current_app.logger.exception("build_dvla_file threw exception") current_app.logger.exception("build_dvla_file threw exception")
raise e self.retry(queue=QueueNames.RETRY)
@notify_celery.task(bind=True, name='update-letter-job-to-sent') @notify_celery.task(bind=True, name='update-letter-job-to-sent')
@@ -520,9 +523,7 @@ def send_inbound_sms_to_service(self, inbound_sms_id, service_id):
) )
if not isinstance(e, HTTPError) or e.response.status_code >= 500: if not isinstance(e, HTTPError) or e.response.status_code >= 500:
try: try:
self.retry(queue=QueueNames.RETRY, self.retry(queue=QueueNames.RETRY)
exc='Unable to send_inbound_sms_to_service for service_id: {} and url: {}. \n{}'.format(
service_id, inbound_api.url, e))
except self.MaxRetriesExceededError: except self.MaxRetriesExceededError:
current_app.logger.exception('Retry: send_inbound_sms_to_service has retried the max number of times') current_app.logger.exception('Retry: send_inbound_sms_to_service has retried the max number of times')
+19 -9
View File
@@ -1,16 +1,17 @@
import codecs
import json import json
import uuid import uuid
from datetime import datetime, timedelta from datetime import datetime, timedelta
from unittest.mock import Mock from unittest.mock import Mock
import pytest import pytest
import requests_mock import requests_mock
from flask import current_app from flask import current_app
from freezegun import freeze_time from freezegun import freeze_time
from requests import RequestException from requests import RequestException
from sqlalchemy.exc import SQLAlchemyError from sqlalchemy.exc import SQLAlchemyError
from notifications_utils.template import SMSMessageTemplate, WithSubjectTemplate, LetterDVLATemplate
from celery.exceptions import Retry from celery.exceptions import Retry
from botocore.exceptions import ClientError
from notifications_utils.template import SMSMessageTemplate, WithSubjectTemplate, LetterDVLATemplate
from app import (encryption, DATETIME_FORMAT) from app import (encryption, DATETIME_FORMAT)
from app.celery import provider_tasks from app.celery import provider_tasks
@@ -1024,9 +1025,7 @@ def test_build_dvla_file(sample_letter_template, mocker):
create_notification(template=job.template, job=job) create_notification(template=job.template, job=job)
mocked_upload = mocker.patch("app.celery.tasks.s3upload") mocked_upload = mocker.patch("app.celery.tasks.s3upload")
mocked_send_task = mocker.patch("app.celery.tasks.notify_celery.send_task") mocked_send_task = mocker.patch("app.celery.tasks.notify_celery.send_task")
mocked_letter_template = mocker.patch("app.celery.tasks.LetterDVLATemplate") mocker.patch("app.celery.tasks.LetterDVLATemplate", return_value='dvla|string')
mocked_letter_template_instance = mocked_letter_template.return_value
mocked_letter_template_instance.__str__.return_value = "dvla|string"
build_dvla_file(job.id) build_dvla_file(job.id)
mocked_upload.assert_called_once_with( mocked_upload.assert_called_once_with(
@@ -1049,12 +1048,25 @@ def test_build_dvla_file_retries_if_all_notifications_are_not_created(sample_let
build_dvla_file(job.id) build_dvla_file(job.id)
mocked.assert_not_called() mocked.assert_not_called()
tasks.build_dvla_file.retry.assert_called_with(queue="retry-tasks", tasks.build_dvla_file.retry.assert_called_with(queue="retry-tasks")
exc="All notifications for job {} are not persisted".format(job.id))
assert Job.query.get(job.id).job_status == 'in progress' assert Job.query.get(job.id).job_status == 'in progress'
mocked_send_task.assert_not_called() mocked_send_task.assert_not_called()
def test_build_dvla_file_retries_if_s3_err(sample_letter_template, mocker):
job = create_job(sample_letter_template, notification_count=1)
create_notification(job.template, job=job)
mocker.patch('app.celery.tasks.LetterDVLATemplate', return_value='dvla|string')
mocker.patch('app.celery.tasks.s3upload', side_effect=ClientError({}, 'operation_name'))
retry_mock = mocker.patch('app.celery.tasks.build_dvla_file.retry', side_effect=Retry)
with pytest.raises(Retry):
build_dvla_file(job.id)
retry_mock.assert_called_once_with(queue='retry-tasks')
def test_create_dvla_file_contents(notify_db_session, mocker): def test_create_dvla_file_contents(notify_db_session, mocker):
service = create_service(service_permissions=SERVICE_PERMISSION_TYPES) service = create_service(service_permissions=SERVICE_PERMISSION_TYPES)
create_letter_contact(service=service, contact_block='London,\nNW1A 1AA') create_letter_contact(service=service, contact_block='London,\nNW1A 1AA')
@@ -1157,7 +1169,6 @@ def test_send_inbound_sms_to_service_retries_if_request_returns_500(notify_api,
) )
assert mocked.call_count == 1 assert mocked.call_count == 1
assert mocked.call_args[1]['queue'] == 'retry-tasks' assert mocked.call_args[1]['queue'] == 'retry-tasks'
assert exc_msg in mocked.call_args[1]['exc']
def test_send_inbound_sms_to_service_retries_if_request_throws_unknown(notify_api, sample_service, mocker): def test_send_inbound_sms_to_service_retries_if_request_throws_unknown(notify_api, sample_service, mocker):
@@ -1177,7 +1188,6 @@ def test_send_inbound_sms_to_service_retries_if_request_throws_unknown(notify_ap
) )
assert mocked.call_count == 1 assert mocked.call_count == 1
assert mocked.call_args[1]['queue'] == 'retry-tasks' assert mocked.call_args[1]['queue'] == 'retry-tasks'
assert exc_msg in mocked.call_args[1]['exc']
def test_send_inbound_sms_to_service_does_not_retries_if_request_returns_404(notify_api, sample_service, mocker): def test_send_inbound_sms_to_service_does_not_retries_if_request_returns_404(notify_api, sample_service, mocker):