fix(anchoring): 3방향 적대 리뷰 반영 — 테스트 격리·Redis 방어·운영 견고화

모듈·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>
This commit is contained in:
민헌 2026-07-02 20:55:46 +09:00
parent 24cd74dda0
commit ba996f47c1
20 changed files with 202 additions and 67 deletions

View File

@ -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 시 자동 갱신)

View File

@ -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)

View File

@ -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개는 우리 스키마에 데이터 소스가 없어

View File

@ -29,10 +29,11 @@ def _parse_price(text_):
class _AnchorAgent(IAgentClient):
"""결정론적 더블: 서비스안내(오프닝) → 기존가격제시(앵커 표시) → 합의 종료.
"""결정론적 더블: 서비스안내(오프닝) → 가격 입력 요청 → 합의 종료.
앵커보다 높은 가격이면 기존가격제시 step 을 반복(노출 기록 1회성 검증용).
앵커보다 높은 가격이면 같은 step 을 반복(마지막 제시가 덮어쓰기 검증용).
매 턴 수신한 ctx.anchor_price 를 기록해 backend 의 앵커 해석을 관찰한다.
(script 에 앵커가 보이는 건 테스트 관찰 편의일 뿐 — 실제 agent 는 비노출.)
"""
def __init__(self):

View File

@ -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, -- 입찰 시각

View File

@ -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; -- 앵커링 배치 소비 마킹

View File

@ -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` 하나뿐.

View File

@ -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` 설정을 적용할 것.

View File

@ -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

View File

@ -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 앵커링가 계산 (내림 검증)

View File

@ -63,7 +63,7 @@ NegoWiz는 여러 회사가 함께 쓰는 플랫폼이므로, 시스템은 협
- 반대로 **거절이 잦아지면**, "너무 셌구나" 하고 **한발 물러섭니다**.
- 받아주는 비율이 **적당한 수준이면 그대로 유지**합니다. 굳이 건드리지 않습니다.
이 시스템은 정확히 이 상인의 감각을 규칙으로 만든 것입니다. 다만 사람과 달리 회사별 수만 개의 칸을 전부 동시에, 감정 없이, 데이터로만 판단합니다.
이 시스템은 정확히 이 상인의 감각을 규칙으로 만든 것입니다. 다만 사람과 달리 회사별 모든 칸(유형×가격대 조합 138개)을 전부 동시에, 감정 없이, 데이터로만 판단합니다.
중요한 특징 하나: **어디까지 깎을 수 있을지는 시스템이 정하는 게 아니라 시장(협력사들)이 정합니다.** 시스템은 상대가 받아주는 한계선을 더듬어 찾아갈 뿐입니다. 그래서 이 값은 "우리가 정한 목표"가 아니라 "시장이 알려준 답"에 가깝습니다.

View File

@ -58,7 +58,7 @@ cd schedules/anchoring
psql -h <DB호스트> -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. 로그 읽는 법

View File

@ -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) 형태 재사용 불가)

View File

@ -1,6 +1,6 @@
"""정적 기본 테이블 — 칸 시작값의 유일한 소스. 규범: §2.
resources/anchoring_base.json(33,334행, 불변)을 서비스 기동 시 메모리에 로드한다.
resources/anchoring_base.json(46행 사다리, 불변)을 서비스 기동 시 메모리에 로드한다.
DB 에 저장하지 않으며 런타임에 절대 수정하지 않는다. 검증 실패 시 기동 중단(§13-7).
"""
import json

View File

@ -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=<uuid>`).
"""
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

View File

@ -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__":

View File

@ -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

View File

@ -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)

View File

@ -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)

View File

@ -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