o2o-negosium-original/agent/services/chat_service.py
jwkim 51b3820cc9 [feat] agent: 완전 자율 협상 모드 (AUTONOMY_MODE) + LLM 멘트
판정 룰(앵커타결/와일드존/3라운드결렬)과 카드 선택을 학습 정책으로 대체:
수락/역제안 금액/압박 화법/결렬 전부 행동 30개(수락1+결렬1+역제안 6단x화법4+압박4)에서 선택.

- autonomy_actions: 행동 공간·특징 인코딩 (학습/서빙 공유)
- autonomy_store: numpy 서빙 + 행동 봉투 7개(수락<=목표가 / 역제안 단조 / 역제시·결렬은
  설득 2회 후 해금 / 마무리국면 압박 금지 / 첫 역제안 앵커 이하 / 최종제안 1회 보장)
- chat_engine: 자율 스텝(역제안/최종제안/압박1~4), 최종제안 금액=목표가, 턴캡 12
- ment_generator: Gemini 멘트 생성 + 가드(숫자 화이트리스트·목표가 비공개·금지어·문장완결,
  실패시 템플릿 폴백, 6초 컷), 인터넷최저가 근거 인용(수집됨+제시가 초과시만), 대화 기억
- chat_service: 자율 행동 experience_logs 로깅(AUT|종류|위치|전략), 대화기억 ctx 관리
- context loader/CRUD: 인터넷최저가·견적기간·협력사 이력 로드 (v3 상태 21차원)
- train_full_autonomy: 시뮬 15k ep — 앵커율 0.8~6% 정합, 협력사 현실화(컷반발·반복짜증·
  양보 상호성), 관측성 마스크(마감 40% 미관측·15% 전부미상 — 서빙 중립값 분포 정합),
  보상 수정(목표가 초과 타결=결렬 취급)
- 서빙 v3.5 (v3.3 목표가 즉시지르기 퇴화, v3.4 소액지형 첫턴 통보 퇴화 — 게이트 반려 이력 보관)

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-10 16:57:38 +09:00

376 lines
24 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""ChatService — 대화형 /chat 오케스트레이션 (P7 슬라이스).
ChatEngine(동기 step 전이) + UCB Q-Table(가격협상 카드선택·학습) + DB(experience_logs) 결합.
세션 상태는 인메모리 스토어(PoC; 단일 워커). 종료 outcome 은 마지막 (state,action)에 종료보상을 역전파.
"""
import uuid
from typing import Optional
import numpy as np
from common.enums import DBType, ErrorType
from common.database.db_session_manager import DB_SESSION_MNG
from common.logger import LOG
from config.server_configs import agent_config
from negotiation.chat.service import ment_generator
from negotiation.chat.service.chat_engine import ChatEngine, ChatSession, StepView
from negotiation.chat.service.indicator import compute_indicator
from negotiation.chat.service.chat_session_repository import ChatSessionRepository
from negotiation.chat.service.negotiation_context_loader import NegotiationContextLoader
from negotiation.chat.service.script_repository import ScriptRepository
from negotiation.policies.base import EpisodeState, PolicyContext, Transition
from negotiation.policies.autonomy_actions import ACTIONS as AUTONOMY_ACTIONS
from negotiation.policy.autonomy_store import AutonomyStore
from negotiation.policy.dqn_store import DQNServingStore
from negotiation.policy.model_store import QTablePolicyStore
from negotiation.qtable.domain.model.snapshot import NegotiationOutcome, NegotiationSnapshot, PartnerType
from negotiation.qtable.domain.service.reward_calculator import RewardCalculator
from negotiation.qtable.domain.service.state_calculator import state_index
from negotiation.qtable.infra.repository.learning_repository import LearningRepository
from router.v1.chat.protocol import Req_Chat, Res_Chat, Res_ChatSession
from tenancy.registry import TenantEngine
# 직접 호출(데모/테스트) 폴백 기본 컨텍스트 — 운영 경로는 NegotiationContextLoader 가
# DB(negotiation.sessions·partner.items·partner.suppliers·quotation.quotations)에서 조회한다.
_DEFAULT_RQ_TYPE = "재협상"
_DEFAULT_TARGET_PRICE = 10000 # KT 목표 매입가
_DEFAULT_ANCHOR_PRICE = 9900 # 앵커링가(목표가보다 낮음). 제시가 ≤ anchor → 우선협상
_DEFAULT_REVENUE_AMOUNT = 20_000_000 # 매출액(원) — suppliers.total_revenue 미기재 시 폴백
_DEFAULT_DISTRIBUTION_CODE = "A" # 유통 코드 — supplier_items.supply_type 미지정 시 폴백
class ChatService:
async def chat(self, engine: TenantEngine, req: Req_Chat) -> Res_Chat:
res = Res_Chat()
repo = ScriptRepository(engine.config, agent_config.tenants_dir)
sess_repo = ChatSessionRepository(engine.company_id)
# 1) 세션 확보 / 시작 (DB 영속 — 재시작/멀티워커 안전, P8-A)
session = await sess_repo.get(req.session_id) if req.session_id else None
# 새 세션 컨텍스트: 요청 페이로드 대신 DB(negotiation.sessions 등)에서 1회 조회.
# 행이 없으면(데모/테스트 직접 호출) 기본값 폴백.
db_ctx = None if session else await NegotiationContextLoader().load(req.session_id, engine.company_id)
rq_type = session.rq_type if session else (db_ctx.rq_type if db_ctx else _DEFAULT_RQ_TYPE)
chat_engine = ChatEngine(repo, rq_type=rq_type)
# 완전 자율 모드(AUTONOMY_MODE=1 + 번들 존재): 가격협상 판정 룰·카드 선택을 정책 행동으로 대체.
# decider 를 감싸 결정을 컨텍스트에 기록 → advance() 후 experience_logs 에 적재(_autonomy_learn).
autonomy = AutonomyStore.policy_for(engine)
if autonomy is not None:
def _decide(ctx, _p=autonomy):
act = _p.decide(ctx)
ctx["autonomy_pending"] = {"idx": AUTONOMY_ACTIONS.index(act), "kind": act.kind,
"q": act.counter_q, "s": act.strategy}
return act
chat_engine.autonomy_decider = _decide
# ① step desync 감지: backend 가 본 직전 봇 step(client_step)이 agent 세션 step 과 다르면 경고.
# agent 가 자기 step 을 정답으로 보고 진행하고(응답의 step/client_step 으로 backend 가 따라옴),
# 추적/리싱크 트리거를 위해 로깅 + 응답 desynced 플래그로 알린다.
desynced = False
if session is not None and req.client_step:
agent_step = session.step
agent_client_step = chat_engine.step_map.get(agent_step, agent_step)
if req.client_step not in (agent_step, agent_client_step):
desynced = True
LOG.w(
f"[ChatService] step desync: backend client_step={req.client_step!r} != "
f"agent step={agent_step!r}(client={agent_client_step!r}) session={session.session_id}"
)
if session is None:
selected_nego_cards = db_ctx.selected_nego_card_numbers if db_ctx else []
selected_wild_cards = db_ctx.selected_wild_card_numbers if db_ctx else []
# 운영 DB 세션은 견적 version_id 에 묶인 카드만 사용한다. 직접 호출/데모(DB context 없음)는
# 기존 테넌트 기본 action mapping 으로 폴백해 로컬 테스트와 콘솔 데모를 유지한다.
action_space_size = (
min(len(selected_nego_cards), engine.action_space_size)
if db_ctx is not None
else engine.action_space_size
)
# session_id honoring: backend 가 보낸 session_id(= negotiation.sessions.session_id)를
# 새 uuid 발급 없이 그대로 세션 키로 쓴다. 없으면(직접 호출/데모) 생성.
session = ChatSession(
session_id=req.session_id or str(uuid.uuid4()), tenant_id=engine.tenant_id, company_id=engine.company_id,
rq_type=rq_type, action_space_size=action_space_size,
context={
# 매출액 = suppliers.total_revenue, 유통코드 = supplier_items.supply_type 매핑 (loader).
# 미기재/미지정이면 기본값 폴백.
"revenue_amount": db_ctx.revenue_amount if db_ctx and db_ctx.revenue_amount > 0 else _DEFAULT_REVENUE_AMOUNT,
"distribution_code": db_ctx.distribution_code if db_ctx and db_ctx.distribution_code else _DEFAULT_DISTRIBUTION_CODE,
# 파트너 유형(PartnerType 값 0/1/2) — 상품별 협력사 수 DB 조회로 세션 시작 시 1회 확정.
# snapshot.partner_count 로 그대로 사용(값 호환).
# 가격 수용률은 컨텍스트에 두지 않는다 — _snapshot 이 라운드별 제시가로 동적 계산.
"partner_count": int(db_ctx.partner_type) if db_ctx else int(PartnerType.SINGLE),
# 목표가/앵커링가: sessions 행(생성 시 박제된 anchoring_price) → 박제 ‰ → 1% 폴백 (loader).
"anchor_price": db_ctx.anchor_price if db_ctx else _DEFAULT_ANCHOR_PRICE,
"target_price": db_ctx.target_price if db_ctx else _DEFAULT_TARGET_PRICE,
"round": 0,
# 기존 공급가(품목 기준가) — 가격협상_확인 인하율 산출용.
"item_price": db_ctx.item_price if db_ctx else 0,
# 견적 생성 시 선택한 카드. 일반카드는 action_id 0..N-1 에 그대로 매핑한다.
# 1% 인하는 기본 와일드카드로 항상 열고, 재원부족 등 선택형 와일드카드는
# 선택된 와일드카드가 있을 때만 허용한다.
"db_context_loaded": db_ctx is not None,
"selected_nego_card_numbers": selected_nego_cards,
"selected_wild_card_numbers": selected_wild_cards,
"allow_selected_wildcards": True if db_ctx is None else bool(selected_wild_cards),
# ---- 자율 에이전트 v3 특징 소스 (미상이면 키 자체를 중립값으로 — JSON 직렬화 안전) ----
"internet_lowest_price": db_ctx.internet_lowest_price if db_ctx else 0,
"deadline_end_ts": db_ctx.deadline_end_ts if db_ctx else None,
"deadline_total_s": db_ctx.deadline_total_s if db_ctx else None,
"hist_n": db_ctx.hist_n if db_ctx else 0,
"hist_success": db_ctx.hist_success if db_ctx else None,
"hist_settle_ratio": db_ctx.hist_settle_ratio if db_ctx else None,
},
)
view = chat_engine.start(session)
else:
view = chat_engine.advance(session, req.user_input)
# 2) 응답 기본 채움 (학습 블록이 가격협상 턴에서 script/indicator 를 덮어쓸 수 있어 먼저 채운다)
res.session_id = session.session_id
res.step = view.step
res.client_step = view.client_step
res.script = view.script
res.input_mode = view.input_mode
res.input_options = view.input_options
res.chat_end = view.chat_end
res.outcome = view.outcome
res.desynced = desynced
# 합의가: 성공 확정 이후 턴(협상완료 요약 → 협상종료)에 내려준다. 와일드카드 수락처럼
# 유저가 직접 입력하지 않은 가격으로 타결될 수 있어 backend 요약/입찰가는 이 값을 최우선으로 쓴다.
if session.context.get("final_outcome") == "success" and session.context.get("input_price"):
res.settled_price = int(session.context["input_price"])
# 3) 학습 결합 (가격협상 카드선택 → 카드 스크립트·협상지표 / 종료 보상)
if view.error is None and session.action_space_size > 0:
if view.needs_card_selection:
await self._select_and_learn(engine, chat_engine, repo, session, res)
elif view.outcome is not None:
await self._terminal_learn(engine, session, view.outcome, res)
# 3-b) 완전 자율 모드: 정책 결정·종료 결과를 experience_logs 에 적재 (실로그 재학습 재료).
if autonomy is not None and view.error is None:
await self._autonomy_learn(engine, session, view, res)
# 자율 스텝 멘트를 LLM 으로 생성 (행동은 RL, 문장은 LLM). 실패/미설정 → 템플릿 유지.
if view.step.startswith("자율_"):
llm_ment = await ment_generator.generate(view.step, session.context)
if llm_ment:
res.script = llm_ment
# 직전 봇 멘트 보존 — 다음 생성에서 같은 문장 구조 반복을 금지하는 힌트.
session.context["autonomy_last_ment"] = (res.script or "")[:200]
if view.error:
res.result.SetResult(ErrorType.NEGO_INVALID_STEP)
res.msg = view.error
await sess_repo.save(session) # 진행 상태 영속화 (재시작/멀티워커 안전)
return res
# ---- 세션 상태 조회 (① desync 리싱크용) ----------------------------
async def get_session_state(self, engine: TenantEngine, session_id: str) -> "Res_ChatSession":
"""backend 가 타임아웃/재진입 시 agent 의 현재 step 을 읽어 정합을 맞춘다."""
res = Res_ChatSession()
repo = ScriptRepository(engine.config, agent_config.tenants_dir)
session = await ChatSessionRepository(engine.company_id).get(session_id)
if session is None:
res.result.SetResult(ErrorType.NEGO_SESSION_NOT_FOUND)
return res
step_map = repo.client_step_mapping()
res.session_id = session.session_id
res.step = session.step
res.client_step = step_map.get(session.step, session.step)
res.rq_type = session.rq_type
res.ended = session.ended
res.found = True
return res
# ---- 학습 ----------------------------------------------------------
@staticmethod
def _acceptance_ratio(context: dict) -> float:
"""가격 수용률 동적 계산 — 기준가 대비 현재 제시가의 양보율.
기준가 = 기존 공급가(item_price) 우선. 스크립트 구조상 협력사의 첫 제시가부터
기존 공급가 대비 인하가 반영되므로(가격협상_확인 멘트의 discount_rate 와 동일 기준)
첫 라운드부터 실값이 나온다 — 첫 제시가 기준 0 아님.
기존 공급가가 없으면(신규 협상 등) 협력사 첫 제시가 기준 폴백(첫 라운드 0).
acceptance = max(0, (기준가 − 현재 제시가) / 기준가). 가격 미입력이면 0(low 버킷).
"""
base = context.get("item_price") or context.get("first_offer_price") or 0
current = context.get("input_price") or 0
if base <= 0 or current <= 0:
return 0.0
return max(0.0, (base - current) / base)
def _snapshot(self, session: ChatSession, outcome: NegotiationOutcome) -> NegotiationSnapshot:
c = session.context
return NegotiationSnapshot(
revenue_amount=c["revenue_amount"], distribution_code=c["distribution_code"],
partner_count=c["partner_count"], acceptance_ratio=self._acceptance_ratio(c),
input_price=c.get("input_price", c["anchor_price"]), anchor_price=c["anchor_price"],
target_price=c["target_price"], round_number=c.get("round", 0), outcome=outcome,
)
async def _select_and_learn(self, engine: TenantEngine, chat_engine: ChatEngine,
scripts: ScriptRepository, session: ChatSession, res: Res_Chat):
snap = self._snapshot(session, NegotiationOutcome.ONGOING)
try:
idx = state_index(snap, engine.config.state)
except ValueError as ex:
LOG.e_no_callstack(f"[ChatService] state error: {ex}")
return
policy, version_id, repo = await QTablePolicyStore.load(engine)
# action space 는 카탈로그 전체(engine.action_space_size)로 고정 — action_id↔카드 대응을
# 견적마다 일정하게 유지해 Q-table 학습 일관성을 지킨다. 견적 선택은 축소가 아니라
# available_mask 로 처리한다(선택 카드만 pickable, 사용분 제외).
ctx = PolicyContext(state_index=idx, snapshot=snap, action_space_size=engine.action_space_size,
available_mask=self._selection_mask(engine, session),
episode=EpisodeState(used_action_ids=set(session.used_action_ids)))
# 카드 '선택'은 DQN 서빙(활성 시), 학습/영속은 아래 Q-table 경로 그대로(오프폴리시 갱신).
# DQN 불가(비활성/번들 없음/후보 특징 없음)면 None → 기존 UCB 선택 폴백.
dqn = DQNServingStore.policy_for(engine)
decision = dqn.select(ctx) if dqn is not None else None
selector_name = dqn.name if decision is not None else policy.name
if decision is None:
decision = policy.select(ctx)
session.used_action_ids.add(decision.action_id)
card_id = self._card_id_for_action(engine, session, decision.action_id)
reward = RewardCalculator(engine.config.reward, engine.config.state).calculate(snap)
policy.update(Transition(state_index=idx, action_id=decision.action_id, reward=reward.total, done=False))
await QTablePolicyStore.persist_cell(repo, version_id, policy, idx, decision.action_id)
session.context["last_state"] = idx
session.context["last_action"] = decision.action_id
await self._log(repo, session, idx, decision.action_id, card_id, snap, reward, decision.propensity, done=False)
res.card_id = card_id
res.policy = selector_name
res.q_value = decision.q_value
res.updated_q = float(policy.qtable.q[idx, decision.action_id])
res.visit_count = int(policy.qtable.visits[idx, decision.action_id])
res.reward_total = reward.total
# 가격협상 턴 연출(선행 chat_server 의 dynamic step type=indicator 재현):
# ① 선택된 카드의 스크립트를 봇 메시지(script)로 출력 ② 협상지표 게이지(indicator_value) 동봉.
# backend/front 가 indicator/bot_chat_type 패스스루·게이지 렌더 준비 완료 → 값만 채우면 표시된다.
# 카드 멘트: backoffice_db 모드면 card.nego_cards.script(negodata 편집 정본), 아니면 파일 폴백.
card_script = await scripts.resolve_card_script(
decision.action_id,
card_id,
chat_engine.vars_for(session),
prefer_db=bool(session.context.get("selected_nego_card_numbers")),
)
if card_script:
res.script = card_script
c = session.context
ind = compute_indicator(c.get("anchor_price", 0), c.get("input_price", 0), c.get("target_price", 0))
if ind is not None:
res.indicator_value = float(ind[0])
res.indicator_range = ind[1]
res.bot_chat_type = "indicator"
async def _autonomy_learn(self, engine: TenantEngine, session: ChatSession, view: StepView, res: Res_Chat):
"""완전 자율 행동 로깅 — Q-table 은 건드리지 않고 experience_logs 만 적재한다.
action_id = autonomy_actions.ACTIONS 인덱스, card_id = "AUT|종류|위치|전략" (카드 재학습
파이프라인이 임베딩 매칭에서 자동 제외하도록 프리픽스로 구분). 종료 시 최종 보상 행(done=True)을
남겨 retrain 의 에피소드 재구성 규약(카드턴 N + 종료 1)과 정합을 맞춘다.
"""
def _card_id(d) -> str:
return f"AUT|{d['kind']}|{d['q']:g}|{d['s']}"[:40]
ctx = session.context
lrepo = LearningRepository(engine.company_id)
pending = ctx.pop("autonomy_pending", None)
if pending is not None:
snap = self._snapshot(session, NegotiationOutcome.ONGOING)
try:
idx = state_index(snap, engine.config.state) # 로깅 호환용 이산 인덱스
except ValueError:
idx = 0 # 자율 모드는 이산 상태를 쓰지 않으므로 폴백해도 학습 오염 없음
if ctx.get("autonomy_last"):
ctx["autonomy_prev"] = ctx["autonomy_last"] # 직전 결정 보존 — 멘트 생성 힌트(양보 언급)용
ctx["autonomy_last"] = dict(pending, state_index=idx)
if pending.get("kind") == "counter":
# 역제안 기억은 별도 키로 보존 — autonomy_last 는 '마지막 행동'이라 사이에 낀
# 설득이 덮어쓴다. 단조 봉투·탄약소진 판정이 이 기억을 기준으로 해야
# counter→press→counter 에서 제안 철회가 새지 않는다 (게이트가 잡은 실버그).
ctx["autonomy_last_counter"] = dict(pending)
if pending.get("kind") == "press":
# 설득 횟수 누적 — 역제시 해금 조건(autonomy_store ③)의 카운터.
ctx["autonomy_press_n"] = int(ctx.get("autonomy_press_n") or 0) + 1
reward = RewardCalculator(engine.config.reward, engine.config.state).calculate(snap)
await self._log(lrepo, session, idx, pending["idx"], _card_id(pending), snap,
reward, (1.0 - 0.1) + 0.1 / len(AUTONOMY_ACTIONS), done=False)
res.policy = "full_autonomy"
last = ctx.get("autonomy_last")
if view.outcome is not None and last is not None:
oc = NegotiationOutcome.SUCCESS if view.outcome == "success" else NegotiationOutcome.FAILURE
snap = self._snapshot(session, oc)
reward = RewardCalculator(engine.config.reward, engine.config.state).calculate(snap)
res.reward_total = reward.total
await self._log(lrepo, session, last["state_index"], last["idx"], _card_id(last), snap,
reward, None, done=True)
async def _terminal_learn(self, engine: TenantEngine, session: ChatSession, outcome: str, res: Res_Chat):
oc = NegotiationOutcome.SUCCESS if outcome == "success" else NegotiationOutcome.FAILURE
snap = self._snapshot(session, oc)
reward = RewardCalculator(engine.config.reward, engine.config.state).calculate(snap)
res.reward_total = reward.total
last_state = session.context.get("last_state")
last_action = session.context.get("last_action")
if last_state is None or last_action is None:
return # 카드선택 없이 종료된 경우(예: 담당자확인 단계 이탈)
policy, version_id, repo = await QTablePolicyStore.load(engine)
policy.update(Transition(state_index=last_state, action_id=last_action, reward=reward.total, done=True))
await QTablePolicyStore.persist_cell(repo, version_id, policy, last_state, last_action)
await self._log(repo, session, last_state, last_action,
self._card_id_for_action(engine, session, last_action), snap, reward, None, done=True)
res.updated_q = float(policy.qtable.q[last_state, last_action])
@staticmethod
def _card_id_for_action(engine: TenantEngine, session: ChatSession, action_id: int) -> Optional[str]:
# action_id ↔ 카드는 테넌트 매핑(action_to_card)으로 고정한다. 견적 선택은 available_mask 로
# 걸러지므로 여기서 selected 리스트를 인덱싱하지 않는다 — 인덱싱하면 견적마다 action_id 의미가
# 달라져(같은 action_id 가 다른 카드) Q-table 학습이 오염된다.
return engine.mapper.get_card_id(action_id)
@staticmethod
def _selection_mask(engine: TenantEngine, session: ChatSession) -> Optional[np.ndarray]:
"""견적에서 선택한 카드(selected_nego_card_numbers)만 pickable 로 하는 available_mask.
action space 전체(engine.action_space_size) 크기의 bool 배열. 선택 카드의 action_id 만 True,
이미 사용한 action 은 False. 선택이 없거나(직접호출/데모) 매핑 불가면 None → 전체 허용(폴백).
번호(card.nego_cards.number)와 action_to_card 값이 일치해야 매핑된다.
"""
selected = session.context.get("selected_nego_card_numbers") or []
if not selected:
return None
used = set(session.used_action_ids)
selected_ids = {engine.mapper.get_action_id(str(n)) for n in selected}
selected_ids.discard(None)
if not selected_ids:
return None # 매핑에 없는 번호뿐 → 폴백(전체 허용)
return np.array(
[(a in selected_ids and a not in used) for a in range(engine.action_space_size)],
dtype=bool,
)
async def _log(self, repo: LearningRepository, session, state_index, action_id, card_id, snap, reward, propensity, done):
data = {
"session_id": session.session_id, "state_index": state_index, "action_id": action_id,
"card_id": card_id, "snapshot": snap.to_dict(), "propensity": propensity,
"turn": snap.round_number, "reward": reward.total, "done": done,
"settled_price": int(snap.input_price) if snap.outcome == NegotiationOutcome.SUCCESS else None,
}
try:
await DB_SESSION_MNG.execute_lambda_run([DBType.MAIN.value], [lambda s: repo.log_transition(s, data)])
except Exception as ex:
LOG.e_no_callstack(f"[ChatService] log failed: {ex}")
def reset_sessions():
"""세션은 DB(learning.chat_sessions)에 영속화된다(P8-A). 테스트는 db_engine 픽스처가 TRUNCATE 하므로 no-op."""
pass