diff --git a/backend/common/database/model/models.py b/backend/common/database/model/models.py index 80fe568..befef90 100644 --- a/backend/common/database/model/models.py +++ b/backend/common/database/model/models.py @@ -165,7 +165,7 @@ class quotations(MAIN_BASE): class quotation_settings(MAIN_BASE): - # quotation.quotation_settings (견적 설정). 앵커링값(anchoring_value) 조회용 — agent RL state 입력. + # quotation.quotation_settings (견적 설정). 견적 설정 스냅샷 — anchoring_value 는 구(舊) 앵커 산출용으로 채팅 경로에서는 더 이상 사용하지 않음(앵커는 sessions.target_anchoring_price 박제값). @staticmethod def DBType(): return DBType.QUOTATION.value @@ -176,7 +176,7 @@ class quotation_settings(MAIN_BASE): qt_setting_id = Column(UUID(as_uuid=True), primary_key=True, server_default=text("gen_random_uuid()")) # 견적 설정 식별자(PK) user_id = Column(UUID(as_uuid=True), nullable=False) # 생성 유저(company.users.user_id) target_margin_rate = Column(Numeric(8, 6), nullable=False) # 목표 마진율 - anchoring_value = Column(Numeric(8, 6), nullable=False, server_default=text("0.01")) # 앵커링 값(비율) — anchor=round(target*(1-value)) + anchoring_value = Column(Numeric(8, 6), nullable=False, server_default=text("0.01")) # 앵커링 값(비율) — 구 방식 앵커 비율(현행 앵커 산출에는 미사용 — schedules/anchoring 참조) card_count = Column(Integer, nullable=False, server_default=text("3")) # 협상 내 협상카드 사용 횟수 created_at = Column(DateTime(timezone=True), nullable=False, server_default=text("(now() AT TIME ZONE 'utc')")) # 생성 시각(UTC) updated_at = Column(DateTime(timezone=True), nullable=False, server_default=text("(now() AT TIME ZONE 'utc')"), onupdate=text("(now() AT TIME ZONE 'utc')")) # 수정 시각(UTC, UPDATE 시 자동 갱신) diff --git a/backend/crud/chat_crud.py b/backend/crud/chat_crud.py index 2d97bf8..4c8978b 100644 --- a/backend/crud/chat_crud.py +++ b/backend/crud/chat_crud.py @@ -7,7 +7,7 @@ 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 +from common.enums import ErrorType, SessionStatus from common.logger import LOG @@ -125,7 +125,12 @@ class ChatCRUD(IChatCRUD): values["reject_reason"] = reject_reason[:255] if reject_price is not None: values["reject_price"] = reject_price - query = update(sessions).where(sessions.session_id == session_id).values(**values) + # 진행중일 때만 전이 — 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) @@ -134,11 +139,16 @@ class ChatCRUD(IChatCRUD): 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).values(last_offered_price=price) + 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) diff --git a/backend/services/agent_client.py b/backend/services/agent_client.py index 6c1435d..409e047 100644 --- a/backend/services/agent_client.py +++ b/backend/services/agent_client.py @@ -43,7 +43,7 @@ class AgentChatContext: tenant_id: str # X-Tenant-ID = 견적(갑) 회사 company_id rq_type: str = "재협상" # 재협상 | 재견적 target_price: int = 0 # 갑 목표 매입가(원) - anchor_price: int = 0 # 앵커링가(목표가보다 낮음). quotation_settings.anchoring_value 로 계산. + anchor_price: int = 0 # 앵커링가(목표가보다 낮음). 세션 생성 시 박제된 sessions.target_anchoring_price. item_price: int = 0 # 기존 공급가(품목 기준가). agent 가격협상_확인 인하율 산출용. # 핸드오프 #4: agent 의 RL 상태(state) 계산 입력. # partner_count 는 견적당 세션 수로 산출(실데이터). 나머지 3개는 우리 스키마에 데이터 소스가 없어 diff --git a/backend/tests/test_anchoring_chat.py b/backend/tests/test_anchoring_chat.py index bcfac85..a9ab911 100644 --- a/backend/tests/test_anchoring_chat.py +++ b/backend/tests/test_anchoring_chat.py @@ -29,10 +29,11 @@ def _parse_price(text_): class _AnchorAgent(IAgentClient): - """결정론적 더블: 서비스안내(오프닝) → 기존가격제시(앵커 표시) → 합의 종료. + """결정론적 더블: 서비스안내(오프닝) → 가격 입력 요청 → 합의 종료. - 앵커보다 높은 가격이면 기존가격제시 step 을 반복(노출 기록 1회성 검증용). + 앵커보다 높은 가격이면 같은 step 을 반복(마지막 제시가 덮어쓰기 검증용). 매 턴 수신한 ctx.anchor_price 를 기록해 backend 의 앵커 해석을 관찰한다. + (script 에 앵커가 보이는 건 테스트 관찰 편의일 뿐 — 실제 agent 는 비노출.) """ def __init__(self): diff --git a/postgres-init/01-schema.sql b/postgres-init/01-schema.sql index 6ab5c50..8dd1289 100644 --- a/postgres-init/01-schema.sql +++ b/postgres-init/01-schema.sql @@ -295,7 +295,10 @@ CREATE TABLE IF NOT EXISTS negotiation.sessions ( qt_round INTEGER NOT NULL, -- 견적 라운드(스냅샷) qt_type SMALLINT NOT NULL, -- 견적 유형(스냅샷, QuotationType): 1=renego(재협상 1:1), 2=requote(재견적 1:N), 3=new_nego(신규협상 1:1), 4=new_quote(신규견적 1:N) target_price BIGINT NOT NULL, -- 목표가(원) - target_anchoring_price BIGINT NULL, -- 앵커링가(원) + target_anchoring_price BIGINT NULL, -- 앵커링가(원) — 생성 시 박제(schedules/anchoring 참조) + anchor_rate_permille SMALLINT NULL, -- 제안 당시 앵커링 값(천분율) 박제 + last_offered_price BIGINT NULL, -- 협력사 마지막 제시가(원) — 앵커링 표본 판정의 "가격 흔적" + anchoring_adjustment_id BIGINT NULL, -- 앵커링 배치 소비 마킹(NULL=미처리 0=제외 >0=조정 id) status SMALLINT NOT NULL, -- 진행 상태(SessionStatus): 1=created(생성), 2=in_progress(진행중), 3=done(완료), 4=not_participated(미참여), 5=rejected(거부) bid_price BIGINT NULL, -- 입찰가(원) bid_at TIMESTAMPTZ NULL, -- 입찰 시각 diff --git a/postgres-init/04-alter.sql b/postgres-init/04-alter.sql index 2afaf64..a2087c5 100644 --- a/postgres-init/04-alter.sql +++ b/postgres-init/04-alter.sql @@ -51,3 +51,12 @@ CREATE TABLE IF NOT EXISTS company.notifications ( CREATE INDEX IF NOT EXISTS idx_notifications_user_id ON company.notifications (user_id); CREATE INDEX IF NOT EXISTS idx_notifications_user_unread ON company.notifications (user_id, created_at) WHERE deleted = FALSE AND read_at IS NULL; CREATE INDEX IF NOT EXISTS idx_notifications_ref_qt_id ON company.notifications (ref_qt_id); + +-- ───────────────────────────────────────────────────────────── +-- [2026-07-02] 앵커링 v1.2 — sessions 판정·마킹 컬럼 3종 +-- (신규 DB 는 01-schema.sql 에 반영됨. anchoring 스키마 자체(rate_adjustments·뷰)는 +-- 모듈 소유 DDL schedules/anchoring/schema.sql 로 적용 — 여기엔 두지 않는다.) +ALTER TABLE negotiation.sessions + ADD COLUMN IF NOT EXISTS anchor_rate_permille SMALLINT NULL, -- 제안 당시 앵커링 값(천분율) 박제 + ADD COLUMN IF NOT EXISTS last_offered_price BIGINT NULL, -- 협력사 마지막 제시가(가격 흔적) + ADD COLUMN IF NOT EXISTS anchoring_adjustment_id BIGINT NULL; -- 앵커링 배치 소비 마킹 diff --git a/schedules/anchoring/README.md b/schedules/anchoring/README.md index db25b46..920e089 100644 --- a/schedules/anchoring/README.md +++ b/schedules/anchoring/README.md @@ -6,7 +6,7 @@ > 규범 문서: **`docs/개발용.md`** (정책: `docs/기획용.md`, 흐름 해설: `docs/워크플로우.md`, 타 팀 적용: `docs/인수인계.md`) > **처음 오신 분 / 운영 담당자** → **`docs/운영및유지보수.md`** 부터 보세요 (설치·실행·로그 읽기·트러블슈팅). -> 후속 과제(가격구간 계단식 재설계, anchoring_records 조회 테이블) → **`TODO.md`** +> 과제 이력·백로그 → **`TODO.md`** (주요 과제는 전부 종결) ## 경계 @@ -22,7 +22,7 @@ schema.sql # 모듈 소유 DDL (rate_adjustments + sessions 3컬럼) — psql 수동 적용 src/anchoring/ constants.py # 상수·enum (δ={1:20, 2:10, 3:15} — 제조/총판 스왑 주의) - resources/anchoring_base.json # 정적 기본 테이블(33,334칸, 전부 10‰) — 불변, 시작값의 유일한 소스 + resources/anchoring_base.json # 정적 기본 테이블(46칸 사다리, 전부 10‰) — 불변, 시작값의 유일한 소스 base_table.py # 로드+검증(실패 시 기동 중단) service.py # 순수 계산 (구간·앵커가·판정·평가) — negodata 이식 대상 reader.py # 현재 rate 조회: Redis → 조정 이력 → 정적 테이블 — negodata 이식 대상 @@ -71,5 +71,5 @@ docker logs anchoring | grep -E "WARNING|ERROR" # 이상 신호만 재기동 후 `--once` 1회 실행으로 즉시 캐치업(격주 게이트만 무시, 정책 파라미터 불변). - **Redis 유실/재기동**: 캐시는 파생값 — 매 실행(매주, 게이트 무관) 시작 시 조정 보유 칸 전체를 re-SET 하고 TTL 7일이 보조하므로 자가 회복된다. 수동 복구가 필요하면 `--once`. -- **노출률 0% WARN**: agent 스크립트 스텝명(`기존가격제시`) 변경이나 backend 노출 기록 배선 유실 신호 — 즉시 점검. +- **가격 제시율 0% WARN**: backend 의 `last_offered_price` 기록 배선 유실 신호(학습 무증상 동결) — 즉시 점검. - 조정 이력은 append-only — UPDATE/DELETE 금지. 배치가 sessions 에 쓰는 컬럼은 `anchoring_adjustment_id` 하나뿐. diff --git a/schedules/anchoring/TODO.md b/schedules/anchoring/TODO.md index 4a5d876..5e64295 100644 --- a/schedules/anchoring/TODO.md +++ b/schedules/anchoring/TODO.md @@ -25,3 +25,21 @@ 상세: `docs/개발용.md` §6.3, 사용법: `docs/운영및유지보수.md` §8. 추후 대시보드에서 "전체 칸 나열(무조정 칸 포함)·페이징" 요구가 생기면 그때 스냅샷 테이블로 승격을 재검토한다. + +--- + +## 백로그 (저우선 — 리뷰에서 식별, 착수 조건 명시) + +- [ ] **percent 입력 모드 대비**: `chat_service.send` 는 `user_input_type == "price"` 만 가격으로 + 파싱한다. agent 에 percent 스크립트가 도입되면 percent 턴이 가격 흔적 없이 지나가 + 학습에서 조용히 빠진다(현재 agent 스크립트에 percent 없음 — 잠복). 도입 시 + percent→price 변환(`target*(100-pct)//100`) 후 동일 경로로 태울 것. +- [ ] **스캔 스트리밍**: 배치 스캔이 pending 전량을 메모리에 올린다. 레거시 수백만 행 + 규모 DB 에 첫 적용할 때는 keyset 페이지네이션으로 전환 검토(제외 마킹은 이미 청크 + 커밋이라 트랜잭션 장기화 없음). +- [ ] **Redis 통합 테스트**: 자동 스위트는 무Redis(폴백 경로)로 돈다. CI 에 redis 컨테이너가 + 생기면 §11.5 의 re-SET 회복·TTL·오염 값 방어(get_rate 범위 검증) 케이스를 자동화. +- [ ] **config 오류 메시지**: config.toml 의 오타 키가 TypeError 로 죽는다 — 파일/섹션명을 + 알려주는 검증 메시지로 개선. +- [ ] **운영 Redis 인증**: compose 는 127.0.0.1 바인딩으로 방어했지만, 운영 네트워크에서 + negodata 가 원격 접속하는 구성이면 `requirepass` + `REDIS_PASSWORD` 설정을 적용할 것. diff --git a/schedules/anchoring/docker-compose.yml b/schedules/anchoring/docker-compose.yml index 012fd54..de276b1 100644 --- a/schedules/anchoring/docker-compose.yml +++ b/schedules/anchoring/docker-compose.yml @@ -13,7 +13,7 @@ services: image: redis:7-alpine container_name: anchoring-redis ports: - - "6379:6379" + - "127.0.0.1:6379:6379" # 호스트 로컬만 — 무인증 Redis 를 외부에 열지 않는다(앵커 값 오염 방지) restart: unless-stopped logging: *default-logging diff --git a/schedules/anchoring/docs/개발용.md b/schedules/anchoring/docs/개발용.md index 5b8f9cd..3fde2dc 100644 --- a/schedules/anchoring/docs/개발용.md +++ b/schedules/anchoring/docs/개발용.md @@ -2,7 +2,7 @@ > **문서 성격**: 이 문서만 보고 앵커링 시스템을 구현·유지보수할 수 있도록 작성된 규범 문서(최종 확정본). > **규범 언어**: `MUST` = 반드시 준수, `MUST NOT` = 금지, `SHOULD` = 권장, `MAY` = 선택. -> **스택**: FastAPI(async) + PostgreSQL(SQLAlchemy async / SQL 수동 적용, Alembic 없음) + Redis + APScheduler — **`schedules/anchoring` 자립 컨테이너**(backend 내장 아님). +> **스택**: Python(asyncio) + PostgreSQL(SQLAlchemy async / SQL 수동 적용, Alembic 없음) + Redis + APScheduler — **`schedules/anchoring` 자립 컨테이너**(backend 내장 아님, FastAPI 미사용). > **버전**: v1.2 (2026-07-02 확정) — 정책 배경은 `기획용.md`, 흐름 해설은 `워크플로우.md`, 타 팀(negodata) 적용 명세는 `인수인계.md` 참조. --- @@ -85,7 +85,7 @@ | 경계값 소속 | `target_price`가 정확히 `upper_bound`와 같으면 **다음 idx** 소속. 예: 30,000원 → "3만 원대" 칸(idx 13) | | 내부 인덱스 변환 | `bracket_index = idx − 1` = `bisect_right(UPPER_BOUNDS, price)` (0-기반). DB·Redis·코드 내부는 `bracket_index` 사용 | | 상한 클램프 | `target_price ≥ 90,000,000` → 마지막 구간(idx 46, `bracket_index` 45). **1억 초과도 예외 없이 마지막 인덱스** | -| 시작값 | 칸의 시작 앵커링 값 = 해당 idx의 `anchoring_value` 천분율 변환 정수: `int(anchoring_value * 1000)`. 현재 전 구간 10‰ | +| 시작값 | 칸의 시작 앵커링 값 = 해당 idx의 `anchoring_value` 천분율 변환 정수: `int(round(anchoring_value * 1000))`. 현재 전 구간 10‰ | | 기동 검증 | 로드 시 46행·idx 연속(1..46)·`upper_bound == constants.UPPER_BOUNDS[i]`(사다리 대조)·`0.01 ≤ anchoring_value ≤ 0.20` 검증, 실패 시 **기동 중단** (§13) | - 시작값은 **정적 테이블에서만** 읽는다. 코드에 `0.01`/`10` 하드코딩 **MUST NOT** (테이블이 유일한 소스). @@ -357,7 +357,7 @@ anchoring.current_rates -- 칸별 현재값(최신 조정 행). 여기 없는 주의사항: - `sessions.target_anchoring_price`는 negodata 가 이미 생성 시 채우는 기존 컬럼 — 앵커가 박제로 그대로 활용(신규 컬럼 아님). -- 신규 DB 구축 시 적용 순서: `postgres-init/01~04` → `schedules/anchoring/schema.sql` (IF NOT EXISTS 라 재적용 안전). +- 신규 DB 구축 시 적용 순서: `postgres-init/01~04` → `schedules/anchoring/schema.sql` (IF NOT EXISTS 라 재적용 안전). **sessions 3컬럼은 backend ORM 이 참조하므로 `postgres-init/01-schema.sql`·`04-alter.sql` 에도 반영돼 있다**(backend 가 모듈 DDL 없이도 기동) — anchoring 스키마 자체(테이블·뷰)는 모듈 파일만이 소유. - backend 모델(`models.py`)에는 **sessions 3컬럼만 추가**한다 — `rate_adjustments` 모델은 backend 에 만들지 않는다(무의존). 배치용 ORM 은 모듈이 자체 보유(읽기전용 sessions/quotations/items 매핑 포함). --- @@ -390,11 +390,11 @@ anchoring.current_rates -- 칸별 현재값(최신 조정 행). 여기 없는 ``` 절차 (run_evaluation_batch(force=False)): 0. force 아니고 격주 게이트 미충족 → 절차 0.5 만 수행 후 종료 -0.5. 캐시 정합(매주, 게이트 무관): 조정 이력 보유 칸 전체의 최신 rate 를 Redis 일괄 re-SET +0.5. 캐시 정합(매주, 게이트 무관): Redis ping 확인 후 조정 이력 보유 칸 전체의 최신 rate 를 일괄 re-SET + (Redis 미가용이면 WARN 후 즉시 건너뜀 — 셀마다 timeout 을 태우며 지연되지 않게) (SET 실패·Redis 옛 스냅샷 재기동으로 인한 stale 을 최대 1주 내 회복 — §7) -1. 미처리 종료 재협상 세션 스캔: - sessions s JOIN quotation.quotations q ON q.qt_id = s.quotation_id - JOIN partner.items i ON i.item_id = s.item_id +1. 미처리 종료 재협상 세션 스캔 (LEFT JOIN + ON 절 deleted 필터 — §10 SQL 참조): + sessions s LEFT JOIN quotations q (deleted=false) LEFT JOIN items i (deleted=false) WHERE s.anchoring_adjustment_id IS NULL AND s.deleted = false AND s.qt_type = 1 AND s.status IN (3, 4, 5) 2. 세션별 파생 판정(§4.3): @@ -467,6 +467,8 @@ anchoring.current_rates -- 칸별 현재값(최신 조정 행). 여기 없는 3. agent 컨텍스트로 anchor_price 전달 (기존 AgentChatContext.anchor_price 그대로) ``` +- 퇴화 케이스: `target_price` 가 0/NULL 인 세션은 anchor 0 을 반환한다(기존 동작 보존) — 정상 데이터에서는 발생하지 않는다. + **가격 흔적 기록** (MUST): - agent 는 **변경하지 않는다**. 앵커가는 협력사에게 표시하지 않고(비노출 전략 — 정보 비대칭·상대 선제안 유도) 엔진 내부 체결 임계로만 쓴다. @@ -546,17 +548,19 @@ def get_current_rate(latest_adjusted_rate: int | None, bracket_index: int) -> in """현재 앵커링 값. §4.5 — 조정 이력 없으면 정적 테이블 시작값.""" if latest_adjusted_rate is not None: return latest_adjusted_rate - return get_base_rate_permille(bracket_index) # int(anchoring_value * 1000) + return get_base_rate_permille(bracket_index) # int(round(anchoring_value * 1000)) ``` ```sql -- 배치의 미처리 세션 스캔 (§8 절차 1) -SELECT s.session_id, s.status, s.bid_price, - s.target_price, s.target_anchoring_price, s.anchor_rate_permille, - s.last_offered_price, q.supplier_type, i.company_id +-- LEFT JOIN + ON 절 deleted 필터: 삭제·소실된 견적/상품의 세션은 칸 해석이 NULL 이 되어 +-- 제외 마킹(0)으로 정리된다 — 철회된 거래를 학습에 쓰지 않으면서 영구 재스캔도 방지. +SELECT s.session_id, s.status, s.bid_price, s.target_price, + s.target_anchoring_price, s.last_offered_price, + q.supplier_type, i.company_id FROM negotiation.sessions s -JOIN quotation.quotations q ON q.qt_id = s.quotation_id AND q.deleted = false -JOIN partner.items i ON i.item_id = s.item_id +LEFT JOIN quotation.quotations q ON q.qt_id = s.quotation_id AND q.deleted = false +LEFT JOIN partner.items i ON i.item_id = s.item_id AND i.deleted = false WHERE s.anchoring_adjustment_id IS NULL AND s.deleted = false AND s.qt_type = 1 @@ -576,7 +580,7 @@ LIMIT 1 ## 11. 검증 벡터 (Golden Tests) -아래 케이스가 전부 통과해야 한다 (pytest 고정). 위치: 골든 벡터·배치 통합 = **`schedules/anchoring/tests/`**, 가격 흔적 기록·NULL 폴백 = `backend/tests/`. +아래 케이스가 전부 성립해야 한다. 위치: 골든 벡터·배치 통합 = **`schedules/anchoring/tests/`**, 가격 흔적 기록·NULL 폴백 = `backend/tests/`. 단 Redis 의존 케이스(re-SET 회복)와 다중 프로세스 동시 실행·타이밍 케이스는 자동 스위트(무Redis·단일 프로세스)가 아닌 **수동/후속 검증** 대상이다 — 자동화된 것은 pytest 로 고정돼 있다. ### 11.1 앵커링가 계산 (내림 검증) diff --git a/schedules/anchoring/docs/기획용.md b/schedules/anchoring/docs/기획용.md index f0cc612..bb8715d 100644 --- a/schedules/anchoring/docs/기획용.md +++ b/schedules/anchoring/docs/기획용.md @@ -63,7 +63,7 @@ NegoWiz는 여러 회사가 함께 쓰는 플랫폼이므로, 시스템은 협 - 반대로 **거절이 잦아지면**, "너무 셌구나" 하고 **한발 물러섭니다**. - 받아주는 비율이 **적당한 수준이면 그대로 유지**합니다. 굳이 건드리지 않습니다. -이 시스템은 정확히 이 상인의 감각을 규칙으로 만든 것입니다. 다만 사람과 달리 회사별 수만 개의 칸을 전부 동시에, 감정 없이, 데이터로만 판단합니다. +이 시스템은 정확히 이 상인의 감각을 규칙으로 만든 것입니다. 다만 사람과 달리 회사별 모든 칸(유형×가격대 조합 138개)을 전부 동시에, 감정 없이, 데이터로만 판단합니다. 중요한 특징 하나: **어디까지 깎을 수 있을지는 시스템이 정하는 게 아니라 시장(협력사들)이 정합니다.** 시스템은 상대가 받아주는 한계선을 더듬어 찾아갈 뿐입니다. 그래서 이 값은 "우리가 정한 목표"가 아니라 "시장이 알려준 답"에 가깝습니다. diff --git a/schedules/anchoring/docs/운영및유지보수.md b/schedules/anchoring/docs/운영및유지보수.md index 96887e8..b279eed 100644 --- a/schedules/anchoring/docs/운영및유지보수.md +++ b/schedules/anchoring/docs/운영및유지보수.md @@ -58,7 +58,7 @@ cd schedules/anchoring psql -h -U <계정> -d negosium_db -f schema.sql ``` -- 테이블 1개(`anchoring.rate_adjustments`)와 `negotiation.sessions` 컬럼 3개를 추가합니다. +- 테이블 1개(`anchoring.rate_adjustments`)·조회용 뷰 2개(`rate_history`, `current_rates`)와 `negotiation.sessions` 컬럼 3개를 추가합니다. - `IF NOT EXISTS` 라 **여러 번 실행해도 안전**합니다. ### STEP 2 — 설정 채우기 @@ -69,7 +69,7 @@ cp config.toml.example config.toml ``` 환경변수로 덮어쓸 수도 있습니다(우선순위: env > config.toml > 기본값): -`DB_HOST` `DB_PORT` `DB_USER` `DB_PASSWORD` `DB_NAME` / `REDIS_HOST` `REDIS_PORT` `REDIS_PASSWORD` / `LOG_LEVEL` +`DB_HOST` `DB_PORT` `DB_USER` `DB_PASSWORD` `DB_NAME` / `REDIS_HOST` `REDIS_PORT` `REDIS_DB` `REDIS_PASSWORD` / `LOG_LEVEL` ### STEP 3-A — 도커로 실행 (운영 권장) @@ -106,13 +106,13 @@ PYTHONPATH=src .venv/bin/python -m anchoring.main --once # 도커: docker exec anchoring python -m anchoring.main --once ``` -마지막 줄에 `종료 {'run_id': ..., 'status': 'done', ...}` 가 나오면 성공입니다. +끝부분에 `종료 {'run_id': ..., 'status': 'done', ...}` 와 `[main] 결과: {...}` 가 나오면 성공입니다(부분 실패면 종료코드 1). (협상 데이터가 없으면 `scanned: 0` — 이것도 정상) 테스트 스위트로 확인하려면: ```bash -PYTHONPATH=src .venv/bin/python -m pytest tests/ -q # 15 passed 기대 +PYTHONPATH=src .venv/bin/python -m pytest tests/ -q # 전부 passed 기대(현재 18개) ``` ## 5. 로그 읽는 법 diff --git a/schedules/anchoring/docs/인수인계.md b/schedules/anchoring/docs/인수인계.md index 642230d..987787e 100644 --- a/schedules/anchoring/docs/인수인계.md +++ b/schedules/anchoring/docs/인수인계.md @@ -13,7 +13,7 @@ | 전달물 | 내용 | |---|---| | `schedules/anchoring/src/anchoring/` 모듈 | `constants.py`(상수·enum) · `base_table.py`(정적 테이블 로더) · `service.py`(순수 계산 함수) · `redis_client.py` · `reader.py`(rate 조회) — **전부 async(SQLAlchemy async + redis.asyncio) 자립형이라 negodata 에 그대로 복사/이식 가능** | -| `schedules/anchoring/src/anchoring/resources/anchoring_base.json` | 정적 기본 테이블 (33,334행, 불변) | +| `schedules/anchoring/src/anchoring/resources/anchoring_base.json` | 정적 기본 테이블 (46행 자릿수 사다리, 불변) | | `schedules/anchoring/schema.sql` | `anchoring.rate_adjustments` 테이블 + `negotiation.sessions` 컬럼 3개 ALTER — 모듈 소유 DDL, psql 수동 적용 (적용 시점 협의) | | `schedules/anchoring/docker-compose.yml` | anchoring 서비스 + redis 동봉 — **negodata 는 이 redis 인스턴스를 바라본다** (`REDIS_HOST` 환경변수) | | 이 문서 | 적용 위치·변경 전후 명세 | @@ -49,7 +49,7 @@ else: # supplier_type = quotations.supplier_type (이번 견적의 유형 코드 1/2/3) # bracket = calc_bracket_index(tp) # 자릿수 사다리(46칸) — service 모듈 함수 그대로 이식 # ② rate 조회 — 전달받은 reader 모듈 사용 -rate = await get_anchor_rate(company_id, supplier_type, bracket) +rate = await get_anchor_rate(db, company_id, supplier_type, bracket) # db = AsyncSession # 내부 동작: Redis GET → miss 시 anchoring.rate_adjustments 최신 행 → 없으면 정적 테이블(10‰) # supplier_type ∉ {1,2,3} 이면 get_base_rate_permille(bracket) 사용 (정적 테이블 시작값) # ③ 앵커링가 — 정수 연산만 (float 곱셈 금지: int(tp * 0.99) 형태 재사용 불가) diff --git a/schedules/anchoring/src/anchoring/base_table.py b/schedules/anchoring/src/anchoring/base_table.py index d44915d..3b41b32 100644 --- a/schedules/anchoring/src/anchoring/base_table.py +++ b/schedules/anchoring/src/anchoring/base_table.py @@ -1,6 +1,6 @@ """정적 기본 테이블 — 칸 시작값의 유일한 소스. 규범: §2. -resources/anchoring_base.json(33,334행, 불변)을 서비스 기동 시 메모리에 로드한다. +resources/anchoring_base.json(46행 사다리, 불변)을 서비스 기동 시 메모리에 로드한다. DB 에 저장하지 않으며 런타임에 절대 수정하지 않는다. 검증 실패 시 기동 중단(§13-7). """ import json diff --git a/schedules/anchoring/src/anchoring/batch.py b/schedules/anchoring/src/anchoring/batch.py index def61e5..68209f8 100644 --- a/schedules/anchoring/src/anchoring/batch.py +++ b/schedules/anchoring/src/anchoring/batch.py @@ -3,7 +3,7 @@ 절차: (매주, 게이트 무관) 캐시 re-SET → 격주 게이트 → 미처리 종료 재협상 세션 스캔 → 파생 판정 → EXCLUDED/칸 불가 마킹 0 → 칸별 [조정 INSERT + 소비 마킹 한 트랜잭션, rowcount ≠ n 이면 전체 롤백(MUST — 유니크 가드 없는 구조에서 이중 조정의 유일한 방어선)] -→ 커밋 후 Redis SET → 요약 로그(노출률 0% 면 WARN). +→ 커밋 후 Redis SET → 요약 로그(가격 제시율 0% 면 WARN). """ from collections import defaultdict from datetime import datetime @@ -28,7 +28,7 @@ from anchoring.db import session_scope from anchoring.log import LOG from anchoring.models import Item, Quotation, RateAdjustment, Session from anchoring.reader import get_latest_adjusted_rate -from anchoring.redis_client import consume_failure_counts, set_rate +from anchoring.redis_client import consume_failure_counts, ping, set_rate from anchoring.service import calc_bracket_index, evaluate_pending, judge_sample_type KST = ZoneInfo("Asia/Seoul") @@ -77,8 +77,13 @@ async def _reconcile_cache() -> int: return ok -async def _scan_pending(db) -> list: - """미처리 종료 재협상 세션 + 칸 해석 소스(supplier_type/company_id) 조인. §8 절차 1""" +async def _scan_pending(db, company_ids: list | None = None) -> list: + """미처리 종료 재협상 세션 + 칸 해석 소스(supplier_type/company_id) 조인. §8 절차 1 + + 조인 ON 절에 deleted 필터 — 소프트 삭제된 견적/상품의 세션은 칸 해석이 NULL 이 되어 + 제외 마킹(0)으로 정리된다(철회된 거래를 학습에 쓰지 않으면서 영구 재스캔도 방지). + company_ids: 대상 회사 한정(테스트·표적 수동 실행용). None = 전체. + """ stmt = ( select( Session.session_id, @@ -90,8 +95,16 @@ async def _scan_pending(db) -> list: Quotation.supplier_type, Item.company_id, ) - .join(Quotation, Quotation.qt_id == Session.quotation_id, isouter=True) - .join(Item, Item.item_id == Session.item_id, isouter=True) + .join( + Quotation, + (Quotation.qt_id == Session.quotation_id) & Quotation.deleted.is_(False), + isouter=True, + ) + .join( + Item, + (Item.item_id == Session.item_id) & Item.deleted.is_(False), + isouter=True, + ) .where( Session.anchoring_adjustment_id.is_(None), Session.deleted.is_(False), @@ -99,6 +112,8 @@ async def _scan_pending(db) -> list: Session.status.in_(TERMINAL_SESSION_STATUSES), ) ) + if company_ids: + stmt = stmt.where(Item.company_id.in_(company_ids)) return (await db.execute(stmt)).all() @@ -168,20 +183,29 @@ def _new_company_agg() -> dict: return {"evaluated": 0, "up": 0, "hold": 0, "down": 0, "carryover": 0, "failed": 0, "excluded": 0} -async def run_evaluation_batch(force: bool = False) -> dict: +async def run_evaluation_batch(force: bool = False, company_ids: list | None = None) -> dict: """배치 1회. force=True 면 격주 게이트만 무시(정책 파라미터는 불변). + company_ids: 대상 회사 한정 — 테스트가 공유 DB 의 실데이터를 소비하지 않게 하는 + 격리 장치이자, 특정 테넌트만 표적 수동 실행하는 운영 옵션. None = 전체(운영 기본). + 로그 규약: 모든 라인에 `[batch {run_id}]` 태그(회차 grep), 칸/회사 단위 라인은 `company=` `type=` `bracket=` key=value 형식(회사별 grep — `grep company=`). """ now = datetime.now(KST) run_id = now.strftime("%Y%m%d-%H%M%S") tag = f"[batch {run_id}]" - LOG.info(f"{tag} 시작 — ISO 주차 {now.isocalendar().week}, force={force}") + scope = f", 대상 회사 {len(company_ids)}곳" if company_ids else "" + LOG.info(f"{tag} 시작 — ISO 주차 {now.isocalendar().week}, force={force}{scope}") - # 절차 0.5 — 캐시 정합(매주, 게이트 무관) - reconciled = await _reconcile_cache() - LOG.info(f"{tag} 캐시 re-SET {reconciled}칸") + # 절차 0.5 — 캐시 정합(매주, 게이트 무관). Redis 다운이면 즉시 건너뜀 + # (셀마다 timeout 을 태우며 수십 분 지연되는 것 방지 — TTL·다음 주 re-SET 이 회복) + if await ping(): + reconciled = await _reconcile_cache() + LOG.info(f"{tag} 캐시 re-SET {reconciled}칸") + else: + reconciled = 0 + LOG.warning(f"{tag} Redis 미가용 — 캐시 re-SET 건너뜀(읽기는 DB 폴백으로 동작)") if not force and not is_evaluation_week(now): _log_redis_failures(tag) @@ -190,7 +214,7 @@ async def run_evaluation_batch(force: bool = False) -> dict: # 절차 1~2 — 스캔 + 파생 판정 async with session_scope() as db: - rows = await _scan_pending(db) + rows = await _scan_pending(db, company_ids) excluded_ids: list = [] cells: dict[tuple, list] = defaultdict(list) @@ -222,11 +246,15 @@ async def run_evaluation_batch(force: bool = False) -> dict: LOG.warning(f"{tag} 가격 제시 흔적 0% (종료 재협상 {len(rows)}건 중 last_offered_price 전무) " f"— backend 가격 입력 기록 배선 점검 필요") - # 절차 2 — 제외 확정 마킹(재스캔 방지) + # 절차 2 — 제외 확정 마킹(재스캔 방지). 청크별 개별 커밋 — 판정이 결정적이라 + # 원자성이 불필요하고(중단 시 다음 회차가 이어서 마킹), 첫 실행의 레거시 대량 + # 마킹이 장시간 단일 트랜잭션(WAL·락)을 만드는 것을 방지한다. if excluded_ids: - async with session_scope() as db: - await _mark_sessions(db, excluded_ids, MARK_EXCLUDED) - LOG.info(f"{tag} 제외 확정 마킹 {len(excluded_ids)}건") + marked_total = 0 + for i in range(0, len(excluded_ids), _MARK_CHUNK): + async with session_scope() as db: + marked_total += await _mark_sessions(db, excluded_ids[i:i + _MARK_CHUNK], MARK_EXCLUDED) + LOG.info(f"{tag} 제외 확정 마킹 {marked_total}건") # 절차 3~4 — 칸별 평가(칸 단위 독립 트랜잭션 — 한 칸 실패가 전파되지 않음) evaluated = up = hold = down = clamped = failed = 0 diff --git a/schedules/anchoring/src/anchoring/main.py b/schedules/anchoring/src/anchoring/main.py index fe95b24..11084fe 100644 --- a/schedules/anchoring/src/anchoring/main.py +++ b/schedules/anchoring/src/anchoring/main.py @@ -8,6 +8,7 @@ 기동 시 정적 테이블 검증 실패 → 예외로 즉시 중단(§13-7 MUST). """ import asyncio +import signal import sys from anchoring.base_table import load_base_table @@ -34,20 +35,32 @@ async def _run(once: bool) -> None: LOG.info("[main] 수동 1회 실행(--once, 격주 게이트 무시)") result = await run_evaluation_batch(force=True) LOG.info(f"[main] 결과: {result}") - return + return result scheduler = build_scheduler() scheduler.start() job = scheduler.get_job("anchoring_biweekly_evaluation") LOG.info(f"[main] 스케줄러 상주 시작 — 다음 실행 예정: {job.next_run_time}") - await asyncio.Event().wait() # 컨테이너 메인 — SIGTERM 까지 대기 + + # SIGTERM(docker stop)/SIGINT 를 받아 정상 종료 — finally(리소스 정리)가 반드시 실행되게 한다 + stop = asyncio.Event() + loop = asyncio.get_running_loop() + for sig in (signal.SIGTERM, signal.SIGINT): + loop.add_signal_handler(sig, stop.set) + await stop.wait() + LOG.info("[main] 종료 신호 수신 — 정리 후 종료") + scheduler.shutdown(wait=False) + return None finally: await close_redis() await dispose_engine() def main() -> None: - asyncio.run(_run(once="--once" in sys.argv)) + result = asyncio.run(_run(once="--once" in sys.argv)) + # --once 가 부분 실패(partial)로 끝나면 비정상 종료코드 — 런북/cron 에서 감지 가능해야 한다 + if result is not None and result.get("status") not in ("done", "skipped"): + sys.exit(1) if __name__ == "__main__": diff --git a/schedules/anchoring/src/anchoring/reader.py b/schedules/anchoring/src/anchoring/reader.py index 6ae0d2e..797f9f2 100644 --- a/schedules/anchoring/src/anchoring/reader.py +++ b/schedules/anchoring/src/anchoring/reader.py @@ -38,5 +38,5 @@ async def get_anchor_rate(db: AsyncSession, company_id, supplier_type: int, brac rate = await get_latest_adjusted_rate(db, company_id, supplier_type, bracket_index) if rate is None: rate = get_base_rate_permille(bracket_index) - await set_rate(company_id, supplier_type, bracket_index, rate) + await set_rate(company_id, supplier_type, bracket_index, rate, nx=True) return rate diff --git a/schedules/anchoring/src/anchoring/redis_client.py b/schedules/anchoring/src/anchoring/redis_client.py index 4b3f129..fa33b04 100644 --- a/schedules/anchoring/src/anchoring/redis_client.py +++ b/schedules/anchoring/src/anchoring/redis_client.py @@ -7,7 +7,7 @@ import redis.asyncio as aioredis from anchoring.config import RedisConfig -from anchoring.constants import CACHE_TTL_SECONDS, REDIS_SOCKET_TIMEOUT +from anchoring.constants import ANCHOR_RATE_MAX, ANCHOR_RATE_MIN, CACHE_TTL_SECONDS, REDIS_SOCKET_TIMEOUT from anchoring.log import LOG _client: aioredis.Redis | None = None @@ -56,24 +56,45 @@ def anchor_key(company_id, supplier_type: int, bracket_index: int) -> str: return f"anchor:{company_id}:{supplier_type}:{bracket_index}" +async def ping() -> bool: + """Redis 가용성 확인 — 대량 re-SET 전에 1회 확인해 다운 시 즉시 건너뛴다.""" + if _client is None: + return False + try: + return bool(await _client.ping()) + except Exception: + return False + + async def get_rate(company_id, supplier_type: int, bracket_index: int) -> int | None: """캐시 조회. 미스·에러·클라이언트 미초기화 → None(호출측이 DB 폴백).""" if _client is None: return None try: raw = await _client.get(anchor_key(company_id, supplier_type, bracket_index)) - return int(raw) if raw is not None else None + if raw is None: + return None + rate = int(raw) + # 방어: 캐시 오염(외부 SET 등)으로 정책 범위 밖 값이 오면 미스로 취급 → DB 폴백 + 재적재로 자가 교정 + if not (ANCHOR_RATE_MIN <= rate <= ANCHOR_RATE_MAX): + LOG.warning(f"[redis] 범위 밖 캐시 값 무시(오염 의심) key={anchor_key(company_id, supplier_type, bracket_index)} value={raw}") + return None + return rate except Exception as ex: _note_failure("get", anchor_key(company_id, supplier_type, bracket_index), ex) return None -async def set_rate(company_id, supplier_type: int, bracket_index: int, rate: int) -> bool: - """캐시 적재(best effort, TTL 7일). 실패해도 예외를 밖으로 던지지 않는다.""" +async def set_rate(company_id, supplier_type: int, bracket_index: int, rate: int, nx: bool = False) -> bool: + """캐시 적재(best effort, TTL 7일). 실패해도 예외를 밖으로 던지지 않는다. + + nx=True: 키가 없을 때만 적재 — 읽기 경로의 미스 백필용(배치가 방금 쓴 새 값을 + 구값으로 덮어쓰는 write-after-read 경합 방지). + """ if _client is None: return False try: - await _client.set(anchor_key(company_id, supplier_type, bracket_index), str(rate), ex=CACHE_TTL_SECONDS) + await _client.set(anchor_key(company_id, supplier_type, bracket_index), str(rate), ex=CACHE_TTL_SECONDS, nx=nx) return True except Exception as ex: _note_failure("set", anchor_key(company_id, supplier_type, bracket_index), ex) diff --git a/schedules/anchoring/tests/test_batch.py b/schedules/anchoring/tests/test_batch.py index c47dccc..95c194b 100644 --- a/schedules/anchoring/tests/test_batch.py +++ b/schedules/anchoring/tests/test_batch.py @@ -15,6 +15,11 @@ from conftest import requires_db pytestmark = requires_db + +async def _run_batch(*seeders): + """테스트 전용: 시드한 회사로 스코프 — 공유 dev DB 의 실데이터를 소비하지 않는다.""" + return await run_evaluation_batch(force=True, company_ids=[s.company_id for s in seeders]) + # 시드 기본값: target 30,000 / rate 10‰ / anchor 29,700 → bracket 12 ("3만 원대" 칸) BRACKET = 12 SUCCESS_BID = 29_000 # ≤ anchor → BID_SUCCESS @@ -55,7 +60,7 @@ async def test_full_cycle_and_idempotency(seeder, caplog): # 합계 13건, 성공 8 → r≈0.615 → +20 with caplog.at_level(logging.INFO, logger="anchoring"): - result = await run_evaluation_batch(force=True) + result = await run_evaluation_batch(force=True, company_ids=[seeder.company_id]) # 로그 규약: run_id 태그 + 회사별 grep 가능한 칸별 조정 라인 + 회사요약 라인 assert result["run_id"] @@ -74,7 +79,7 @@ async def test_full_cycle_and_idempotency(seeder, caplog): assert all(v == adj.id for v in marks.values()) # 13건 모두 소비 마킹 # 재실행 — 마킹 멱등: 우리 칸 조정은 그대로 1건 - await run_evaluation_batch(force=True) + await _run_batch(seeder) async with adb.session_scope() as db: assert len(await _adjustments(db, seeder)) == 1 @@ -96,7 +101,7 @@ async def test_carryover(seeder): async with adb.session_scope() as db: first = await _seed_mixed(db, seeder, success=5, fail=2) # 7건 < 10 - await run_evaluation_batch(force=True) + await _run_batch(seeder) async with adb.session_scope() as db: assert await _adjustments(db, seeder) == [] marks = await _marks(db, first) @@ -104,7 +109,7 @@ async def test_carryover(seeder): second = await _seed_mixed(db, seeder, success=3, fail=3) # 누적 13건 (8S/5F) - await run_evaluation_batch(force=True) + await _run_batch(seeder) async with adb.session_scope() as db: adjustments = await _adjustments(db, seeder) assert len(adjustments) == 1 @@ -120,7 +125,7 @@ async def test_company_isolation_and_reader(seeder): await _seed_mixed(db, seeder, success=10, fail=0, supplier_type=1) # 유통 → +20 await _seed_mixed(db, seeder, success=10, fail=0, supplier_type=2) # 제조 → +10 (δ 스왑 가드) - await run_evaluation_batch(force=True) + await _run_batch(seeder) other_company = uuid.uuid4() async with adb.session_scope() as db: @@ -140,8 +145,9 @@ async def test_excluded_and_unsampleable(seeder): excluded += [await seeder.seed_session(db, status=4, last_offered_price=None) for _ in range(7)] valid = await _seed_mixed(db, seeder, success=5, fail=0) # 유효 5 < 10 untyped = [await seeder.seed_session(db, supplier_type=0, bid_price=SUCCESS_BID)] # 칸 구성 불가 + untyped.append(await seeder.seed_session(db, supplier_type=None, bid_price=SUCCESS_BID)) # NULL 도 동일(§13-6) - await run_evaluation_batch(force=True) + await _run_batch(seeder) async with adb.session_scope() as db: assert await _adjustments(db, seeder) == [] # 유효 5 < 10 → 평가 없음 @@ -172,3 +178,14 @@ async def test_marking_conflict_rolls_back(seeder): assert await _adjustments(db, seeder) == [] # 롤백 — 이중 조정 없음 marks = await _marks(db, ids[1:]) assert all(v is None for v in marks.values()) # 나머지 9건 마킹도 롤백 + + +# ── 무증상 고장 감지: 가격 제시 흔적 0% → WARN ──────────── +async def test_priced_rate_zero_warns(seeder, caplog): + async with adb.session_scope() as db: + for _ in range(3): # 전부 가격 흔적 없는 종료 → priced_rate 0 + await seeder.seed_session(db, status=5, last_offered_price=None) + + with caplog.at_level(logging.WARNING, logger="anchoring"): + await run_evaluation_batch(force=True, company_ids=[seeder.company_id]) + assert any("가격 제시 흔적 0%" in m for m in caplog.messages) diff --git a/schedules/anchoring/tests/test_core.py b/schedules/anchoring/tests/test_core.py index 9b8dfb3..cc03cea 100644 --- a/schedules/anchoring/tests/test_core.py +++ b/schedules/anchoring/tests/test_core.py @@ -2,10 +2,13 @@ 실행: cd schedules/anchoring && PYTHONPATH=src python -m pytest tests/test_core.py -q """ +from datetime import datetime + import pytest from anchoring.base_table import BaseTableError, _validate, get_base_rate_permille, load_base_table from anchoring.constants import BRACKET_COUNT, UPPER_BOUNDS, AnchoringSampleType +from anchoring.batch import is_evaluation_week from anchoring.service import ( calc_anchor_price, calc_bracket_index, @@ -130,3 +133,11 @@ def test_judge_sample_type(): def test_get_current_rate(): assert get_current_rate(70, 0) == 70 assert get_current_rate(None, 0) == 10 # 이력 없으면 정적 테이블 시작값 + + +# ── §8 격주 게이트 (ISO 주차 짝수 토요일만 평가) ────────── +def test_is_evaluation_week(): + assert datetime(2026, 7, 4).isocalendar().week == 27 # 홀수 주 토요일 + assert is_evaluation_week(datetime(2026, 7, 4)) is False + assert datetime(2026, 7, 11).isocalendar().week == 28 # 짝수 주 토요일 + assert is_evaluation_week(datetime(2026, 7, 11)) is True