"""Base classes module.""" import abc import csv from dataclasses import dataclass from datetime import datetime from io import StringIO import logging import pathlib from typing import Any, cast, Collection, Dict, List, Literal, Optional, Self from sentry_sdk import capture_message from config import app_logger as logger from src.connectors import s3 from src.processors.exceptions import ProcessorConfigurationError @dataclass class LogEntry: """Log entry class.""" message: str group_condition: str | None = None additional_data: dict[str, Any] | None = None class Processor(abc.ABC): """Base processor.""" _input_file_buffer: StringIO _bucket_name: str _file_path: str _logs: Optional[Dict[int, List[LogEntry]]] = None def __init__(self, bucket_name: str, file_path: str) -> None: """Init.""" self._bucket_name = bucket_name self._file_path = file_path @classmethod def execute(cls, bucket_name: str, file_path: str) -> Self: """Make some config and run the logic.""" instance = cls(bucket_name, file_path) instance._download_file() # _report_logs must run even when process() raises, otherwise every # row-level _add_log accumulated before the failure is silently # dropped (no Sentry, no logger output). Swallow any exception from # _report_logs in the failure path so it never masks the original # process() exception. Catch `Exception` (not `BaseException`) so # KeyboardInterrupt / SystemExit / GeneratorExit propagate without # the inner Sentry I/O — flushing logs during interpreter shutdown # has caused deadlocks in the past. try: instance.process() except Exception: try: instance._report_logs() except Exception: bucket = getattr(instance, '_bucket_name', '?') file_path = getattr(instance, '_file_path', '?') logger.exception( 'Failed to flush logs while handling a processor error ' f'(bucket={bucket} file={file_path})' ) raise instance._report_logs() return instance def _download_file(self) -> None: """Download file from S3 bucket.""" logger.info(f'Processing file {self._file_path}') try: self._input_file_buffer = s3.download_file( self._bucket_name, self._file_path ) except Exception as e: logger.error( f'Unable to get the file from bucket: {self._bucket_name}/{self._file_path}, {e}' ) raise e @property def csv_dict_reader(self) -> csv.DictReader: # type: ignore """Get CSV dict reader.""" return csv.DictReader(self._input_file_buffer) @abc.abstractmethod def process(self) -> None: """Run main processing logic.""" def _add_log( self, message: str, log_level: int = logging.ERROR, group_condition: Optional[str] = None, **additional_data: Any, ) -> None: """Add error.""" if self._logs is None: self._logs = {} self._logs.setdefault(log_level, []).append( LogEntry(message, group_condition, additional_data or None) ) def _report_logs(self) -> None: """Report errors.""" if not self._logs: return for log_level, log_entries in self._logs.items(): level_name = logging.getLevelName(log_level).lower() logger.log(log_level, f'{level_name}s count: {len(log_entries)}') log_messages = [] grouped_messages: dict[str, list[LogEntry]] = {} for entry in log_entries: if entry.group_condition: grouped_messages.setdefault(entry.group_condition, []).append(entry) else: log_messages.append(entry.message) for group_condition, group_entries in grouped_messages.items(): first_entry = group_entries[0] if not first_entry.additional_data: raise ValueError('Provide parameters to specify each log entry') if len(first_entry.additional_data) == 1: field_name = list(first_entry.additional_data.keys())[0] group_data = [ entry.additional_data[field_name] for entry in group_entries if entry.additional_data ] data_str = f'{field_name} {", ".join(sorted(group_data))}' else: group_data = [ ', '.join( sorted(f'{k}={v}' for k, v in entry.additional_data.items()) ) for entry in group_entries if entry.additional_data ] data_str = '; '.join(sorted(f'({i})' for i in group_data)) log_messages.append( f'{first_entry.message} | {group_condition}: {data_str}' ) for message in log_messages: logger.log(log_level, message) if log_level >= logging.WARNING: single_error_message = '\n'.join(log_messages) capture_message( single_error_message, level=cast( Literal[ 'fatal', 'critical', 'error', 'warning', 'info', 'debug' ], level_name, ), ) class ReportProcessor(Processor, abc.ABC): """Base processor that returns file result.""" _report_type: str _output_file_buffer: StringIO | None = None _output_columns: Collection[str] _output_data: list[Any] @classmethod def execute(cls, bucket_name: str, file_path: str) -> Self: """Make some config and run the logic.""" instance: Self = super().execute(bucket_name, file_path) instance._generate_report() instance._upload_file() return instance def _upload_file(self) -> None: """Upload result file.""" path = pathlib.PurePath(self._file_path) try: s3.upload_file( self._bucket_name, f'reports/{self.report_file_name}', self.get_output_file_buffer(True), metadata={'original-file-name': path.name}, tags={'report-type': self.report_type}, ) except Exception as e: logger.error( f'Unable to put a file to the bucket: {self._bucket_name}/{self._file_path}, {e}' ) raise e @property def report_type(self) -> str: """Report type.""" return self._report_type @property def report_file_name(self) -> str: """Generate report file name.""" ts = datetime.now().strftime('%Y%m%d_%H%M%S') return f'{self._report_type}_report_{ts}.csv' def get_output_file_buffer(self, for_reading: bool = False) -> StringIO: """Get result memory buffer.""" if not self._output_file_buffer: self._output_file_buffer = StringIO() elif for_reading: self._output_file_buffer.seek(0) return self._output_file_buffer def _generate_report(self) -> None: """Generate report.""" if self._output_data is None: raise ProcessorConfigurationError('Please set result data') if not self._output_data: return if not hasattr(self, '_output_columns') or not self._output_columns: self._output_columns = self._output_data[0].keys() writer = csv.DictWriter( self.get_output_file_buffer(), fieldnames=self._output_columns ) writer.writeheader() for item in self._output_data: writer.writerow(dict(item))