"""Payments Generate Lambda Invocation.""" from datetime import datetime import json from airflow import DAG from airflow.operators.python import PythonOperator from hooks.lambda_hook import OrchLambdaHook from lib import constants from lib.config import PAYMENTS_GENERATE_PROCESSOR_LAMBDA_NAME from lib.utils import event from lib.utils.ows import get_abacus_state_id, update_abacus_state dag = DAG( 'payments_generate', description='Invoke payments generate lambda', start_date=datetime(2019, 1, 1), default_args={'provide_context': True}, schedule_interval=None ) def get_event_from_params(dag_run, **kwargs): """Get event from params passed to task callbacks.""" abacus_event = event.get_abacus_event(dag_run, **kwargs) event.validate_event_for_handler( abacus_event, target_type=constants.DAG_PAYMENTS_GENERATE_TARGET_TYPE, event_name=constants.DAG_PAYMENTS_GENERATE_EVENT_NAME ) return abacus_event def invoke_payments_generate(dag_run, *args, **kwargs): """Invoke payments generate lambda.""" hook = OrchLambdaHook(PAYMENTS_GENERATE_PROCESSOR_LAMBDA_NAME) event = dag_run.conf payload = json.dumps(event) response = hook.invoke_lambda(payload) if not response.function_response.succeeded: raise Exception(response.error_message) def notify_failure(dag_run, *args, **kwargs): """Update payment group payment status on error.""" abacus_event = get_event_from_params(dag_run, **kwargs) payment_group_payment_id = abacus_event.target_id abacus_state_id = get_abacus_state_id( 'payment-group-payment', payment_group_payment_id, constants.PAYMENT_GROUP_PAYMENT_ACTIONS.GENERATE_PAYMENTS ) return update_abacus_state( abacus_state_id, body=dict( action_status=constants.ABACUS_STATE_STATUSES.ERROR ) ) def notify_success(dag_run, *args, **kwargs): """Update payment group payment status on success.""" abacus_event = get_event_from_params(dag_run, **kwargs) payment_group_payment_id = abacus_event.target_id abacus_state_id = get_abacus_state_id( 'payment-group-payment', payment_group_payment_id, constants.PAYMENT_GROUP_PAYMENT_ACTIONS.GENERATE_PAYMENTS ) return update_abacus_state( abacus_state_id, body=dict( action_status=constants.ABACUS_STATE_STATUSES.COMPLETE ) ) invoke_payments_generate_operator = PythonOperator( task_id='invoke_payments_generate', python_callable=invoke_payments_generate, dag=dag ) notify_failure_operator = PythonOperator( dag=dag, python_callable=notify_failure, task_id='notify_failure', trigger_rule='one_failed' ) notify_success_operator = PythonOperator( dag=dag, python_callable=notify_success, task_id='notify_success' ) invoke_payments_generate_operator >> notify_success_operator >> \ notify_failure_operator