from abc import ABC, abstractmethod from typing import Tuple from sqlalchemy import and_, select, func, update from sqlalchemy.ext.asyncio import AsyncSession from common.database.db_session_manager import DB_SESSION_MNG from common.database.model.models import places, users from common.enums import ErrorType from common.logger import LOG from common.utils.gtime import GTime def _place_count_subquery(): """계정 하나가 가진 사업장 수 — users 에 상관 서브쿼리로 얹는다(N+1 회피).""" return ( select(func.count()) .select_from(places) .where(places.owner_user_id == users.user_id, places.deleted == False) # noqa: E712 .correlate(users) .scalar_subquery() ) # CRUD 는 인터페이스(I*) 와 구현(*) 으로 분리한다. class IUserCRUD(ABC): @abstractmethod async def get_user_by_login_id(self, cdb: AsyncSession, login_id: str) -> Tuple[ErrorType, users]: pass @abstractmethod async def get_user_by_provider_uid(self, cdb: AsyncSession, provider: int, provider_uid: str) -> Tuple[ErrorType, users]: pass @abstractmethod async def get_user_by_email(self, cdb: AsyncSession, email: str) -> Tuple[ErrorType, users]: pass @abstractmethod async def is_user(self, cdb: AsyncSession, login_id: str) -> ErrorType: pass @abstractmethod async def add_user(self, cdb: AsyncSession, user: users) -> ErrorType: pass @abstractmethod async def update_last_accessed(self, cdb: AsyncSession, user_id) -> ErrorType: pass @abstractmethod async def get_by_user_id(self, cdb: AsyncSession, user_id) -> Tuple[ErrorType, users]: pass @abstractmethod async def update_user(self, cdb: AsyncSession, user_id, data: dict) -> ErrorType: pass @abstractmethod async def list_users( self, cdb: AsyncSession, roles: list, skip: int, limit: int, search: str | None = None ) -> Tuple[ErrorType, list, int]: """내부 운영(DEVELOPER) 전용 전체 계정 목록.""" pass class UserCRUD(IUserCRUD): async def get_user_by_login_id(self, cdb: AsyncSession, login_id: str) -> Tuple[ErrorType, users]: try: query = select(users).where(users.id == login_id, users.deleted == False).limit(1) # noqa: E712 err_type, row_list = await DB_SESSION_MNG.execute(cdb, query, f"get_user_by_login_id(ID:{login_id}) failed.") if err_type != ErrorType.SUCCESS: return err_type, None if len(row_list) != 1: return ErrorType.DB_INVALID_KEY, None return ErrorType.SUCCESS, row_list[0] except Exception as ex: LOG.e_no_callstack(ex) return ErrorType.DB_RUN_FAILED, None async def get_user_by_provider_uid(self, cdb: AsyncSession, provider: int, provider_uid: str) -> Tuple[ErrorType, users]: """소셜 계정 조회 키는 provider_uid(구글 sub) 다 — 이메일이 아니다.""" try: query = ( select(users) .where(users.provider == provider, users.provider_uid == provider_uid, users.deleted == False) # noqa: E712 .limit(1) ) err_type, row_list = await DB_SESSION_MNG.execute(cdb, query, "get_user_by_provider_uid failed.") if err_type != ErrorType.SUCCESS: return err_type, None if len(row_list) != 1: return ErrorType.DB_INVALID_KEY, None return ErrorType.SUCCESS, row_list[0] except Exception as ex: LOG.e_no_callstack(ex) return ErrorType.DB_RUN_FAILED, None async def get_user_by_email(self, cdb: AsyncSession, email: str) -> Tuple[ErrorType, users]: """이메일로 1건.""" try: query = ( select(users) .where(func.lower(users.email) == email.strip().lower(), users.deleted == False) # noqa: E712 .order_by(users.created_at.asc()) .limit(1) ) err_type, row_list = await DB_SESSION_MNG.execute(cdb, query, "get_user_by_email failed.") if err_type != ErrorType.SUCCESS: return err_type, None if len(row_list) != 1: return ErrorType.DB_INVALID_KEY, None return ErrorType.SUCCESS, row_list[0] except Exception as ex: LOG.e_no_callstack(ex) return ErrorType.DB_RUN_FAILED, None async def is_user(self, cdb: AsyncSession, login_id: str) -> ErrorType: try: query = select(users).where(users.id == login_id, users.deleted == False).limit(1) # noqa: E712 err_type, row_list = await DB_SESSION_MNG.execute(cdb, query) if err_type != ErrorType.SUCCESS: return err_type if row_list: return ErrorType.DB_ALREADY_SAME_KEY return ErrorType.SUCCESS except Exception as ex: LOG.e_no_callstack(ex) return ErrorType.DB_RUN_FAILED async def add_user(self, cdb: AsyncSession, user: users) -> ErrorType: try: return await DB_SESSION_MNG.insert(cdb, user) except Exception as ex: LOG.e_no_callstack(ex) return ErrorType.DB_RUN_FAILED async def update_last_accessed(self, cdb: AsyncSession, user_id) -> ErrorType: try: query = update(users).where(users.user_id == user_id).values(last_accessed_at=GTime.UTC()) return await DB_SESSION_MNG.add(cdb, query) except Exception as ex: LOG.e_no_callstack(ex) return ErrorType.DB_RUN_FAILED async def get_by_user_id(self, cdb: AsyncSession, user_id) -> Tuple[ErrorType, users]: try: query = select(users).where(users.user_id == user_id, users.deleted == False).limit(1) # noqa: E712 err_type, row_list = await DB_SESSION_MNG.execute(cdb, query) if err_type != ErrorType.SUCCESS: return err_type, None if len(row_list) != 1: return ErrorType.DB_INVALID_KEY, None return ErrorType.SUCCESS, row_list[0] except Exception as ex: LOG.e_no_callstack(ex) return ErrorType.DB_RUN_FAILED, None async def update_user(self, cdb: AsyncSession, user_id, data: dict) -> ErrorType: try: if not data: return ErrorType.SUCCESS query = update(users).where(users.user_id == user_id).values(**data) return await DB_SESSION_MNG.add(cdb, query) except Exception as ex: LOG.e_no_callstack(ex) return ErrorType.DB_RUN_FAILED async def list_users( self, cdb: AsyncSession, roles: list, skip: int, limit: int, search: str | None = None ) -> Tuple[ErrorType, list, int]: """role 이 roles 안에 있는 계정만 본다 — 개발자 계정은 호출측이 roles 에서 뺀다 (UserRole 주석: "개발자 계정은 고객사에 존재를 노출하지 않는다" 원칙을 내부 화면에서도 지킨다).""" try: where = and_(users.deleted == False, users.role.in_(roles)) # noqa: E712 if search: like = f"%{search.strip()}%" where = and_(where, (users.email.ilike(like) | users.name.ilike(like) | users.id.ilike(like))) cnt_err, cnt_rows = await DB_SESSION_MNG.execute(cdb, select(func.count()).select_from(users).where(where)) if cnt_err != ErrorType.SUCCESS: return cnt_err, [], 0 total = int(cnt_rows[0] or 0) if cnt_rows else 0 query = ( select(users, _place_count_subquery()) .where(where) .order_by(users.created_at.desc()) .offset(skip) .limit(limit) ) list_err, rows = await DB_SESSION_MNG.execute(cdb, query) if list_err != ErrorType.SUCCESS: return list_err, [], 0 return ErrorType.SUCCESS, list(rows), total except Exception as ex: LOG.e_no_callstack(ex) return ErrorType.DB_RUN_FAILED, [], 0