""" KCSC MCP 서버 애플리케이션 한국건설기준센터(KCSC) API와 Model Context Protocol(MCP)을 연동한 자연어 검색 및 질의응답 API 서버 """ import logging import os import json import time import uuid from datetime import datetime from typing import List, Optional, Dict, Any, Union from fastapi import FastAPI, HTTPException, Depends, Request, BackgroundTasks from fastapi.middleware.cors import CORSMiddleware from fastapi.responses import JSONResponse from fastapi.openapi.docs import get_swagger_ui_html from fastapi.openapi.utils import get_openapi from pydantic import BaseModel, Field, validator from src import config from src.kcsc_api_client import KCSCApiClient from src.vector_db_client import KCSCVectorDB from src.mcp_processor import MCPProcessor from src.utils import ( format_error_response, KCSCBaseException, sanitize_string ) # 설정 초기화 config.init() logger = config.setup_logger("app") # FastAPI 앱 초기화 app = FastAPI( title="KCSC MCP API", description="한국건설기준센터(KCSC) API와 Model Context Protocol(MCP)을 연동한 자연어 검색 및 질의응답 API", version="1.0.1", docs_url=None, # 기본 문서 경로 비활성화 (커스텀 경로 사용) redoc_url=None, # Redoc 비활성화 ) # CORS 설정 origins = [ "http://localhost", "http://localhost:8000", "http://127.0.0.1", "http://127.0.0.1:8000", # 필요에 따라 추가 도메인 허용 ] app.add_middleware( CORSMiddleware, allow_origins=origins, allow_credentials=True, allow_methods=["*"], allow_headers=["*"], ) # 모델 정의 class SearchQuery(BaseModel): """검색 쿼리 모델""" query: str = Field(..., description="검색어") doc_types: Optional[List[str]] = Field(None, description="문서 유형 목록") limit: Optional[int] = Field(5, description="검색 결과 수", ge=1, le=20) @validator('query') def validate_query(cls, v): if not v or not v.strip(): raise ValueError('검색어는 비워둘 수 없습니다.') return sanitize_string(v.strip()) class MCPQuery(BaseModel): """MCP 쿼리 모델""" query: str = Field(..., description="질의 내용") conversation_id: Optional[str] = Field(None, description="대화 ID") enable_rag: bool = Field(True, description="RAG(검색 증강 생성) 사용 여부") system_prompt: Optional[str] = Field(None, description="시스템 프롬프트") temperature: Optional[float] = Field(0.2, description="생성 온도(창의성)", ge=0.0, le=1.0) max_tokens: Optional[int] = Field(2000, description="최대 생성 토큰 수", ge=1, le=4096) @validator('query') def validate_query(cls, v): if not v or not v.strip(): raise ValueError('질의 내용은 비워둘 수 없습니다.') return sanitize_string(v.strip()) class CodeDetailRequest(BaseModel): """코드 상세 정보 요청 모델""" code: str = Field(..., description="코드 번호") doc_type: str = Field(..., description="문서 유형") @validator('code', 'doc_type') def validate_fields(cls, v): if not v or not v.strip(): raise ValueError('필드는 비워둘 수 없습니다.') return v.strip() class SearchResult(BaseModel): """검색 결과 모델""" id: str = Field(..., description="결과 ID") text: str = Field(..., description="결과 텍스트") metadata: Dict[str, Any] = Field(..., description="메타데이터") relevance: float = Field(..., description="관련도 점수", ge=0.0, le=1.0) class SearchResponse(BaseModel): """검색 응답 모델""" results: List[SearchResult] = Field(..., description="검색 결과 목록") query: str = Field(..., description="원본 검색어") timestamp: str = Field(..., description="타임스탬프") class CollectRequest(BaseModel): """데이터 수집 요청 모델""" doc_types: Optional[List[str]] = Field(["KDS", "KCS"], description="수집할 문서 유형 목록") class IndexParams(BaseModel): """인덱스 구축 파라미터 모델""" chunk_size: Optional[int] = Field(None, description="청크 크기") overlap: Optional[int] = Field(None, description="청크 간 겹치는 문자 수") reset: Optional[bool] = Field(True, description="기존 데이터 리셋 여부") # 종속성 주입 - 싱글턴 객체 def get_kcsc_api(): """KCSC API 클라이언트 반환""" return KCSCApiClient() def get_vector_db(): """벡터 데이터베이스 클라이언트 반환""" return KCSCVectorDB() def get_mcp_processor(): """MCP 프로세서 반환""" return MCPProcessor() # 백그라운드 작업 def background_data_collection(api_client, doc_types): """백그라운드 데이터 수집 작업""" try: logger.info(f"백그라운드 데이터 수집 시작: {doc_types}") all_codes = api_client.collect_codes(doc_types) for doc_type, codes in all_codes.items(): logger.info(f"{doc_type} 상세 정보 수집 시작 (총 {len(codes)}개 코드)") api_client.fetch_code_details(doc_type, codes) logger.info(f"백그라운드 데이터 수집 완료") except Exception as e: logger.error(f"백그라운드 데이터 수집 중 오류 발생: {str(e)}") def background_index_building(vector_db, params): """백그라운드 인덱스 구축 작업""" try: logger.info("백그라운드 인덱스 구축 시작") vector_db.build_index_from_files( chunk_size=params.get("chunk_size"), overlap=params.get("overlap"), reset=params.get("reset", True) ) logger.info("백그라운드 인덱스 구축 완료") except Exception as e: logger.error(f"백그라운드 인덱스 구축 중 오류 발생: {str(e)}") # 에러 핸들러 @app.exception_handler(KCSCBaseException) async def kcsc_exception_handler(request: Request, exc: KCSCBaseException): """KCSC 예외 처리기""" return JSONResponse( status_code=500, content=format_error_response(exc) ) @app.exception_handler(Exception) async def general_exception_handler(request: Request, exc: Exception): """일반 예외 처리기""" logger.error(f"요청 처리 중 예외 발생: {str(exc)}") return JSONResponse( status_code=500, content={ "status": "error", "error": { "type": exc.__class__.__name__, "message": str(exc), "timestamp": datetime.now().isoformat() } } ) # 서버 시작 시 초기화 @app.on_event("startup") async def startup_event(): """서버 시작 이벤트 핸들러""" logger.info("서버 시작 중...") try: # 필요한 디렉토리 생성 os.makedirs(config.DATA_DIR, exist_ok=True) os.makedirs(config.LOGS_DIR, exist_ok=True) os.makedirs(config.VECTOR_DB_DIR, exist_ok=True) # 모듈 초기화 kcsc_api = get_kcsc_api() vector_db = get_vector_db() mcp_processor = get_mcp_processor() # 벡터 DB 초기화 상태 확인 stats = vector_db.get_collection_stats() has_data = any(count > 0 for count in stats.values()) if not has_data: logger.info("벡터 DB에 데이터가 없습니다. 데이터 수집 및 인덱싱이 필요합니다.") logger.info("서버 시작 완료") except Exception as e: logger.error(f"서버 시작 중 오류 발생: {str(e)}") # 오류가 있더라도 서버는 실행 (필요시 일부 기능만 제한) # API 엔드포인트 구현 @app.get("/") async def root(): """루트 엔드포인트""" return { "status": "online", "message": "KCSC MCP API 서버 실행 중", "version": app.version, "timestamp": datetime.now().isoformat() } # 문서 커스텀 엔드포인트 @app.get("/docs", include_in_schema=False) async def custom_swagger_ui_html(): """커스텀 Swagger UI""" return get_swagger_ui_html( openapi_url="/openapi.json", title=app.title + " - API 문서", swagger_js_url="https://cdn.jsdelivr.net/npm/swagger-ui-dist@4/swagger-ui-bundle.js", swagger_css_url="https://cdn.jsdelivr.net/npm/swagger-ui-dist@4/swagger-ui.css", ) @app.get("/openapi.json", include_in_schema=False) async def get_open_api_endpoint(): """OpenAPI 스키마""" return get_openapi( title=app.title, version=app.version, description=app.description, routes=app.routes, ) # MCP 프로토콜 메타데이터 엔드포인트 @app.get("/mcp") async def mcp_metadata(): """MCP 서버 메타데이터 - Claude Desktop에서 서버 검증에 사용""" return { "name": "KCSC API", "description": "한국건설기준센터(KCSC) API를 통한 검색 및 컨텍스트 제공", "version": app.version, "actions": [ { "name": "search", "description": "건설기준 정보 검색", "parameters": { "query": "검색어", "doc_types": "문서 유형 (선택적)", "limit": "결과 수 (선택적)" } }, { "name": "code_detail", "description": "특정 코드의 상세 정보 조회", "parameters": { "code": "코드 번호", "doc_type": "문서 유형" } } ] } # MCP 메시지 처리 엔드포인트 - Claude Desktop과 통신하는 핵심 부분 @app.post("/mcp/messages") async def handle_mcp_message( request: Request, vector_db: KCSCVectorDB = Depends(get_vector_db), kcsc_api: KCSCApiClient = Depends(get_kcsc_api), mcp_processor: MCPProcessor = Depends(get_mcp_processor) ): """ Claude Desktop과 호환되는 MCP 메시지 처리 이 엔드포인트는 MCP 프로토콜을 통해 Claude Desktop과 통신합니다. """ try: # 요청 데이터 파싱 data = await request.json() message = data.get("message", "") conversation_id = data.get("conversation_id") action = data.get("action") action_params = data.get("action_params", {}) logger.info(f"MCP 메시지 수신: action={action}, message_prefix={message[:30] if message else 'None'}...") # 액션이 없으면 기본 검색 수행 if not action: # 검색 쿼리로 처리 search_results = vector_db.search( query=message, limit=3 ) # 컨텍스트 형식으로 변환 context = [] for i, result in enumerate(search_results): context.append({ "id": f"result-{i}", "content": result["text"], "metadata": { "source": "KCSC", "code": result["metadata"]["code"], "name": result["metadata"]["name"], "doc_type": result["metadata"]["doc_type"], "relevance": result["relevance"] } }) logger.info(f"기본 검색 완료: {len(context)}개 결과") # 응답 생성 return { "content": "다음은 한국건설기준센터(KCSC)에서 검색한 관련 정보입니다:", "context": context, "conversation_id": conversation_id } # 액션별 처리 elif action == "search": # 검색 파라미터 추출 query = action_params.get("query", message) doc_types = action_params.get("doc_types") limit = int(action_params.get("limit", 5)) # 검색 수행 search_results = vector_db.search( query=query, doc_types=doc_types, limit=limit ) # 컨텍스트 형식으로 변환 context = [] for i, result in enumerate(search_results): context.append({ "id": f"result-{i}", "content": result["text"], "metadata": { "source": "KCSC", "code": result["metadata"]["code"], "name": result["metadata"]["name"], "doc_type": result["metadata"]["doc_type"], "relevance": result["relevance"] } }) logger.info(f"검색 액션 완료: '{query}'에 대해 {len(context)}개 결과") # 응답 생성 return { "content": f"'{query}'에 대한 검색 결과입니다:", "context": context, "conversation_id": conversation_id } elif action == "code_detail": # 코드 파라미터 추출 code = action_params.get("code") doc_type = action_params.get("doc_type") if not code or not doc_type: return { "content": "코드 번호와 문서 유형이 필요합니다.", "conversation_id": conversation_id } # 코드 상세 정보 조회 detail = kcsc_api.get_code_details(doc_type, code) if not detail: return { "content": f"코드 {doc_type}/{code}에 대한 정보를 찾을 수 없습니다.", "conversation_id": conversation_id } # 포맷팅된 상세 정보 formatted_detail = f""" 코드: {detail.get('Code')} 이름: {detail.get('Name')} 버전: {detail.get('Version')} 업데이트: {detail.get('UpdateDate')} 내용: {detail.get('Contents')} """ # 응답 생성 return { "content": formatted_detail, "context": [{ "id": f"{doc_type}-{code}", "content": detail.get('Contents', ''), "metadata": { "source": "KCSC", "code": detail.get('Code'), "name": detail.get('Name'), "doc_type": doc_type } }], "conversation_id": conversation_id } else: return { "content": f"지원하지 않는 액션입니다: {action}", "conversation_id": conversation_id } except Exception as e: logger.error(f"MCP 메시지 처리 중 오류 발생: {str(e)}") return { "content": f"오류 발생: {str(e)}", "conversation_id": conversation_id } @app.get("/status") async def get_status( kcsc_api: KCSCApiClient = Depends(get_kcsc_api), vector_db: KCSCVectorDB = Depends(get_vector_db) ): """서버 상태 확인""" try: # 벡터 DB 상태 확인 db_stats = vector_db.get_collection_stats() return { "status": "online", "vector_db": db_stats, "timestamp": datetime.now().isoformat(), "version": app.version } except Exception as e: logger.error(f"상태 확인 중 오류 발생: {str(e)}") raise HTTPException(status_code=500, detail=str(e)) @app.post("/collect") async def collect_data( request: CollectRequest, background_tasks: BackgroundTasks, kcsc_api: KCSCApiClient = Depends(get_kcsc_api) ): """ KCSC API에서 데이터 수집 시작 백그라운드에서 데이터를 수집하여 서버 응답을 즉시 반환합니다. """ try: doc_types = request.doc_types logger.info(f"데이터 수집 요청: {doc_types}") # 백그라운드에서 데이터 수집 시작 background_tasks.add_task(background_data_collection, kcsc_api, doc_types) return { "status": "accepted", "message": "데이터 수집이 백그라운드에서 시작되었습니다.", "doc_types": doc_types, "timestamp": datetime.now().isoformat() } except Exception as e: logger.error(f"데이터 수집 요청 처리 중 오류 발생: {str(e)}") raise HTTPException(status_code=500, detail=str(e)) @app.post("/index") async def build_index( params: IndexParams = None, background_tasks: BackgroundTasks = None, vector_db: KCSCVectorDB = Depends(get_vector_db) ): """ 벡터 데이터베이스 인덱스 구축 백그라운드에서 인덱스를 구축하여 서버 응답을 즉시 반환합니다. """ try: logger.info("벡터 인덱스 구축 요청") params_dict = params.dict() if params else {} # 백그라운드에서 인덱스 구축 시작 background_tasks.add_task(background_index_building, vector_db, params_dict) return { "status": "accepted", "message": "벡터 인덱스 구축이 백그라운드에서 시작되었습니다.", "parameters": params_dict, "timestamp": datetime.now().isoformat() } except Exception as e: logger.error(f"인덱스 구축 요청 처리 중 오류 발생: {str(e)}") raise HTTPException(status_code=500, detail=str(e)) @app.post("/search", response_model=SearchResponse) async def search( query: SearchQuery, vector_db: KCSCVectorDB = Depends(get_vector_db) ): """ 자연어 검색 수행 문서 내용을 자연어 쿼리로 검색합니다. """ try: logger.info(f"검색 요청: {query.query}") # 검색 수행 results = vector_db.search( query=query.query, doc_types=query.doc_types, limit=query.limit ) # 응답 형식 변환 formatted_results = [] for result in results: formatted_results.append(SearchResult( id=result["id"], text=result["text"], metadata=result["metadata"], relevance=result["relevance"] )) return SearchResponse( results=formatted_results, query=query.query, timestamp=datetime.now().isoformat() ) except Exception as e: logger.error(f"검색 중 오류 발생: {str(e)}") raise HTTPException(status_code=500, detail=str(e)) @app.post("/code_detail") async def get_code_detail( request: CodeDetailRequest, kcsc_api: KCSCApiClient = Depends(get_kcsc_api) ): """ 특정 코드의 상세 정보 조회 코드 번호와 문서 유형으로 상세 정보를 조회합니다. """ try: logger.info(f"코드 상세 정보 요청: {request.doc_type}/{request.code}") # 코드 상세 정보 조회 detail = kcsc_api.get_code_details(request.doc_type, request.code) if not detail: raise HTTPException(status_code=404, detail="해당 코드 정보를 찾을 수 없습니다.") return detail except HTTPException: raise except Exception as e: logger.error(f"코드 상세 정보 조회 중 오류 발생: {str(e)}") raise HTTPException(status_code=500, detail=str(e)) @app.post("/mcp_query") async def process_mcp_query( query: MCPQuery, vector_db: KCSCVectorDB = Depends(get_vector_db), mcp_processor: MCPProcessor = Depends(get_mcp_processor) ): """ MCP를 통한 쿼리 처리 RAG(검색 증강 생성)을 활용하여 질의에 답변합니다. """ try: logger.info(f"MCP 쿼리 요청: {query.query[:50]}...") # 컨텍스트 데이터 준비 (RAG 사용하는 경우) context_data = None if query.enable_rag: search_results = vector_db.search(query.query, limit=3) if search_results: context_data = [ { "id": result["id"], "content": result["text"], "metadata": { "title": result["metadata"]["name"], "code": result["metadata"]["code"], "doc_type": result["metadata"]["doc_type"], "relevance": result["relevance"] } } for result in search_results ] logger.debug(f"컨텍스트 데이터 준비 완료: {len(context_data)}개 항목") # MCP를 통한 처리 result = mcp_processor.process_query( query=query.query, context_data=context_data, conversation_id=query.conversation_id, system_prompt=query.system_prompt, temperature=query.temperature, max_tokens=query.max_tokens ) return { "response": result["content"], "conversation_id": result["conversation_id"], "used_context": result["used_context"], "timestamp": datetime.now().isoformat() } except Exception as e: logger.error(f"MCP 쿼리 처리 중 오류 발생: {str(e)}") raise HTTPException(status_code=500, detail=str(e)) @app.get("/health") async def health_check(): """ 헬스 체크 엔드포인트 서버 상태의 간단한 헬스 체크를 제공합니다. """ return { "status": "healthy", "timestamp": datetime.now().isoformat(), "service": "KCSC MCP API", "version": app.version } # 메인 실행 if __name__ == "__main__": import uvicorn uvicorn.run( "src.app:app", host=config.HOST, port=config.PORT, reload=config.DEBUG )