"""EarningsTransferProcessor — orchestrates the full pipeline.""" from __future__ import annotations import os from lambdacommon.common_config import logger from src.calculator import evaluate from src.connectors.s3 import S3Connector, get_s3_connector from src.constants import OUTPUT_PREFIX from src.errors import InputValidationError from src.schemas.responses import ProcessResult from src.types import TransferRecord from src.writer import write_output_to_buffer class EarningsTransferProcessor: """Orchestrates calculate, write, and upload.""" def __init__(self, bucket: str, s3: S3Connector | None = None) -> None: """Set up the processor with bucket name.""" self.bucket = bucket self.s3 = s3 or get_s3_connector() def process_records( self, records: list[TransferRecord], source: str = 'ows-royalties' ) -> ProcessResult: """Process pre-loaded records (manual trigger, no S3 download).""" if not records: raise InputValidationError(f'No records provided from "{source}".') logger.info( 'Processing pre-loaded records', extra={'source': source, 'record_count': len(records)}, ) return self._run_pipeline(records, source, source) def _run_pipeline( self, records: list[TransferRecord], fmt: str, input_key: str ) -> ProcessResult: """Evaluate, write, and upload.""" results = evaluate(records) logger.info( 'Calculation complete', extra={'record_count': len(results)}, ) output_buffer = write_output_to_buffer(results) base_name = os.path.splitext(os.path.basename(input_key))[0] output_key = f'{OUTPUT_PREFIX}{base_name}_output.xlsx' logger.info( 'Uploading output file', extra={'bucket': self.bucket, 'key': output_key}, ) self.s3.upload_buffer(self.bucket, output_key, output_buffer) return ProcessResult( input_key=input_key, output_key=output_key, format=fmt, record_count=len(records), )