"""Application logic layer.""" import asyncio from typing import List from typing import Optional from uuid import uuid4 from dataexport.conf import constants from dataexport.conf.connectors import cursor from dataexport.conf.connectors import get_kafka_consumer from dataexport.dtos import Job from dataexport.dtos import JobQuery from dataexport.helpers.logging_helpers import app_logger from dataexport.helpers.query_helpers import format_query from dataexport.helpers.query_helpers import format_status_query __all__ = ( "create_job", "format_query", "consume_kafka_msgs", "get_all_jobs", ) async def update_status(job: Job) -> Optional[str]: query = format_status_query(job) cursor.execute_async(query) return cursor.sfqid async def create_job(query: JobQuery) -> Job: """Create a job based on a query.""" job_id = str(uuid4()) formatted_query = format_query(query.query_string, job_id) cursor.execute_async(formatted_query) query_id = cursor.sfqid if not query_id: ValueError("Snowflake have not returned any query id") job = Job( id=job_id, status=constants.JOB_STATUS_STARTED, callback=None, **query.dict() ) update_query_id = await update_status(job) app_logger.debug(f"sent status update: {update_query_id}") return job def get_all_jobs() -> List[Job]: """Fetch all jobs.""" # just return some mock data return [ Job( id="abc", query_string="bcd", status=constants.JOB_STATUS_STARTED, callback=None, ) ] async def consume_kafka_msgs(): """Method to consume the messages from Kafka and Update SF table.""" kafka_consumer = get_kafka_consumer() await kafka_consumer.start() try: async for msg in kafka_consumer: try: job_id, job_status = msg.value.split(",") job = Job(id=job_id, status=job_status, query_string="NULL") update_query_id = await update_status(job) app_logger.info(f"sent status update: {update_query_id}") except Exception as e: app_logger.debug(f"Kafka msg '{msg}' failed with exception: {str(e)}") finally: await kafka_consumer.stop() asyncio.sleep(20)