import boto.swf import os from datetime import datetime, timedelta from swf.event import Event import time def fetch_events(domain, run_id, workflow_id): events = {} completed = False aws_swf_l1 = boto.swf.layer1.Layer1( aws_access_key_id=os.environ.get('AWS_ACCESS_KEY_ID'), aws_secret_access_key=os.environ.get('AWS_SECRET_ACCESS_KEY')) while not completed: history = aws_swf_l1.get_workflow_execution_history( domain=domain, run_id=run_id, workflow_id=workflow_id) for event in history['events']: if event['eventId'] not in events: yield event events.update({event['eventId']: Event(event)}) if event['eventType'] == 'WorkflowExecutionCompleted': completed = True def find_all_types(key_id, key_secret): l1 = boto.swf.layer1.Layer1( aws_access_key_id=key_id, aws_secret_access_key=key_secret) workflow_types = [] for domain in l1.list_domains('REGISTERED')['domainInfos']: w_types = l1.list_workflow_types( domain['name'], 'REGISTERED')['typeInfos'] for flow_type in w_types: flow_type['domain'] = domain['name'] workflow_types.append(flow_type) return workflow_types def find_top_3_executions(key_id, key_secret, workflow_type): l1 = boto.swf.layer1.Layer1( aws_access_key_id=key_id, aws_secret_access_key=key_secret) try: executions = l1.list_open_workflow_executions( oldest_date=(datetime.now() - timedelta(days=90)).timestamp(), domain=workflow_type['domain'], workflow_name=workflow_type['workflowType']['name'], workflow_version=workflow_type['workflowType']['version'], maximum_page_size=5)['executionInfos'] executions += l1.list_closed_workflow_executions( start_oldest_date=(datetime.now() - timedelta(days=90)).timestamp(), domain=workflow_type['domain'], workflow_name=workflow_type['workflowType']['name'], workflow_version=workflow_type['workflowType']['version'], maximum_page_size=5)['executionInfos'] except Exception: time.sleep(0.5) return find_top_3_executions(key_id, key_secret, workflow_type) executions.sort(key=lambda x: x['startTimestamp']) return executions def all_regions(): return boto.swf.get_regions('swf')