Added statsd integration into the API

- new client for statsd, follows conventions used elsewhere for configuration
- client wraps underlying library so we can use a config property to send/not send statsd

Added statsd metrics for:
- count of API successful calls SMS/Email
- count of successful task execution for SMS/Email
- count of errors from Client libraries
- timing of API calls to third party clients
- timing of how long messages live on the SQS queue
This commit is contained in:
Martyn Inglis
2016-05-13 17:15:39 +01:00
parent 3c8e45093c
commit 3f7559b286
17 changed files with 234 additions and 74 deletions

View File

@@ -1,5 +1,7 @@
import uuid
import os
import statsd
from flask import request, url_for, g
from flask import Flask, _request_ctx_stack
from flask.ext.sqlalchemy import SQLAlchemy
@@ -10,10 +12,10 @@ from notifications_utils import logging
from app.celery.celery import NotifyCelery
from app.clients import Clients
from app.clients.sms.mmg import MMGClient
from app.clients.sms.twilio import TwilioClient
from app.clients.sms.firetext import FiretextClient
from app.clients.sms.loadtesting import LoadtestingClient
from app.clients.email.aws_ses import AwsSesClient
from app.clients.statsd.statsd_client import StatsdClient
from app.encryption import Encryption
DATETIME_FORMAT = "%Y-%m-%dT%H:%M:%S.%f"
@@ -22,12 +24,12 @@ DATE_FORMAT = "%Y-%m-%d"
db = SQLAlchemy()
ma = Marshmallow()
notify_celery = NotifyCelery()
twilio_client = TwilioClient()
firetext_client = FiretextClient()
loadtest_client = LoadtestingClient()
mmg_client = MMGClient()
aws_ses_client = AwsSesClient()
encryption = Encryption()
statsd_client = StatsdClient()
clients = Clients()
@@ -47,11 +49,11 @@ def create_app(app_name=None):
ma.init_app(application)
init_app(application)
logging.init_app(application)
twilio_client.init_app(application)
firetext_client.init_app(application)
loadtest_client.init_app(application)
mmg_client.init_app(application.config)
aws_ses_client.init_app(application.config['AWS_REGION'])
statsd_client.init_app(application)
firetext_client.init_app(application, statsd_client=statsd_client)
loadtest_client.init_app(application, statsd_client=statsd_client)
mmg_client.init_app(application.config, statsd_client=statsd_client)
aws_ses_client.init_app(application.config['AWS_REGION'], statsd_client=statsd_client)
notify_celery.init_app(application)
encryption.init_app(application)
clients.init_app(sms_clients=[firetext_client, mmg_client, loadtest_client], email_clients=[aws_ses_client])

View File

@@ -3,7 +3,7 @@ from datetime import datetime
from flask import current_app
from sqlalchemy.exc import SQLAlchemyError
from app import clients
from app import clients, statsd_client
from app.clients.email import EmailClientException
from app.clients.sms import SmsClientException
from app.dao.services_dao import dao_fetch_service_by_id
@@ -190,6 +190,7 @@ def process_job(job_id):
current_app.logger.info(
"Job {} created at {} started at {} finished at {}".format(job_id, job.created_at, start, finished)
)
statsd_client.incr("notifications.tasks.process-job")
@notify_celery.task(name="remove-job")
@@ -235,7 +236,11 @@ def send_sms(service_id, notification_id, encrypted_notification, created_at):
sent_by=provider.get_name(),
content_char_count=template.replaced_content_count
)
statsd_client.timing_with_dates(
"notifications.tasks.send-sms.queued-for",
sent_at,
datetime.strptime(created_at, DATETIME_FORMAT)
)
dao_create_notification(notification_db_object, TEMPLATE_TYPE_SMS, provider.get_name())
if restricted:
@@ -259,6 +264,7 @@ def send_sms(service_id, notification_id, encrypted_notification, created_at):
current_app.logger.info(
"SMS {} created at {} sent at {}".format(notification_id, created_at, sent_at)
)
statsd_client.incr("notifications.tasks.send-sms")
except SQLAlchemyError as e:
current_app.logger.exception(e)
@@ -293,6 +299,11 @@ def send_email(service_id, notification_id, from_address, encrypted_notification
)
dao_create_notification(notification_db_object, TEMPLATE_TYPE_EMAIL, provider.get_name())
statsd_client.timing_with_dates(
"notifications.tasks.send-email.queued-for",
sent_at,
datetime.strptime(created_at, DATETIME_FORMAT)
)
if restricted:
return
@@ -309,15 +320,18 @@ def send_email(service_id, notification_id, from_address, encrypted_notification
body=template.replaced_govuk_escaped,
html_body=template.as_HTML_email,
)
update_notification_reference_by_id(notification_id, reference)
except EmailClientException as e:
current_app.logger.exception(e)
notification_db_object.status = 'failed'
dao_update_notification(notification_db_object)
dao_update_notification(notification_db_object)
current_app.logger.info(
"Email {} created at {} sent at {}".format(notification_id, created_at, sent_at)
)
statsd_client.incr("notifications.tasks.send-email")
except SQLAlchemyError as e:
current_app.logger.exception(e)

View File

@@ -39,10 +39,11 @@ class AwsSesClient(EmailClient):
Amazon SES email client.
'''
def init_app(self, region, *args, **kwargs):
def init_app(self, region, statsd_client, *args, **kwargs):
self._client = boto3.client('ses', region_name=region)
super(AwsSesClient, self).__init__(*args, **kwargs)
self.name = 'ses'
self.statsd_client = statsd_client
def get_name(self):
return self.name
@@ -88,7 +89,9 @@ class AwsSesClient(EmailClient):
ReplyToAddresses=reply_to_addresses)
elapsed_time = monotonic() - start_time
current_app.logger.info("AWS SES request finished in {}".format(elapsed_time))
self.statsd_client.timing("notifications.clients.ses.request-time", elapsed_time)
return response['MessageId']
except Exception as e:
# TODO logging exceptions
self.statsd_client.incr("notifications.clients.ses.error")
raise AwsSesClientException(str(e))

View File

@@ -50,11 +50,12 @@ class FiretextClient(SmsClient):
FireText sms client.
'''
def init_app(self, config, *args, **kwargs):
def init_app(self, config, statsd_client, *args, **kwargs):
super(SmsClient, self).__init__(*args, **kwargs)
self.api_key = config.config.get('FIRETEXT_API_KEY')
self.from_number = config.config.get('FIRETEXT_NUMBER')
self.name = 'firetext'
self.statsd_client = statsd_client
def get_name(self):
return self.name
@@ -90,8 +91,10 @@ class FiretextClient(SmsClient):
api_error.message
)
)
self.statsd_client.incr("notifications.clients.firetext.error")
raise api_error
finally:
elapsed_time = monotonic() - start_time
current_app.logger.info("Firetext request finished in {}".format(elapsed_time))
self.statsd_client.timing("notifications.clients.firetext.request-time", elapsed_time)
return response

View File

@@ -11,8 +11,9 @@ class LoadtestingClient(FiretextClient):
Loadtest sms client.
'''
def init_app(self, config, *args, **kwargs):
def init_app(self, config, statsd_client, *args, **kwargs):
super(FiretextClient, self).__init__(*args, **kwargs)
self.api_key = config.config.get('LOADTESTING_API_KEY')
self.from_number = config.config.get('LOADTESTING_NUMBER')
self.name = 'loadtesting'
self.statsd_client = statsd_client

View File

@@ -39,11 +39,12 @@ class MMGClient(SmsClient):
MMG sms client
'''
def init_app(self, config, *args, **kwargs):
def init_app(self, config, statsd_client, *args, **kwargs):
super(SmsClient, self).__init__(*args, **kwargs)
self.api_key = config.get('MMG_API_KEY')
self.from_number = config.get('MMG_FROM_NUMBER')
self.name = 'mmg'
self.statsd_client = statsd_client
def get_name(self):
return self.name
@@ -78,8 +79,10 @@ class MMGClient(SmsClient):
api_error.message
)
)
self.statsd_client.incr("notifications.clients.mmg.error")
raise api_error
finally:
elapsed_time = monotonic() - start_time
self.statsd_client.timing("notifications.clients.mmg.request-time", elapsed_time)
current_app.logger.info("MMG request finished in {}".format(elapsed_time))
return response

View File

@@ -1,54 +0,0 @@
from monotonic import monotonic
from app.clients.sms import (
SmsClient, SmsClientException)
from twilio.rest import TwilioRestClient
from twilio import TwilioRestException
from flask import current_app
class TwilioClientException(SmsClientException):
pass
class TwilioClient(SmsClient):
'''
Twilio sms client.
'''
def init_app(self, config, *args, **kwargs):
super(TwilioClient, self).__init__(*args, **kwargs)
self.client = TwilioRestClient(
config.config.get('TWILIO_ACCOUNT_SID'),
config.config.get('TWILIO_AUTH_TOKEN'))
self.from_number = config.config.get('TWILIO_NUMBER')
self.name = 'twilio'
def get_name(self):
return self.name
def send_sms(self, to, content):
start_time = monotonic()
try:
response = self.client.messages.create(
body=content,
to=to,
from_=self.from_number
)
return response.sid
except TwilioRestException as e:
current_app.logger.exception(e)
raise TwilioClientException(e)
finally:
elapsed_time = monotonic() - start_time
current_app.logger.info("Twilio request finished in {}".format(elapsed_time))
def status(self, message_id):
try:
response = self.client.messages.get(message_id)
if response.status in ('delivered', 'failed'):
return response.status
elif response.status == 'undelivered':
return 'sending'
return None
except TwilioRestException as e:
current_app.logger.exception(e)
raise TwilioClientException(e)

View File

View File

@@ -0,0 +1,26 @@
from statsd import StatsClient
class StatsdClient(StatsClient):
def init_app(self, app, *args, **kwargs):
StatsClient.__init__(
self,
app.config.get('STATSD_HOST'),
app.config.get('STATSD_PORT'),
prefix=app.config.get('STATSD_PREFIX')
)
self.active = app.config.get('STATSD_ENABLED')
def incr(self, stat, count=1, rate=1):
if self.active:
super(StatsClient, self).incr(stat, count, rate)
def timing(self, stat, delta, rate=1):
if self.active:
super(StatsClient, self).timing(stat, delta, rate)
def timing_with_dates(self, stat, start, end, rate=1):
if self.active:
delta = (start - end).total_seconds() * 1000
super(StatsClient, self).timing(stat, delta, rate)

View File

@@ -1,5 +1,7 @@
import uuid
from flask import current_app
from app import statsd_client
from app.dao import notifications_dao
from app.clients.sms.firetext import get_firetext_responses
from app.clients.sms.mmg import get_mmg_responses
@@ -67,5 +69,6 @@ def process_sms_client_response(status, reference, client_name):
reference,
notification_status_message))
statsd_client.incr('notifications.callback.{}.{}'.format(client_name.lower(), notification_statistics_status))
success = "{} callback succeeded. reference {} updated".format(client_name, reference)
return success, errors

View File

@@ -1,5 +1,5 @@
from datetime import datetime
import statsd
import itertools
from flask import (
Blueprint,
@@ -12,7 +12,7 @@ from flask import (
from notifications_utils.recipients import allowed_to_send_to, first_column_heading
from notifications_utils.template import Template
from app.clients.email.aws_ses import get_aws_responses
from app import api_user, encryption, create_uuid, DATETIME_FORMAT, DATE_FORMAT
from app import api_user, encryption, create_uuid, DATETIME_FORMAT, DATE_FORMAT, statsd_client
from app.authentication.auth import require_admin
from app.dao import (
templates_dao,
@@ -103,6 +103,7 @@ def process_ses_response():
)
)
statsd_client.incr('notifications.callback.ses.{}'.format(notification_statistics_status))
return jsonify(
result="success", message="SES callback succeeded"
), 200
@@ -374,4 +375,5 @@ def send_notification(notification_type):
datetime.utcnow().strftime(DATETIME_FORMAT)
), queue='email')
statsd_client.incr('notifications.api.{}'.format(notification_type))
return jsonify(data={"notification": {"id": notification_id}}), 201

View File

@@ -87,12 +87,20 @@ class Config(object):
CSV_UPLOAD_BUCKET_NAME = 'local-notifications-csv-upload'
NOTIFICATIONS_ALERT = 5 # five mins
STATSD_ENABLED = False
STATSD_HOST = "localhost"
STATSD_PORT = None
STATSD_PREFIX = None
class Development(Config):
DEBUG = True
MMG_API_KEY = os.environ['MMG_API_KEY']
CSV_UPLOAD_BUCKET_NAME = 'development-notifications-csv-upload'
STATSD_ENABLED = True
STATSD_HOST = os.getenv('STATSD_HOST')
STATSD_PORT = os.getenv('STATSD_PORT')
STATSD_PREFIX = os.getenv('STATSD_PREFIX')
class Preview(Config):
MMG_API_KEY = os.environ['MMG_API_KEY']

View File

@@ -18,6 +18,10 @@ class Live(Config):
TWILIO_AUTH_TOKEN = os.getenv('LIVE_TWILIO_AUTH_TOKEN')
MMG_API_KEY = os.environ['LIVE_MMG_API_KEY']
CSV_UPLOAD_BUCKET_NAME = 'live-notifications-csv-upload'
STATSD_ENABLED = True
STATSD_HOST = os.getenv('LIVE_STATSD_HOST')
STATSD_PORT = os.getenv('LIVE_STATSD_PORT')
STATSD_PREFIX = os.getenv('LIVE_STATSD_PREFIX')
BROKER_TRANSPORT_OPTIONS = {
'region': 'eu-west-1',

View File

@@ -22,3 +22,7 @@ export MMG_API_KEY='mmg-secret-key'
export MMG_FROM_NUMBER='test'
export LOADTESTING_API_KEY="loadtesting"
export LOADTESTING_NUMBER="loadtesting"
export STATSD_ENABLED=True
export STATSD_HOST="somehost"
export STATSD_PORT=1000
export STATSD_PREFIX="stats-prefix"

View File

@@ -17,7 +17,7 @@ boto==2.39.0
celery==3.1.20
twilio==4.6.0
monotonic==0.3
statsd==3.2.1
git+https://github.com/alphagov/notifications-python-client.git@1.0.0#egg=notifications-python-client==1.0.0

View File

@@ -16,7 +16,7 @@ from app.celery.tasks import (
delete_successful_notifications,
provider_to_use
)
from app import (aws_ses_client, encryption, DATETIME_FORMAT, mmg_client)
from app import (aws_ses_client, encryption, DATETIME_FORMAT, mmg_client, statsd_client)
from app.clients.email.aws_ses import AwsSesClientException
from app.clients.sms.mmg import MMGClientException
from app.dao import notifications_dao, jobs_dao, provider_details_dao
@@ -103,6 +103,7 @@ def test_should_call_delete_invotations_on_delete_invitations_task(notify_api, m
@freeze_time("2016-01-01 11:09:00.061258")
def test_should_process_sms_job(sample_job, mocker, mock_celery_remove_job):
mocker.patch('app.statsd_client.incr')
mocker.patch('app.celery.tasks.s3.get_job_from_s3', return_value=load_example_csv('sms'))
mocker.patch('app.celery.tasks.send_sms.apply_async')
mocker.patch('app.encryption.encrypt', return_value="something_encrypted")
@@ -124,6 +125,7 @@ def test_should_process_sms_job(sample_job, mocker, mock_celery_remove_job):
)
job = jobs_dao.dao_get_job_by_id(sample_job.id)
assert job.status == 'finished'
statsd_client.incr.assert_called_once_with("notifications.tasks.process-job")
@freeze_time("2016-01-01 11:09:00.061258")
@@ -329,21 +331,37 @@ def test_should_send_template_to_correct_sms_provider_and_persist(sample_templat
mocker.patch('app.encryption.decrypt', return_value=notification)
mocker.patch('app.mmg_client.send_sms')
mocker.patch('app.mmg_client.get_name', return_value="mmg")
mocker.patch('app.statsd_client.incr')
mocker.patch('app.statsd_client.timing_with_dates')
notification_id = uuid.uuid4()
freezer = freeze_time("2016-01-01 11:09:00.00000")
freezer.start()
now = datetime.utcnow()
freezer.stop()
freezer = freeze_time("2016-01-01 11:10:00.00000")
freezer.start()
send_sms(
sample_template_with_placeholders.service_id,
notification_id,
"encrypted-in-reality",
now.strftime(DATETIME_FORMAT)
)
freezer.stop()
statsd_client.timing_with_dates.assert_called_once_with(
"notifications.tasks.send-sms.queued-for", datetime(2016, 1, 1, 11, 10, 0, 00000), datetime(2016, 1, 1, 11, 9, 0, 00000)
)
mmg_client.send_sms.assert_called_once_with(
to=format_phone_number(validate_phone_number("+447234123123")),
content="Sample service: Hello Jo",
reference=str(notification_id)
)
statsd_client.incr.assert_called_once_with("notifications.tasks.send-sms")
persisted_notification = notifications_dao.get_notification(
sample_template_with_placeholders.service_id, notification_id
)
@@ -539,11 +557,21 @@ def test_should_use_email_template_and_persist(sample_email_template_with_placeh
"personalisation": {"name": "Jo"}
}
mocker.patch('app.encryption.decrypt', return_value=notification)
mocker.patch('app.aws_ses_client.send_email')
mocker.patch('app.statsd_client.incr')
mocker.patch('app.statsd_client.timing_with_dates')
mocker.patch('app.aws_ses_client.get_name', return_value='ses')
mocker.patch('app.aws_ses_client.send_email', return_value='ses')
notification_id = uuid.uuid4()
freezer = freeze_time("2016-01-01 11:09:00.00000")
freezer.start()
now = datetime.utcnow()
freezer.stop()
freezer = freeze_time("2016-01-01 11:10:00.00000")
freezer.start()
send_email(
sample_email_template_with_placeholders.service_id,
notification_id,
@@ -551,6 +579,8 @@ def test_should_use_email_template_and_persist(sample_email_template_with_placeh
"encrypted-in-reality",
now.strftime(DATETIME_FORMAT)
)
freezer.stop()
aws_ses_client.send_email.assert_called_once_with(
"email_from",
"my_email@my_email.com",
@@ -558,9 +588,16 @@ def test_should_use_email_template_and_persist(sample_email_template_with_placeh
body="Hello Jo",
html_body=AnyStringWith("Hello Jo")
)
statsd_client.incr.assert_called_once_with("notifications.tasks.send-email")
statsd_client.timing_with_dates.assert_called_once_with(
"notifications.tasks.send-email.queued-for", datetime(2016, 1, 1, 11, 10, 0, 00000), datetime(2016, 1, 1, 11, 9, 0, 00000)
)
persisted_notification = notifications_dao.get_notification(
sample_email_template_with_placeholders.service_id, notification_id
)
assert persisted_notification.id == notification_id
assert persisted_notification.to == 'my_email@my_email.com'
assert persisted_notification.template_id == sample_email_template_with_placeholders.id

View File

@@ -1510,6 +1510,110 @@ def test_should_handle_validation_code_callbacks(notify_api, notify_db, notify_d
assert json_resp['message'] == 'SES callback succeeded'
def test_should_record_email_request_in_statsd(notify_api, notify_db, notify_db_session, sample_email_template, mocker):
with notify_api.test_request_context():
with notify_api.test_client() as client:
mocker.patch('app.statsd_client.incr')
mocker.patch('app.celery.tasks.send_email.apply_async')
mocker.patch('app.encryption.encrypt', return_value="something_encrypted")
data = {
'to': 'ok@ok.com',
'template': str(sample_email_template.id)
}
auth_header = create_authorization_header(service_id=sample_email_template.service_id)
response = client.post(
path='/notifications/email',
data=json.dumps(data),
headers=[('Content-Type', 'application/json'), auth_header])
assert response.status_code == 201
app.statsd_client.incr.assert_called_once_with("notification.api.email")
def test_should_record_sms_request_in_statsd(notify_api, notify_db, notify_db_session, sample_template, mocker):
with notify_api.test_request_context():
with notify_api.test_client() as client:
mocker.patch('app.statsd_client.incr')
mocker.patch('app.celery.tasks.send_sms.apply_async')
mocker.patch('app.encryption.encrypt', return_value="something_encrypted")
data = {
'to': '07123123123',
'template': str(sample_template.id)
}
auth_header = create_authorization_header(service_id=sample_template.service_id)
response = client.post(
path='/notifications/sms',
data=json.dumps(data),
headers=[('Content-Type', 'application/json'), auth_header])
assert response.status_code == 201
app.statsd_client.incr.assert_called_once_with("notification.api.sms")
def test_ses_callback_should_update_record_statsd(
notify_api,
notify_db,
notify_db_session,
sample_email_template,
mocker):
with notify_api.test_request_context():
with notify_api.test_client() as client:
mocker.patch('app.statsd_client.incr')
notification = create_sample_notification(
notify_db,
notify_db_session,
template=sample_email_template,
reference='ref'
)
assert get_notification_by_id(notification.id).status == 'sending'
client.post(
path='/notifications/email/ses',
data=ses_notification_callback(),
headers=[('Content-Type', 'text/plain; charset=UTF-8')]
)
app.statsd_client.incr.assert_called_once_with("notifications.callback.ses.delivered")
def test_process_mmg_response_records_statsd(notify_api, sample_notification, mocker):
with notify_api.test_client() as client:
mocker.patch('app.statsd_client.incr')
data = json.dumps({"reference": "mmg_reference",
"CID": str(sample_notification.id),
"MSISDN": "447777349060",
"status": "3",
"deliverytime": "2016-04-05 16:01:07"})
client.post(path='notifications/sms/mmg',
data=data,
headers=[('Content-Type', 'application/json')])
app.statsd_client.incr.assert_called_once_with("notifications.callback.mmg.delivered")
def test_firetext_callback_should_record_statsd(notify_api, notify_db, notify_db_session, mocker):
with notify_api.test_request_context():
with notify_api.test_client() as client:
mocker.patch('app.statsd_client.incr')
notification = create_sample_notification(notify_db, notify_db_session, status='delivered')
client.post(
path='/notifications/sms/firetext',
data='mobile=441234123123&status=2&time=2016-03-10 14:17:00&reference={}'.format(
notification.id
),
headers=[('Content-Type', 'application/x-www-form-urlencoded')])
app.statsd_client.incr.assert_called_once_with("notifications.callback.firetext.delivered")
def ses_validation_code_callback():
return b'{\n "Type" : "Notification",\n "MessageId" : "ref",\n "TopicArn" : "arn:aws:sns:eu-west-1:123456789012:testing",\n "Message" : "{\\"notificationType\\":\\"Delivery\\",\\"mail\\":{\\"timestamp\\":\\"2016-03-14T12:35:25.909Z\\",\\"source\\":\\"valid-code@test.com\\",\\"sourceArn\\":\\"arn:aws:ses:eu-west-1:123456789012:identity/testing-notify\\",\\"sendingAccountId\\":\\"123456789012\\",\\"messageId\\":\\"ref\\",\\"destination\\":[\\"testing@digital.cabinet-office.gov.uk\\"]},\\"delivery\\":{\\"timestamp\\":\\"2016-03-14T12:35:26.567Z\\",\\"processingTimeMillis\\":658,\\"recipients\\":[\\"testing@digital.cabinet-office.gov.u\\"],\\"smtpResponse\\":\\"250 2.0.0 OK 1457958926 uo5si26480932wjc.221 - gsmtp\\",\\"reportingMTA\\":\\"a6-238.smtp-out.eu-west-1.amazonses.com\\"}}",\n "Timestamp" : "2016-03-14T12:35:26.665Z",\n "SignatureVersion" : "1",\n "Signature" : "X8d7eTAOZ6wlnrdVVPYanrAlsX0SMPfOzhoTEBnQqYkrNWTqQY91C0f3bxtPdUhUtOowyPAOkTQ4KnZuzphfhVb2p1MyVYMxNKcBFB05/qaCX99+92fjw4x9LeUOwyGwMv5F0Vkfi5qZCcEw69uVrhYLVSTFTrzi/yCtru+yFULMQ6UhbY09GwiP6hjxZMVr8aROQy5lLHglqQzOuSZ4KeD85JjifHdKzlx8jjQ+uj+FLzHXPMAPmPU1JK9kpoHZ1oPshAFgPDpphJe+HwcJ8ezmk+3AEUr3wWli3xF+49y8Z2anASSVp6YI2YP95UT8Rlh3qT3T+V9V8rbSVislxA==",\n "SigningCertURL" : "https://sns.eu-west-1.amazonaws.com/SimpleNotificationService-bb750dd426d95ee9390147a5624348ee.pem",\n "UnsubscribeURL" : "https://sns.eu-west-1.amazonaws.com/?Action=Unsubscribe&SubscriptionArn=arn:aws:sns:eu-west-1:302763885840:preview-emails:d6aad3ef-83d6-4cf3-a470-54e2e75916da"\n}' # noqa