import logging from dataclasses import dataclass from pathlib import Path from typing import Any from src.connectors.auth0_mng import auth0_client from src.connectors.cypher import Neo4jContext, db_client from src.errors import UserProcessError from src.file_processor.base import BaseAsyncFileProcessor from src.file_processor.reader import CsvDictReader, CsvDictReaderT from src.file_processor.writer import CsvDictWriter, CsvDictWriterT from src.result import BaseResult, StatusEnum logger = logging.getLogger('users_cleanup') @dataclass(kw_only=True) class PrepareUsersResult(BaseResult): email: str auth0_id: str | None = None identity_id: str | None = None class PrepareUsers(BaseAsyncFileProcessor[CsvDictReaderT, CsvDictWriterT]): def __init__(self, *, email_key: str, fi: Path, fo: Path, concurrency: int) -> None: self.email_key = email_key super().__init__(fi=fi, fo=fo, concurrency=concurrency) def get_reader(self) -> CsvDictReader: return CsvDictReader(path=self.fi, ensure_keys=(self.email_key,)) def get_writer(self) -> CsvDictWriter: return CsvDictWriter(path=self.fo, queue=self.result_queue, header_keys=PrepareUsersResult.get_fields()) def validate_user_data(self, identity: dict[str, Any], auth0_id: str) -> None: if auth0_id.startswith('auth0|'): auth0_id = auth0_id[6:] if identity['active'] != 'Y': raise UserProcessError('User is not active') if 'auth0UserId' not in identity: raise UserProcessError('identity missing "auth0UserId"') if identity['auth0UserId'] != auth0_id: raise UserProcessError('identity.auth0UserId does not match auth0.user_id') def validate_email(self, email: str) -> None: if not email: raise UserProcessError('Email is empty') async def process_row(self, _index: int, row: CsvDictReaderT) -> CsvDictWriterT: email = row[self.email_key] logger.debug(f'Processing user "{email}"') result = PrepareUsersResult( email=email, status=StatusEnum.SUCCESS, ) try: self.validate_email(email) auth0_id = await auth0_client.get_id_by_email(email) result.auth0_id = auth0_id identity = await db_client.get_identity_by_email(email) result.identity_id = identity['id'] self.validate_user_data(identity, auth0_id) except UserProcessError as e: result.status = StatusEnum.SKIPPED result.message = str(e) except Exception as e: logger.exception('Unhandled exception') result.status = StatusEnum.ERROR result.message = str(e) return result.to_dict() async def process(self) -> None: # Add neo4j context async with Neo4jContext(): return await super().process()