용어 체계: 값=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>
156 lines
7.0 KiB
Python
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
|