import csv import datetime import io from typing import BinaryIO from email_validator import EmailNotValidError from email_validator import validate_email as _validate_email from snowflake.snowpark import Session from backend.dtos import CcpaRow, CcpaUploadResult from backend.queries import upsert_ccpa_emails def parse_and_validate_ccpa_csv( file: BinaryIO, ) -> tuple[list[CcpaRow], list[str]]: raw_content = file.read() if isinstance(raw_content, bytes): content = raw_content.decode("utf-8-sig") else: content = raw_content.lstrip("\ufeff") reader = csv.reader(io.StringIO(content)) rows = list(reader) if not rows: return [], ["CSV file is empty."] header = [col.strip().lower() for col in rows[0]] if "email" not in header or "requested_date" not in header: return [], [ 'CSV must contain "email" and "requested_date" columns. ' f"Found columns: {', '.join(rows[0])}" ] email_idx = header.index("email") date_idx = header.index("requested_date") data_rows = rows[1:] if not data_rows: return [], ["CSV file contains only a header row with no data."] valid_rows: list[CcpaRow] = [] errors: list[str] = [] for row_num, row in enumerate(data_rows, start=2): if not row or all(cell.strip() == "" for cell in row): continue if len(row) <= max(email_idx, date_idx): errors.append(f"Row {row_num}: missing columns.") continue raw_email = row[email_idx].strip().lower() raw_date = row[date_idx].strip() if not raw_email: errors.append(f"Row {row_num}: email is empty.") continue try: _validate_email(raw_email, check_deliverability=False) except EmailNotValidError: errors.append(f"Row {row_num}: invalid email '{raw_email}'.") continue if not raw_date: errors.append(f"Row {row_num}: requested_date is empty.") continue try: parsed_date = datetime.datetime.strptime(raw_date, "%m/%d/%y").date() except ValueError: errors.append( f"Row {row_num}: invalid date '{raw_date}'. Use MM/DD/YY format." ) continue valid_rows.append(CcpaRow(email=raw_email, requested_date=parsed_date)) # Deduplicate: keep only the oldest (minimum) date per email email_min_date: dict[str, datetime.date] = {} for ccpa_row in valid_rows: if ( ccpa_row.email not in email_min_date or ccpa_row.requested_date < email_min_date[ccpa_row.email] ): email_min_date[ccpa_row.email] = ccpa_row.requested_date deduped_rows = [ CcpaRow(email=email, requested_date=date) for email, date in email_min_date.items() ] return deduped_rows, errors def upload_ccpa_emails_from_csv( session: Session, file: BinaryIO, user: str, env: str = "QA" ) -> CcpaUploadResult: parsed_rows, errors = parse_and_validate_ccpa_csv(file) if errors: return CcpaUploadResult( inserted_count=0, updated_count=0, skipped_count=0, errors=errors, ) if not parsed_rows: return CcpaUploadResult( inserted_count=0, updated_count=0, skipped_count=0, errors=[], ) inserted_count, updated_count, skipped_count = upsert_ccpa_emails( session, parsed_rows, user, env ) return CcpaUploadResult( inserted_count=inserted_count, updated_count=updated_count, skipped_count=skipped_count, errors=[], )