"""DAG for calculating sales against contracts.""" from datetime import datetime from airflow import DAG from airflow.operators.python import PythonOperator from lib import constants from lib.utils.slack import slack_failure_callback, slack_success_callback from tasks.accounting_run_calculate.invoke_ledger_accounting_run_balance_lambda import \ invoke_ledger_accounting_run_balance_lambda_task from tasks.accounting_run_calculate.snapshot_contract_mechanical_deductions import \ snapshot_contract_mechanical_deductions_task from tasks.accounting_run_calculate.snapshot_contracts import snapshot_contracts from tasks.accounting_run_calculate.snowflake_calculate_mechanical_deductions import \ calculate_mech_deductions_task from tasks.accounting_run_calculate.snowflake_calculate_totals import \ calculate_run_totals_task from tasks.accounting_run_calculate.snowflake_flatten_contract_json import \ flatten_contract_json_task from tasks.accounting_run_calculate.snowflake_match_contract_terms import \ match_contract_terms_to_sales_task from tasks.accounting_run_calculate.update_run_status_complete import \ status_complete from tasks.accounting_run_calculate.update_run_status_error import status_error from tasks.accounting_run_calculate.update_run_status_started import \ status_started dag = DAG( constants.DAG_CALC_EVENT_NAME, description='Calculate revenue based on contract terms', start_date=datetime(2022, 1, 1), default_args={ 'provide_context': True, 'on_failure_callback': slack_failure_callback, 'on_success_callback': slack_success_callback, }, schedule_interval=None ) """ ###################################################################### ############################# OPERATORS ############################## ###################################################################### """ invoke_ledger_accounting_run_balance_lambda_operator = PythonOperator( dag=dag, python_callable=invoke_ledger_accounting_run_balance_lambda_task, task_id='invoke_ledger_accounting_run_balance_lambda', trigger_rule='none_failed' ) snapshot_contracts_operator = PythonOperator( dag=dag, python_callable=snapshot_contracts, task_id='snapshot_contracts', ) snapshot_contract_mechanical_deductions_operator = PythonOperator( dag=dag, python_callable=snapshot_contract_mechanical_deductions_task, task_id='snapshot_contract_mechanical_deductions', ) snowflake_calculate_mech_deductions_operator = PythonOperator( dag=dag, python_callable=calculate_mech_deductions_task, task_id='snowflake_calculate_mech_deductions' ) snowflake_calculate_totals_operator = PythonOperator( dag=dag, python_callable=calculate_run_totals_task, task_id='snowflake_calculate_run_totals' ) snowflake_flatten_contract_json_operator = PythonOperator( dag=dag, python_callable=flatten_contract_json_task, task_id='snowflake_flatten_contract_json' ) snowflake_match_contract_terms_operator = PythonOperator( dag=dag, python_callable=match_contract_terms_to_sales_task, task_id='snowflake_match_contract_terms_to_sales' ) status_complete_operator = PythonOperator( dag=dag, python_callable=status_complete, task_id='update_status_complete', trigger_rule='none_failed' ) status_error_operator = PythonOperator( dag=dag, python_callable=status_error, task_id='status_error', trigger_rule='one_failed', ) status_started_operator = PythonOperator( dag=dag, python_callable=status_started, task_id='status_started' ) """ ###################################################################### ################################ DAG ################################# ###################################################################### """ status_started_operator >> \ snapshot_contracts_operator >> \ snapshot_contract_mechanical_deductions_operator >> \ snowflake_flatten_contract_json_operator >> \ snowflake_match_contract_terms_operator >> \ snowflake_calculate_totals_operator >> \ snowflake_calculate_mech_deductions_operator >> \ invoke_ledger_accounting_run_balance_lambda_operator >> \ status_complete_operator >> \ status_error_operator