mirror of
https://github.com/GSA/notifications-api.git
synced 2026-08-01 12:18:50 -04:00
Merge branch 'main' into 2199-add-pending-message-data-to-daily-and-user_daily-stats
This commit is contained in:
@@ -1,7 +1,9 @@
|
||||
import csv
|
||||
import datetime
|
||||
import re
|
||||
import time
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
from io import StringIO
|
||||
|
||||
import botocore
|
||||
from boto3 import Session
|
||||
@@ -395,31 +397,25 @@ def get_job_from_s3(service_id, job_id):
|
||||
|
||||
|
||||
def extract_phones(job):
|
||||
job = job.split("\r\n")
|
||||
first_row = job[0]
|
||||
job.pop(0)
|
||||
first_row = first_row.split(",")
|
||||
job_csv_data = StringIO(job)
|
||||
csv_reader = csv.reader(job_csv_data)
|
||||
first_row = next(csv_reader)
|
||||
|
||||
phone_index = 0
|
||||
for item in first_row:
|
||||
# Note: may contain a BOM and look like \ufeffphone number
|
||||
if item.lower() in [
|
||||
"phone number",
|
||||
"\\ufeffphone number",
|
||||
"\\ufeffphone number\n",
|
||||
"phone number\n",
|
||||
]:
|
||||
for i, item in enumerate(first_row):
|
||||
if item.lower().lstrip("\ufeff") == "phone number":
|
||||
phone_index = i
|
||||
break
|
||||
phone_index = phone_index + 1
|
||||
|
||||
phones = {}
|
||||
job_row = 0
|
||||
for row in job:
|
||||
row = row.split(",")
|
||||
for row in csv_reader:
|
||||
|
||||
if phone_index >= len(row):
|
||||
phones[job_row] = "Unavailable"
|
||||
current_app.logger.error(
|
||||
"Corrupt csv file, missing columns or possibly a byte order mark in the file",
|
||||
f"Corrupt csv file, missing columns or\
|
||||
possibly a byte order mark in the file, row looks like {row}",
|
||||
)
|
||||
|
||||
else:
|
||||
|
||||
@@ -1,107 +1,19 @@
|
||||
import json
|
||||
import os
|
||||
from datetime import timedelta
|
||||
|
||||
from botocore.exceptions import ClientError
|
||||
from flask import current_app
|
||||
from sqlalchemy.orm.exc import NoResultFound
|
||||
|
||||
from app import aws_cloudwatch_client, notify_celery, redis_store
|
||||
from app import notify_celery, redis_store
|
||||
from app.clients.email import EmailClientNonRetryableException
|
||||
from app.clients.email.aws_ses import AwsSesClientThrottlingSendRateException
|
||||
from app.clients.sms import SmsClientResponseException
|
||||
from app.config import Config, QueueNames
|
||||
from app.dao import notifications_dao
|
||||
from app.dao.notifications_dao import (
|
||||
sanitize_successful_notification_by_id,
|
||||
update_notification_status_by_id,
|
||||
)
|
||||
from app.dao.notifications_dao import update_notification_status_by_id
|
||||
from app.delivery import send_to_providers
|
||||
from app.enums import NotificationStatus
|
||||
from app.exceptions import NotificationTechnicalFailureException
|
||||
from app.utils import utc_now
|
||||
|
||||
# This is the amount of time to wait after sending an sms message before we check the aws logs and look for delivery
|
||||
# receipts
|
||||
DELIVERY_RECEIPT_DELAY_IN_SECONDS = 30
|
||||
|
||||
|
||||
@notify_celery.task(
|
||||
bind=True,
|
||||
name="check_sms_delivery_receipt",
|
||||
max_retries=48,
|
||||
default_retry_delay=300,
|
||||
)
|
||||
def check_sms_delivery_receipt(self, message_id, notification_id, sent_at):
|
||||
"""
|
||||
This is called after deliver_sms to check the status of the message. This uses the same number of
|
||||
retries and the same delay period as deliver_sms. In addition, this fires five minutes after
|
||||
deliver_sms initially. So the idea is that most messages will succeed and show up in the logs quickly.
|
||||
Other message will resolve successfully after a retry or to. A few will fail but it will take up to
|
||||
4 hours to know for sure. The call to check_sms will raise an exception if neither a success nor a
|
||||
failure appears in the cloudwatch logs, so this should keep retrying until the log appears, or until
|
||||
we run out of retries.
|
||||
"""
|
||||
# TODO the localstack cloudwatch doesn't currently have our log groups. Possibly create them with awslocal?
|
||||
if aws_cloudwatch_client.is_localstack():
|
||||
status = "success"
|
||||
provider_response = "this is a fake successful localstack sms message"
|
||||
carrier = "unknown"
|
||||
else:
|
||||
try:
|
||||
status, provider_response, carrier = aws_cloudwatch_client.check_sms(
|
||||
message_id, notification_id, sent_at
|
||||
)
|
||||
except NotificationTechnicalFailureException as ntfe:
|
||||
provider_response = "Unable to find carrier response -- still looking"
|
||||
status = "pending"
|
||||
carrier = ""
|
||||
update_notification_status_by_id(
|
||||
notification_id,
|
||||
status,
|
||||
carrier=carrier,
|
||||
provider_response=provider_response,
|
||||
)
|
||||
raise self.retry(exc=ntfe)
|
||||
except ClientError as err:
|
||||
# Probably a ThrottlingException but could be something else
|
||||
error_code = err.response["Error"]["Code"]
|
||||
provider_response = (
|
||||
f"{error_code} while checking sms receipt -- still looking"
|
||||
)
|
||||
status = "pending"
|
||||
carrier = ""
|
||||
update_notification_status_by_id(
|
||||
notification_id,
|
||||
status,
|
||||
carrier=carrier,
|
||||
provider_response=provider_response,
|
||||
)
|
||||
raise self.retry(exc=err)
|
||||
|
||||
if status == "success":
|
||||
status = NotificationStatus.DELIVERED
|
||||
elif status == "failure":
|
||||
status = NotificationStatus.FAILED
|
||||
# if status is not success or failure the client raised an exception and this method will retry
|
||||
|
||||
if status == NotificationStatus.DELIVERED:
|
||||
sanitize_successful_notification_by_id(
|
||||
notification_id, carrier=carrier, provider_response=provider_response
|
||||
)
|
||||
current_app.logger.info(
|
||||
f"Sanitized notification {notification_id} that was successfully delivered"
|
||||
)
|
||||
else:
|
||||
update_notification_status_by_id(
|
||||
notification_id,
|
||||
status,
|
||||
carrier=carrier,
|
||||
provider_response=provider_response,
|
||||
)
|
||||
current_app.logger.info(
|
||||
f"Updated notification {notification_id} with response '{provider_response}'"
|
||||
)
|
||||
|
||||
|
||||
@notify_celery.task(
|
||||
@@ -127,17 +39,8 @@ def deliver_sms(self, notification_id):
|
||||
ansi_green + f"AUTHENTICATION CODE: {notification.content}" + ansi_reset
|
||||
)
|
||||
# Code branches off to send_to_providers.py
|
||||
message_id = send_to_providers.send_sms_to_provider(notification)
|
||||
send_to_providers.send_sms_to_provider(notification)
|
||||
|
||||
# DEPRECATED
|
||||
# We have to put it in UTC. For other timezones, the delay
|
||||
# will be ignored and it will fire immediately (although this probably only affects developer testing)
|
||||
my_eta = utc_now() + timedelta(seconds=DELIVERY_RECEIPT_DELAY_IN_SECONDS)
|
||||
check_sms_delivery_receipt.apply_async(
|
||||
[message_id, notification_id, notification.created_at],
|
||||
eta=my_eta,
|
||||
queue=QueueNames.CHECK_SMS,
|
||||
)
|
||||
except Exception as e:
|
||||
update_notification_status_by_id(
|
||||
notification_id,
|
||||
|
||||
@@ -11,6 +11,7 @@ from app.celery.tasks import (
|
||||
process_job,
|
||||
process_row,
|
||||
)
|
||||
from app.clients.cloudwatch.aws_cloudwatch import AwsCloudwatchClient
|
||||
from app.config import QueueNames
|
||||
from app.dao.invited_org_user_dao import (
|
||||
delete_org_invitations_created_more_than_two_days_ago,
|
||||
@@ -22,7 +23,10 @@ from app.dao.jobs_dao import (
|
||||
find_jobs_with_missing_rows,
|
||||
find_missing_row_for_job,
|
||||
)
|
||||
from app.dao.notifications_dao import notifications_not_yet_sent
|
||||
from app.dao.notifications_dao import (
|
||||
dao_update_delivery_receipts,
|
||||
notifications_not_yet_sent,
|
||||
)
|
||||
from app.dao.services_dao import (
|
||||
dao_find_services_sending_to_tv_numbers,
|
||||
dao_find_services_with_high_failure_rates,
|
||||
@@ -32,6 +36,7 @@ from app.enums import JobStatus, NotificationType
|
||||
from app.models import Job
|
||||
from app.notifications.process_notifications import send_notification_to_queue
|
||||
from app.utils import utc_now
|
||||
from notifications_utils import aware_utcnow
|
||||
from notifications_utils.clients.zendesk.zendesk_client import NotifySupportTicket
|
||||
|
||||
MAX_NOTIFICATION_FAILS = 10000
|
||||
@@ -95,27 +100,29 @@ def check_job_status():
|
||||
select
|
||||
from jobs
|
||||
where job_status == 'in progress'
|
||||
and processing started between 30 and 35 minutes ago
|
||||
and processing started some time ago
|
||||
OR where the job_status == 'pending'
|
||||
and the job scheduled_for timestamp is between 30 and 35 minutes ago.
|
||||
and the job scheduled_for timestamp is some time ago.
|
||||
if any results then
|
||||
update the job_status to 'error'
|
||||
process the rows in the csv that are missing (in another task) just do the check here.
|
||||
"""
|
||||
thirty_minutes_ago = utc_now() - timedelta(minutes=30)
|
||||
thirty_five_minutes_ago = utc_now() - timedelta(minutes=35)
|
||||
START_MINUTES = 245
|
||||
END_MINUTES = 240
|
||||
end_minutes_ago = utc_now() - timedelta(minutes=END_MINUTES)
|
||||
start_minutes_ago = utc_now() - timedelta(minutes=START_MINUTES)
|
||||
|
||||
incomplete_in_progress_jobs = Job.query.filter(
|
||||
Job.job_status == JobStatus.IN_PROGRESS,
|
||||
between(Job.processing_started, thirty_five_minutes_ago, thirty_minutes_ago),
|
||||
between(Job.processing_started, start_minutes_ago, end_minutes_ago),
|
||||
)
|
||||
incomplete_pending_jobs = Job.query.filter(
|
||||
Job.job_status == JobStatus.PENDING,
|
||||
Job.scheduled_for.isnot(None),
|
||||
between(Job.scheduled_for, thirty_five_minutes_ago, thirty_minutes_ago),
|
||||
between(Job.scheduled_for, start_minutes_ago, end_minutes_ago),
|
||||
)
|
||||
|
||||
jobs_not_complete_after_30_minutes = (
|
||||
jobs_not_complete_after_allotted_time = (
|
||||
incomplete_in_progress_jobs.union(incomplete_pending_jobs)
|
||||
.order_by(Job.processing_started, Job.scheduled_for)
|
||||
.all()
|
||||
@@ -124,7 +131,7 @@ def check_job_status():
|
||||
# temporarily mark them as ERROR so that they don't get picked up by future check_job_status tasks
|
||||
# if they haven't been re-processed in time.
|
||||
job_ids = []
|
||||
for job in jobs_not_complete_after_30_minutes:
|
||||
for job in jobs_not_complete_after_allotted_time:
|
||||
job.job_status = JobStatus.ERROR
|
||||
dao_update_job(job)
|
||||
job_ids.append(str(job.id))
|
||||
@@ -169,9 +176,7 @@ def check_for_missing_rows_in_completed_jobs():
|
||||
for row_to_process in missing_rows:
|
||||
row = recipient_csv[row_to_process.missing_row]
|
||||
current_app.logger.info(
|
||||
"Processing missing row: {} for job: {}".format(
|
||||
row_to_process.missing_row, job.id
|
||||
)
|
||||
f"Processing missing row: {row_to_process.missing_row} for job: {job.id}"
|
||||
)
|
||||
process_row(row, template, job, job.service, sender_id=sender_id)
|
||||
|
||||
@@ -231,3 +236,45 @@ def check_for_services_with_high_failure_rates_or_sending_to_tv_numbers():
|
||||
technical_ticket=True,
|
||||
)
|
||||
zendesk_client.send_ticket_to_zendesk(ticket)
|
||||
|
||||
|
||||
@notify_celery.task(
|
||||
bind=True, max_retries=7, default_retry_delay=3600, name="process-delivery-receipts"
|
||||
)
|
||||
def process_delivery_receipts(self):
|
||||
"""
|
||||
Every eight minutes or so (see config.py) we run this task, which searches the last ten
|
||||
minutes of logs for delivery receipts and batch updates the db with the results. The overlap
|
||||
is intentional. We don't mind re-updating things, it is better than losing data.
|
||||
|
||||
We also set this to retry with exponential backoff in the case of failure. The only way this would
|
||||
fail is if, for example the db went down, or redis filled causing the app to stop processing. But if
|
||||
it does fail, we need to go back over at some point when things are running again and process those results.
|
||||
"""
|
||||
try:
|
||||
batch_size = 1000 # in theory with postgresql this could be 10k to 20k?
|
||||
|
||||
cloudwatch = AwsCloudwatchClient()
|
||||
cloudwatch.init_app(current_app)
|
||||
start_time = aware_utcnow() - timedelta(minutes=3)
|
||||
end_time = aware_utcnow()
|
||||
delivered_receipts, failed_receipts = cloudwatch.check_delivery_receipts(
|
||||
start_time, end_time
|
||||
)
|
||||
delivered_receipts = list(delivered_receipts)
|
||||
for i in range(0, len(delivered_receipts), batch_size):
|
||||
batch = delivered_receipts[i : i + batch_size]
|
||||
dao_update_delivery_receipts(batch, True)
|
||||
failed_receipts = list(failed_receipts)
|
||||
for i in range(0, len(failed_receipts), batch_size):
|
||||
batch = failed_receipts[i : i + batch_size]
|
||||
dao_update_delivery_receipts(batch, False)
|
||||
except Exception as ex:
|
||||
retry_count = self.request.retries
|
||||
wait_time = 3600 * 2**retry_count
|
||||
try:
|
||||
raise self.retry(ex=ex, countdown=wait_time)
|
||||
except self.MaxRetriesExceededError:
|
||||
current_app.logger.error(
|
||||
"Failed process delivery receipts after max retries"
|
||||
)
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import json
|
||||
from time import sleep
|
||||
|
||||
from celery.signals import task_postrun
|
||||
from flask import current_app
|
||||
@@ -38,9 +39,7 @@ def process_job(job_id, sender_id=None):
|
||||
start = utc_now()
|
||||
job = dao_get_job_by_id(job_id)
|
||||
current_app.logger.info(
|
||||
"Starting process-job task for job id {} with status: {}".format(
|
||||
job_id, job.job_status
|
||||
)
|
||||
f"Starting process-job task for job id {job_id} with status: {job.job_status}"
|
||||
)
|
||||
|
||||
if job.job_status != JobStatus.PENDING:
|
||||
@@ -56,7 +55,7 @@ def process_job(job_id, sender_id=None):
|
||||
job.job_status = JobStatus.CANCELLED
|
||||
dao_update_job(job)
|
||||
current_app.logger.warning(
|
||||
"Job {} has been cancelled, service {} is inactive".format(
|
||||
f"Job {job_id} has been cancelled, service {service.id} is inactive".format(
|
||||
job_id, service.id
|
||||
)
|
||||
)
|
||||
@@ -70,13 +69,21 @@ def process_job(job_id, sender_id=None):
|
||||
)
|
||||
|
||||
current_app.logger.info(
|
||||
"Starting job {} processing {} notifications".format(
|
||||
job_id, job.notification_count
|
||||
)
|
||||
f"Starting job {job_id} processing {job.notification_count} notifications"
|
||||
)
|
||||
|
||||
# notify-api-1495 we are going to sleep periodically to give other
|
||||
# jobs running at the same time a chance to get some of their messages
|
||||
# sent. Sleep for 1 second after every 3 sends, which gives us throughput
|
||||
# of about 3600*3 per hour and would keep the queue clear assuming only one sender.
|
||||
# It will also hopefully eliminate throttling when we send messages which we are
|
||||
# currently seeing.
|
||||
count = 0
|
||||
for row in recipient_csv.get_rows():
|
||||
process_row(row, template, job, service, sender_id=sender_id)
|
||||
count = count + 1
|
||||
if count % 3 == 0:
|
||||
sleep(1)
|
||||
|
||||
# End point/Exit point for message send flow.
|
||||
job_complete(job, start=start)
|
||||
@@ -206,9 +213,7 @@ def save_sms(self, service_id, notification_id, encrypted_notification, sender_i
|
||||
f"service not allowed to send for job_id {notification.get('job', None)}, aborting"
|
||||
)
|
||||
)
|
||||
current_app.logger.debug(
|
||||
"SMS {} failed as restricted service".format(notification_id)
|
||||
)
|
||||
current_app.logger.debug(f"SMS {notification_id} failed as restricted service")
|
||||
return
|
||||
|
||||
try:
|
||||
@@ -218,22 +223,30 @@ def save_sms(self, service_id, notification_id, encrypted_notification, sender_i
|
||||
job = dao_get_job_by_id(job_id)
|
||||
created_by_id = job.created_by_id
|
||||
|
||||
saved_notification = persist_notification(
|
||||
template_id=notification["template"],
|
||||
template_version=notification["template_version"],
|
||||
recipient=notification["to"],
|
||||
service=service,
|
||||
personalisation=notification.get("personalisation"),
|
||||
notification_type=NotificationType.SMS,
|
||||
api_key_id=None,
|
||||
key_type=KeyType.NORMAL,
|
||||
created_at=utc_now(),
|
||||
created_by_id=created_by_id,
|
||||
job_id=notification.get("job", None),
|
||||
job_row_number=notification.get("row_number", None),
|
||||
notification_id=notification_id,
|
||||
reply_to_text=reply_to_text,
|
||||
)
|
||||
try:
|
||||
saved_notification = persist_notification(
|
||||
template_id=notification["template"],
|
||||
template_version=notification["template_version"],
|
||||
recipient=notification["to"],
|
||||
service=service,
|
||||
personalisation=notification.get("personalisation"),
|
||||
notification_type=NotificationType.SMS,
|
||||
api_key_id=None,
|
||||
key_type=KeyType.NORMAL,
|
||||
created_at=utc_now(),
|
||||
created_by_id=created_by_id,
|
||||
job_id=notification.get("job", None),
|
||||
job_row_number=notification.get("row_number", None),
|
||||
notification_id=notification_id,
|
||||
reply_to_text=reply_to_text,
|
||||
)
|
||||
except IntegrityError:
|
||||
current_app.logger.warning(
|
||||
f"{NotificationType.SMS}: {notification_id} already exists."
|
||||
)
|
||||
# If we don't have the return statement here, we will fall through and end
|
||||
# up retrying because IntegrityError is a subclass of SQLAlchemyError
|
||||
return
|
||||
|
||||
# Kick off sns process in provider_tasks.py
|
||||
sn = saved_notification
|
||||
@@ -247,11 +260,8 @@ def save_sms(self, service_id, notification_id, encrypted_notification, sender_i
|
||||
)
|
||||
|
||||
current_app.logger.debug(
|
||||
"SMS {} created at {} for job {}".format(
|
||||
saved_notification.id,
|
||||
saved_notification.created_at,
|
||||
notification.get("job", None),
|
||||
)
|
||||
f"SMS {saved_notification.id} created at {saved_notification.created_at} "
|
||||
f"for job {notification.get('job', None)}"
|
||||
)
|
||||
|
||||
except SQLAlchemyError as e:
|
||||
|
||||
@@ -1,15 +1,12 @@
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
from datetime import timedelta
|
||||
|
||||
from boto3 import client
|
||||
from flask import current_app
|
||||
|
||||
from app.clients import AWS_CLIENT_CONFIG, Client
|
||||
from app.cloudfoundry_config import cloud_config
|
||||
from app.exceptions import NotificationTechnicalFailureException
|
||||
from app.utils import hilite, utc_now
|
||||
|
||||
|
||||
class AwsCloudwatchClient(Client):
|
||||
@@ -49,48 +46,32 @@ class AwsCloudwatchClient(Client):
|
||||
def is_localstack(self):
|
||||
return self._is_localstack
|
||||
|
||||
def _get_log(self, my_filter, log_group_name, sent_at):
|
||||
def _get_log(self, log_group_name, start, end):
|
||||
# Check all cloudwatch logs from the time the notification was sent (currently 5 minutes previously) until now
|
||||
now = utc_now()
|
||||
beginning = sent_at
|
||||
next_token = None
|
||||
all_log_events = []
|
||||
current_app.logger.info(f"START TIME {beginning} END TIME {now}")
|
||||
# There has been a change somewhere and the time range we were previously using has become too
|
||||
# narrow or wrong in some way, so events can't be found. For the time being, adjust by adding
|
||||
# a buffer on each side of 12 hours.
|
||||
TWELVE_HOURS = 12 * 60 * 60 * 1000
|
||||
|
||||
while True:
|
||||
if next_token:
|
||||
response = self._client.filter_log_events(
|
||||
logGroupName=log_group_name,
|
||||
filterPattern=my_filter,
|
||||
nextToken=next_token,
|
||||
startTime=int(beginning.timestamp() * 1000) - TWELVE_HOURS,
|
||||
endTime=int(now.timestamp() * 1000) + TWELVE_HOURS,
|
||||
startTime=int(start.timestamp() * 1000),
|
||||
endTime=int(end.timestamp() * 1000),
|
||||
)
|
||||
else:
|
||||
response = self._client.filter_log_events(
|
||||
logGroupName=log_group_name,
|
||||
filterPattern=my_filter,
|
||||
startTime=int(beginning.timestamp() * 1000) - TWELVE_HOURS,
|
||||
endTime=int(now.timestamp() * 1000) + TWELVE_HOURS,
|
||||
startTime=int(start.timestamp() * 1000),
|
||||
endTime=int(end.timestamp() * 1000),
|
||||
)
|
||||
log_events = response.get("events", [])
|
||||
all_log_events.extend(log_events)
|
||||
if len(log_events) > 0:
|
||||
# We found it
|
||||
|
||||
break
|
||||
next_token = response.get("nextToken")
|
||||
if not next_token:
|
||||
break
|
||||
return all_log_events
|
||||
|
||||
def _extract_account_number(self, ses_domain_arn):
|
||||
account_number = ses_domain_arn.split(":")
|
||||
return account_number
|
||||
|
||||
def warn_if_dev_is_opted_out(self, provider_response, notification_id):
|
||||
if (
|
||||
"is opted out" in provider_response.lower()
|
||||
@@ -108,60 +89,108 @@ class AwsCloudwatchClient(Client):
|
||||
return logline
|
||||
return None
|
||||
|
||||
def check_sms(self, message_id, notification_id, created_at):
|
||||
def _extract_account_number(self, ses_domain_arn):
|
||||
account_number = ses_domain_arn.split(":")
|
||||
return account_number
|
||||
|
||||
def event_to_db_format(self, event):
|
||||
|
||||
# massage the data into the form the db expects. When we switch
|
||||
# from filter_log_events to log insights this will be convenient
|
||||
if isinstance(event, str):
|
||||
event = json.loads(event)
|
||||
|
||||
# Don't trust AWS to always send the same JSON structure back
|
||||
# However, if we don't get message_id and status we might as well blow up
|
||||
# because it's pointless to continue
|
||||
phone_carrier = self._aws_value_or_default(event, "delivery", "phoneCarrier")
|
||||
provider_response = self._aws_value_or_default(
|
||||
event, "delivery", "providerResponse"
|
||||
)
|
||||
my_timestamp = self._aws_value_or_default(event, "notification", "timestamp")
|
||||
return {
|
||||
"notification.messageId": event["notification"]["messageId"],
|
||||
"status": event["status"],
|
||||
"delivery.phoneCarrier": phone_carrier,
|
||||
"delivery.providerResponse": provider_response,
|
||||
"@timestamp": my_timestamp,
|
||||
}
|
||||
|
||||
# Here is an example of how to get the events with log insights
|
||||
# def do_log_insights():
|
||||
# query = """
|
||||
# fields @timestamp, status, message, recipient
|
||||
# | filter status = "DELIVERED"
|
||||
# | sort @timestamp asc
|
||||
# """
|
||||
# temp_client = boto3.client(
|
||||
# "logs",
|
||||
# region_name="us-gov-west-1",
|
||||
# aws_access_key_id=AWS_ACCESS_KEY_ID,
|
||||
# aws_secret_access_key=AWS_SECRET_ACCESS_KEY,
|
||||
# config=AWS_CLIENT_CONFIG,
|
||||
# )
|
||||
# start = utc_now()
|
||||
# end = utc_now - timedelta(hours=1)
|
||||
# response = temp_client.start_query(
|
||||
# logGroupName = LOG_GROUP_NAME_DELIVERED,
|
||||
# startTime = int(start.timestamp()),
|
||||
# endTime= int(end.timestamp()),
|
||||
# queryString = query
|
||||
|
||||
# )
|
||||
# query_id = response['queryId']
|
||||
# while True:
|
||||
# result = temp_client.get_query_results(queryId=query_id)
|
||||
# if result['status'] == 'Complete':
|
||||
# break
|
||||
# time.sleep(1)
|
||||
|
||||
# delivery_receipts = []
|
||||
# for log in result['results']:
|
||||
# receipt = {field['field']: field['value'] for field in log}
|
||||
# delivery_receipts.append(receipt)
|
||||
# print(receipt)
|
||||
|
||||
# print(len(delivery_receipts))
|
||||
|
||||
# In the long run we want to use Log Insights because it is more efficient
|
||||
# that filter_log_events. But we are blocked by a permissions issue in the broker.
|
||||
# So for now, use filter_log_events and grab all log_events over a 10 minute interval,
|
||||
# and run this on a schedule.
|
||||
def check_delivery_receipts(self, start, end):
|
||||
region = cloud_config.sns_region
|
||||
# TODO this clumsy approach to getting the account number will be fixed as part of notify-api #258
|
||||
account_number = self._extract_account_number(cloud_config.ses_domain_arn)
|
||||
|
||||
time_now = utc_now()
|
||||
log_group_name = f"sns/{region}/{account_number[4]}/DirectPublishToPhoneNumber"
|
||||
filter_pattern = '{$.notification.messageId="XXXXX"}'
|
||||
filter_pattern = filter_pattern.replace("XXXXX", message_id)
|
||||
all_log_events = self._get_log(filter_pattern, log_group_name, created_at)
|
||||
if all_log_events and len(all_log_events) > 0:
|
||||
event = all_log_events[0]
|
||||
message = json.loads(event["message"])
|
||||
self.warn_if_dev_is_opted_out(
|
||||
message["delivery"]["providerResponse"], notification_id
|
||||
)
|
||||
# Here we map the answer from aws to the message_id.
|
||||
# Previously, in send_to_providers, we mapped the job_id and row number
|
||||
# to the message id. And on the admin side we mapped the csv filename
|
||||
# to the job_id. So by tracing through all the logs we can go:
|
||||
# filename->job_id->message_id->what really happened
|
||||
current_app.logger.info(
|
||||
hilite(f"DELIVERED: {message} for message_id {message_id}")
|
||||
)
|
||||
return (
|
||||
"success",
|
||||
message["delivery"]["providerResponse"],
|
||||
message["delivery"].get("phoneCarrier", "Unknown Carrier"),
|
||||
)
|
||||
|
||||
delivered_event_set = self._get_receipts(log_group_name, start, end)
|
||||
current_app.logger.info(
|
||||
(f"Delivered message count: {len(delivered_event_set)}")
|
||||
)
|
||||
log_group_name = (
|
||||
f"sns/{region}/{account_number[4]}/DirectPublishToPhoneNumber/Failure"
|
||||
)
|
||||
all_failed_events = self._get_log(filter_pattern, log_group_name, created_at)
|
||||
if all_failed_events and len(all_failed_events) > 0:
|
||||
event = all_failed_events[0]
|
||||
message = json.loads(event["message"])
|
||||
self.warn_if_dev_is_opted_out(
|
||||
message["delivery"]["providerResponse"], notification_id
|
||||
)
|
||||
failed_event_set = self._get_receipts(log_group_name, start, end)
|
||||
current_app.logger.info((f"Failed message count: {len(failed_event_set)}"))
|
||||
|
||||
current_app.logger.info(
|
||||
hilite(f"FAILED: {message} for message_id {message_id}")
|
||||
)
|
||||
return (
|
||||
"failure",
|
||||
message["delivery"]["providerResponse"],
|
||||
message["delivery"].get("phoneCarrier", "Unknown Carrier"),
|
||||
)
|
||||
return delivered_event_set, failed_event_set
|
||||
|
||||
if time_now > (created_at + timedelta(hours=3)):
|
||||
# see app/models.py Notification. This message corresponds to "permanent-failure",
|
||||
# but we are copy/pasting here to avoid circular imports.
|
||||
return "failure", "Unable to find carrier response."
|
||||
raise NotificationTechnicalFailureException(
|
||||
f"No event found for message_id {message_id} notification_id {notification_id}"
|
||||
)
|
||||
def _get_receipts(self, log_group_name, start, end):
|
||||
event_set = set()
|
||||
all_events = self._get_log(log_group_name, start, end)
|
||||
for event in all_events:
|
||||
try:
|
||||
actual_event = self.event_to_db_format(event["message"])
|
||||
event_set.add(json.dumps(actual_event))
|
||||
except Exception:
|
||||
current_app.logger.exception(
|
||||
f"Could not format delivery receipt {event} for db insert"
|
||||
)
|
||||
return event_set
|
||||
|
||||
def _aws_value_or_default(self, event, top_level, second_level):
|
||||
if event.get(top_level) is None or event[top_level].get(second_level) is None:
|
||||
my_var = ""
|
||||
else:
|
||||
my_var = event[top_level][second_level]
|
||||
|
||||
return my_var
|
||||
|
||||
@@ -198,6 +198,11 @@ class Config(object):
|
||||
"schedule": timedelta(minutes=63),
|
||||
"options": {"queue": QueueNames.PERIODIC},
|
||||
},
|
||||
"process-delivery-receipts": {
|
||||
"task": "process-delivery-receipts",
|
||||
"schedule": timedelta(minutes=2),
|
||||
"options": {"queue": QueueNames.PERIODIC},
|
||||
},
|
||||
"expire-or-delete-invitations": {
|
||||
"task": "expire-or-delete-invitations",
|
||||
"schedule": timedelta(minutes=66),
|
||||
|
||||
@@ -1,7 +1,21 @@
|
||||
import json
|
||||
from datetime import timedelta
|
||||
from time import time
|
||||
|
||||
from flask import current_app
|
||||
from sqlalchemy import asc, delete, desc, func, or_, select, text, union, update
|
||||
from sqlalchemy import (
|
||||
TIMESTAMP,
|
||||
asc,
|
||||
cast,
|
||||
delete,
|
||||
desc,
|
||||
func,
|
||||
or_,
|
||||
select,
|
||||
text,
|
||||
union,
|
||||
update,
|
||||
)
|
||||
from sqlalchemy.orm import joinedload
|
||||
from sqlalchemy.orm.exc import NoResultFound
|
||||
from sqlalchemy.sql import functions
|
||||
@@ -52,6 +66,12 @@ def dao_get_last_date_template_was_used(template_id, service_id):
|
||||
return last_date
|
||||
|
||||
|
||||
def dao_notification_exists(notification_id) -> bool:
|
||||
stmt = select(Notification).where(Notification.id == notification_id)
|
||||
result = db.session.execute(stmt).scalar()
|
||||
return result is not None
|
||||
|
||||
|
||||
@autocommit
|
||||
def dao_create_notification(notification):
|
||||
if not notification.id:
|
||||
@@ -73,9 +93,7 @@ def dao_create_notification(notification):
|
||||
notification.normalised_to = "1"
|
||||
|
||||
# notify-api-1454 insert only if it doesn't exist
|
||||
stmt = select(Notification).where(Notification.id == notification.id)
|
||||
result = db.session.execute(stmt).scalar()
|
||||
if result is None:
|
||||
if not dao_notification_exists(notification.id):
|
||||
db.session.add(notification)
|
||||
|
||||
|
||||
@@ -707,3 +725,58 @@ def get_service_ids_with_notifications_on_date(notification_type, date):
|
||||
union(notification_table_query, ft_status_table_query).subquery()
|
||||
).distinct()
|
||||
}
|
||||
|
||||
|
||||
def dao_update_delivery_receipts(receipts, delivered):
|
||||
start_time_millis = time() * 1000
|
||||
new_receipts = []
|
||||
for r in receipts:
|
||||
if isinstance(r, str):
|
||||
r = json.loads(r)
|
||||
new_receipts.append(r)
|
||||
|
||||
receipts = new_receipts
|
||||
|
||||
id_to_carrier = {
|
||||
r["notification.messageId"]: r["delivery.phoneCarrier"] for r in receipts
|
||||
}
|
||||
id_to_provider_response = {
|
||||
r["notification.messageId"]: r["delivery.providerResponse"] for r in receipts
|
||||
}
|
||||
id_to_timestamp = {r["notification.messageId"]: r["@timestamp"] for r in receipts}
|
||||
|
||||
status_to_update_with = NotificationStatus.DELIVERED
|
||||
if not delivered:
|
||||
status_to_update_with = NotificationStatus.FAILED
|
||||
stmt = (
|
||||
update(Notification)
|
||||
.where(Notification.message_id.in_(id_to_carrier.keys()))
|
||||
.values(
|
||||
carrier=case(
|
||||
*[
|
||||
(Notification.message_id == key, value)
|
||||
for key, value in id_to_carrier.items()
|
||||
]
|
||||
),
|
||||
status=status_to_update_with,
|
||||
sent_at=case(
|
||||
*[
|
||||
(Notification.message_id == key, cast(value, TIMESTAMP))
|
||||
for key, value in id_to_timestamp.items()
|
||||
]
|
||||
),
|
||||
provider_response=case(
|
||||
*[
|
||||
(Notification.message_id == key, value)
|
||||
for key, value in id_to_provider_response.items()
|
||||
]
|
||||
),
|
||||
)
|
||||
)
|
||||
db.session.execute(stmt)
|
||||
db.session.commit()
|
||||
elapsed_time = (time() * 1000) - start_time_millis
|
||||
current_app.logger.info(
|
||||
f"#loadtestperformance batch update query time: \
|
||||
updated {len(receipts)} notification in {elapsed_time} ms"
|
||||
)
|
||||
|
||||
@@ -577,7 +577,16 @@ class Service(db.Model, Versioned):
|
||||
return self.inbound_number.number
|
||||
|
||||
def get_default_sms_sender(self):
|
||||
default_sms_sender = [x for x in self.service_sms_senders if x.is_default]
|
||||
# notify-api-1513 let's try a minimalistic fix
|
||||
# to see if we can get the right numbers back
|
||||
default_sms_sender = [
|
||||
x
|
||||
for x in self.service_sms_senders
|
||||
if x.is_default and x.service_id == self.id
|
||||
]
|
||||
current_app.logger.info(
|
||||
f"#notify-api-1513 senders for service {self.name} are {self.service_sms_senders}"
|
||||
)
|
||||
return default_sms_sender[0].sms_sender
|
||||
|
||||
def get_default_reply_to_email_address(self):
|
||||
|
||||
@@ -8,6 +8,7 @@ from app.config import QueueNames
|
||||
from app.dao.notifications_dao import (
|
||||
dao_create_notification,
|
||||
dao_delete_notifications_by_id,
|
||||
dao_notification_exists,
|
||||
get_notification_by_id,
|
||||
)
|
||||
from app.enums import KeyType, NotificationStatus, NotificationType
|
||||
@@ -153,6 +154,10 @@ def persist_notification(
|
||||
return notification
|
||||
|
||||
|
||||
def notification_exists(notification_id):
|
||||
return dao_notification_exists(notification_id)
|
||||
|
||||
|
||||
def send_notification_to_queue_detached(
|
||||
key_type, notification_type, notification_id, queue=None
|
||||
):
|
||||
|
||||
@@ -54,7 +54,7 @@ def _create_service_invite(invited_user, nonce, state):
|
||||
data["invited_user_email"] = invited_user.email_address
|
||||
|
||||
invite_redis_key = f"invite-data-{unquote(state)}"
|
||||
redis_store.set(invite_redis_key, get_user_data_url_safe(data))
|
||||
redis_store.set(invite_redis_key, get_user_data_url_safe(data), ex=2 * 24 * 60 * 60)
|
||||
|
||||
url = os.environ["LOGIN_DOT_GOV_REGISTRATION_URL"]
|
||||
|
||||
|
||||
Reference in New Issue
Block a user