mirror of
https://github.com/GSA/notifications-api.git
synced 2026-02-01 07:35:34 -05:00
The send_delivery_status_to_service task was refactor to take the details of the notification and service api callback such that the task no longer needed to go to the database to provide the status update.
This PR removes the code that is no longer used. This extra step was necessary to keep the tasks backward compatible.
This commit is contained in:
@@ -1,119 +1,52 @@
|
||||
import json
|
||||
|
||||
from flask import current_app
|
||||
from notifications_utils.statsd_decorators import statsd
|
||||
|
||||
from app import (
|
||||
db,
|
||||
DATETIME_FORMAT,
|
||||
notify_celery,
|
||||
encryption
|
||||
)
|
||||
from app.dao.notifications_dao import (
|
||||
get_notification_by_id,
|
||||
)
|
||||
|
||||
from app.dao.service_callback_api_dao import get_service_callback_api_for_service
|
||||
from requests import (
|
||||
HTTPError,
|
||||
request,
|
||||
RequestException
|
||||
)
|
||||
from flask import current_app
|
||||
|
||||
from app import (
|
||||
notify_celery,
|
||||
encryption
|
||||
)
|
||||
from app.config import QueueNames
|
||||
|
||||
|
||||
@notify_celery.task(bind=True, name="send-delivery-status", max_retries=5, default_retry_delay=300)
|
||||
@statsd(namespace="tasks")
|
||||
def send_delivery_status_to_service(self, notification_id,
|
||||
encrypted_status_update=None
|
||||
encrypted_status_update
|
||||
):
|
||||
if not encrypted_status_update:
|
||||
process_update_with_notification_id(self, notification_id=notification_id)
|
||||
else:
|
||||
try:
|
||||
status_update = encryption.decrypt(encrypted_status_update)
|
||||
|
||||
data = {
|
||||
"id": str(notification_id),
|
||||
"reference": status_update['notification_client_reference'],
|
||||
"to": status_update['notification_to'],
|
||||
"status": status_update['notification_status'],
|
||||
"created_at": status_update['notification_created_at'],
|
||||
"completed_at": status_update['notification_updated_at'],
|
||||
"sent_at": status_update['notification_sent_at'],
|
||||
"notification_type": status_update['notification_type']
|
||||
}
|
||||
|
||||
response = request(
|
||||
method="POST",
|
||||
url=status_update['service_callback_api_url'],
|
||||
data=json.dumps(data),
|
||||
headers={
|
||||
'Content-Type': 'application/json',
|
||||
'Authorization': 'Bearer {}'.format(status_update['service_callback_api_bearer_token'])
|
||||
},
|
||||
timeout=60
|
||||
)
|
||||
current_app.logger.info('send_delivery_status_to_service sending {} to {}, response {}'.format(
|
||||
notification_id,
|
||||
status_update['service_callback_api_url'],
|
||||
response.status_code
|
||||
))
|
||||
response.raise_for_status()
|
||||
except RequestException as e:
|
||||
current_app.logger.warning(
|
||||
"send_delivery_status_to_service request failed for notification_id: {} and url: {}. exc: {}".format(
|
||||
notification_id,
|
||||
status_update['service_callback_api_url'],
|
||||
e
|
||||
)
|
||||
)
|
||||
if not isinstance(e, HTTPError) or e.response.status_code >= 500:
|
||||
try:
|
||||
self.retry(queue=QueueNames.RETRY)
|
||||
except self.MaxRetriesExceededError:
|
||||
current_app.logger.exception(
|
||||
"""Retry: send_delivery_status_to_service has retried the max num of times
|
||||
for notification: {}""".format(notification_id)
|
||||
)
|
||||
|
||||
|
||||
def process_update_with_notification_id(self, notification_id):
|
||||
retry = False
|
||||
try:
|
||||
notification = get_notification_by_id(notification_id)
|
||||
service_callback_api = get_service_callback_api_for_service(service_id=notification.service_id)
|
||||
if not service_callback_api:
|
||||
# No delivery receipt API info set
|
||||
return
|
||||
|
||||
# Release DB connection before performing an external HTTP request
|
||||
db.session.close()
|
||||
status_update = encryption.decrypt(encrypted_status_update)
|
||||
|
||||
data = {
|
||||
"id": str(notification_id),
|
||||
"reference": str(notification.client_reference),
|
||||
"to": notification.to,
|
||||
"status": notification.status,
|
||||
"created_at": notification.created_at.strftime(DATETIME_FORMAT),
|
||||
"completed_at": notification.updated_at.strftime(DATETIME_FORMAT),
|
||||
"sent_at": notification.sent_at.strftime(DATETIME_FORMAT),
|
||||
"notification_type": notification.notification_type
|
||||
"reference": status_update['notification_client_reference'],
|
||||
"to": status_update['notification_to'],
|
||||
"status": status_update['notification_status'],
|
||||
"created_at": status_update['notification_created_at'],
|
||||
"completed_at": status_update['notification_updated_at'],
|
||||
"sent_at": status_update['notification_sent_at'],
|
||||
"notification_type": status_update['notification_type']
|
||||
}
|
||||
|
||||
response = request(
|
||||
method="POST",
|
||||
url=service_callback_api.url,
|
||||
url=status_update['service_callback_api_url'],
|
||||
data=json.dumps(data),
|
||||
headers={
|
||||
'Content-Type': 'application/json',
|
||||
'Authorization': 'Bearer {}'.format(service_callback_api.bearer_token)
|
||||
'Authorization': 'Bearer {}'.format(status_update['service_callback_api_bearer_token'])
|
||||
},
|
||||
timeout=60
|
||||
)
|
||||
current_app.logger.info('send_delivery_status_to_service sending {} to {}, response {}'.format(
|
||||
notification_id,
|
||||
service_callback_api.url,
|
||||
status_update['service_callback_api_url'],
|
||||
response.status_code
|
||||
))
|
||||
response.raise_for_status()
|
||||
@@ -121,26 +54,18 @@ def process_update_with_notification_id(self, notification_id):
|
||||
current_app.logger.warning(
|
||||
"send_delivery_status_to_service request failed for notification_id: {} and url: {}. exc: {}".format(
|
||||
notification_id,
|
||||
service_callback_api.url,
|
||||
status_update['service_callback_api_url'],
|
||||
e
|
||||
)
|
||||
)
|
||||
if not isinstance(e, HTTPError) or e.response.status_code >= 500:
|
||||
retry = True
|
||||
except Exception as e:
|
||||
current_app.logger.exception(
|
||||
'Unhandled exception when sending callback for notification {}'.format(notification_id)
|
||||
)
|
||||
retry = True
|
||||
|
||||
if retry:
|
||||
try:
|
||||
self.retry(queue=QueueNames.RETRY)
|
||||
except self.MaxRetriesExceededError:
|
||||
current_app.logger.exception(
|
||||
"""Retry: send_delivery_status_to_service has retried the max num of times
|
||||
for notification: {}""".format(notification_id)
|
||||
)
|
||||
try:
|
||||
self.retry(queue=QueueNames.RETRY)
|
||||
except self.MaxRetriesExceededError:
|
||||
current_app.logger.exception(
|
||||
"""Retry: send_delivery_status_to_service has retried the max num of times
|
||||
for notification: {}""".format(notification_id)
|
||||
)
|
||||
|
||||
|
||||
def create_encrypted_callback_data(notification, service_callback_api):
|
||||
|
||||
Reference in New Issue
Block a user