o2o-negosium-original/backend/crud/chat_crud.py
민헌 a2c299aa14 refactor(anchoring): 도메인 이름 전면 개편 — adjustments·anchoring_value/price·price_range·sample 용어 통일
용어 체계: 값=anchoring_value(정수‰)·가격=anchoring_price·조정=adjustment·구간=price_range·표본=sample

- DB: rate_adjustments→anchoring.adjustments (id→adjustment_id, price_bracket_index→price_range_index,
  nego_count→sample_count, anchor_rate_before/after→anchoring_value_before/after,
  consumed_session_ids→used_session_ids)
- sessions: target_anchoring_price→anchoring_price, anchor_rate_permille→anchoring_value,
  last_offered_price→last_offer_price, anchoring_adjustment_id→used_by_adjustment_id
- 뷰: rate_history/current_rates→value_history/current_values, delta_permille→value_change
- 코드: calc_price_range_index·calc_anchoring_price·evaluate_samples·get_current_value·
  get_latest_adjusted_value·get_current_anchoring_value·fetch_current_values·get_base_anchoring_value·
  Adjustment(ORM)·update_last_offer_price, 상수 ANCHORING_VALUE_MIN/MAX·ADJUSTMENT_STEP·
  PRICE_RANGE_COUNT/INDEX_MAX, 배치 로그 키 bracket=→price_range=
- API: negodata protocol 필드 target_anchoring_price→anchoring_price (front 생성 모델·컴포넌트 동반)
- 기존 DB 마이그레이션 신설: schedules/anchoring/migrations/20260706_rename_anchoring.sql
  (멱등 DO 블록 — 테이블·컬럼·뷰·인덱스·PK 제약. 코드 배포와 동시 적용 필요)
- postgres-init 01·04, 문서 6종 동기화
- 실배포 전 수정 포함: main.py argparse 화(--dry-run 단독·오타 플래그 기동 전 차단),
  박제 정합식 calc_anchoring_price 재사용, clamped 지표가 실제 포화만 집계(경계값 유지 제외)

주의: sessions.anchoring_value(정수‰)와 quotation_settings.anchoring_value(구 float 비율)는
같은 이름·다른 단위 — 구 컬럼은 미변경.

검증: 모듈 20·negodata 50·backend 57 테스트 통과, front tsc·vite build 통과,
로컬 DB 마이그레이션 적용 후 배치 dry-run·상주 기동·양 서버 부팅 확인.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-06 11:17:01 +09:00

156 lines
7.0 KiB
Python

from abc import ABC, abstractmethod
from datetime import datetime, timezone
from typing import Optional, Tuple
from sqlalchemy import asc, desc, select, update
from sqlalchemy.ext.asyncio import AsyncSession
from common.database.db_session_manager import DB_SESSION_MNG
from common.database.model.models import chats, items, sessions
from common.enums import ErrorType, SessionStatus
from common.logger import LOG
# 협상 채팅 CRUD. 메시지 로그(negotiation.chats)와 종료 시 세션 입찰 확정(negotiation.sessions)을 다룬다.
# chats / sessions 모두 NEGOTIATION 논리 DB 라 한 트랜잭션(execute_lambda_run)으로 묶을 수 있다.
class IChatCRUD(ABC):
@abstractmethod
async def list_by_session(self, cdb: AsyncSession, session_id) -> Tuple[ErrorType, list]:
pass
@abstractmethod
async def get_last(self, cdb: AsyncSession, session_id) -> Tuple[ErrorType, Tuple[int, Optional[int], Optional[dict]]]:
"""마지막 메시지의 (seq, sender, meta). 없으면 (0, None, None).
동시전송 가드 + seq 채번 + 입력-모드 검증(직전 봇 meta.input_mode)에 사용."""
pass
@abstractmethod
async def insert_message(self, cdb: AsyncSession, message: chats) -> ErrorType:
pass
@abstractmethod
async def soft_delete_message(self, cdb: AsyncSession, chat_id) -> ErrorType:
pass
@abstractmethod
async def get_item_by_id(self, cdb: AsyncSession, item_id) -> Tuple[ErrorType, items]:
pass
@abstractmethod
async def finalize_session(
self, cdb: AsyncSession, session_id, status: int,
bid_price: Optional[int] = None, reject_reason: Optional[str] = None, reject_price: Optional[int] = None,
) -> ErrorType:
pass
@abstractmethod
async def update_last_offer_price(self, cdb: AsyncSession, session_id, price: int) -> ErrorType:
pass
class ChatCRUD(IChatCRUD):
async def list_by_session(self, cdb: AsyncSession, session_id) -> Tuple[ErrorType, list]:
try:
# (session_id, seq) 유니크 인덱스가 정렬 스캔을 커버한다.
query = (
select(chats)
.where(chats.session_id == session_id, chats.deleted == False) # noqa: E712
.order_by(asc(chats.seq))
)
err_type, rows = await DB_SESSION_MNG.execute(cdb, query, "list_by_session failed.")
if err_type != ErrorType.SUCCESS:
return err_type, []
return ErrorType.SUCCESS, rows
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED, []
async def get_last(self, cdb: AsyncSession, session_id) -> Tuple[ErrorType, Tuple[int, Optional[int], Optional[dict]]]:
try:
query = (
select(chats.seq, chats.sender, chats.meta)
.where(chats.session_id == session_id, chats.deleted == False) # noqa: E712
.order_by(desc(chats.seq))
.limit(1)
)
err_type, rows = await DB_SESSION_MNG.execute(cdb, query, "get_last failed.")
if err_type != ErrorType.SUCCESS:
return err_type, (0, None, None)
if not rows:
return ErrorType.SUCCESS, (0, None, None)
return ErrorType.SUCCESS, (rows[0][0], rows[0][1], rows[0][2])
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED, (0, None, None)
async def insert_message(self, cdb: AsyncSession, message: chats) -> ErrorType:
try:
return await DB_SESSION_MNG.insert(cdb, message, raise_error=False)
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED
async def soft_delete_message(self, cdb: AsyncSession, chat_id) -> ErrorType:
# agent 실패 시 선점(pre-claim)한 유저 메시지를 되돌린다. 부분 유니크(WHERE deleted=FALSE)라 seq 가 다시 비워진다.
try:
query = update(chats).where(chats.chat_id == chat_id).values(deleted=True)
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_item_by_id(self, cdb: AsyncSession, item_id) -> Tuple[ErrorType, items]:
try:
query = select(items).where(items.item_id == item_id, items.deleted == False).limit(1) # noqa: E712
err_type, row_list = await DB_SESSION_MNG.execute(cdb, query, f"get_item_by_id({item_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 finalize_session(
self, cdb: AsyncSession, session_id, status: int,
bid_price: Optional[int] = None, reject_reason: Optional[str] = None, reject_price: Optional[int] = None,
) -> ErrorType:
try:
values = {"status": status}
if bid_price is not None:
values["bid_price"] = bid_price
values["bid_at"] = datetime.now(timezone.utc)
if reject_reason is not None:
values["reject_reason"] = reject_reason[:255]
if reject_price is not None:
values["reject_price"] = reject_price
# 진행중일 때만 전이 — negodata 일괄마감/중복 전송 경합이 종료된 세션을 되살리지 못하게 가드
query = (
update(sessions)
.where(sessions.session_id == session_id, sessions.status == SessionStatus.IN_PROGRESS.value)
.values(**values)
)
return await DB_SESSION_MNG.add(cdb, query)
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED
async def update_last_offer_price(self, cdb: AsyncSession, session_id, price: int) -> ErrorType:
"""협력사 마지막 제시가 갱신 — 가격 입력 턴의 봇 메시지 저장과 같은 트랜잭션에서 호출.
진행 중엔 매 가격 입력마다 덮어쓰고 종료 후엔 불변. 앵커링 표본 판정에서
"가격을 한 번이라도 써낸 협상"을 가르는 기준값(NULL=가격 흔적 없음 → 집계 제외).
status 가드: 일괄마감 등으로 이미 종료된 세션의 흔적을 사후에 바꾸지 못하게 한다(판정 결정성 보호).
"""
try:
query = (
update(sessions)
.where(sessions.session_id == session_id, sessions.status == SessionStatus.IN_PROGRESS.value)
.values(last_offer_price=price)
)
return await DB_SESSION_MNG.add(cdb, query)
except Exception as ex:
LOG.e_no_callstack(ex)
return ErrorType.DB_RUN_FAILED