Added docx to pdf lambda function
This commit is contained in:
@@ -0,0 +1,250 @@
|
||||
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
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user