250 lines
10 KiB
Python
250 lines
10 KiB
Python
|
|
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
|
||
|
|
}
|
||
|
|
|
||
|
|
|