#!/usr/bin/env python3 """Monitor running workflows and kill task when none are active.""" from datetime import datetime, timedelta import os import sys import time import boto3 import botocore from owslogger import logger import requests import sentry_sdk import config # Global logger log = logger.setup( config.LOGGER_DSN, config.ENVIRONMENT, config.LOGGER_NAME, config.LOGGER_LEVEL, config.SERVICE_NAME, '1.0.0') # Sentry handler sentry_sdk.init(config.SENTRY_DSN) def describe_open_swf_workflows(): """Check for open workflows matching config parameters. Returns: list: A list of active workflows """ client = boto3.client('swf', region_name=config.AWS_DEFAULT_REGION) try: result = client.list_open_workflow_executions( domain=config.SWF_MONITOR_DOMAIN, startTimeFilter={ 'oldestDate': datetime.utcnow() - timedelta( days=config.SWF_MONITOR_MAX_RUNNING_DAYS), 'latestDate': datetime.utcnow() } ) executions = result['executionInfos'] open_executions = [] for flow in executions: if flow['workflowType']['name'] == \ config.SWF_MONITOR_WORKFLOW_NAME \ and flow['executionStatus'] == 'OPEN': open_executions.append(flow['execution']['workflowId']) return open_executions except botocore.exceptions.ClientError as error: sentry_sdk.capture_exception(error) log.error('Error polling SWF: {}'.format(error)) sys.exit(1) def get_task_arn(): """ Get the ARN of the currently running task. This only works when running in ECS Fargate, since it uses the task metadata endpoint to retrieve info. Returns: str: cluster of currently running task str: currently running task arn """ try: # First try v4 (platform 1.4.0+) metadata endpoint if os.environ.get('ECS_CONTAINER_METADATA_URI_V4'): url = os.environ.get('ECS_CONTAINER_METADATA_URI_V4') # If that is not available, fall back to the older endpoint elif os.environ.get('ECS_CONTAINER_METADATA_URI'): url = os.environ.get('ECS_CONTAINER_METADATA_URI') else: log.error('Task metadata URI cannot be determined') sys.exit(1) response = requests.get('{}/task'.format(url)) json_response = response.json() return json_response['Cluster'], json_response['TaskARN'] except requests.exceptions.RequestException as error: sentry_sdk.capture_exception(error) log.error('Error retrieving task metadata: {}'.format(error)) def kill_running_fargate_task(cluster, task_arn): """Kill the task in which this code is running. Args: cluster (str): cluster in which the task is running task_arn (str): ARN of task to kill """ client = boto3.client('ecs', region_name=config.AWS_DEFAULT_REGION) client.describe_tasks( cluster=cluster, tasks=[task_arn], ) response = client.stop_task( cluster=cluster, task=task_arn, reason='No active workflows found for {}'.format( config.SWF_MONITOR_WORKFLOW_NAME) ) log.info('Desired task status: {}'.format( response['task']['desiredStatus'])) def main(): """Execute main entrypoint.""" log.info('Polling SWF for status of {} flows'.format( config.SWF_MONITOR_WORKFLOW_NAME)) timeout = time.time() + config.SWF_MONITOR_TASK_RUNTIME_LIMIT_MINUTES * 60 grace_period = time.time() + config.SWF_MONITOR_GRACE_PERIOD_SECONDS while time.time() < timeout: running_flows = describe_open_swf_workflows() if running_flows: log.info('Workflow IDs {} are active'.format(running_flows)) time.sleep(60) else: if time.time() < grace_period: log.info('No workflows active for {}, but task runtime is ' 'within grace period of {} seconds'.format( config.SWF_MONITOR_WORKFLOW_NAME, config.SWF_MONITOR_GRACE_PERIOD_SECONDS ) ) time.sleep(60) else: log.info('No workflows active for {}. Killing task.'.format( config.SWF_MONITOR_WORKFLOW_NAME)) """ Wait a minute for any garcon/swf cleanup tasks to finish after execution is closed """ time.sleep(60) cluster, task_arn = get_task_arn() log.info('Killing task ARN {}.'.format(task_arn)) kill_running_fargate_task(cluster, task_arn) if __name__ == '__main__': main()