""" Snowflake Caching Monitor Screenboard in Datadog. https://app.datadoghq.com/dashboard/pya-m3v-mvb """ import os from garcon import activity from garcon import runner import raven from swf_monitoring import logger from swf_monitoring.flows import base from swf_monitoring.flows.datadog_caching_monitor import config from swf_monitoring.flows.datadog_caching_monitor import tasks class Flow(base.FlowBaseMixin): """Class representing the workflow.""" def __init__(self): """Initialize flow object.""" self.name = self.generate_feed_name(config.FEED_NAME) self.feed_name = config.FEED_NAME self.version = '1.0' self.domain = os.environ.get('SWF_DOMAIN', 'dev') self.timeout = 60 * 35 self.sentry_dsn = config.SENTRY_DSN if self.sentry_dsn: self.sentry_client = raven.Client() self.create = activity.create( self.domain, self.name, version=self.version, on_exception=self.on_exception) def on_exception(self, actor, exception): """Capture an exception that has occurred in the application. Args: actor (ActivityWorker, DeciderWorker): The actor that has received the exception. exception (Exception): The exception to capture. """ if self.sentry_dsn: self.sentry_client.captureException() if isinstance(actor, activity.Activity): actor.logger.error(exception, exc_info=True) else: logger.error(exception, exc_info=True) def decider(self, schedule): """Activity decider. Args: schedule (callable): The scheduler method. """ populate_cache_stats = schedule( 'populate_cache_stats', self.populate_cache_stats) # Stop flow if metadata query took more than 10 mins and was aborted if populate_cache_stats.result.get('populate_cache_stats.stop'): return {'result': 'Metadata query timeouted'} update_metadata_in_cache_stats_table = schedule( 'update_metadata_in_cache_stats_table', self.update_metadata_in_cache_stats_table, requires=[populate_cache_stats]) send_query_level_metrics_to_datadog = schedule( 'send_query_level_metrics_to_datadog', self.send_query_level_metrics_to_datadog, requires=[update_metadata_in_cache_stats_table]) schedule( 'send_hits_to_kw_cache_percentage_to_datadog', self.send_hits_to_kw_cache_percentage_to_datadog, requires=[send_query_level_metrics_to_datadog]) @property def populate_cache_stats(self): """Populate stats table with the fresh queries.""" return self.create( name='populate_cache_stats', retry=0, tasks=runner.Sync( tasks.populate_cache_stats.fill( namespace='populate_cache_stats'))) @property def update_metadata_in_cache_stats_table(self): """Update cache stats table with metadata from Snowflake API.""" return self.create( name='update_metadata_in_cache_stats_table', retry=0, tasks=runner.Sync( tasks.update_metadata_in_cache_stats_table.fill( namespace='update_metadata_in_cache_stats_table'))) @property def send_query_level_metrics_to_datadog(self): """Send query-level metrics to Datadog.""" return self.create( name='send_query_level_metrics_to_datadog', retry=0, tasks=runner.Sync( tasks.send_query_level_metrics_to_datadog.fill( namespace='send_query_level_metrics_to_datadog'))) @property def send_hits_to_kw_cache_percentage_to_datadog(self): """Send aggregated cumulative daily metric to Datadog.""" return self.create( name='send_hits_to_kw_cache_percentage_to_datadog', retry=0, tasks=runner.Sync( tasks.send_hits_to_kw_cache_percentage_to_datadog.fill( namespace=( 'send_hits_to_kw_cache_percentage_to_datadog'))))