Merge branch 'master' into schedule-api-notification

Conflicts:
	app/celery/scheduled_tasks.py
	app/v2/notifications/post_notifications.py
	tests/app/celery/test_scheduled_tasks.py
This commit is contained in:
Rebecca Law
2017-05-22 14:05:57 +01:00
48 changed files with 1917 additions and 291 deletions

View File

@@ -11,7 +11,7 @@ from app.models import (Job,
Template,
JOB_STATUS_SCHEDULED,
JOB_STATUS_PENDING,
LETTER_TYPE)
LETTER_TYPE, JobStatistics)
from app.statsd_decorators import statsd
@@ -109,6 +109,11 @@ def dao_get_future_scheduled_job_by_id_and_service_id(job_id, service_id):
def dao_create_job(job):
job_stats = JobStatistics(
job_id=job.id,
updated_at=datetime.utcnow()
)
db.session.add(job_stats)
db.session.add(job)
db.session.commit()

View File

@@ -0,0 +1,22 @@
from app import db
from app.dao.dao_utils import transactional
from app.models import ServicePermission, SERVICE_PERMISSION_TYPES
def dao_fetch_service_permissions(service_id):
return ServicePermission.query.filter(
ServicePermission.service_id == service_id).all()
@transactional
def dao_add_service_permission(service_id, permission):
service_permission = ServicePermission(service_id=service_id, permission=permission)
db.session.add(service_permission)
def dao_remove_service_permission(service_id, permission):
deleted = ServicePermission.query.filter(
ServicePermission.service_id == service_id,
ServicePermission.permission == permission).delete()
db.session.commit()
return deleted

View File

@@ -25,9 +25,13 @@ from app.models import (
User,
InvitedUser,
Service,
ServicePermission,
KEY_TYPE_TEST,
NOTIFICATION_STATUS_TYPES,
TEMPLATE_TYPES,
JobStatistics,
SMS_TYPE,
EMAIL_TYPE
)
from app.service.statistics import format_monthly_template_notification_stats
from app.statsd_decorators import statsd
@@ -60,6 +64,19 @@ def dao_fetch_service_by_id(service_id, only_active=False):
return query.one()
def dao_fetch_service_by_id_with_api_keys(service_id, only_active=False):
query = Service.query.filter_by(
id=service_id
).options(
joinedload('api_keys')
)
if only_active:
query = query.filter(Service.active)
return query.one()
def dao_fetch_all_services_by_user(user_id, only_active=False):
query = Service.query.filter(
Service.users.any(id=user_id)
@@ -111,13 +128,18 @@ def dao_fetch_service_by_id_and_user(service_id, user_id):
@transactional
@version_class(Service)
def dao_create_service(service, user, service_id=None):
def dao_create_service(service, user, service_id=None, service_permissions=[SMS_TYPE, EMAIL_TYPE]):
from app.dao.permissions_dao import permission_dao
service.users.append(user)
permission_dao.add_default_service_permissions_for_user(user, service)
service.id = service_id or uuid.uuid4() # must be set now so version history model can use same id
service.active = True
service.research_mode = False
for permission in service_permissions:
service_permission = ServicePermission(service_id=service.id, permission=permission)
db.session.add(service_permission)
db.session.add(service)
@@ -160,6 +182,10 @@ def delete_service_and_all_associated_db_objects(service):
query.delete()
db.session.commit()
job_stats = JobStatistics.query.join(Job).filter(Job.service_id == service.id)
list(map(db.session.delete, job_stats))
db.session.commit()
_delete_commit(NotificationStatistics.query.filter_by(service=service))
_delete_commit(TemplateStatistics.query.filter_by(service=service))
_delete_commit(ProviderStatistics.query.filter_by(service=service))
@@ -167,11 +193,12 @@ def delete_service_and_all_associated_db_objects(service):
_delete_commit(Permission.query.filter_by(service=service))
_delete_commit(ApiKey.query.filter_by(service=service))
_delete_commit(ApiKey.get_history_model().query.filter_by(service_id=service.id))
_delete_commit(Job.query.filter_by(service=service))
_delete_commit(NotificationHistory.query.filter_by(service=service))
_delete_commit(Notification.query.filter_by(service=service))
_delete_commit(Job.query.filter_by(service=service))
_delete_commit(Template.query.filter_by(service=service))
_delete_commit(TemplateHistory.query.filter_by(service_id=service.id))
_delete_commit(ServicePermission.query.filter_by(service_id=service.id))
verify_codes = VerifyCode.query.join(User).filter(User.id.in_([x.id for x in service.users]))
list(map(db.session.delete, verify_codes))

149
app/dao/statistics_dao.py Normal file
View File

@@ -0,0 +1,149 @@
from datetime import datetime, timedelta
from itertools import groupby
from flask import current_app
from sqlalchemy import func
from sqlalchemy.exc import IntegrityError, SQLAlchemyError
from app import db
from app.dao.dao_utils import transactional
from app.models import (
JobStatistics,
Notification,
EMAIL_TYPE,
SMS_TYPE,
LETTER_TYPE,
NOTIFICATION_STATUS_TYPES_FAILED,
NOTIFICATION_STATUS_SUCCESS,
NOTIFICATION_DELIVERED,
NOTIFICATION_SENT)
from app.statsd_decorators import statsd
@transactional
def timeout_job_counts(notifications_type, timeout_start):
total_updated = 0
sent = columns(notifications_type, 'sent')
delivered = columns(notifications_type, 'delivered')
failed = columns(notifications_type, 'failed')
results = db.session.query(
JobStatistics.job_id.label('job_id'),
func.count(Notification.status).label('count'),
Notification.status.label('status')
).filter(
Notification.notification_type == notifications_type,
JobStatistics.job_id == Notification.job_id,
JobStatistics.created_at < timeout_start,
sent != failed + delivered
).group_by(Notification.status, JobStatistics.job_id).order_by(JobStatistics.job_id).all()
sort = sorted(results, key=lambda result: result.job_id)
groups = []
for k, g in groupby(sort, key=lambda result: result.job_id):
groups.append(list(g))
for job in groups:
sent_count = 0
delivered_count = 0
failed_count = 0
for notification_status in job:
if notification_status.status in NOTIFICATION_STATUS_SUCCESS:
delivered_count += notification_status.count
else:
failed_count += notification_status.count
sent_count += notification_status.count
total_updated += JobStatistics.query.filter_by(
job_id=notification_status.job_id
).update({
sent: sent_count,
failed: failed_count,
delivered: delivered_count
}, synchronize_session=False)
return total_updated
@statsd(namespace="dao")
def dao_timeout_job_statistics(timeout_period):
timeout_start = datetime.utcnow() - timedelta(seconds=timeout_period)
sms_count = timeout_job_counts(SMS_TYPE, timeout_start)
email_count = timeout_job_counts(EMAIL_TYPE, timeout_start)
return sms_count + email_count
@statsd(namespace="dao")
def create_or_update_job_sending_statistics(notification):
if __update_job_stats_sent_count(notification) == 0:
try:
__insert_job_stats(notification)
except IntegrityError as e:
current_app.logger.exception(e)
if __update_job_stats_sent_count(notification) == 0:
raise SQLAlchemyError("Failed to create job statistics for {}".format(notification.job_id))
@transactional
def __update_job_stats_sent_count(notification):
column = columns(notification.notification_type, 'sent')
return db.session.query(JobStatistics).filter_by(
job_id=notification.job_id,
).update({
column: column + 1
})
@transactional
def __insert_job_stats(notification):
stats = JobStatistics(
job_id=notification.job_id,
emails_sent=1 if notification.notification_type == EMAIL_TYPE else 0,
sms_sent=1 if notification.notification_type == SMS_TYPE else 0,
letters_sent=1 if notification.notification_type == LETTER_TYPE else 0,
updated_at=datetime.utcnow()
)
db.session.add(stats)
def columns(notification_type, status):
keys = {
EMAIL_TYPE: {
'failed': JobStatistics.emails_failed,
'delivered': JobStatistics.emails_delivered,
'sent': JobStatistics.emails_sent
},
SMS_TYPE: {
'failed': JobStatistics.sms_failed,
'delivered': JobStatistics.sms_delivered,
'sent': JobStatistics.sms_sent
},
LETTER_TYPE: {
'failed': JobStatistics.letters_failed,
'sent': JobStatistics.letters_sent
}
}
return keys.get(notification_type).get(status)
@transactional
def update_job_stats_outcome_count(notification):
if notification.status in NOTIFICATION_STATUS_TYPES_FAILED:
column = columns(notification.notification_type, 'failed')
elif notification.status in [NOTIFICATION_DELIVERED,
NOTIFICATION_SENT] and notification.notification_type != LETTER_TYPE:
column = columns(notification.notification_type, 'delivered')
else:
column = None
if column:
return db.session.query(JobStatistics).filter_by(
job_id=notification.job_id,
).update({
column: column + 1
})
else:
return 0