"""File upload helper functions.""" import logging import math from owsresponse import response from abacus_file_upload.connectors.s3 import S3Connector from abacus_file_upload.constants import S3_MAX_PARTS, UPLOAD_STATUSES from abacus_file_upload.models import FileUpload, FileUploadConfig from abacus_file_upload.utils.exception import ( AWSOperationError, ) from abacus_file_upload.utils.s3 import ( calculate_optimal_chunk_size, convert_md5_hex_to_base64, format_bytes, get_file_extension, validate_file_type, ) from core.config import Config logger = logging.getLogger(Config.LOGGER_NAME) def cleanup_s3_upload(s3_connector: S3Connector, file_upload: FileUpload) -> None: """Clean up S3 resources for an upload. Args: s3_connector: S3 connector instance file_upload: FileUpload record to clean up """ if file_upload.multipart_upload_id: # Abort multipart upload try: s3_connector.abort_multipart_upload( file_upload.s3_bucket, file_upload.s3_key, file_upload.multipart_upload_id, ) except Exception as e: logger.warning(f'Failed to abort S3 multipart upload: {e}') else: # Delete single-part uploaded file if it exists try: if s3_connector.object_exists(file_upload.s3_bucket, file_upload.s3_key): s3_connector.delete_object(file_upload.s3_bucket, file_upload.s3_key) except Exception as e: logger.warning(f'Failed to delete S3 object: {e}') def handle_multipart_upload( s3_connector: S3Connector, s3_bucket: str, s3_key: str, file_size_bytes: int, min_chunk_size: int, s3_metadata: dict, mime_type: str | None, ) -> dict: """Handle multipart upload setup and URL generation. Args: s3_connector: S3 connector instance s3_bucket: S3 bucket name s3_key: S3 object key file_size_bytes: File size in bytes min_chunk_size: Minimum chunk size for multipart s3_metadata: S3 metadata dict mime_type: Optional MIME type Returns: Dict with parts, complete_url, chunk_size_bytes, multipart_upload_id, total_parts Raises: ValueError: If file too large for multipart """ # Calculate optimal chunk size chunk_size_bytes = calculate_optimal_chunk_size(file_size_bytes, min_chunk_size) # Validate part count total_parts = math.ceil(file_size_bytes / chunk_size_bytes) if total_parts > S3_MAX_PARTS: raise ValueError( f'File too large: would require {total_parts} parts (max {S3_MAX_PARTS})' ) # Initiate multipart upload with S3 multipart_upload_id = s3_connector.initiate_multipart_upload( s3_bucket, s3_key, metadata=s3_metadata, content_type=mime_type, ) # Generate presigned URLs for all parts presigned_parts = s3_connector.generate_multipart_presigned_urls( s3_bucket, s3_key, multipart_upload_id, total_parts ) parts_response = [ { 'part_number': p['part_number'], 'url': p['url'], 'expires_at': p['expires_at'].isoformat(), } for p in presigned_parts ] # Generate complete URL complete_url = s3_connector.generate_complete_multipart_presigned_url( s3_bucket, s3_key, multipart_upload_id ) return { 'parts': parts_response, 'complete_url': complete_url, 'chunk_size_bytes': chunk_size_bytes, 'multipart_upload_id': multipart_upload_id, 'total_parts': total_parts, } def handle_single_part_upload( s3_connector: S3Connector, s3_bucket: str, s3_key: str, s3_metadata: dict, md5sum: str, mime_type: str | None, ) -> dict: """Handle single-part upload setup and URL generation. Args: s3_connector: S3 connector instance s3_bucket: S3 bucket name s3_key: S3 object key s3_metadata: S3 metadata dict md5sum: MD5 hash mime_type: Optional MIME type Returns: Dict with upload_url and required_headers """ # Generate single-part upload URL with metadata, ContentMD5, and ContentType content_md5 = convert_md5_hex_to_base64(md5sum) upload_url = s3_connector.generate_put_presigned_url( s3_bucket, s3_key, metadata=s3_metadata, content_md5=content_md5, content_type=mime_type, ) # Return required headers that client must send with PUT request # Include Content-MD5, Content-Type, and all x-amz-meta-* headers required_headers = {'Content-MD5': content_md5} if mime_type: required_headers['Content-Type'] = mime_type # Add metadata as x-amz-meta-* headers (required for signature validation) for key, value in s3_metadata.items(): # Keys are already lowercase, just add the x-amz-meta- prefix metadata_key = f'x-amz-meta-{key}' required_headers[metadata_key] = value return {'upload_url': upload_url, 'required_headers': required_headers} def handle_upload_operation_error( file_upload: FileUpload, file_key: str, error: Exception, operation: str, status_code: int, ) -> response.Response: """Handle errors during upload operations (e.g., complete, cancel). Marks the upload status as ERROR and returns an error response. Args: file_upload: FileUpload record to update file_key: File key for logging error: The exception that occurred operation: Operation being performed (e.g., complete, cancel) status_code: HTTP status code to return Returns: Error response """ try: file_upload.update_attributes( upload_status=UPLOAD_STATUSES.ERROR, error_message=str(error), ) FileUpload.commit_changes() except Exception as db_error: logger.error(f'Failed marking upload {file_key} as ERROR: {db_error}') logger.error(f'Failed to {operation} upload {file_key}: {error}') return response.Response( message=f'Failed to {operation} upload: {str(error)}', status=status_code ) def validate_with_upload_config( config: FileUploadConfig, filename: str, file_size_bytes: int, ) -> None: """Validate file type and size against config. Args: config: Upload configuration filename: Original filename file_size_bytes: File size in bytes Raises: ValueError: If validation fails """ # Validate file type file_type = get_file_extension(filename) if not validate_file_type(file_type, config.allowed_file_types): allowed = ', '.join(config.allowed_file_types or []) raise ValueError(f"File type not allowed. Allowed types: '{allowed}'") # Validate file size if file_size_bytes > config.max_file_size_bytes: max_size = format_bytes(config.max_file_size_bytes) raise ValueError(f"File size exceeds maximum allowed size of '{max_size}'") def verify_s3_upload(s3_connector: S3Connector, file_upload: FileUpload) -> None: """Verify S3 upload is complete and valid. Args: s3_connector: S3 connector instance file_upload: FileUpload record to verify Raises: ValueError: If verification fails """ # Verify file exists in S3 if not s3_connector.object_exists(file_upload.s3_bucket, file_upload.s3_key): raise ValueError('File not found in S3. Upload may not be complete.') # Check file size matches expected metadata = s3_connector.get_object_metadata( file_upload.s3_bucket, file_upload.s3_key ) actual_size = metadata.get('size', 0) if actual_size != file_upload.file_size_bytes: raise ValueError( f'File size mismatch: expected {file_upload.file_size_bytes}, got {actual_size}' ) # For single-part uploads, verify MD5 matches if not file_upload.multipart_upload_id and file_upload.md5sum: s3_etag = metadata.get('etag', '').strip('"').lower() expected_md5 = file_upload.md5sum.lower() if s3_etag != expected_md5: raise ValueError( f'MD5 verification failed: expected {expected_md5}, got {s3_etag}' ) def quarantine_s3_upload(s3_connector: S3Connector, file_upload: FileUpload) -> None: """Move infected uploads to the quarantine bucket. Args: s3_connector: S3 connector instance file_upload: FileUpload record to move """ does_s3_object_exist = s3_connector.object_exists( file_upload.s3_bucket, file_upload.s3_key ) if does_s3_object_exist is False: raise AWSOperationError( f"File doesn't exist in source bucket {Config.S3_ABACUS_ADJUSTMENTS_BUCKET}" ) copy_response = s3_connector.copy_object( file_upload.s3_bucket, Config.S3_ABACUS_QUARANTINE_BUCKET, file_upload.s3_key, ) copy_status = copy_response['ResponseMetadata']['HTTPStatusCode'] if copy_status not in [200, 201]: logging.error(f'Copy failed with status: {copy_status}') raise AWSOperationError( f'Failed to copy object to {Config.S3_ABACUS_QUARANTINE_BUCKET}' ) delete_response = s3_connector.delete_object( file_upload.s3_bucket, file_upload.s3_key ) del_status = delete_response['ResponseMetadata']['HTTPStatusCode'] if del_status not in [200, 204]: logging.error(f'Delete failed with status: {copy_status}') raise AWSOperationError(f'Failed to delete S3 object.')