import os from io import BytesIO import tarfile import boto3 import subprocess import brotli from botocore.exceptions import ClientError import logging from configparser import ConfigParser from urllib.parse import unquote_plus import json # Initialize logger logger = logging.getLogger(__name__) logger.setLevel(logging.INFO) # Initialize S3 & Textract clients s3_client = boto3.client('s3') LIBRE_OFFICE_INSTALL_DIR = '/tmp/instdir' # Function to retrieve configuration values from S3 def load_config_from_s3(bucket_name, file_key): # Download the config file from S3 response = s3_client.get_object(Bucket=bucket_name, Key=file_key) config_content = response['Body'].read().decode('utf-8') # Parse the config file config_parser = ConfigParser() config_parser.read_string(config_content) # Convert the configuration to a dictionary config_dict = {} for section in config_parser.sections(): config_dict[section] = {key.upper(): value for key, value in config_parser.items(section)} return config_dict # Function to move file within S3 def move_file_within_s3(source_bucket, source_path, destination_path): try: # Copy the file to the destination folder s3_client.copy_object(Bucket=source_bucket, CopySource={'Bucket': source_bucket, 'Key': source_path}, Key=destination_path) # Delete the file from the source folder s3_client.delete_object(Bucket=source_bucket, Key=source_path) logger.info(f"File moved from {source_path} to {destination_path}") except Exception as e: logger.error(f"Error moving file: {e}") def load_libre_office(): if os.path.exists(LIBRE_OFFICE_INSTALL_DIR) and os.path.isdir(LIBRE_OFFICE_INSTALL_DIR): print('We have a cached copy of LibreOffice, skipping extraction') else: print('No cached copy of LibreOffice, extracting tar stream from Brotli file.') buffer = BytesIO() with open('/opt/lo.tar.br', 'rb') as brotli_file: d = brotli.Decompressor() while True: chunk = brotli_file.read(1024) buffer.write(d.decompress(chunk)) if len(chunk) < 1024: break buffer.seek(0) print('Extracting tar stream to /tmp for caching.') with tarfile.open(fileobj=buffer) as tar: tar.extractall('/tmp') print('Done caching LibreOffice!') return f'{LIBRE_OFFICE_INSTALL_DIR}/program/soffice.bin' def download_from_s3(bucket, key, download_path): s3 = boto3.client("s3") s3.download_file(bucket, key, download_path) def upload_to_s3(file_path, bucket, key): s3 = boto3.client("s3") s3.upload_file(file_path, bucket, key) def convert_word_to_pdf(soffice_path, word_file_path, output_dir): print(word_file_path) conv_cmd = f"{soffice_path} --headless --norestore --invisible --nodefault --nofirststartwizard --nolockcheck --nologo --convert-to pdf:writer_pdf_Export --outdir {output_dir} {word_file_path}" print(conv_cmd) response = subprocess.run(conv_cmd.split(), stdout=subprocess.PIPE, stderr=subprocess.PIPE) if response.returncode != 0: response = subprocess.run(conv_cmd.split(), stdout=subprocess.PIPE, stderr=subprocess.PIPE) print(response.returncode, response.stdout, response.stderr) if len(response.stderr) > 0: print(response.stderr) if len(response.stdout) > 0: print(response.stdout) if response.returncode != 0: return False return True def lambda_handler(event, context): try: # Read environment variables property_file_path = os.environ.get('PROPERTY_FILE_S3_PATH', '') batch_id = os.environ.get('BATCH_ID', '') # Read config.properties file_path_array = property_file_path.split("/") # Valid if file_path_array has more than 2 elements if len(file_path_array) > 1: # Extract BUCKET_NAME and config_file_path S3_BUCKET_NAME = file_path_array[0] CONFIG_FILE_PATH = "/".join(file_path_array[1:]) logger.info(f'S3_BUCKET_NAME: {S3_BUCKET_NAME}') logger.info(f'CONFIG_FILE_PATH: {CONFIG_FILE_PATH}') # Load config file config_dict = load_config_from_s3(S3_BUCKET_NAME, CONFIG_FILE_PATH) #logger.info('## CONFIG DICTIONARY\r' + str(config_dict)) # Extract configuration values SOURCE_PDF_LOCATION = config_dict['FOLDER_LOCATIONS']['SOURCE_LOCATION'].format(batch_id) # SOURCE_LOCATION SOURCE_DOCX_LOCATION = config_dict['FOLDER_LOCATIONS']['SOURCE_DOCX_LOCATION'].format(batch_id) # SOURCE_DOCX_LOCATION SOURCE_DOCX_PROCESSED_LOCATION = config_dict['FOLDER_LOCATIONS']['SOURCE_DOCX_PROCESSED_LOCATION'].format(batch_id) # SOURCE_DOCX_LOCATION SOURCE_DOCX_UNPROCESSED_LOCATION = config_dict['FOLDER_LOCATIONS']['SOURCE_DOCX_UNPROCESSED_LOCATION'].format(batch_id) # SOURCE_DOCX_LOCATION logger.info('SOURCE_PDF_LOCATION: ' + SOURCE_PDF_LOCATION) logger.info('SOURCE_DOCX_LOCATION: ' + SOURCE_DOCX_LOCATION) logger.info('SOURCE_DOCX_PROCESSED_LOCATION: ' + SOURCE_DOCX_PROCESSED_LOCATION) logger.info('SOURCE_DOCX_UNPROCESSED_LOCATION: ' + SOURCE_DOCX_UNPROCESSED_LOCATION) files = {} files_success = 0 files_failure = 0 # Process each message from the SQS event for record in event['Records']: # Extract the message body from the record record_body = json.loads(record['body']) print('SQS Message Count: ',len(record_body['Records'])) for sqs_record in record_body['Records']: print(sqs_record) # decode source path source_path = unquote_plus(sqs_record['s3']['object']['key']) logging.info("SOURCE_PATH: {source_path} ") # Verify source path have valid file extension if os.path.splitext(source_path)[1].lower() not in ['.doc','.docx']: print('File type not supported: ', os.path.splitext(source_path)[1]) files_failure+=1 continue # Construct destination path #destination_path = SOURCE_PDF_LOCATION + source_path.replace(SOURCE_DOCX_LOCATION,"") processed_destination_path = SOURCE_DOCX_PROCESSED_LOCATION + source_path.replace(SOURCE_DOCX_LOCATION,"") unprocessed_destination_path = SOURCE_DOCX_UNPROCESSED_LOCATION + source_path.replace(SOURCE_DOCX_LOCATION,"") destination_path = "pdf_converted_files/" temporary_filename = "conversion_file" key_prefix, base_name = os.path.split(source_path) file_name, file_ext = os.path.splitext(base_name) download_path = f"/tmp/{temporary_filename}{file_ext}" output_dir = "/tmp" files[source_path] = False # Load Libreoffice library libreoffice_exec_path = "" if os.path.isfile('/opt/lo.tar.br'): logging.info('compressed Libreoffice found!') libreoffice_exec_path = load_libre_office() else: print('libreoffice Layer not found!') return {'body': 'libreoffice Layer not found!'} logging.info("DOWNLOAD_PATH: {download_path}") download_from_s3(S3_BUCKET_NAME, source_path, download_path) logging.info('Downloading Finished!') print("Files list after download: ",os.listdir('/tmp')) logger.info('Starting Conversion') is_converted = convert_word_to_pdf(libreoffice_exec_path, download_path, output_dir) output_filepath = f"{output_dir}/{temporary_filename}.pdf" print("Files list after conversion: ",os.listdir('/tmp')) if is_converted and os.path.isfile(output_filepath): logger.info('Conversion Success!') file_name, _ = os.path.splitext(base_name) files_success+=1 files[source_path] = True logger.info(f'Saving converted file in Bucket {S3_BUCKET_NAME} with Key : {destination_path}{file_name}.pdf') # Upload file to S3 location upload_to_s3(output_filepath, S3_BUCKET_NAME, f"{destination_path}{file_name}.pdf") # Move file to processed folder move_file_within_s3(S3_BUCKET_NAME, source_path, processed_destination_path) logger.info("File converted: " + source_path) else: # Move file to unprocessed folder move_file_within_s3(S3_BUCKET_NAME, source_path, unprocessed_destination_path) files_failure+=1 logger.error("File not converted: " + source_path) files_converted = { 'statusCode': 200, 'body': { 'Files_Processed': len(files), 'Success' : files_success, 'Failure': files_failure } } logger.info(files_converted) return files_converted except ClientError as e: # Handle specific Textract client errors error_message = f"Error in pdf conversion operation: {e}" logger.error(error_message) return { 'statusCode': 500, 'body': error_message } except Exception as e: # Handle other exceptions error_message = f"Unexpected error: {e}" logger.error(error_message) return { 'statusCode': 500, 'body': error_message }