import json import validictory from fpcapture.connectors import sentry from fpcapture.logic import codegen from fpcapture.validation_schema import message_schema DEFAULT_NUM_MESSAGES = 10 # We generally want to use a long visibility timeout so we have enough time # to process and ACK a message before it becomes available to consumers again. DEFAULT_VISIBILITY_TIMEOUT = 60 DEFAULT_WAIT_TIME_SECONDS = 20 def poll_in_loop( queue, db_session, timer, num_messages=DEFAULT_NUM_MESSAGES, visibility_timeout=DEFAULT_VISIBILITY_TIMEOUT, wait_time_seconds=DEFAULT_WAIT_TIME_SECONDS): """Polls the given queue until the timer expires. Args: queue (boto.sqs.queue.Queue): queue object to poll db_session (sqlalchemy.orm.session.Session): DB session timer (utils.timer.Timer): timer object that indicates when to stop num_messages (int): number of messages to fetch from queue on each read visibility_timeout (int): number of seconds for messages to remain in flight after being read from the queue wait_time_seconds (int): number of seconds for each SQS read to wait before returning, if no messages are immediately available """ while not timer.is_expired(): poll_queue( db_session=db_session, num_messages=num_messages, queue=queue, visibility_timeout=visibility_timeout, wait_time_seconds=wait_time_seconds) def poll_queue( queue, db_session, num_messages=DEFAULT_NUM_MESSAGES, visibility_timeout=DEFAULT_VISIBILITY_TIMEOUT, wait_time_seconds=DEFAULT_WAIT_TIME_SECONDS): """Polls the queue and processes any messages that are returned Args: queue (boto.sqs.queue.Queue): queue object to poll db_session (sqlalchemy.orm.session.Session): DB session num_messages (int): number of messages to fetch from queue on each read visibility_timeout (int): number of seconds for messages to remain in flight after being read from the queue wait_time_seconds (int): number of seconds for each SQS read to wait before returning, if no messages are immediately available """ messages = queue.get_messages( num_messages=num_messages, wait_time_seconds=wait_time_seconds, visibility_timeout=visibility_timeout) for message in messages: if process_message(message, db_session): queue.delete_message(message) def process_message(message, db_session): """Parses a message from the queue and passes it along to the codegen logic Args: message (boto.sqs.message.Message): message from the queue db_session (sqlalchemy.orm.session.Session): DB session Returns: bool: True on success, False on error """ try: payload = json.loads(message.get_body()) validictory.validate(payload, message_schema.schema) except ValueError: sentry.sentry_client.captureException() return False codegen.encode( db_session=db_session, tuid=payload.get('tuid'), upc=payload.get('upc'), filename=payload.get('filename'), file_upload_time=payload.get('file_upload_time'), correlation_id=payload.get('correlation_id'), track_source=payload.get('track_source')) return True