Agent 2.0 exdev 서버 배포 스택

- server-dev start/stop/deploy 및 Gitea push 자동 배포
- local-dev 로컬 개발 환경

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
Macbook
2026-07-21 22:57:30 +09:00
commit 4b86b2a660
344 changed files with 45787 additions and 0 deletions
@@ -0,0 +1,907 @@
"""
OpenAI-style tool calling agent loop.
"""
import json
import os
import re
from datetime import datetime, timezone
from typing import Any, Dict, List, Optional
from agent.pending_store import AgentPendingStore
from agent.tool_executor import ToolExecutor
AGENT_SYSTEM_PROMPT = """당신은 한국도로공사 AI 챗봇 에이전트입니다.
사용자 질문에 답하기 위해 제공된 tool을 적절히 호출하세요.
규칙:
1. 실시간 DB/운영 데이터(통행료, 미납, IC전화, 도로정체, 휴게소 주유·음식·매장 등)는 해당 domain tool을 사용하세요.
2. 일반 FAQ/절차/안내는 rag_search tool로 지식베이스를 검색하세요.
3. 필수 정보가 부족하면 ask_user tool로 사용자에게 되물으세요.
4. tool 결과를 바탕으로 한국어로 정확하고 친절하게 최종 답변을 작성하세요.
5. 반드시 tool(rag_search 또는 domain tool)로 얻은 결과에 있는 내용만 사용하세요.
tool 결과에 없는 제도·요금·정책·수치·날짜는 절대 추측하거나 만들어내지 마세요.
6. 업무·정보성 질문은 반드시 먼저 적절한 tool을 호출하세요. tool 없이 임의로 답변하지 마세요.
- 이전 대화 이력에 비슷한 내용이 있어 보여도, 지식·제도·정책·요금 등 정보성 질문이면 매 턴마다 다시 rag_search(또는 domain tool)를 호출하세요.
- 대화 이력만 근거로 사실 답변을 재생성하지 마세요. 근거는 항상 이번 턴 tool 결과에서 가져와야 합니다.
7. tool 결과에서 근거를 찾지 못하면, 정확한 정보를 확인하기 어렵다고 안내하고 한국도로공사 콜센터(1588-2504)로 문의하도록 하세요.
8. 이전 턴 pending intent가 있고 사용자가 누락 파라미터만 짧게 답한 경우, pending intent의 파라미터로 해석하세요.
"""
class AgentService:
"""Runs LLM tool-calling loop with local RAG + remote domain tools."""
MAX_ROUNDS = int(os.getenv("AGENT_MAX_ROUNDS", "6"))
def __init__(
self,
llm_client,
tool_executor: ToolExecutor,
config,
pending_store: Optional[AgentPendingStore] = None,
prompt_builder=None,
llm_handler=None,
intent_detector=None,
greeting_handler=None,
emotion_detector=None,
emotion_handler=None,
suggestion_handler=None,
chat_manager=None,
response_handler=None,
query_rewriter=None,
):
self.llm_client = llm_client
self.tool_executor = tool_executor
self.config = config
self.pending_store = pending_store or AgentPendingStore()
self.prompt_builder = prompt_builder
self.llm_handler = llm_handler
# Legacy parity 핸들러 (선택 주입)
self.intent_detector = intent_detector
self.query_rewriter = query_rewriter
self.greeting_handler = greeting_handler
self.emotion_detector = emotion_detector
self.emotion_handler = emotion_handler
self.suggestion_handler = suggestion_handler
self.chat_manager = chat_manager
self.response_handler = response_handler
def chat(
self,
query: str,
bot_id: Optional[str] = None,
*,
pending_intent_type: Optional[str] = None,
pending_params: Optional[Dict[str, Any]] = None,
) -> Dict[str, Any]:
self.tool_executor.reset_state()
ts = datetime.now(timezone.utc).isoformat()
stored_pending = self.pending_store.get(bot_id)
effective_pending_type = pending_intent_type or (
stored_pending.get("pendingIntentType") if stored_pending else None
)
effective_pending_params = pending_params or (
stored_pending.get("pendingParams") if stored_pending else None
) or {}
had_pending = effective_pending_type is not None
# ⓪ 인사/종료 특별의도 선처리 (Legacy /ask와 동일). pending 진행 중에는 건너뜀.
if not had_pending:
greeting_response = self._handle_special_intent(query, bot_id, ts)
if greeting_response is not None:
return greeting_response
# 대화 이력 (멀티턴 맥락) — tool 선택·슬롯필링에 활용
conversation_messages = self._load_history_messages(bot_id, ts)
tools = self.tool_executor.all_tools()
if len(tools) <= 2:
print(
"[AgentService] WARN: remote domain tools unavailable "
f"(only {len(tools)} tools: rag_search, ask_user)"
)
system_content = AGENT_SYSTEM_PROMPT
if effective_pending_type:
system_content += (
"\n\n【이전 턴 pending】\n"
f"- intent: {effective_pending_type}\n"
f"- 이미 수집된 파라미터: {json.dumps(effective_pending_params, ensure_ascii=False)}\n"
"- 사용자의 이번 발화가 누락 슬롯을 채우는 답이면 해당 intent tool을 재호출하세요."
)
messages: List[Dict[str, Any]] = [
{"role": "system", "content": system_content},
]
if conversation_messages:
messages.extend(conversation_messages)
messages.append({"role": "user", "content": query})
final_content = ""
for round_idx in range(self.MAX_ROUNDS):
try:
message = self.llm_client.chat_completion_message(
messages=messages,
tools=tools,
max_tokens=min(self.config.llm_max_tokens, 1024),
temperature=0.2,
)
except Exception as exc:
print(f"[AgentService] {ts} tool 선택 LLM 호출 실패 → 안전망 진행: {exc}")
break
tool_calls = message.get("tool_calls") or []
content = (message.get("content") or "").strip()
if not tool_calls:
parsed = self._parse_json_tool_call(content)
if parsed:
tool_calls = [parsed]
else:
final_content = self._strip_think_tags(content)
break
messages.append(
{
"role": "assistant",
"content": content or None,
"tool_calls": tool_calls,
}
)
for call in tool_calls:
fn = call.get("function") or {}
tool_name = fn.get("name") or call.get("name")
raw_args = fn.get("arguments") or call.get("arguments") or "{}"
arguments = raw_args if isinstance(raw_args, dict) else json.loads(raw_args)
if tool_name == "ask_user":
question = arguments.get("question") or "조회에 필요한 정보를 조금 더 알려주세요."
# intentType 누락 시: 이번 턴에 시도한 domain tool로 폴백 → 멀티턴 슬롯필링 유지
intent_for_pending = (
arguments.get("intentType")
or arguments.get("pendingIntentType")
or effective_pending_type
or self._last_domain_intent_from_trace()
)
return self._finalize_clarify(
bot_id=bot_id,
answer=question,
intent_type=intent_for_pending,
params=effective_pending_params,
missing_params=arguments.get("missingParams"),
user_query=query,
ts=ts,
)
tool_result_raw = self.tool_executor.execute(
tool_name,
arguments,
bot_id=bot_id,
user_input=query,
)
tool_result = json.loads(tool_result_raw) if isinstance(tool_result_raw, str) else tool_result_raw
if tool_name not in ("rag_search", "ask_user") and isinstance(tool_result, dict):
if tool_result.get("needsClarification"):
question = tool_result.get("clarificationQuestion") or "추가 정보가 필요합니다."
return self._finalize_clarify(
bot_id=bot_id,
answer=question,
intent_type=tool_result.get("intentType"),
fetch_owner=tool_result.get("fetchOwner"),
ui_type=tool_result.get("uiType"),
domain_data=tool_result.get("domainData"),
params=tool_result.get("params"),
ic_candidates=tool_result.get("icCandidates"),
missing_params=tool_result.get("missingParams"),
user_query=query,
ts=ts,
)
if tool_result.get("status") is False:
msg = (
tool_result.get("statusMsg")
or tool_result.get("error")
or "요청을 처리하지 못했습니다."
)
return self._build_response(
answer=msg,
route_type="agent",
references=[],
faq_urls=[],
)
messages.append(
{
"role": "tool",
"tool_call_id": call.get("id") or f"call_{round_idx}_{tool_name}",
"name": tool_name,
"content": tool_result_raw if isinstance(tool_result_raw, str) else json.dumps(tool_result, ensure_ascii=False),
}
)
continue
if not final_content:
final_content = "죄송합니다. 답변을 생성하지 못했습니다."
domain_payload = self.tool_executor.last_domain_result or {}
rag_payload = self.tool_executor.last_rag_result or {}
# 사용 가능한 domain 결과 수집 (복합 질의 시 여러 tool 결과 누적)
usable_domain_results = [
d
for d in (self.tool_executor.domain_results or [])
if self._has_usable_domain_data(d.get("domainData"))
]
if not usable_domain_results and self._has_usable_domain_data(domain_payload.get("domainData")):
usable_domain_results = [domain_payload]
has_rag = bool(rag_payload.get("status") and rag_payload.get("references"))
# [안전망] 어떤 경우라도 정보성 질문은 Qdrant를 조회한다.
# LLM이 대화이력만 보고 rag_search를 스킵하면 동일 질문에 답이 달라지는 문제가 발생하므로,
# domain 근거가 없고 rag_search가 한 번도 실행되지 않았다면 강제로 Qdrant를 조회한다.
# (domain tool이 데이터를 가져온 경우엔 그 자체가 근거이므로 강제 조회하지 않는다.)
#
# 단, pending(슬롯필링) 진행 중에는 이 안전망을 건너뛴다:
# - 강제 rag가 rag_payload.status=True를 만들면 아래 stale-pending 정리가 빈 결과에도 발동해
# 슬롯 채우는 중인 유효 pending이 지워진다.
# - domain 가드가 params={}로 clarify를 반환하면 누적된 pending 파라미터가 덮어써진다.
# pending 중 LLM이 tool 재호출에 실패하면 guidance로 폴백하되 pending은 그대로 보존한다.
if not had_pending and not usable_domain_results and self.tool_executor.last_rag_result is None:
# 강제 rag 전에 명백한 domain 의도(요금 계산/미납 조회 등)면 FAQ 대신 되물어 정확 흐름으로 유도.
# FAQ성 질문(할인/방법/절차 등)은 가드하지 않아 기존 rag 경로를 유지한다(다자녀할인 등 회귀 방지).
domain_guard = self._domain_intent_guard(query)
if domain_guard is not None:
print(
f"[AgentService] {ts} 강제 rag 대신 domain 되물음: "
f"intent={domain_guard['intent_type']}"
)
return self._finalize_clarify(
bot_id=bot_id,
answer=domain_guard["message"],
intent_type=domain_guard["intent_type"],
params={},
user_query=query,
ts=ts,
)
self._force_rag_search(query, bot_id, ts)
rag_payload = self.tool_executor.last_rag_result or {}
has_rag = bool(rag_payload.get("status") and rag_payload.get("references"))
# 감정 분석 (부정 감정 시 공감 톤 지시) — Legacy /ask parity
emotion_instruction, emotion_name = self._detect_emotion(query, ts)
grounded_used = False
# 환각 방지: 최종 답변은 반드시 tool/qdrant 근거로만 생성한다.
# 근거가 있으면 PromptBuilder로 재작성, 근거가 없으면 guidance(콜센터 안내)로 폴백.
guidance_used = False
if usable_domain_results:
combined_domain_data = self._combine_domain_data(usable_domain_results)
grounded = self._generate_grounded_answer(
query=query,
rag_payload=rag_payload,
domain_data=combined_domain_data,
conversation_history=conversation_messages,
emotion_instruction=emotion_instruction,
emotion_name=emotion_name,
)
if grounded:
final_content = grounded
grounded_used = True
elif has_rag:
# rag-only(FAQ) 경로도 Legacy/Admin과 동일한 PromptBuilder로 최종 답변 생성
grounded = self._generate_grounded_answer(
query=query,
rag_payload=rag_payload,
domain_data=None,
conversation_history=conversation_messages,
emotion_instruction=emotion_instruction,
emotion_name=emotion_name,
)
if grounded:
final_content = grounded
grounded_used = True
else:
# tool/qdrant 근거 없음 → agent LLM 자유 답변 폐기, 안내 프롬프트로 폴백
guidance = self._guidance_answer(query, bot_id, conversation_messages)
if guidance:
final_content = guidance
guidance_used = True
references = self._resolve_references(domain_payload, rag_payload)
faq_urls = self._extract_faq_urls(references)
# 낮은 신뢰도 시 대안 질문 제안(💡) — rag 근거가 있을 때만
if grounded_used and has_rag:
final_content = self._apply_suggestions(final_content, rag_payload)
if domain_payload.get("intentType"):
self.pending_store.clear(bot_id)
elif (
had_pending
and rag_payload.get("status")
and not domain_payload.get("intentType")
):
# 무관 FAQ 등 rag_search만으로 답한 경우 stale pending 제거
self.pending_store.clear(bot_id)
history_metadata = {
"type": "no_match" if guidance_used else "agent",
"routeType": self._resolve_route_type(domain_payload),
"intentType": domain_payload.get("intentType"),
"num_references": len(references),
"searchMode": rag_payload.get("searchMode"),
"candidateCount": rag_payload.get("candidateCount"),
}
if guidance_used:
# chatbotAdmin 통계가 guidance를 정상(success)으로 보지 않도록 명시한다.
history_metadata.update(
{
"answer_confidence": "low",
"reason": "no_grounding",
"statusMsg": "no_match",
"searchMode": rag_payload.get("searchMode"),
"keywordRetryUsed": rag_payload.get("keywordRetryUsed"),
"keywordRetryAccepted": rag_payload.get("keywordRetryAccepted"),
}
)
# 대화 이력 저장 (멀티턴 컨텍스트 유지) — Legacy /ask parity
self._save_history(
bot_id=bot_id,
user_query=query,
answer=final_content,
references=references,
metadata=history_metadata,
ts=ts,
)
return self._build_response(
answer=final_content,
route_type=self._resolve_route_type(domain_payload),
intent_type=domain_payload.get("intentType"),
fetch_owner=domain_payload.get("fetchOwner"),
ui_type=domain_payload.get("uiType"),
domain_data=domain_payload.get("domainData"),
params=domain_payload.get("params"),
references=references,
faq_urls=faq_urls,
)
def _finalize_clarify(
self,
*,
bot_id: Optional[str],
answer: str,
intent_type: Optional[str] = None,
fetch_owner: Optional[str] = None,
ui_type: Optional[str] = None,
domain_data: Optional[Dict[str, Any]] = None,
params: Optional[Dict[str, Any]] = None,
ic_candidates: Optional[List[Any]] = None,
missing_params: Optional[List[Any]] = None,
user_query: Optional[str] = None,
ts: Optional[str] = None,
) -> Dict[str, Any]:
if intent_type:
self.pending_store.save(
bot_id,
pending_intent_type=intent_type,
pending_params=params or {},
missing_params=missing_params,
)
# clarify(되물음) 턴도 대화 이력에 저장 → 멀티턴 맥락/재작성기가 참조 가능
if user_query is not None:
self._save_history(
bot_id=bot_id,
user_query=user_query,
answer=answer,
references=[],
metadata={"type": "clarify", "intentType": intent_type},
ts=ts or datetime.now(timezone.utc).isoformat(),
)
return self._build_response(
answer=answer,
route_type="clarify",
intent_type=intent_type,
fetch_owner=fetch_owner,
ui_type=ui_type,
domain_data=domain_data,
params=params,
needs_clarification=True,
clarification_question=answer,
ic_candidates=ic_candidates,
pending_intent_type=intent_type,
missing_params=missing_params,
)
def _resolve_references(
self,
domain_payload: Dict[str, Any],
rag_payload: Dict[str, Any],
) -> List[Any]:
domain_refs = domain_payload.get("references")
if isinstance(domain_refs, list) and domain_refs:
return domain_refs
rag_refs = rag_payload.get("references")
if isinstance(rag_refs, list):
return rag_refs
return []
def _extract_faq_urls(self, references: List[Any]) -> List[str]:
urls: List[str] = []
seen = set()
for ref in references:
if not isinstance(ref, dict):
continue
url = ref.get("url")
if url and url not in seen:
seen.add(url)
urls.append(url)
if len(urls) >= 3:
break
return urls
def _resolve_route_type(self, domain_payload: Dict[str, Any]) -> str:
fetch_owner = domain_payload.get("fetchOwner")
if fetch_owner == "WEB":
return "web_domain"
if domain_payload.get("intentType"):
return "domain"
return "agent"
def _has_usable_domain_data(self, domain_data: Optional[Dict[str, Any]]) -> bool:
if self.prompt_builder:
return self.prompt_builder._has_usable_domain_data(domain_data)
if not domain_data:
return False
status = domain_data.get("status")
if status is False:
return False
if isinstance(status, str) and status.lower() == "false":
return False
return True
def _rag_refs_to_prompt_format(
self,
rag_payload: Dict[str, Any],
) -> tuple[List[Dict[str, Any]], List[float]]:
references: List[Dict[str, Any]] = []
scores: List[float] = []
for ref in rag_payload.get("references") or []:
if not isinstance(ref, dict):
continue
references.append(
{
"q": ref.get("q") or ref.get("question") or "",
"a": ref.get("a") or ref.get("answer") or "",
"url": ref.get("url"),
"category": ref.get("category"),
"source": "faq",
}
)
score = ref.get("score")
scores.append(float(score) if score is not None else 0.0)
return references, scores
def _generate_grounded_answer(
self,
*,
query: str,
rag_payload: Dict[str, Any],
domain_data: Optional[Dict[str, Any]] = None,
conversation_history: Optional[List[Dict[str, str]]] = None,
emotion_instruction: Optional[str] = None,
emotion_name: Optional[str] = None,
) -> Optional[str]:
"""Legacy /ask와 동일한 PromptBuilder로 최종 문장을 생성.
- 도메인 tool 성공: domain_data(llmSummary 포함)를 【DB 조회 결과】 컨텍스트로 사용
- rag-only(FAQ): domain_data=None, FAQ references만으로 답변 (Admin/Legacy와 동일 경로)
- 대화 이력/감정 지시사항을 함께 반영 (Legacy parity)
"""
if not self.prompt_builder or not self.llm_handler:
return None
references, scores = self._rag_refs_to_prompt_format(rag_payload)
messages = self.prompt_builder.build_answer_prompt_messages(
original_query=query,
rewritten_query=None,
references=references,
scores=scores,
conversation_history=conversation_history or [],
emotion_instruction=emotion_instruction,
emotion_name=emotion_name,
domain_data=domain_data,
)
ts = datetime.now(timezone.utc).isoformat()
try:
return self.llm_handler.generate_answer_from_messages(messages, ts)
except Exception as exc:
print(f"[AgentService] {ts} 최종 LLM 답변 실패, 폴백 사용: {exc}")
if domain_data:
return self._domain_fallback_answer(domain_data)
return None
def _combine_domain_data(
self,
domain_results: List[Dict[str, Any]],
) -> Optional[Dict[str, Any]]:
"""복합 질의 시 여러 domain 결과의 llmSummary/필드를 프롬프트용으로 병합."""
usable = [d.get("domainData") for d in domain_results if d.get("domainData")]
if not usable:
return None
if len(usable) == 1:
return usable[0]
summaries: List[str] = []
merged_fields: Dict[str, Any] = {"status": True}
for dd in usable:
summary = dd.get("llmSummary")
if summary is not None and str(summary).strip():
summaries.append(str(summary).strip())
for key, value in dd.items():
if key in ("status", "statusMsg", "errorMsg", "llmSummary"):
continue
if value is not None and key not in merged_fields:
merged_fields[key] = value
if summaries:
merged_fields["llmSummary"] = "\n\n".join(summaries)
return merged_fields
# 강제 rag 폴백 시 명백한 domain 의도만 되묻기로 유도(고정밀). FAQ성 질문은 가드하지 않음.
_FAQ_HINT_TOKENS = (
"할인", "감면", "방법", "절차", "어떻게", "안내", "신청", "등록",
"해지", "종류", "자격", "대상", "무엇", "인가요", "되나요", "가능",
)
_CAR_NO_RE = re.compile(r"\d{2,3}[가-힣]\d{4}")
_ROUTE_RE = re.compile(r"(에서|부터).{0,15}(까지)")
def _last_domain_intent_from_trace(self) -> Optional[str]:
"""이번 턴 tool_trace에서 마지막으로 시도한 domain tool명(=intentType) 반환.
LLM이 ask_user를 intentType 없이 호출했을 때 pending 저장용 폴백으로 사용.
rag_search/ask_user는 domain intent가 아니므로 제외한다.
"""
try:
for entry in reversed(self.tool_executor.tool_trace or []):
name = entry.get("tool")
if name and name not in ("rag_search", "ask_user"):
return name
except Exception:
pass
return None
def _domain_intent_guard(self, query: str) -> Optional[Dict[str, str]]:
"""강제 rag 직전, 명백한 domain 의도면 rag(FAQ) 대신 되물음으로 유도.
- FAQ성 표현(할인/방법/절차 등)이 있으면 가드하지 않는다 → 기존 rag 경로 유지(회귀 방지).
- 전체 차량번호 패턴 → 미납/환불 조회 의도.
- 'A에서 B까지' 경로 + '얼마' → 통행요금 조회 의도.
"""
text = (query or "").strip()
if not text:
return None
if any(tok in text for tok in self._FAQ_HINT_TOKENS):
return None
compact = re.sub(r"\s+", "", text)
if self._CAR_NO_RE.search(compact):
return {
"intent_type": "FARE_UNPAID",
"message": (
"챗봇에서 미납 통행료 조회가 가능합니다. "
"조회하려면 전체 차량번호를 입력해 주세요. 예: 12가3456 미납 조회"
),
}
if self._ROUTE_RE.search(text) and "얼마" in text:
return {
"intent_type": "FARE_SEARCH",
"message": (
"통행요금 조회를 위해 출발 IC와 도착 IC를 알려주세요. "
"예: 판교에서 신갈까지"
),
}
return None
def _contextualize_search_query(
self, query: str, bot_id: Optional[str], ts: str
) -> str:
"""후속 질문('얼마야?' 등)은 대화 이력을 반영해 완결형 검색어로 재작성.
강제 rag는 LLM이 인자를 만들지 않고 원문 발화를 그대로 검색하므로,
맥락이 필요한 후속 질문은 여기서 재작성해 검색 리콜을 보전한다.
재작성이 불필요(새 주제)하면 원문을 그대로 사용한다.
"""
if not self.query_rewriter or not self.chat_manager or not bot_id:
return query
if not getattr(self.config, "query_rewrite_enabled", True):
return query
try:
history = self.chat_manager.get_recent_history(
bot_id=bot_id,
hours=getattr(self.config, "chat_history_hours", 24),
limit=getattr(self.config, "chat_history_limit", 10),
)
if not history:
return query
rewritten = self.query_rewriter.rewrite_query(query, history, ts)
if rewritten:
print(f"[AgentService] {ts} 강제 rag 검색어 맥락화: '{query}' → '{rewritten}'")
return rewritten
except Exception as exc:
print(f"[AgentService] {ts} 검색어 맥락화 실패(원문 사용): {exc}")
return query
def _force_rag_search(self, query: str, bot_id: Optional[str], ts: str) -> None:
"""LLM이 tool을 스킵해도 정보성 질문은 Qdrant를 반드시 조회한다(일관성/환각방지 안전망)."""
try:
search_query = self._contextualize_search_query(query, bot_id, ts)
print(f"[AgentService] {ts} 근거 없음 → rag_search 강제 실행: {search_query}")
self.tool_executor.execute(
"rag_search", {"query": search_query}, user_input=query
)
except Exception as exc:
print(f"[AgentService] {ts} 강제 rag_search 실패: {exc}")
def _guidance_answer(
self,
query: str,
bot_id: Optional[str],
conversation_messages: Optional[List[Dict[str, str]]] = None,
) -> Optional[str]:
"""근거를 찾지 못했을 때 Legacy /ask no-match와 동일한 안내(콜센터) 답변 생성."""
if not self.prompt_builder or not self.llm_handler:
return None
ts = datetime.now(timezone.utc).isoformat()
conversation_context = self._history_messages_to_text(conversation_messages)
try:
system_prompt, user_prompt = self.prompt_builder.build_guidance_prompt(
query, conversation_context
)
return self.llm_handler.generate_answer(system_prompt, user_prompt, ts)
except Exception as exc:
print(f"[AgentService] {ts} guidance 답변 실패, 기본 안내 사용: {exc}")
return (
"문의하신 내용은 현재 정확한 정보를 확인하기 어렵습니다.\n"
"정확한 확인이 필요한 경우 한국도로공사 콜센터(1588-2504)로 문의해 주세요."
)
# ── Legacy parity helpers ─────────────────────────────────
def _handle_special_intent(
self,
query: str,
bot_id: Optional[str],
ts: str,
) -> Optional[Dict[str, Any]]:
"""인사/종료 등 특별의도를 Legacy /ask와 동일하게 고정 응답으로 처리."""
if not self.intent_detector or not self.greeting_handler:
return None
try:
intent = self.intent_detector.detect(query, ts)
except Exception as exc:
print(f"[AgentService] {ts} intent 감지 실패: {exc}")
return None
if not getattr(intent, "is_special", None) or not intent.is_special():
return None
try:
greeting = self.greeting_handler.generate_response(
intent_name=intent.name,
query=query,
matched_keywords=getattr(intent, "matched_keywords", []),
bot_id=bot_id,
)
except Exception as exc:
print(f"[AgentService] {ts} 특별의도 응답 생성 실패: {exc}")
return None
answer = greeting.get("answer") or ""
self._save_history(
bot_id=bot_id,
user_query=query,
answer=answer,
references=[],
metadata={"type": "special_intent", "intent": intent.name},
ts=ts,
)
return self._build_response(
answer=answer,
route_type="greeting",
references=[],
faq_urls=[],
quick_replies=greeting.get("quick_replies"),
)
def _load_history_messages(self, bot_id: Optional[str], ts: str) -> List[Dict[str, str]]:
"""MongoDB 대화 이력을 messages 포맷으로 로드 (멀티턴 컨텍스트)."""
if not self.chat_manager or not bot_id:
return []
if not getattr(self.config, "chat_history_always_include", True):
return []
try:
messages = self.chat_manager.get_messages_for_llm(
bot_id=bot_id,
hours=getattr(self.config, "chat_history_hours", 24),
max_conversations=getattr(self.config, "chat_history_limit", 10),
)
return messages or []
except Exception as exc:
print(f"[AgentService] {ts} 대화 이력 조회 실패: {exc}")
return []
def _history_messages_to_text(
self,
conversation_messages: Optional[List[Dict[str, str]]],
) -> Optional[str]:
if not conversation_messages:
return None
lines = []
for msg in conversation_messages:
role = msg.get("role")
content = (msg.get("content") or "").strip()
if not content:
continue
speaker = "고객" if role == "user" else "상담원"
lines.append(f"{speaker}: {content}")
if not lines:
return None
return "【이전 대화 이력】\n" + "\n".join(lines)
def _detect_emotion(self, query: str, ts: str) -> tuple[Optional[str], Optional[str]]:
"""감정 분석 → (emotion_instruction, emotion_name). Legacy /ask parity."""
if not self.emotion_detector or not self.emotion_handler:
return None, None
try:
emotion = self.emotion_detector.detect(query, ts)
instruction = self.emotion_handler.get_emotion_instruction(emotion.primary)
return (instruction or None), emotion.primary
except Exception as exc:
print(f"[AgentService] {ts} 감정 분석 실패: {exc}")
return None, None
def _apply_suggestions(self, answer: str, rag_payload: Dict[str, Any]) -> str:
"""낮은 신뢰도 시 대안 질문 제안(💡) 추가. Legacy /ask parity."""
if not self.suggestion_handler:
return answer
top_results, top_scores = self._rag_refs_to_prompt_format(rag_payload)
if not top_results:
return answer
try:
return self.suggestion_handler.enhance_answer_with_suggestions(
answer=answer,
top_results=top_results,
top_scores=top_scores,
max_suggestions=3,
)
except Exception as exc:
print(f"[AgentService] 제안 문구 생성 실패: {exc}")
return answer
def _save_history(
self,
*,
bot_id: Optional[str],
user_query: str,
answer: str,
references: List[Any],
metadata: Dict[str, Any],
ts: str,
) -> None:
"""Agent 응답을 MongoDB 대화 이력에 저장 (멀티턴 유지)."""
if not self.response_handler or not bot_id:
return
try:
matched_questions = []
scores = []
for ref in references or []:
if not isinstance(ref, dict):
continue
q = ref.get("question") or ref.get("q")
if q:
matched_questions.append(q)
score = ref.get("score")
if score is not None:
scores.append(score)
self.response_handler.save_to_mongodb(
bot_id=bot_id,
user_query=user_query,
ai_response=answer,
matched_questions=matched_questions,
scores=scores,
metadata=metadata,
ts=ts,
)
except Exception as exc:
print(f"[AgentService] {ts} 대화 이력 저장 실패: {exc}")
def _domain_fallback_answer(self, domain_data: Dict[str, Any]) -> str:
status_msg = domain_data.get("statusMsg")
if status_msg is not None and str(status_msg).strip():
return str(status_msg).strip()
summary = domain_data.get("llmSummary")
if summary is not None and str(summary).strip():
lines = [
line.strip()
for line in str(summary).splitlines()
if line.strip() and "안내하세요" not in line
]
if lines:
return "\n".join(lines)
return "조회 결과를 안내드리지 못했습니다. 잠시 후 다시 시도해 주세요."
def _build_response(
self,
*,
answer: str,
route_type: str,
intent_type: Optional[str] = None,
fetch_owner: Optional[str] = None,
ui_type: Optional[str] = None,
domain_data: Optional[Dict[str, Any]] = None,
params: Optional[Dict[str, Any]] = None,
references: Optional[List[Any]] = None,
faq_urls: Optional[List[str]] = None,
needs_clarification: bool = False,
clarification_question: Optional[str] = None,
ic_candidates: Optional[List[Any]] = None,
pending_intent_type: Optional[str] = None,
missing_params: Optional[List[Any]] = None,
quick_replies: Optional[List[Any]] = None,
) -> Dict[str, Any]:
return {
"status": True,
"routeType": route_type,
"answer": answer,
"llmAnswer": answer,
"intentType": intent_type,
"fetchOwner": fetch_owner,
"uiType": ui_type,
"domainData": domain_data,
"params": params or {},
"faqUrls": faq_urls or [],
"references": references or [],
"needsClarification": needs_clarification,
"clarificationQuestion": clarification_question,
"pendingIntentType": pending_intent_type,
"missingParams": missing_params,
"icCandidates": ic_candidates,
"quickReplies": quick_replies or [],
"toolTrace": self.tool_executor.tool_trace,
"reason": "agent_tool_calling",
}
def _parse_json_tool_call(self, content: str) -> Optional[Dict[str, Any]]:
if not content:
return None
match = re.search(r"\{[\s\S]*\}", content)
if not match:
return None
try:
payload = json.loads(match.group(0))
except json.JSONDecodeError:
return None
tool_name = payload.get("tool") or payload.get("name")
if not tool_name:
return None
arguments = payload.get("arguments") or payload.get("params") or {}
return {
"id": "parsed_json_call",
"type": "function",
"function": {
"name": tool_name,
"arguments": json.dumps(arguments, ensure_ascii=False),
},
}
def _strip_think_tags(self, text: str) -> str:
if "<think>" in text and "</think>" in text:
return re.sub(
r"<think>.*?</think>\s*",
"",
text,
flags=re.DOTALL,
).strip()
return text
@@ -0,0 +1,55 @@
"""
chatbotApi tool registry / execution HTTP client.
"""
import os
from typing import Any, Dict, List, Optional
import httpx
CHATBOT_API_BASE_URL = os.getenv(
"CHATBOT_API_BASE_URL",
os.getenv("TOOL_API_BASE_URL", "http://127.0.0.1:8086/api"),
).rstrip("/")
INTERNAL_TOOL_API_KEY = os.getenv("INTERNAL_TOOL_API_KEY", "")
TOOL_TIMEOUT = float(os.getenv("CHATBOT_TOOL_TIMEOUT", "30"))
class ChatbotToolClient:
"""Fetch tool schemas and execute domain tools via chatbotApi."""
def __init__(self, base_url: Optional[str] = None, timeout: float = TOOL_TIMEOUT):
self.base_url = (base_url or CHATBOT_API_BASE_URL).rstrip("/")
headers = {}
if INTERNAL_TOOL_API_KEY:
headers["X-Internal-Tool-Key"] = INTERNAL_TOOL_API_KEY
self.client = httpx.Client(timeout=timeout, headers=headers)
def list_tool_definitions(self) -> List[Dict[str, Any]]:
url = f"{self.base_url}/v1/tools/definitions"
response = self.client.get(url)
response.raise_for_status()
payload = response.json()
return payload.get("tools") or []
def execute_tool(
self,
tool_name: str,
arguments: Dict[str, Any],
*,
bot_id: Optional[str] = None,
user_input: Optional[str] = None,
) -> Dict[str, Any]:
url = f"{self.base_url}/v1/tools/execute"
body = {
"toolName": tool_name,
"arguments": arguments or {},
"botId": bot_id,
"userInput": user_input,
}
response = self.client.post(url, json=body)
response.raise_for_status()
return response.json()
def close(self) -> None:
self.client.close()
@@ -0,0 +1,117 @@
"""
Agent pending intent state — MongoDB with in-memory fallback.
"""
import os
from datetime import datetime, timezone
from typing import Any, Dict, Optional
try:
from pymongo import MongoClient
from pymongo.errors import PyMongoError
except ImportError: # pragma: no cover
MongoClient = None
PyMongoError = Exception
MONGO_HOST = os.getenv("MONGO_HOST", "localhost")
MONGO_PORT = int(os.getenv("MONGO_PORT", "27017"))
MONGO_USER = os.getenv("MONGO_USER", "")
MONGO_PASSWORD = os.getenv("MONGO_PASSWORD", "")
MONGO_DATABASE = os.getenv("MONGO_DATABASE", "chat_history")
PENDING_COLLECTION = os.getenv("AGENT_PENDING_COLLECTION", "agent_pending")
PENDING_TTL_SECONDS = int(os.getenv("AGENT_PENDING_TTL_SECONDS", str(30 * 60)))
class AgentPendingStore:
"""Stores clarify/pending context per botId for multi-turn slot filling."""
def __init__(self):
self._memory: Dict[str, Dict[str, Any]] = {}
self._collection = None
self._init_mongo()
def _init_mongo(self) -> None:
if MongoClient is None:
print("[AgentPending] pymongo unavailable — using in-memory store")
return
try:
if MONGO_USER and MONGO_PASSWORD:
uri = (
f"mongodb://{MONGO_USER}:{MONGO_PASSWORD}@{MONGO_HOST}:{MONGO_PORT}/"
f"{MONGO_DATABASE}?authSource={MONGO_DATABASE}"
)
else:
uri = f"mongodb://{MONGO_HOST}:{MONGO_PORT}/"
client = MongoClient(uri, serverSelectionTimeoutMS=3000)
client.admin.command("ping")
db = client[MONGO_DATABASE]
self._collection = db[PENDING_COLLECTION]
self._collection.create_index("bot_id", unique=True, background=True)
self._collection.create_index(
[("updated_at", 1)],
expireAfterSeconds=PENDING_TTL_SECONDS,
background=True,
name="idx_pending_ttl",
)
print(f"[AgentPending] Mongo connected: {MONGO_HOST}:{MONGO_PORT}/{PENDING_COLLECTION}")
except PyMongoError as exc:
print(f"[AgentPending] Mongo unavailable ({exc}) — using in-memory store")
self._collection = None
def get(self, bot_id: Optional[str]) -> Optional[Dict[str, Any]]:
key = self._normalize_bot_id(bot_id)
if self._collection is not None:
doc = self._collection.find_one({"bot_id": key}, {"_id": 0})
if doc:
return {
"pendingIntentType": doc.get("pending_intent_type"),
"pendingParams": doc.get("pending_params") or {},
"missingParams": doc.get("missing_params") or [],
}
return None
return self._memory.get(key)
def save(
self,
bot_id: Optional[str],
*,
pending_intent_type: str,
pending_params: Optional[Dict[str, Any]] = None,
missing_params: Optional[list] = None,
) -> None:
key = self._normalize_bot_id(bot_id)
payload = {
"pendingIntentType": pending_intent_type,
"pendingParams": pending_params or {},
"missingParams": missing_params or [],
}
if self._collection is not None:
now = datetime.now(timezone.utc)
self._collection.update_one(
{"bot_id": key},
{
"$set": {
"bot_id": key,
"pending_intent_type": pending_intent_type,
"pending_params": pending_params or {},
"missing_params": missing_params or [],
"updated_at": now,
}
},
upsert=True,
)
return
self._memory[key] = payload
def clear(self, bot_id: Optional[str]) -> None:
key = self._normalize_bot_id(bot_id)
if self._collection is not None:
self._collection.delete_one({"bot_id": key})
return
self._memory.pop(key, None)
@staticmethod
def _normalize_bot_id(bot_id: Optional[str]) -> str:
if bot_id and str(bot_id).strip():
return str(bot_id).strip()
return "anonymous"
@@ -0,0 +1,324 @@
"""
Local + remote tool execution for the agent loop.
"""
import json
import re
from datetime import datetime, timezone
from typing import Any, Dict, List, Optional
from agent.chatbot_tool_client import ChatbotToolClient
RAG_SEARCH_TOOL = {
"type": "function",
"function": {
"name": "rag_search",
"description": (
"한국도로공사 FAQ/상담 지식베이스(Qdrant)에서 질문과 유사한 Q&A를 검색합니다. "
"통행료 안내, Hi-pass, 환불 절차, 민원, 일반 상담 FAQ 등 정적 지식 질문에 사용하세요."
),
"parameters": {
"type": "object",
"properties": {
"query": {"type": "string", "description": "검색할 질문 문장"},
},
"required": ["query"],
},
},
}
ASK_USER_TOOL = {
"type": "function",
"function": {
"name": "ask_user",
"description": (
"조회에 필요한 정보(차량번호, IC명, 휴게소명 등)가 부족할 때 사용자에게 "
"추가 질문을 합니다. 최종 답변 대신 clarification이 필요할 때만 사용하세요."
),
"parameters": {
"type": "object",
"properties": {
"question": {"type": "string", "description": "사용자에게 되물을 질문"},
"intentType": {
"type": "string",
"description": "되묻는 대상 intent (예: FARE_SEARCH, FARE_UNPAID). 첫 턴 clarify 시 필수.",
},
},
"required": ["question"],
},
},
}
class ToolExecutor:
"""Executes rag_search locally and domain tools via chatbotApi."""
_KEYWORD_RETRY_LIMIT = 3
_KEYWORD_RETRY_FIELDS = ("q", "a", "question", "answer", "category", "source")
_KEYWORD_RETRY_MAX_SCAN = 10000
_GENERIC_KEYWORD_QUERIES = {
"요금",
"통행료",
"할인",
"감면",
"환불",
"신청",
"방법",
"절차",
"문의",
"안내",
"고속도로",
"휴게소",
"하이패스",
}
def __init__(
self,
search_handler,
config,
chatbot_tool_client: Optional[ChatbotToolClient] = None,
):
self.search_handler = search_handler
self.config = config
self.chatbot_tool_client = chatbot_tool_client or ChatbotToolClient()
self.last_domain_result: Optional[Dict[str, Any]] = None
self.last_rag_result: Optional[Dict[str, Any]] = None
# 한 턴에 여러 domain tool이 성공한 경우(복합 질의) 모두 누적
self.domain_results: List[Dict[str, Any]] = []
self.tool_trace: List[Dict[str, Any]] = []
def reset_state(self) -> None:
"""Clear per-turn domain result and tool trace."""
self.last_domain_result = None
self.last_rag_result = None
self.domain_results = []
self.tool_trace = []
def load_remote_tools(self) -> List[Dict[str, Any]]:
try:
return self.chatbot_tool_client.list_tool_definitions()
except Exception as exc:
print(f"[ToolExecutor] remote tool definitions unavailable: {exc}")
return []
def all_tools(self) -> List[Dict[str, Any]]:
return [RAG_SEARCH_TOOL, ASK_USER_TOOL] + self.load_remote_tools()
def execute(
self,
tool_name: str,
arguments: Dict[str, Any],
*,
bot_id: Optional[str] = None,
user_input: Optional[str] = None,
) -> str:
started = datetime.now(timezone.utc).isoformat()
try:
if tool_name == "rag_search":
result = self._execute_rag_search(arguments)
elif tool_name == "ask_user":
result = {"status": True, "clarificationQuestion": arguments.get("question")}
else:
payload = self.chatbot_tool_client.execute_tool(
tool_name,
arguments,
bot_id=bot_id,
user_input=user_input,
)
if isinstance(payload, dict):
if payload.get("needsClarification") or payload.get("status") is True:
self.last_domain_result = payload
if payload.get("status") is True and payload.get("intentType"):
self.domain_results.append(payload)
result = payload
success = not (
isinstance(result, dict)
and result.get("status") is False
and not result.get("needsClarification")
)
self.tool_trace.append(
{
"tool": tool_name,
"arguments": arguments,
"startedAt": started,
"status": success,
}
)
return json.dumps(result, ensure_ascii=False)
except Exception as exc:
self.tool_trace.append(
{
"tool": tool_name,
"arguments": arguments,
"startedAt": started,
"status": False,
"error": str(exc),
}
)
return json.dumps({"status": False, "error": str(exc)}, ensure_ascii=False)
def _execute_rag_search(self, arguments: Dict[str, Any]) -> Dict[str, Any]:
query = (arguments or {}).get("query") or ""
ts = datetime.now().strftime("%H:%M:%S")
query_vec = self.search_handler.embed_query(query, ts)
if query_vec is None:
return {"status": False, "error": "embedding_failed", "references": []}
search_results = self.search_handler.search(query_vec, self.config.threshold, ts)
# no-match 시 완화된 threshold로 재검색 (Legacy /ask recall parity)
if not search_results:
retry_threshold = getattr(self.config, "threshold_rewrite", None)
if retry_threshold is not None and retry_threshold < self.config.threshold:
print(f"[ToolExecutor] rag_search no-match → threshold {retry_threshold} 재검색")
search_results = self.search_handler.search(query_vec, retry_threshold, ts)
# search 결과는 {"meta": {...}, "score": ...} 형태 → legacy /ask와 동일하게 meta 언랩
candidates = [r["meta"] for r in search_results if isinstance(r, dict) and r.get("meta")]
top_results, scores, _, _ = self.search_handler.rerank(query, candidates, ts)
references = self._build_references(top_results, scores)
search_info = getattr(self.search_handler, "last_search_info", {}) or {}
search_mode = search_info.get("searchMode") or "vector"
keyword_retry_used = False
keyword_retry_accepted = False
# 최종 guidance 직전 보조 검색: 벡터/완화 재검색이 모두 실패한 경우에만,
# 문자열 포함 결과 최대 3건을 reranker로 검증해 충분히 맞을 때만 근거로 채택한다.
if not references:
keyword_retry_used = self._keyword_retry_allowed(query)
if keyword_retry_used:
keyword_references = self._keyword_retry_with_rerank(query, ts)
if keyword_references:
references = keyword_references
search_mode = "keyword_retry"
keyword_retry_accepted = True
result = {
"status": True,
"query": query,
"references": references,
"referenceCount": len(references),
"searchMode": search_mode,
"candidateCount": search_info.get("candidateCount"),
"keywordRetryUsed": keyword_retry_used,
"keywordRetryAccepted": keyword_retry_accepted,
}
self.last_rag_result = result
return result
def _build_references(
self, top_results: List[Dict[str, Any]], scores: List[Any], limit: int = 5
) -> List[Dict[str, Any]]:
references = []
for item, score in zip((top_results or [])[:limit], (scores or [])[:limit]):
references.append(
{
"question": item.get("q") or item.get("question"),
"answer": item.get("a") or item.get("answer"),
"score": score,
"category": item.get("category"),
"url": item.get("url"),
"searchSource": item.get("_search_source"),
"denseScore": item.get("_dense_score"),
"sparseScore": item.get("_sparse_score"),
}
)
return references
def _keyword_retry_allowed(self, query: str) -> bool:
text = (query or "").strip()
if not text:
return False
normalized = re.sub(r"\s+", "", text).lower()
if len(normalized) < 3:
return False
if not re.search(r"[0-9a-zA-Z가-힣]", normalized):
return False
tokens = re.findall(r"[0-9a-zA-Z가-힣]+", text.lower())
if not tokens:
return False
compact_tokens = [re.sub(r"\s+", "", token) for token in tokens if token.strip()]
if len(compact_tokens) == 1 and compact_tokens[0] in self._GENERIC_KEYWORD_QUERIES:
return False
compact_query = "".join(compact_tokens)
if compact_query in self._GENERIC_KEYWORD_QUERIES:
return False
return True
def _keyword_retry_with_rerank(self, query: str, ts: str) -> List[Dict[str, Any]]:
vector_store = getattr(self.search_handler, "vector_store", None)
keyword_search = getattr(vector_store, "keyword_search", None)
if not callable(keyword_search):
return []
try:
keyword_items, total, exhausted = keyword_search(
query,
fields=self._KEYWORD_RETRY_FIELDS,
limit=self._KEYWORD_RETRY_LIMIT,
max_scan=self._KEYWORD_RETRY_MAX_SCAN,
count_total=False,
)
except TypeError:
# 이전 시그니처/테스트 더블 호환: count_total을 받지 못하면 기본 호출로 재시도
keyword_items, total, exhausted = keyword_search(
query,
fields=self._KEYWORD_RETRY_FIELDS,
limit=self._KEYWORD_RETRY_LIMIT,
)
except Exception as exc:
print(f"[ToolExecutor] keyword retry 실패: {exc}")
return []
if not isinstance(keyword_items, list) or not keyword_items:
print(f"[ToolExecutor] keyword retry no-match: query={query}")
return []
candidates = [
item.get("meta")
for item in keyword_items[: self._KEYWORD_RETRY_LIMIT]
if isinstance(item, dict) and isinstance(item.get("meta"), dict)
]
if not candidates:
return []
print(
f"[ToolExecutor] keyword retry → rerank 검증: "
f"query={query}, candidates={len(candidates)}, total_seen={total}, exhausted={exhausted}"
)
top_results, scores, rerank_used, _ = self.search_handler.rerank(query, candidates, ts)
if not top_results or not scores or scores[0] is None:
print("[ToolExecutor] keyword retry 거부: rerank 점수 없음")
return []
threshold = self._low_confidence_threshold()
try:
top_score = float(scores[0])
except (TypeError, ValueError):
print(f"[ToolExecutor] keyword retry 거부: 잘못된 rerank 점수={scores[0]}")
return []
if not rerank_used or top_score < threshold:
print(
f"[ToolExecutor] keyword retry 거부: score={top_score:.4f}, "
f"threshold={threshold:.4f}, rerank_used={rerank_used}"
)
return []
print(f"[ToolExecutor] keyword retry 채택: score={top_score:.4f}")
return self._build_references(
top_results[: self._KEYWORD_RETRY_LIMIT],
scores[: self._KEYWORD_RETRY_LIMIT],
limit=self._KEYWORD_RETRY_LIMIT,
)
def _low_confidence_threshold(self) -> float:
value = getattr(self.config, "low_confidence_threshold", 0.65)
try:
return float(value)
except (TypeError, ValueError):
return 0.65