"""Distribution Fee ETL Decider.""" from garcon.param import StaticParam from flows import runner from flows.distribution_fee import queries from flows.distribution_fee import tasks from flows.flow import DatabaseParam from flows.flow import FlowBase class Flow(FlowBase): """Distribution Fee SWF flow class.""" def decider(self, schedule): """Schedule the next appropriate activity for the SWF execution.""" create_temp_tables = schedule( 'create_temp_tables', self.create_temp_tables_activity) load_vendor_contract = schedule( 'load_vendor_contract', self.load_vendor_contract_activity, requires=[create_temp_tables]) load_dist_fee = schedule( 'load_dist_fee', self.load_dist_fee_activity, requires=[load_vendor_contract]) calculate_client_amount = schedule( 'calculate_client_amount', self.calculate_client_amount_activity, requires=[load_dist_fee]) update_distribution_fee_table = schedule( 'update_distribution_fee_table', self.update_distribution_fee_tables_activity, requires=[calculate_client_amount]) drop_temp_tables = schedule( 'drop_temp_tables', self.drop_temp_tables_activity, # TODO: change requirements when the other activities are done requires=[update_distribution_fee_table]) schedule( 'set_dynamo_status', self.set_dynamo_status_activity, requires=[drop_temp_tables]) schedule( 'send_sns_notification', self.send_sns_notification_activity, requires=[drop_temp_tables]) @property def create_temp_tables_activity(self): """Create temp tables for storing distribution fees. Returns: garcon.activity.Activity instance. """ return self.create( name='create_temp_tables', retry=10, tasks=runner.Sync( tasks.create_temp_table.fill( correlation_id='correlation_id', create_statement=StaticParam( queries.CREATE_TEMP_TABLE_FEE_DMS), namespace='temp_table_fee_dms', table_name=StaticParam(queries.TEMP_TABLE_NAME_FEE_DMS)), tasks.create_temp_table.fill( correlation_id='correlation_id', create_statement=StaticParam( queries.CREATE_TEMP_TABLE_FEE_REGULAR), namespace='temp_table_fee_regular', table_name=StaticParam( queries.TEMP_TABLE_NAME_FEE_REGULAR)), tasks.create_temp_table.fill( correlation_id='correlation_id', create_statement=StaticParam( queries.CREATE_TEMP_TABLE_FEE_TERRITORY), namespace='temp_table_fee_territory', table_name=StaticParam( queries.TEMP_TABLE_NAME_FEE_TERRITORY)), tasks.create_temp_table.fill( correlation_id='correlation_id', create_statement=StaticParam( queries.CREATE_TEMP_TABLE_VENDOR_CONTRACT), namespace='temp_table_vendor_contract', table_name=StaticParam( queries.TEMP_TABLE_NAME_VENDOR_CONTRACT)))) @property def load_vendor_contract_activity(self): """Load vendor and contract IDs. Returns: garcon.activity.Activity instance. """ return self.create( name='load_vendor_contract', retry=10, tasks=runner.Sync( tasks.load_vendor_contract.fill( call_query=StaticParam(queries.GET_ACTIVE_VENDOR_CONTRACT), correlation_id='correlation_id', insert_query=StaticParam(queries.INSERT_VENDOR_CONTRACT), table_name='temp_table_vendor_contract.name', upcs=DatabaseParam('upcs')))) @property def load_dist_fee_activity(self): """Load vendor and contract IDs. Returns: garcon.activity.Activity instance. """ return self.create( name='load_dist_fee', retry=10, tasks=runner.Sync( tasks.load_dist_fee_regular.fill( contract_table='temp_table_vendor_contract.name', correlation_id='correlation_id', fee_table='temp_table_fee_regular.name'), tasks.load_dist_fee_territory.fill( contract_table='temp_table_vendor_contract.name', correlation_id='correlation_id', fee_table='temp_table_fee_territory.name'))) @property def calculate_client_amount_activity(self): """Apply distribution fee to separate client_amount table. Returns: garcon.activity.Activity instance. """ return self.create( name='calculate_client_amount', retry=10, tasks=runner.Sync( tasks.calculate_client_amount.fill( correlation_id='correlation_id', fee_dms_table='temp_table_fee_dms.name', fee_ter_table='temp_table_fee_territory.name', fee_reg_table='temp_table_fee_regular.name', upcs=DatabaseParam('upcs')))) @property def update_distribution_fee_tables_activity(self): """Update distribution fee table. Returns: garcon.activity.Activity instance. """ return self.create( name='update_distribution_fee_table', retry=10, tasks=runner.Sync( tasks.update_distribution_fee_table.fill( correlation_id='correlation_id', fee_dms_table='temp_table_fee_dms.name', fee_ter_table='temp_table_fee_territory.name', fee_reg_table='temp_table_fee_regular.name', upcs=DatabaseParam('upcs')))) @property def drop_temp_tables_activity(self): """Drop temp tables. Returns: garcon.activity.Activity instance. """ table_names = map(lambda ns: '{}.name'.format(ns), [ 'temp_table_fee_dms', 'temp_table_fee_regular', 'temp_table_fee_territory', 'temp_table_vendor_contract']) task_list = [ tasks.drop_temp_table.fill( correlation_id='correlation_id', drop_statement=StaticParam(queries.DROP_TEMP_TABLE), table_name=table_name) for table_name in table_names] return self.create( name='drop_temp_tables', retry=10, tasks=runner.Sync(*task_list)) @property def set_dynamo_status_activity(self): """Set DynamoDB status. Returns: Activity: garcon.activity.Activity instance. """ return self.create( name='set_status', retry=10, tasks=runner.Sync(tasks.set_dynamo_status.fill( correlation_id='correlation_id'))) @property def send_sns_notification_activity(self): """Task description for notifying any listening apps. Returns: garcon.activity.Activity instance. """ return self.create( name='send_sns_notification', retry=10, tasks=runner.Sync( tasks.send_sns_notification.fill( correlation_id='correlation_id', upcs='upcs')))