협력사가 단종·품절을 대화 도중 알아채도 봇의 결렬 선언을 기다려야 했고, 거부로 끝난 협상은 가격을 남겨도 낙찰 후보에서 빠져 계약으로 이어지지 않았다. negosium - 채팅 액션바에 협상 거부 진입점 — 대화가 끝나지 않고 입력을 기다리는 동안만 노출, 주 CTA 와 붙지 않게 넓은 화면은 우측 끝 고정·좁은 화면은 wrap - 거부 팝업은 목록 거부 팝업과 같은 어휘·규격, 대화 중이라 공급 희망 가격·의견을 더 받는다 - /reject 에 reject_price·opinion 추가 — sessions.reject_price 저장, 의견은 custom 병합 - 화면 문구 '거절' → '거부' 통일 (버튼·배지·탭·토스트·안내 팝업) negodata - 개찰 견적 직접 낙찰 후보 = 가격을 써낸 세션 — 투찰한 협상완료 + 공급 희망가를 남긴 협상거부 - 계약가 파생 _award_price/awardPrice — coalesce(투찰가, 거부 시 공급 희망가) - 통계 낙찰 세션 조인도 같은 기준 — 안 고치면 거부가로 낙찰한 건이 절감 집계에서 빠진다 - 세션 상태 탭 라벨 '거절사유/거절가격/거절배송방식' → '거부…' 자동 마감 판정(close_and_decide)은 그대로 — 자동 낙찰은 투찰가만 본다.
876 lines
41 KiB
Python
876 lines
41 KiB
Python
from abc import ABC, abstractmethod
|
|
from datetime import datetime
|
|
from typing import Optional, Tuple
|
|
|
|
from sqlalchemy import select, func, and_, or_, update
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
from common.database.db_session_manager import DB_SESSION_MNG
|
|
from common.database.model.models import (
|
|
quotations, sessions, chats, nego_cards, wild_cards, items, suppliers, quotation_settings,
|
|
version_nego_cards, version_wild_cards, users, supplier_items, companies,
|
|
)
|
|
from common.enums import CloseReason, ErrorType, QuotationStatus, SessionStatus
|
|
from common.logger import LOG
|
|
from common.utils.gtime import GTime
|
|
|
|
|
|
# 견적 CRUD.
|
|
class IQuotationCRUD(ABC):
|
|
@abstractmethod
|
|
async def search(
|
|
self, cdb: AsyncSession, company_id, owner, search, status, type_, start_from, start_to, skip, limit
|
|
) -> Tuple[ErrorType, list, int]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def get_by_id(self, cdb: AsyncSession, qt_id, company_id=None) -> Tuple[ErrorType, quotations]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def add_quotation(self, cdb: AsyncSession, quotation: quotations) -> ErrorType:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def add_sessions(self, cdb: AsyncSession, session_list: list) -> ErrorType:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def get_item_prices(self, cdb: AsyncSession, item_ids) -> Tuple[ErrorType, dict]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def get_item_companies(self, cdb: AsyncSession, item_ids) -> Tuple[ErrorType, dict]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def get_supply_types(self, cdb: AsyncSession, item_ids, supplier_ids) -> Tuple[ErrorType, dict]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def get_setting_rates(self, cdb: AsyncSession, qt_setting_id) -> Tuple[ErrorType, dict]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def get_company_settings(self, cdb: AsyncSession, user_id) -> dict:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def get_company_brand(self, cdb: AsyncSession, user_id) -> Tuple[str, dict]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def add_rows(self, cdb: AsyncSession, obj_list: list) -> ErrorType:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def classify_card_ids(self, cdb: AsyncSession, card_ids) -> Tuple[ErrorType, dict]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def get_version_cards(self, cdb: AsyncSession, version_id) -> Tuple[ErrorType, list]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def update_quotation(self, cdb: AsyncSession, qt_id, data: dict) -> ErrorType:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def update_sessions_status(self, cdb: AsyncSession, qt_id, from_statuses: list[int], to_status: int) -> ErrorType:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def soft_delete(self, cdb: AsyncSession, qt_id) -> ErrorType:
|
|
"""견적과 연결된 하위 데이터(세션·대화)를 함께 소프트 삭제."""
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def list_sessions(self, cdb: AsyncSession, qt_id) -> Tuple[ErrorType, list]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def list_chats(self, cdb: AsyncSession, session_id) -> Tuple[ErrorType, list]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def list_used_cards(self, cdb: AsyncSession, qt_id) -> Tuple[ErrorType, list]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def list_sessions_with_supplier(self, cdb: AsyncSession, qt_id) -> Tuple[ErrorType, list]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def get_session_with_supplier(self, cdb: AsyncSession, session_id) -> Tuple[ErrorType, Optional[tuple]]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def mark_sessions_emailed(self, cdb: AsyncSession, session_ids, ts) -> ErrorType:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def session_counts(self, cdb: AsyncSession, qt_ids) -> Tuple[ErrorType, dict]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def item_map(self, cdb: AsyncSession, qt_ids) -> Tuple[ErrorType, dict]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def user_name_map(self, cdb: AsyncSession, user_ids) -> Tuple[ErrorType, dict]:
|
|
pass
|
|
|
|
# ----- 스케줄러(크론) 전용 -----
|
|
@abstractmethod
|
|
async def list_due_for_close(self, cdb: AsyncSession, now) -> Tuple[ErrorType, list]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def list_all_sessions_ended(self, cdb: AsyncSession) -> Tuple[ErrorType, list]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def list_chain_close_reasons(self, cdb: AsyncSession, number, current_round) -> Tuple[ErrorType, list]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def chain_max_round(self, cdb: AsyncSession, number) -> Tuple[ErrorType, int]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def list_done_sessions(self, cdb: AsyncSession, qt_id) -> Tuple[ErrorType, list]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def list_sessions_status(self, cdb: AsyncSession, qt_id) -> Tuple[ErrorType, list]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def bulk_update_quotation_status(self, cdb: AsyncSession, qt_ids, status: int) -> ErrorType:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def bulk_update_sessions_status(self, cdb: AsyncSession, qt_ids, from_statuses: list[int], to_status: int) -> ErrorType:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def claim_for_close(self, cdb: AsyncSession, qt_id) -> Tuple[ErrorType, int]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
async def claim_for_award(self, cdb: AsyncSession, qt_id, supplier_id, supplier_name) -> Tuple[ErrorType, int]:
|
|
pass
|
|
|
|
|
|
class QuotationCRUD(IQuotationCRUD):
|
|
async def search(
|
|
self,
|
|
cdb: AsyncSession,
|
|
company_id,
|
|
owner,
|
|
search: Optional[str],
|
|
status: Optional[str],
|
|
type_: Optional[str],
|
|
start_from: Optional[datetime],
|
|
start_to: Optional[datetime],
|
|
skip: int,
|
|
limit: int,
|
|
) -> Tuple[ErrorType, list, int]:
|
|
try:
|
|
# 회사 스코프(멀티테넌트): quotations 엔 company_id 가 없어 작성자(user_id)→users.company_id 로 건다.
|
|
conditions = [
|
|
quotations.deleted == False, # noqa: E712
|
|
quotations.user_id.in_(select(users.user_id).where(users.company_id == company_id)),
|
|
]
|
|
if owner:
|
|
conditions.append(quotations.user_id == owner) # '내 견적만' — 작성자(user_id)=로그인 유저
|
|
if search:
|
|
conditions.append(or_(quotations.name.ilike(f"%{search}%"), quotations.number.ilike(f"%{search}%")))
|
|
if status:
|
|
conditions.append(quotations.status == int(status)) # status/type 는 SMALLINT 코드 — 문자열 쿼리값을 정수로
|
|
if type_:
|
|
conditions.append(quotations.type == int(type_))
|
|
if start_from:
|
|
conditions.append(quotations.start_time >= start_from)
|
|
if start_to:
|
|
conditions.append(quotations.start_time <= start_to)
|
|
where = and_(*conditions)
|
|
|
|
cnt_err, cnt_rows = await DB_SESSION_MNG.execute(cdb, select(func.count()).select_from(quotations).where(where))
|
|
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(quotations).where(where).order_by(quotations.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 session_counts(self, cdb: AsyncSession, qt_ids) -> Tuple[ErrorType, dict]:
|
|
"""견적 id 목록에 대해 참여 협력사 수(distinct supplier)를 한 번에 센다. {qt_id: count}."""
|
|
try:
|
|
if not qt_ids:
|
|
return ErrorType.SUCCESS, {}
|
|
query = (
|
|
select(sessions.quotation_id, func.count(func.distinct(sessions.supplier_id)))
|
|
.where(sessions.quotation_id.in_(qt_ids), sessions.deleted == False) # noqa: E712
|
|
.group_by(sessions.quotation_id)
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, {}
|
|
return ErrorType.SUCCESS, {r[0]: int(r[1] or 0) for r in rows}
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, {}
|
|
|
|
async def item_map(self, cdb: AsyncSession, qt_ids) -> Tuple[ErrorType, dict]:
|
|
"""견적 id 목록에 대해 대표 상품(세션의 첫 item) {qt_id: (item_id, item_name)} 을 한 번에 가져온다."""
|
|
try:
|
|
if not qt_ids:
|
|
return ErrorType.SUCCESS, {}
|
|
query = (
|
|
select(sessions.quotation_id, sessions.item_id, items.name)
|
|
.join(items, items.item_id == sessions.item_id)
|
|
.where(
|
|
sessions.quotation_id.in_(qt_ids),
|
|
sessions.deleted == False, # noqa: E712
|
|
items.deleted == False, # noqa: E712
|
|
)
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, {}
|
|
result = {}
|
|
for qid, iid, iname in rows:
|
|
if qid not in result: # 견적당 대표 1개(첫 세션 상품)
|
|
result[qid] = (iid, iname)
|
|
return ErrorType.SUCCESS, result
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, {}
|
|
|
|
async def user_name_map(self, cdb: AsyncSession, user_ids) -> Tuple[ErrorType, dict]:
|
|
"""user_id 목록 → {user_id: name}. 견적 목록 '작성자(등록자)' 표기용(company.users 조인)."""
|
|
try:
|
|
if not user_ids:
|
|
return ErrorType.SUCCESS, {}
|
|
query = select(users.user_id, users.name).where(users.user_id.in_(user_ids))
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, {}
|
|
return ErrorType.SUCCESS, {uid: name for uid, name in rows}
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, {}
|
|
|
|
async def get_by_id(self, cdb: AsyncSession, qt_id, company_id=None) -> Tuple[ErrorType, quotations]:
|
|
try:
|
|
# company_id 가 주어지면 회사 스코프(작성자 회사)로 좁힌다 — 남의 회사 견적은 '없음'으로 떨어진다.
|
|
conds = [quotations.qt_id == qt_id, quotations.deleted == False] # noqa: E712
|
|
if company_id is not None:
|
|
conds.append(quotations.user_id.in_(select(users.user_id).where(users.company_id == company_id)))
|
|
query = select(quotations).where(*conds).limit(1)
|
|
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 add_quotation(self, cdb: AsyncSession, quotation: quotations) -> ErrorType:
|
|
try:
|
|
return await DB_SESSION_MNG.insert(cdb, quotation)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED
|
|
|
|
async def add_sessions(self, cdb: AsyncSession, session_list: list) -> ErrorType:
|
|
"""견적 생성 시 만들어진 협상 세션들을 한 번에 insert. 빈 목록이면 그냥 통과."""
|
|
try:
|
|
if not session_list:
|
|
return ErrorType.SUCCESS
|
|
return await DB_SESSION_MNG.insert(cdb, session_list)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED
|
|
|
|
async def add_rows(self, cdb: AsyncSession, obj_list: list) -> ErrorType:
|
|
"""임의 ORM 행 묶음 insert(버전/버전-카드 매핑 등). 빈 목록이면 통과."""
|
|
try:
|
|
if not obj_list:
|
|
return ErrorType.SUCCESS
|
|
return await DB_SESSION_MNG.insert(cdb, obj_list)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED
|
|
|
|
async def classify_card_ids(self, cdb: AsyncSession, card_ids) -> Tuple[ErrorType, dict]:
|
|
"""선택 카드 id 를 협상(1)/와일드(2)로 분류. {card_id: card_type}."""
|
|
try:
|
|
if not card_ids:
|
|
return ErrorType.SUCCESS, {}
|
|
out = {}
|
|
n_err, n_rows = await DB_SESSION_MNG.execute(
|
|
cdb, select(nego_cards.nego_card_id).where(nego_cards.nego_card_id.in_(card_ids), nego_cards.deleted == False) # noqa: E712
|
|
)
|
|
if n_err != ErrorType.SUCCESS:
|
|
return n_err, {}
|
|
for r in n_rows:
|
|
out[r] = 1
|
|
w_err, w_rows = await DB_SESSION_MNG.execute(
|
|
cdb, select(wild_cards.wild_card_id).where(wild_cards.wild_card_id.in_(card_ids), wild_cards.deleted == False) # noqa: E712
|
|
)
|
|
if w_err != ErrorType.SUCCESS:
|
|
return w_err, {}
|
|
for r in w_rows:
|
|
out[r] = 2
|
|
return ErrorType.SUCCESS, out
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, {}
|
|
|
|
async def get_version_cards(self, cdb: AsyncSession, version_id) -> Tuple[ErrorType, list]:
|
|
"""견적 버전에 묶인 카드. version_nego_cards/version_wild_cards 조인.
|
|
반환: [(card_type, card_pk, number, name, script, edit_script, condition, memo), ...]."""
|
|
try:
|
|
out = []
|
|
n_q = (
|
|
select(
|
|
nego_cards.nego_card_id, nego_cards.number, nego_cards.name,
|
|
nego_cards.script, nego_cards.edit_script,
|
|
)
|
|
.join(version_nego_cards, version_nego_cards.nego_card_id == nego_cards.nego_card_id)
|
|
.where(version_nego_cards.version_id == version_id, version_nego_cards.deleted == False, nego_cards.deleted == False) # noqa: E712
|
|
)
|
|
n_err, n_rows = await DB_SESSION_MNG.execute(cdb, n_q)
|
|
if n_err != ErrorType.SUCCESS:
|
|
return n_err, []
|
|
for pk, number, name, script, edit in n_rows:
|
|
out.append((1, pk, number, name, script, edit, None, None))
|
|
w_q = (
|
|
select(
|
|
wild_cards.wild_card_id, wild_cards.number, wild_cards.name,
|
|
wild_cards.script, wild_cards.edit_script, wild_cards.condition, wild_cards.memo,
|
|
)
|
|
.join(version_wild_cards, version_wild_cards.wild_card_id == wild_cards.wild_card_id)
|
|
.where(version_wild_cards.version_id == version_id, version_wild_cards.deleted == False, wild_cards.deleted == False) # noqa: E712
|
|
)
|
|
w_err, w_rows = await DB_SESSION_MNG.execute(cdb, w_q)
|
|
if w_err != ErrorType.SUCCESS:
|
|
return w_err, []
|
|
for pk, number, name, script, edit, condition, memo in w_rows:
|
|
out.append((2, pk, number, name, script, edit, condition, memo))
|
|
return ErrorType.SUCCESS, out
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, []
|
|
|
|
async def get_item_prices(self, cdb: AsyncSession, item_ids) -> Tuple[ErrorType, dict]:
|
|
"""item_id -> (internet_lowest_price, purchase_price, selling_price)(원, NULL 가능) 매핑. 세션 목표가 계산 입력."""
|
|
try:
|
|
if not item_ids:
|
|
return ErrorType.SUCCESS, {}
|
|
query = select(
|
|
items.item_id, items.internet_lowest_price, items.purchase_price, items.selling_price
|
|
).where(
|
|
items.item_id.in_(item_ids), items.deleted == False # noqa: E712
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, {}
|
|
return ErrorType.SUCCESS, {r[0]: (r[1], r[2], r[3]) for r in rows}
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, {}
|
|
|
|
async def get_company_settings(self, cdb: AsyncSession, user_id):
|
|
# 이 유저가 속한 회사의 settings(JSONB) 전체를 돌려준다. 회사가 없거나 조회에 실패하면 빈 dict.
|
|
# 어떤 키를 어떻게 해석할지(목표가 정책·메일 브랜딩 등)는 호출하는 쪽 몫이다.
|
|
try:
|
|
query = (
|
|
select(companies.settings)
|
|
.select_from(users)
|
|
.join(companies, companies.company_id == users.company_id)
|
|
.where(users.user_id == user_id, companies.deleted == False) # noqa: E712
|
|
.limit(1)
|
|
)
|
|
err, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err != ErrorType.SUCCESS or not rows:
|
|
return {}
|
|
return rows[0] or {}
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return {}
|
|
|
|
async def get_company_brand(self, cdb: AsyncSession, user_id):
|
|
# 이 유저가 속한 회사의 (이름, settings). 초청 메일 헤더 기본값이 회사명이라 이름까지 같이 읽는다.
|
|
try:
|
|
query = (
|
|
select(companies.name, companies.settings)
|
|
.select_from(users)
|
|
.join(companies, companies.company_id == users.company_id)
|
|
.where(users.user_id == user_id, companies.deleted == False) # noqa: E712
|
|
.limit(1)
|
|
)
|
|
err, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err != ErrorType.SUCCESS or not rows:
|
|
return "", {}
|
|
return rows[0][0] or "", rows[0][1] or {}
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return "", {}
|
|
|
|
async def get_item_companies(self, cdb: AsyncSession, item_ids) -> Tuple[ErrorType, dict]:
|
|
"""item_id -> company_id(소유 회사) 매핑. 앵커링 칸(회사×유형×가격구간) 해석 입력."""
|
|
try:
|
|
if not item_ids:
|
|
return ErrorType.SUCCESS, {}
|
|
query = select(items.item_id, items.company_id).where(
|
|
items.item_id.in_(item_ids), items.deleted == False # noqa: E712
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, {}
|
|
return ErrorType.SUCCESS, {r[0]: r[1] for r in rows}
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, {}
|
|
|
|
async def get_supply_types(self, cdb: AsyncSession, item_ids, supplier_ids) -> Tuple[ErrorType, dict]:
|
|
"""(item_id, supplier_id) -> supplier_items.supply_type 매핑.
|
|
|
|
매핑이 없거나 supply_type 이 0/NULL 이면 호출측이 정적 앵커링 값으로 폴백한다.
|
|
"""
|
|
try:
|
|
if not item_ids or not supplier_ids:
|
|
return ErrorType.SUCCESS, {}
|
|
query = select(
|
|
supplier_items.item_id,
|
|
supplier_items.supplier_id,
|
|
supplier_items.supply_type,
|
|
).where(
|
|
supplier_items.item_id.in_(item_ids),
|
|
supplier_items.supplier_id.in_(supplier_ids),
|
|
supplier_items.deleted == False, # noqa: E712
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, {}
|
|
return ErrorType.SUCCESS, {(r[0], r[1]): r[2] for r in rows}
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, {}
|
|
|
|
async def get_setting_rates(self, cdb: AsyncSession, qt_setting_id) -> Tuple[ErrorType, dict]:
|
|
"""견적 세팅의 목표 마진율·협상 완료 상한율: {margin, done_ceiling_rate}.
|
|
margin·수수료는 목표가 산정 입력, done_ceiling_rate(‰)는 완료 상한 = 목표가×(1+값/1000).
|
|
(낙찰 정책은 견적 단위 이관, 앵커링은 칸 rate v1.2 → 세팅 컬럼 제거됨.)"""
|
|
try:
|
|
query = select(
|
|
quotation_settings.target_margin_rate,
|
|
quotation_settings.done_ceiling_rate,
|
|
).where(quotation_settings.qt_setting_id == qt_setting_id).limit(1)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, {}
|
|
if not rows:
|
|
return ErrorType.SUCCESS, {}
|
|
row = rows[0]
|
|
return ErrorType.SUCCESS, {
|
|
"margin": float(row.target_margin_rate) if row.target_margin_rate is not None else None,
|
|
"done_ceiling_rate": int(row.done_ceiling_rate) if row.done_ceiling_rate is not None else None,
|
|
}
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, {}
|
|
|
|
async def update_quotation(self, cdb: AsyncSession, qt_id, data: dict) -> ErrorType:
|
|
try:
|
|
if not data:
|
|
return ErrorType.SUCCESS
|
|
query = update(quotations).where(quotations.qt_id == qt_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 claim_for_close(self, cdb: AsyncSession, qt_id) -> Tuple[ErrorType, int]:
|
|
"""[동시 마감 가드] 아직 안 닫힌(status != CLOSED, not deleted) 견적만 CLOSED 로 선점 전이.
|
|
반환: (ErrorType, 적용행수). 동시 호출 시 Postgres 행 잠금으로 직렬화되어
|
|
실제로 CLOSED 로 바꾼 호출자만 1, 이미 닫혀 있던(진 호출자/재처리) 경우는 0 을 받는다.
|
|
close_and_decide 가 이 결과로 '마감 판정 권한'을 단 한 번만 갖도록 한다."""
|
|
try:
|
|
query = (
|
|
update(quotations)
|
|
.where(
|
|
quotations.qt_id == qt_id,
|
|
quotations.status != QuotationStatus.CLOSED.value,
|
|
quotations.deleted == False, # noqa: E712
|
|
)
|
|
.values(status=QuotationStatus.CLOSED.value, updated_at=GTime.UTC())
|
|
)
|
|
return await DB_SESSION_MNG.add_with_rowcount(cdb, query)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, 0
|
|
|
|
async def claim_for_award(self, cdb: AsyncSession, qt_id, supplier_id, supplier_name) -> Tuple[ErrorType, int]:
|
|
"""[동시 직접낙찰 가드] 개찰(마감·낙찰자 미정, close_reason ∈ OPEN_*)인 견적만 낙찰(AWARDED)로 선점 전이.
|
|
선택 협력사를 낙찰자(preferred_sp_*)로 박고 동가 플래그는 내린다. status 는 이미 CLOSED 라 유지.
|
|
반환: (ErrorType, 적용행수). 이미 낙찰됐거나(재클릭) 개찰이 아니면 0 → 서비스가 한 번만 통과시킨다."""
|
|
try:
|
|
query = (
|
|
update(quotations)
|
|
.where(
|
|
quotations.qt_id == qt_id,
|
|
quotations.status == QuotationStatus.CLOSED.value,
|
|
quotations.close_reason.in_(
|
|
[CloseReason.OPEN_PRICE.value, CloseReason.OPEN_EQUAL.value,
|
|
CloseReason.OPEN_NOSHOW.value, CloseReason.OPEN_REJECT.value]
|
|
),
|
|
quotations.deleted == False, # noqa: E712
|
|
)
|
|
.values(
|
|
close_reason=CloseReason.AWARDED.value,
|
|
preferred_sp_yn=True,
|
|
preferred_sp_id=supplier_id,
|
|
preferred_sp_name=supplier_name,
|
|
equal_bid_yn=False,
|
|
updated_at=GTime.UTC(),
|
|
)
|
|
)
|
|
return await DB_SESSION_MNG.add_with_rowcount(cdb, query)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, 0
|
|
|
|
async def update_sessions_status(self, cdb: AsyncSession, qt_id, from_statuses: list[int], to_status: int) -> ErrorType:
|
|
# 견적에 딸린 세션 중 from_statuses 에 속한 것만 to_status 로 일괄 전이(삭제 제외). 다른 상태는 건드리지 않는다.
|
|
try:
|
|
query = (
|
|
update(sessions)
|
|
.where(
|
|
sessions.quotation_id == qt_id,
|
|
sessions.status.in_(from_statuses),
|
|
sessions.deleted == False, # noqa: E712
|
|
)
|
|
.values(status=to_status, updated_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 soft_delete(self, cdb: AsyncSession, qt_id) -> ErrorType:
|
|
"""견적 + 연결된 하위 데이터(협상 세션·대화)를 한 트랜잭션으로 소프트 삭제한다.
|
|
대화(chats)→세션(sessions)→견적(quotations) 순서로 deleted=True. 한 스텝이라도 실패하면 execute_lambda_run 이 롤백한다."""
|
|
try:
|
|
now = GTime.UTC()
|
|
# 이 견적에 매달린 세션들 — chats 삭제 범위를 잡는 서브쿼리(quotation_id 기준, deleted 무관).
|
|
session_ids = select(sessions.session_id).where(sessions.quotation_id == qt_id)
|
|
for query in (
|
|
update(chats).where(chats.session_id.in_(session_ids)).values(deleted=True, updated_at=now),
|
|
update(sessions).where(sessions.quotation_id == qt_id).values(deleted=True, updated_at=now),
|
|
update(quotations).where(quotations.qt_id == qt_id).values(deleted=True, updated_at=now),
|
|
):
|
|
err_type = await DB_SESSION_MNG.add(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type
|
|
return ErrorType.SUCCESS
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED
|
|
|
|
# ----- 스케줄러(크론) 전용 -----
|
|
async def list_due_for_close(self, cdb: AsyncSession, now) -> Tuple[ErrorType, list]:
|
|
"""[잡①] 마감시각이 지났는데 아직 안 닫힌 견적 qt_id 목록.
|
|
조건: end_time < now AND status != 견적마감 AND not deleted."""
|
|
try:
|
|
query = select(quotations.qt_id).where(
|
|
quotations.end_time < now,
|
|
quotations.status != QuotationStatus.CLOSED.value,
|
|
quotations.deleted == False, # noqa: E712
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, []
|
|
return ErrorType.SUCCESS, list(rows)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, []
|
|
|
|
async def list_all_sessions_ended(self, cdb: AsyncSession) -> Tuple[ErrorType, list]:
|
|
"""[잡②] 아직 안 닫혔는데 모든 세션이 종결된 견적 qt_id 목록(견적 타입 무관 — KTC 와 동일).
|
|
종결 = 협상완료/협상거부/미참여 → 즉 진행중(IN_PROGRESS)·미시작(CREATED) 세션이 하나도 없음.
|
|
세션이 1개 이상 있어야 하며, 마감일과 무관하게 협상이 다 끝났으면 즉시 마감 대상."""
|
|
try:
|
|
# 아직 안 끝난(진행중·미시작) 세션이 하나라도 있으면 제외
|
|
pending_exists = (
|
|
select(sessions.session_id)
|
|
.where(
|
|
sessions.quotation_id == quotations.qt_id,
|
|
sessions.status.in_([SessionStatus.CREATED.value, SessionStatus.IN_PROGRESS.value]),
|
|
sessions.deleted == False, # noqa: E712
|
|
)
|
|
.exists()
|
|
)
|
|
# 세션이 최소 1개는 있어야(세션 없는 견적은 대상 아님)
|
|
any_session = (
|
|
select(sessions.session_id)
|
|
.where(sessions.quotation_id == quotations.qt_id, sessions.deleted == False) # noqa: E712
|
|
.exists()
|
|
)
|
|
query = select(quotations.qt_id).where(
|
|
quotations.status != QuotationStatus.CLOSED.value,
|
|
quotations.deleted == False, # noqa: E712
|
|
any_session,
|
|
~pending_exists,
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, []
|
|
return ErrorType.SUCCESS, list(rows)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, []
|
|
|
|
async def list_done_sessions(self, cdb: AsyncSession, qt_id) -> Tuple[ErrorType, list]:
|
|
"""[잡②] 견적의 협상완료(DONE) 세션 → (supplier_id, bid_price, supplier_name) 목록. 낙찰자 판정 입력."""
|
|
try:
|
|
query = (
|
|
select(sessions.supplier_id, sessions.bid_price, suppliers.name)
|
|
.join(suppliers, suppliers.supplier_id == sessions.supplier_id)
|
|
.where(
|
|
sessions.quotation_id == qt_id,
|
|
sessions.status == SessionStatus.DONE.value,
|
|
sessions.deleted == False, # noqa: E712
|
|
suppliers.deleted == False, # noqa: E712
|
|
)
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, []
|
|
return ErrorType.SUCCESS, list(rows)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, []
|
|
|
|
async def list_chain_close_reasons(self, cdb: AsyncSession, number, current_round) -> Tuple[ErrorType, list]:
|
|
"""[재생성 한도] 같은 견적번호(체인)의 이전 라운드(round < current_round)들의 close_reason 코드 목록. 삭제 제외.
|
|
_chain_regen_counts 가 REGEN_* 값(목표초과재협상/동가재입찰/미참여재소집)만 사유별로 센다.
|
|
낙찰(AWARDED)·유찰(FAIL_*)·미마감(NULL)은 재생성 한도에 안 셈."""
|
|
try:
|
|
query = select(quotations.close_reason).where(
|
|
quotations.number == number,
|
|
quotations.round < current_round,
|
|
quotations.deleted == False, # noqa: E712
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, []
|
|
# 단일 컬럼 select → 각 행이 스칼라(close_reason 코드 or None)
|
|
return ErrorType.SUCCESS, [r for r in rows]
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, []
|
|
|
|
async def chain_max_round(self, cdb: AsyncSession, number) -> Tuple[ErrorType, int]:
|
|
"""같은 견적번호(체인)의 최대 round. 다음 라운드 = 이 값 + 1.
|
|
uq_quotations_number(number, round) 충돌 방지 — 원본 round+1이 아니라 체인 최신 기준으로 매긴다."""
|
|
try:
|
|
query = select(func.max(quotations.round)).where(
|
|
quotations.number == number,
|
|
quotations.deleted == False, # noqa: E712
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, 0
|
|
# 단일 컬럼 select → rows[0] 이 스칼라(max) 값. 행 없거나 전부 NULL이면 None.
|
|
mx = rows[0] if rows else None
|
|
return ErrorType.SUCCESS, int(mx or 0)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, 0
|
|
|
|
async def list_sessions_status(self, cdb: AsyncSession, qt_id) -> Tuple[ErrorType, list]:
|
|
"""[마감 판정] 견적의 모든 세션 → (status, supplier_id, bid_price, name, target_price, anchoring_price, reject_price). 삭제 제외.
|
|
공급사가 지워졌어도 세션 집계엔 포함되도록 outerjoin(이때 name 은 None).
|
|
target/anchoring 은 마감 가격게이트 입력(견적당 상품 1개라 세션 공통값).
|
|
reject_price 는 거부 협력사의 공급 희망 가격 — 자동 마감 판정엔 안 쓰고 수동 직접 낙찰 후보에서만 본다."""
|
|
try:
|
|
query = (
|
|
select(sessions.status, sessions.supplier_id, sessions.bid_price, suppliers.name,
|
|
sessions.target_price, sessions.anchoring_price, sessions.reject_price)
|
|
.outerjoin(suppliers, suppliers.supplier_id == sessions.supplier_id)
|
|
.where(sessions.quotation_id == qt_id, sessions.deleted == False) # noqa: E712
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, []
|
|
return ErrorType.SUCCESS, list(rows)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, []
|
|
|
|
async def bulk_update_quotation_status(self, cdb: AsyncSession, qt_ids, status: int) -> ErrorType:
|
|
"""[잡①] 여러 견적의 status 를 한 번에 전이."""
|
|
try:
|
|
if not qt_ids:
|
|
return ErrorType.SUCCESS
|
|
query = (
|
|
update(quotations)
|
|
.where(quotations.qt_id.in_(qt_ids))
|
|
.values(status=status, updated_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 bulk_update_sessions_status(self, cdb: AsyncSession, qt_ids, from_statuses: list[int], to_status: int) -> ErrorType:
|
|
"""[잡①] 여러 견적에 딸린 세션 중 from_statuses 에 속한 것만 to_status 로 일괄 전이(삭제 제외)."""
|
|
try:
|
|
if not qt_ids:
|
|
return ErrorType.SUCCESS
|
|
query = (
|
|
update(sessions)
|
|
.where(
|
|
sessions.quotation_id.in_(qt_ids),
|
|
sessions.status.in_(from_statuses),
|
|
sessions.deleted == False, # noqa: E712
|
|
)
|
|
.values(status=to_status, updated_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 list_sessions(self, cdb: AsyncSession, qt_id) -> Tuple[ErrorType, list]:
|
|
try:
|
|
query = (
|
|
select(sessions)
|
|
.where(sessions.quotation_id == qt_id, sessions.deleted == False) # noqa: E712
|
|
.order_by(sessions.created_at.asc())
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, []
|
|
return ErrorType.SUCCESS, list(rows)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, []
|
|
|
|
async def list_sessions_with_supplier(self, cdb: AsyncSession, qt_id) -> Tuple[ErrorType, list]:
|
|
"""견적의 세션 + 공급사(담당자 이메일/이름) 조인. 초청 메일 발송 대상 조회용.
|
|
반환: [(session, supplier_name, manager_email), ...] (created_at asc). 행은 인덱스로 언팩."""
|
|
try:
|
|
query = (
|
|
select(sessions, suppliers.name, suppliers.manager_email)
|
|
.outerjoin(suppliers, suppliers.supplier_id == sessions.supplier_id)
|
|
.where(sessions.quotation_id == qt_id, sessions.deleted == False) # noqa: E712
|
|
.order_by(sessions.created_at.asc())
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, []
|
|
return ErrorType.SUCCESS, list(rows)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, []
|
|
|
|
async def get_session_with_supplier(self, cdb: AsyncSession, session_id) -> Tuple[ErrorType, Optional[tuple]]:
|
|
"""단일 세션 + 공급사(이름/이메일). 행별 재발송용. 반환: (session, name, email) | None."""
|
|
try:
|
|
query = (
|
|
select(sessions, suppliers.name, suppliers.manager_email)
|
|
.outerjoin(suppliers, suppliers.supplier_id == sessions.supplier_id)
|
|
.where(sessions.session_id == session_id, sessions.deleted == False) # noqa: E712
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, None
|
|
rows = list(rows)
|
|
return ErrorType.SUCCESS, (rows[0] if rows else None)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, None
|
|
|
|
async def mark_sessions_emailed(self, cdb: AsyncSession, session_ids, ts) -> ErrorType:
|
|
"""발송 성공 세션들의 email_sent_at 을 ts 로 기록(write)."""
|
|
try:
|
|
if not session_ids:
|
|
return ErrorType.SUCCESS
|
|
query = (
|
|
update(sessions)
|
|
.where(sessions.session_id.in_(session_ids))
|
|
.values(email_sent_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 list_chats(self, cdb: AsyncSession, session_id) -> Tuple[ErrorType, list]:
|
|
try:
|
|
query = (
|
|
select(chats)
|
|
.where(chats.session_id == session_id, chats.deleted == False) # noqa: E712
|
|
.order_by(chats.seq.asc())
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, []
|
|
return ErrorType.SUCCESS, list(rows)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, []
|
|
|
|
async def list_used_cards(self, cdb: AsyncSession, qt_id) -> Tuple[ErrorType, list]:
|
|
"""견적의 세션들에서 실제 사용된 카드(chats.card_used_yn)를 카드 카탈로그와 조인.
|
|
반환: [(chat_row, card_id, number, name, script, edit_script, condition, memo), ...].
|
|
card_type 1=nego_cards / 2=wild_cards 양쪽을 LEFT JOIN 해서 어느 쪽이든 잡는다.
|
|
condition/memo 는 wild_cards 에만 있는 컬럼이라 nego 카드면 NULL 로 나온다.
|
|
"""
|
|
try:
|
|
query = (
|
|
select(
|
|
chats,
|
|
func.coalesce(nego_cards.nego_card_id, wild_cards.wild_card_id).label("card_pk"),
|
|
func.coalesce(nego_cards.number, wild_cards.number).label("card_number"),
|
|
func.coalesce(nego_cards.name, wild_cards.name).label("card_name"),
|
|
func.coalesce(nego_cards.script, wild_cards.script).label("card_script"),
|
|
func.coalesce(nego_cards.edit_script, wild_cards.edit_script).label("card_edit_script"),
|
|
wild_cards.condition.label("card_condition"),
|
|
wild_cards.memo.label("card_memo"),
|
|
)
|
|
.join(sessions, sessions.session_id == chats.session_id)
|
|
.outerjoin(nego_cards, and_(nego_cards.nego_card_id == chats.card_id, chats.card_type == 1))
|
|
.outerjoin(wild_cards, and_(wild_cards.wild_card_id == chats.card_id, chats.card_type == 2))
|
|
.where(
|
|
sessions.quotation_id == qt_id,
|
|
chats.card_used_yn == True, # noqa: E712
|
|
chats.deleted == False, # noqa: E712
|
|
sessions.deleted == False, # noqa: E712
|
|
)
|
|
.order_by(chats.created_at.asc())
|
|
)
|
|
err_type, rows = await DB_SESSION_MNG.execute(cdb, query)
|
|
if err_type != ErrorType.SUCCESS:
|
|
return err_type, []
|
|
return ErrorType.SUCCESS, list(rows)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED, []
|