"""DAG for accounting run commit.""" from datetime import datetime from airflow import DAG from airflow.operators.dummy import DummyOperator from airflow.operators.python import BranchPythonOperator from airflow.operators.python import PythonOperator from airflow.sensors.python import PythonSensor from lib.utils.ows import update_accounting_run from lib.utils.slack import slack_failure_callback, slack_success_callback from tasks.accounting_run_commit.commit_royalties_lambda_response_check import \ commit_royalties_lambda_response_check from tasks.accounting_run_commit.contract_type_branch import \ choose_branch_based_on_contract_type_task from tasks.accounting_run_commit.debit_reserves import \ debit_reserves_for_accounting_run_task from tasks.accounting_run_commit.invoke_commit_mechanicals_lambda import \ invoke_commit_mechanicals_lambda_task from tasks.accounting_run_commit.invoke_commit_royalties_lambda import \ invoke_commit_royalties_lambda_task from tasks.accounting_run_commit.invoke_reserves_schedule_lambda import \ invoke_reserves_schedule_lambda_task from tasks.accounting_run_commit.invoke_reserves_take_lambda import \ invoke_reserves_take_lambda_task from tasks.accounting_run_commit.load_mech_results_into_publishing_escrow \ import load_mech_results_into_publishing_escrow_task from tasks.accounting_run_commit.snowflake_copy_contract_transactions import \ copy_contract_transactions from tasks.accounting_run_commit.snowflake_copy_mech_deductions import \ copy_mech_deductions from tasks.accounting_run_commit.snowflake_copy_results import \ copy_results_in_snowflake from tasks.accounting_run_commit.snowflake_export_mechanicals import \ export_mechanicals_task dag = DAG( 'accounting_run_commit', description='Move data from staging to final tables, conditionally take and schedule reserves.', # noqa: E501 start_date=datetime(2019, 1, 1), default_args={ 'provide_context': True, 'on_failure_callback': slack_failure_callback, 'on_success_callback': slack_success_callback, }, schedule_interval=None ) def complete_commit(dag_run, *args, **kwargs): """Mark run as 'Committed'.""" config = dag_run.conf accounting_run_id = config['target_id'] return update_accounting_run(accounting_run_id, accounting_run_status='Committed') def notify_failure(dag_run, *args, **kwargs): """Update accounting run status to 'Error'.""" config = dag_run.conf accounting_run_id = config['target_id'] update_accounting_run(accounting_run_id, accounting_run_status='Error') raise ValueError('Commit failed') """ ###################################################################### ############################# OPERATORS ############################## ###################################################################### """ complete_commit_operator = PythonOperator( dag=dag, python_callable=complete_commit, task_id='complete_commit', trigger_rule='none_failed_min_one_success' ) contract_type_branch_operator = BranchPythonOperator( dag=dag, python_callable=choose_branch_based_on_contract_type_task, task_id='contract_type_branch' ) copy_snowflake_txns_operator = PythonOperator( dag=dag, python_callable=copy_results_in_snowflake, task_id='copy_snowflake_txns' ) copy_snowflake_contract_txns_operator = PythonOperator( dag=dag, python_callable=copy_contract_transactions, task_id='copy_snowflake_contract_txns' ) copy_snowflake_mech_deductions_operator = PythonOperator( dag=dag, python_callable=copy_mech_deductions, task_id='copy_snowflake_mech_deductions' ) invoke_reserves_take_lambda_operator = PythonOperator( dag=dag, python_callable=invoke_reserves_take_lambda_task, task_id='invoke_reserves_take_lambda' ) debit_reserves_operator = PythonOperator( dag=dag, python_callable=debit_reserves_for_accounting_run_task, task_id='debit_reserves' ) invoke_reserves_schedule_lambda_operator = PythonOperator( dag=dag, python_callable=invoke_reserves_schedule_lambda_task, task_id='invoke_reserves_schedule_lambda' ) invoke_commit_mechanicals_lambda_operator = PythonOperator( dag=dag, python_callable=invoke_commit_mechanicals_lambda_task, task_id='invoke_commit_mechanicals_lambda' ) load_mech_results_into_publishing_escrow_operator = PythonOperator( dag=dag, python_callable=load_mech_results_into_publishing_escrow_task, task_id='load_mech_into_art_relations_publishing_escrow' ) skip_distro_operator = DummyOperator( dag=dag, task_id='skip_distro' ) snowflake_export_mechanicals_operator = PythonOperator( dag=dag, python_callable=export_mechanicals_task, task_id='snowflake_export_mechanicals' ) invoke_commit_royalties_lambda_operator = PythonOperator( dag=dag, python_callable=invoke_commit_royalties_lambda_task, task_id='invoke_commit_royalties_lambda' ) wait_commit_royalties_lambda_result_operator = PythonSensor( dag=dag, python_callable=commit_royalties_lambda_response_check, timeout=900, poke_interval=30, task_id='wait_commit_royalties_lambda_result' ) notify_failure_operator = PythonOperator( dag=dag, python_callable=notify_failure, task_id='notify_failure', trigger_rule='one_failed' ) """ ###################################################################### ################################ DAG ################################# ###################################################################### """ copy_snowflake_txns_operator >> \ copy_snowflake_contract_txns_operator >> \ invoke_commit_royalties_lambda_operator >> \ wait_commit_royalties_lambda_result_operator >> \ contract_type_branch_operator contract_type_branch_operator >> \ invoke_reserves_take_lambda_operator >> \ debit_reserves_operator >> \ invoke_reserves_schedule_lambda_operator >> \ copy_snowflake_mech_deductions_operator >> \ snowflake_export_mechanicals_operator >> \ load_mech_results_into_publishing_escrow_operator >> \ invoke_commit_mechanicals_lambda_operator contract_type_branch_operator >> \ skip_distro_operator [invoke_commit_mechanicals_lambda_operator, skip_distro_operator] >> \ complete_commit_operator >> \ notify_failure_operator