import io import logging import time import uuid from uuid import UUID from common.connectors import s3 from common.connectors.graphql import obo_graphql as graphql from common.logic import fargate from common.schemas.s3_reference import S3Reference import config from src.connectors import sfn from src.logic import validation logger = logging.getLogger() def main() -> None: """ Fargate task entry point. This function is executed when the task is started, and depends on the environment variables "KEY" and "BUCKET" to be set. In a step function, these will be provided by the previous step. In local development, these are provided by the event.json file. """ # Figure out why .info() logs don't show up in datadog # TODO .warning() logs should be .info() logger.warning("Started validate_spreadsheet") start_time = time.time() s3_reference: S3Reference | None = validation.get_s3_reference() task_token: str | None = fargate.get_task_token() if not task_token: error_msg = "Task token not provided. Cannot proceed without task token." logger.error(error_msg) raise ValueError(error_msg) try: if not s3_reference: error_msg = ( "S3 reference not provided. Cannot proceed without S3 reference." ) logger.error(error_msg) raise ValueError(error_msg) result = s3.get_object(s3_reference) stream = io.BytesIO(result["Body"].read()) metadata = s3.get_object_metadata(s3_reference) bulk_session_id = UUID(metadata["bulk_session_id"]) identity_uuid = metadata["identity_uuid"] graphql.init_graphql_client( environment=config.ENVIRONMENT, service_name=config.M2M_APPLICATION_NAME, graphql_service_name=config.GRAPHQL_SERVICE_NAME, identity_id=identity_uuid, profile_id=config.PROFILE_ID, profile_type=config.PROFILE_TYPE, headers={ "Correlation-Id": str(uuid.uuid4()), "Orchard-Roles": config.ROLE, }, ) response = validation.stream_to_response( task_token=task_token, stream=stream, bulk_session_id=bulk_session_id, identity_uuid=identity_uuid, ) if response["is_valid"]: sfn.send_task_success( token=task_token, payload={"is_valid": True, "is_classical": response["is_classical"]}, ) else: output_key = str(f"errors/{uuid.uuid4()}.json") s3.put_object( s3_reference=S3Reference(bucket=s3_reference.bucket, key=output_key), stream=validation.dict_to_stream(response), ) sfn.send_task_success( token=task_token, payload={"errors": output_key, "is_valid": False} ) except Exception as e: logger.error("validate_spreadsheet failed with exception %s", str(e)) sfn.send_task_failure(token=task_token, cause=str(e)) duration = int(time.time() - start_time) logger.warning("Finished validate_spreadsheet after %d seconds", duration) if __name__ == "__main__": main()