모듈·backend·문서 3개 관점의 적대 리뷰에서 확인된 결함 일괄 수정: [모듈] - 배치에 company_ids 스코프 옵션 추가 — 통합 테스트가 공유 dev DB 의 실세션을 소비/마킹하던 문제 해소(테스트는 시드 회사로 한정), 표적 수동 실행 옵션 겸용 - Redis 방어: compose 포트를 127.0.0.1 바인딩(무인증 공개 차단), get_rate 에 범위([10,200]) 검증 — 오염 캐시값은 미스 취급 후 자가 교정, 미스 백필은 SET NX (배치가 방금 쓴 새 값을 구값으로 덮는 write-after-read 경합 방지) - 배치: Redis ping 후 re-SET(다운 시 셀×timeout 지연 없이 즉시 스킵), 스캔 조인 ON 절에 quotations/items deleted 필터(철회 거래를 학습에서 배제), 제외 마킹을 청크별 커밋(레거시 대량 첫 실행의 장시간 단일 트랜잭션 방지) - main: SIGTERM/SIGINT 핸들러(docker stop 시 정리 로직 보장), --once 부분 실패 시 종료코드 1(런북/cron 감지 가능) [backend] - finalize_session·update_last_offered_price 에 status=IN_PROGRESS 가드 — negodata 일괄마감/중복 전송 경합이 종료된 세션을 되살리거나 가격 흔적을 사후 변경하는 것 차단(파생 판정 결정성 보호) - 신규 DB 부트스트랩: sessions 3컬럼을 postgres-init/01-schema·04-alter 에도 반영(backend 가 모듈 DDL 없이 기동) — anchoring 스키마 자체는 모듈 소유 유지 - 낡은 주석 정리(agent_client·quotation_settings 의 구 앵커 산출 서술) [테스트·문서] - 신규 테스트: 격주 게이트 골든(ISO 주차), supplier_type NULL, 가격 제시율 0% WARN — 모듈 18개·backend 57개 통과 - 문서 정합 감사 20건 반영: 잔존 33,334/노출 문구 제거, §10 SQL 을 실제 코드 (LEFT JOIN+deleted)와 일치, §11 자동/수동 검증 구분, FastAPI 오기 제거, 인수인계 reader 시그니처(db 인자), TODO 백로그 5건 기록 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_offered_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_offered_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_offered_price=price)
|
|
)
|
|
return await DB_SESSION_MNG.add(cdb, query)
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(ex)
|
|
return ErrorType.DB_RUN_FAILED
|