"""analytics_digest workflow.""" from garcon import runner from activity_detector.flows.analytics_digest import config from activity_detector.flows.analytics_digest.tasks.\ create_label_top_new_releases_table import \ create_label_top_new_releases_table from activity_detector.flows.analytics_digest.tasks.\ create_label_top_releases_table import create_label_top_releases_table from activity_detector.flows.analytics_digest.tasks.\ create_label_top_tracks_table import create_label_top_tracks_table from activity_detector.flows.analytics_digest.tasks.\ fetch_last_available_date_per_store \ import fetch_last_available_date_per_store from activity_detector.flows.analytics_digest.tasks.generate_label_digest \ import generate_label_digest from activity_detector.flows.analytics_digest.tasks.get_date_range \ import get_date_range from activity_detector.flows.analytics_digest.tasks.get_growth_date_range \ import get_growth_date_range from activity_detector.flows.base_flow import BaseFlow class Flow(BaseFlow): """Class representing the workflow.""" timeout = 2160 * 60 # 36 hours def __init__(self): """Initialize flow object.""" super(Flow, self).__init__( flow_name=config.SWF_FLOW_NAME, version=config.SWF_FLOW_VERSION) def decider(self, schedule, context): """Activity decider. Args: schedule (callable): The scheduler method. """ if 'label_ids' in context: get_date_range = schedule( 'get_date_range', self.get_date_range) fetch_last_available_date_per_store = schedule( 'fetch_last_available_date_per_store', self.fetch_last_available_date_per_store, requires=[get_date_range]) get_growth_date_range = schedule( 'get_growth_date_range', self.get_growth_date_range, requires=[fetch_last_available_date_per_store]) create_tables = schedule( 'create_tables', self.create_tables, requires=[get_growth_date_range]) schedule( 'generate_label_digests', self.generate_label_digests, requires=[create_tables]) else: get_date_range = schedule( 'get_date_range', self.get_date_range) fetch_last_available_date_per_store = schedule( 'fetch_last_available_date_per_store', self.fetch_last_available_date_per_store, requires=[get_date_range]) get_growth_date_range = schedule( 'get_growth_date_range', self.get_growth_date_range, requires=[fetch_last_available_date_per_store]) create_tables = schedule( 'create_tables', self.create_tables, requires=[get_growth_date_range]) schedule( 'generate_digests', self.generate_digests, requires=[create_tables]) @property def get_date_range(self): """Get date range.""" return self.create( name='get_date_range', tasks=runner.Sync( get_date_range.fill( namespace='get_date_range'))) @property def fetch_last_available_date_per_store(self): """Fetch last available dates per store.""" return self.create( name='fetch_last_available_date_per_store', tasks=runner.Sync( fetch_last_available_date_per_store.fill( namespace='fetch_last_available_date_per_store'))) @property def get_growth_date_range(self): """Get date range for growth percentages.""" return self.create( name='get_growth_date_range', tasks=runner.Sync( get_growth_date_range.fill( namespace='get_growth_date_range', start_date='get_date_range.start_date', end_date='get_date_range.end_date'))) @property def create_tables(self): """Create summary tables.""" return self.create( name='create_tables', tasks=runner.Async( create_label_top_releases_table.fill( namespace='create_label_top_releases_table', start_date='get_date_range.start_date', end_date='get_date_range.end_date', growth_start_date='get_growth_date_range.start_date', growth_end_date='get_growth_date_range.end_date'), create_label_top_tracks_table.fill( namespace='create_label_top_tracks_table', start_date='get_date_range.start_date', end_date='get_date_range.end_date', growth_start_date='get_growth_date_range.start_date', growth_end_date='get_growth_date_range.end_date'), create_label_top_new_releases_table.fill( namespace='create_label_top_new_releases_table', start_date='get_date_range.start_date', end_date='get_date_range.end_date'))) @property def generate_digests(self): """Generate digests.""" return self.create( name='generate_digests', tasks=runner.Sync( generate_label_digest.fill( namespace='generate_label_digest', start_date='get_date_range.start_date', end_date='get_date_range.end_date', last_dates_per_store='fetch_last_available_date_per_store.' 'last_dates_per_store'))) @property def generate_label_digests(self): """Generate label digests.""" return self.create( name='generate_label_digests', tasks=runner.Sync( generate_label_digest.fill( namespace='generate_label_digest', start_date='get_date_range.start_date', end_date='get_date_range.end_date', label_ids='label_ids', last_dates_per_store='fetch_last_available_date_per_store.' 'last_dates_per_store')))