o2o-negosium-original/negodata/backend/crud/notification_crud.py
Mina Choi 1d111a031e [feat] negodata: 알림함 — 견적 마감 결과(낙찰/재생성/결렬) 자동 통지
- company.notifications 테이블 + NotificationType(SUCCESS/REGENERATED/FAILURE)
- 백엔드 알림 조회/읽음 API + close_and_decide 결과 분기마다 알림 생성
- 헤더 알림 벨(안읽음 배지) + 알림 페이지(읽음·딥링크)

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-30 17:24:59 +09:00

108 lines
4.0 KiB
Python

from abc import ABC, abstractmethod
from typing import Tuple
from sqlalchemy import 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 notifications
from common.enums import ErrorType
from common.logger import LOG
# 알림 CRUD. 항상 user_id(수신자)로 스코프한다.
class INotificationCRUD(ABC):
@abstractmethod
async def list_for_user(self, cdb: AsyncSession, user_id, skip, limit) -> Tuple[ErrorType, list, int]:
pass
@abstractmethod
async def count_unread(self, cdb: AsyncSession, user_id) -> Tuple[ErrorType, int]:
pass
@abstractmethod
async def mark_read(self, cdb: AsyncSession, user_id, notification_id, ts) -> ErrorType:
pass
@abstractmethod
async def mark_all_read(self, cdb: AsyncSession, user_id, ts) -> ErrorType:
pass
class NotificationCRUD(INotificationCRUD):
async def list_for_user(self, cdb: AsyncSession, user_id, skip, limit) -> Tuple[ErrorType, list, int]:
try:
cnt_err, cnt_rows = await DB_SESSION_MNG.execute(
cdb,
select(func.count()).select_from(notifications).where(
notifications.user_id == user_id, notifications.deleted == False # noqa: E712
),
)
if cnt_err != ErrorType.SUCCESS:
return cnt_err, [], 0
total = int(cnt_rows[0] or 0) if cnt_rows else 0
list_err, rows = await DB_SESSION_MNG.execute(
cdb,
select(notifications)
.where(notifications.user_id == user_id, notifications.deleted == False) # noqa: E712
.order_by(notifications.created_at.desc())
.offset(skip)
.limit(limit),
)
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
async def count_unread(self, cdb: AsyncSession, user_id) -> Tuple[ErrorType, int]:
try:
err_type, rows = await DB_SESSION_MNG.execute(
cdb,
select(func.count()).select_from(notifications).where(
notifications.user_id == user_id,
notifications.read_at.is_(None),
notifications.deleted == False, # noqa: E712
),
)
if err_type != ErrorType.SUCCESS:
return err_type, 0
return ErrorType.SUCCESS, (int(rows[0] or 0) if rows else 0)
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED, 0
async def mark_read(self, cdb: AsyncSession, user_id, notification_id, ts) -> ErrorType:
try:
query = (
update(notifications)
.where(
notifications.notification_id == notification_id,
notifications.user_id == user_id, # 남의 알림 못 건드리게 수신자 스코프
notifications.read_at.is_(None),
)
.values(read_at=ts, updated_at=ts)
)
return await DB_SESSION_MNG.add(cdb, query)
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED
async def mark_all_read(self, cdb: AsyncSession, user_id, ts) -> ErrorType:
try:
query = (
update(notifications)
.where(
notifications.user_id == user_id,
notifications.read_at.is_(None),
notifications.deleted == False, # noqa: E712
)
.values(read_at=ts, updated_at=ts)
)
return await DB_SESSION_MNG.add(cdb, query)
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED