o2o-negosium-original/negodata/backend/crud/statistics_crud.py

348 lines
16 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

from abc import ABC, abstractmethod
from typing import Tuple
from sqlalchemy import select, func, and_, or_, case
from sqlalchemy.orm import aliased
from sqlalchemy.ext.asyncio import AsyncSession
from common.database.db_session_manager import DB_SESSION_MNG
from common.database.model.models import quotations, sessions, items, chats, users
from common.enums import ErrorType, QuotationStatus, CloseReason, SessionStatus, CardType, ChatSender
from common.logger import LOG
# 통계 유니버스 = 현재 마감사유 5코드로 마감된 견적. 레거시 REGEN_*(2·3·4) 은 제외해
# KPI(낙찰률·마감수)와 유형별/결과분해의 분모를 일치시킨다(프로덕션엔 레거시 없어 전체 마감과 동일).
CURRENT_CLOSE_REASONS = [
CloseReason.AWARDED.value,
CloseReason.OPEN_PRICE.value,
CloseReason.OPEN_EQUAL.value,
CloseReason.OPEN_NOSHOW.value,
CloseReason.OPEN_REJECT.value,
]
# 통계 집계 CRUD. 대시보드와 동일하게 회사 스코프(작성자 user_id→users.company_id)로 건다.
# owner(user_id) 가 주어지면 '내가 만든 견적'으로 더 좁힌다. quotations 엔 company_id 컬럼이 없어 서브쿼리로.
def _company_scope(company_id, owner) -> list:
conds = [
quotations.deleted == False, # noqa: E712
quotations.user_id.in_(select(users.user_id).where(users.company_id == company_id)),
]
if owner is not None:
conds.append(quotations.user_id == owner)
return conds
class IStatisticsCRUD(ABC):
@abstractmethod
async def winning_sessions(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, list]:
pass
@abstractmethod
async def outcome_counts(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, list]:
pass
@abstractmethod
async def type_counts(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, list]:
pass
@abstractmethod
async def participation_counts(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, list]:
pass
@abstractmethod
async def regen_avg_round(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, float]:
pass
@abstractmethod
async def markup_suppression(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, float]:
pass
@abstractmethod
async def markup_suppression_monthly(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, list]:
pass
@abstractmethod
async def card_usage(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, list]:
pass
@abstractmethod
async def card_effect_chats(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, list]:
pass
class StatisticsCRUD(IStatisticsCRUD):
async def winning_sessions(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, list]:
# 낙찰 마감 견적의 '낙찰 세션'(supplier_id=preferred_sp_id) 행 — 절감/추이/유형/카테고리/앵커도달률의 단일 원천.
# 파생: 저장 안 하고 조회 때 조인. category 는 items LEFT JOIN(자유텍스트·NULL 허용).
try:
stmt = (
select(
quotations.updated_at,
quotations.type,
items.category,
sessions.target_price,
sessions.bid_price,
sessions.anchoring_price,
)
.select_from(quotations)
.join(
sessions,
and_(
sessions.quotation_id == quotations.qt_id,
sessions.supplier_id == quotations.preferred_sp_id,
sessions.bid_price.isnot(None),
sessions.deleted == False, # noqa: E712
),
)
.join(items, items.item_id == sessions.item_id, isouter=True)
.where(
and_(
*_company_scope(company_id, owner),
quotations.status == QuotationStatus.CLOSED.value,
quotations.close_reason == CloseReason.AWARDED.value,
quotations.updated_at >= since,
)
)
)
err, rows = await DB_SESSION_MNG.execute(cdb, stmt)
return (err, list(rows) if err == ErrorType.SUCCESS else [])
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED, []
async def outcome_counts(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, list]:
# 마감 결과 분해: close_reason 별 건수. 낙찰률·마감건수도 여기서 파생.
try:
stmt = (
select(quotations.close_reason, func.count())
.where(
and_(
*_company_scope(company_id, owner),
quotations.status == QuotationStatus.CLOSED.value,
quotations.close_reason.in_(CURRENT_CLOSE_REASONS),
quotations.updated_at >= since,
)
)
.group_by(quotations.close_reason)
)
err, rows = await DB_SESSION_MNG.execute(cdb, stmt)
return (err, list(rows) if err == ErrorType.SUCCESS else [])
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED, []
async def type_counts(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, list]:
# 유형별(협상/경매) 마감 건수 + 낙찰 건수 → 유형별 낙찰률.
try:
awarded = func.sum(case((quotations.close_reason == CloseReason.AWARDED.value, 1), else_=0))
stmt = (
select(quotations.type, func.count(), awarded)
.where(
and_(
*_company_scope(company_id, owner),
quotations.status == QuotationStatus.CLOSED.value,
quotations.close_reason.in_(CURRENT_CLOSE_REASONS),
quotations.updated_at >= since,
)
)
.group_by(quotations.type)
)
err, rows = await DB_SESSION_MNG.execute(cdb, stmt)
return (err, list(rows) if err == ErrorType.SUCCESS else [])
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED, []
async def participation_counts(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, list]:
# 협력사 참여: 회사 견적(창 내 생성)의 세션을 status 별 집계(응찰/미응찰/거부).
try:
conds = [
sessions.deleted == False, # noqa: E712
quotations.deleted == False, # noqa: E712
quotations.created_at >= since,
quotations.user_id.in_(select(users.user_id).where(users.company_id == company_id)),
]
if owner is not None:
conds.append(quotations.user_id == owner)
stmt = (
select(sessions.status, func.count())
.select_from(sessions)
.join(quotations, quotations.qt_id == sessions.quotation_id)
.where(and_(*conds))
.group_by(sessions.status)
)
err, rows = await DB_SESSION_MNG.execute(cdb, stmt)
return (err, list(rows) if err == ErrorType.SUCCESS else [])
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED, []
async def regen_avg_round(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, float]:
# 평균 재견적 라운드. TODO: 체인키 없어 avg(round) 단순버전 — 체인당 최대 라운드 정의는 root_qt_id 도입 후.
try:
stmt = select(func.avg(quotations.round)).where(
and_(
*_company_scope(company_id, owner),
quotations.status == QuotationStatus.CLOSED.value,
quotations.close_reason.in_(CURRENT_CLOSE_REASONS),
quotations.updated_at >= since,
)
)
err, rows = await DB_SESSION_MNG.execute(cdb, stmt)
if err != ErrorType.SUCCESS:
return err, 0.0
# 단일컬럼 집계는 execute 가 스칼라 리스트를 반환한다(대시보드 _count 와 동일). rows[0] 이 곧 avg 값.
val = rows[0] if rows else None
return ErrorType.SUCCESS, float(val) if val is not None else 0.0
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED, 0.0
async def markup_suppression(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, float]:
# 인상억제율(재협상 전용): 같은 견적번호(qt_number)의 직전 라운드 투찰가 대비 이번 라운드 투찰가가
# 얼마나 안 올랐나 = avg((직전투찰 이번투찰) / 직전투찰). 양수=인하(억제 성공), 음수=인상 허용.
# 직전·이번 둘 다 유효 투찰(bid_price)이 있는 재협상 쌍만 대상(직전이 개찰/거부면 비교 불가 → 제외).
# 새 컬럼 없이 sessions.qt_number+qt_round+bid_price 로만 파생.
try:
prev = aliased(sessions)
stmt = (
select(func.avg((prev.bid_price - sessions.bid_price) * 1.0 / prev.bid_price))
.select_from(sessions)
.join(
prev,
and_(
prev.qt_number == sessions.qt_number,
prev.item_id == sessions.item_id,
prev.supplier_id == sessions.supplier_id,
prev.qt_round == sessions.qt_round - 1,
prev.bid_price.isnot(None),
prev.bid_price > 0,
prev.deleted == False, # noqa: E712
),
)
.join(quotations, quotations.qt_id == sessions.quotation_id)
.where(
and_(
*_company_scope(company_id, owner),
sessions.bid_price.isnot(None),
sessions.qt_round >= 2,
sessions.deleted == False, # noqa: E712
quotations.updated_at >= since,
)
)
)
err, rows = await DB_SESSION_MNG.execute(cdb, stmt)
if err != ErrorType.SUCCESS:
return err, 0.0
val = rows[0] if rows else None
return ErrorType.SUCCESS, float(val) if val is not None else 0.0
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED, 0.0
async def markup_suppression_monthly(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, list]:
# 월별 인상억제율: 이번 라운드 마감월(quotations.updated_at)별 avg((직전투찰 이번투찰)/직전투찰).
try:
prev = aliased(sessions)
month = func.to_char(quotations.updated_at, "YYYY-MM")
stmt = (
select(month.label("m"), func.avg((prev.bid_price - sessions.bid_price) * 1.0 / prev.bid_price))
.select_from(sessions)
.join(
prev,
and_(
prev.qt_number == sessions.qt_number,
prev.item_id == sessions.item_id,
prev.supplier_id == sessions.supplier_id,
prev.qt_round == sessions.qt_round - 1,
prev.bid_price.isnot(None),
prev.bid_price > 0,
prev.deleted == False, # noqa: E712
),
)
.join(quotations, quotations.qt_id == sessions.quotation_id)
.where(
and_(
*_company_scope(company_id, owner),
sessions.bid_price.isnot(None),
sessions.qt_round >= 2,
sessions.deleted == False, # noqa: E712
quotations.updated_at >= since,
)
)
.group_by(month)
.order_by(month)
)
err, rows = await DB_SESSION_MNG.execute(cdb, stmt)
return (err, list(rows) if err == ErrorType.SUCCESS else [])
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED, []
async def card_usage(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, list]:
# 카드 유형별 사용 빈도: card_used_yn=True 채팅 + 1% 인하 시스템 카드(meta.step='wild_card_1pct', card_id 미제공)를
# 와일드로 함께 집계. (1% 카드는 카탈로그 카드가 아니라 card_type 로그가 없어 step 으로 잡는다.)
try:
step_1pct = chats.meta["step"].astext == "wild_card_1pct"
ctype = case(
(chats.card_used_yn.is_(True), chats.card_type),
(step_1pct, CardType.WILD.value),
else_=None,
)
conds = [
chats.deleted == False, # noqa: E712
or_(chats.card_used_yn.is_(True), step_1pct),
quotations.deleted == False, # noqa: E712
quotations.created_at >= since,
quotations.user_id.in_(select(users.user_id).where(users.company_id == company_id)),
]
if owner is not None:
conds.append(quotations.user_id == owner)
stmt = (
select(ctype, func.count())
.select_from(chats)
.join(sessions, sessions.session_id == chats.session_id)
.join(quotations, quotations.qt_id == sessions.quotation_id)
.where(and_(*conds))
.group_by(ctype)
)
err, rows = await DB_SESSION_MNG.execute(cdb, stmt)
return (err, list(rows) if err == ErrorType.SUCCESS else [])
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED, []
async def card_effect_chats(self, cdb: AsyncSession, company_id, owner, since) -> Tuple[ErrorType, list]:
# 카드 사용 직후 제시가 하락 산출용 — 유저 제시가 chat + 카드 사용 chat(1% 인하 포함)을 세션·순번 순으로.
try:
step_1pct = chats.meta["step"].astext == "wild_card_1pct"
conds = [
chats.deleted == False, # noqa: E712
quotations.deleted == False, # noqa: E712
quotations.created_at >= since,
quotations.user_id.in_(select(users.user_id).where(users.company_id == company_id)),
or_(
and_(chats.sender == ChatSender.USER.value, chats.target_price > 0),
chats.card_used_yn.is_(True),
step_1pct,
),
]
if owner is not None:
conds.append(quotations.user_id == owner)
stmt = (
select(chats.session_id, chats.seq, chats.sender, chats.target_price, chats.card_used_yn,
chats.card_type, step_1pct, sessions.bid_price, sessions.status)
.select_from(chats)
.join(sessions, sessions.session_id == chats.session_id)
.join(quotations, quotations.qt_id == sessions.quotation_id)
.where(and_(*conds))
.order_by(chats.session_id, chats.seq)
)
err, rows = await DB_SESSION_MNG.execute(cdb, stmt)
return (err, list(rows) if err == ErrorType.SUCCESS else [])
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED, []