""" Reserve payouts workflow. Example call: garcon exec reserve_payouts -c '{"period_id":"199", "validate_only":true}' """ import os from garcon import runner from accounting.flows.base import FlowBase from accounting.flows.reserve_payouts import setting from accounting.flows.reserve_payouts import tasks class Flow(FlowBase): """Reserve Payouts Workflow.""" required_params = ['period_id'] available_params = [ 'period_id', 'skip_prepare_label_data', 'validate_only'] def __init__(self): """Initialize the flow using proper attributes.""" self.timeout = setting.START_TO_CLOSE_TIMEOUT super().__init__('reserve_payouts', '1.0') self.env = os.getenv('Environment') def workflow_id(self, initial_context=None): """Build the workflow id. Args: initial_context (dict): the initial context for the flow. Returns: str: a unique identifier for a workflow being executed In the forms of '-', where period_id is a period from initial_context (e.g. 214). """ return ( '{flow_name}-{period_id}').format( flow_name=self.name, period_id=initial_context.get('period_id')) def decider(self, schedule, context): """Activity decider. Args: schedule (callable): the scheduler method. context (dict): Initial context of workflow. """ bootstrap = schedule('bootstrap', self.bootstrap) if bootstrap.result.get('bootstrap.stop'): return prepare_temp_table = ('prepare_temp_table', ['bootstrap']) prepare_label_data = ('prepare_label_data', ['prepare_temp_table']) calculate_reserve_payouts_full_run = ( 'calculate_reserve_payouts', ['prepare_label_data']) calculate_reserve_payouts_no_deps = ('calculate_reserve_payouts', []) validate_reserve_payout_calculation_full_run = ( 'validate_reserve_payout_calculation', ['calculate_reserve_payouts']) validate_only = ( 'validate_reserve_payout_calculation', []) task_to_schedule = ( prepare_temp_table, prepare_label_data, calculate_reserve_payouts_full_run, validate_reserve_payout_calculation_full_run, ) if bootstrap.result.get('bootstrap.validate_only'): task_to_schedule = (validate_only,) elif bootstrap.result.get('bootstrap.skip_prepare_label_data'): task_to_schedule = ( calculate_reserve_payouts_no_deps, validate_reserve_payout_calculation_full_run) scheduled_tasks = {'bootstrap': bootstrap} for task, requirements in task_to_schedule: requires = [] for required_task in requirements: garcon_activity = scheduled_tasks[required_task] if garcon_activity.result.get('stop'): return requires.append(garcon_activity) scheduled_tasks[task] = schedule( task, getattr(self, task), requires=requires) @property def bootstrap(self): """Bootstrap initial configuration.""" return self.create( name='bootstrap', tasks=runner.Sync( tasks.bootstrap.fill( namespace='bootstrap', period_id='period_id', skip_prepare_label_data='skip_prepare_label_data', validate_only='validate_only'))) @property def prepare_temp_table(self): """Populate the temp table for a given period id.""" return self.create( name='prepare_temp_table', tasks=runner.Sync( tasks.prepare_temp_table.fill( period_id='bootstrap.period_id'))) @property def prepare_label_data(self): """Do the label data preparation.""" return self.create( name='prepare_label_data', tasks=runner.Sync( tasks.prepare_label_data.fill( period_id='bootstrap.period_id'))) @property def calculate_reserve_payouts(self): """Calculate the reserves and payouts.""" return self.create( name='calculate_reserve_payouts', tasks=runner.Sync(tasks.calculate_reserve_payouts.fill( period_id='bootstrap.period_id'))) @property def validate_reserve_payout_calculation(self): """Validate reserve payout calculation flow results.""" return self.create( name='validate_reserve_payout_calculation', tasks=runner.Sync( tasks.validate_reserve_payout_calculation.fill()))